fix: enforce workspace deletion blockers
This commit is contained in:
@@ -948,7 +948,6 @@ CREATE TABLE workspace_deletion_operations (
|
||||
workspace_revision TEXT NOT NULL,
|
||||
owner_account_id TEXT NOT NULL,
|
||||
actor_account_id TEXT NOT NULL,
|
||||
force_delete_dirty_workdirs INTEGER NOT NULL CHECK(force_delete_dirty_workdirs IN (0, 1)),
|
||||
state TEXT NOT NULL CHECK(state IN ('queued', 'running', 'blocked', 'failed', 'succeeded')),
|
||||
resource_counts_json TEXT NOT NULL,
|
||||
child_operation_ids_json TEXT NOT NULL,
|
||||
|
||||
@@ -90,11 +90,11 @@ use workspace_api::{
|
||||
WorkingDirectoryRemovalDisposition, WorkingDirectoryRemovalRequest,
|
||||
WorkingDirectoryRemovalResponse, WorkingDirectoryRepositoryOption,
|
||||
WorkspaceCatalogListResponse, WorkspaceCreateResponse, WorkspaceDeletionBlocker,
|
||||
WorkspaceDeletionBlockerKind, WorkspaceDeletionOperationResponse, WorkspaceDeletionRequest,
|
||||
WorkspaceDeletionState, WorkspaceExtensionPointState, WorkspaceExtensionPoints,
|
||||
WorkspaceMetadataMutationResponse, WorkspaceMetadataSettingsResponse,
|
||||
WorkspacePermissionSummary, WorkspaceRepositoryRecord, WorkspaceResponse,
|
||||
WorkspaceRuntimeDetail, WorkspaceRuntimeResource, WorkspaceSummary,
|
||||
WorkspaceDeletionBlockerKind, WorkspaceDeletionOperationResponse,
|
||||
WorkspaceDeletionPreflightResponse, WorkspaceDeletionRequest, WorkspaceDeletionState,
|
||||
WorkspaceExtensionPointState, WorkspaceExtensionPoints, WorkspaceMetadataMutationResponse,
|
||||
WorkspaceMetadataSettingsResponse, WorkspacePermissionSummary, WorkspaceRepositoryRecord,
|
||||
WorkspaceResponse, WorkspaceRuntimeDetail, WorkspaceRuntimeResource, WorkspaceSummary,
|
||||
WorkspaceWorkerDiscoveryItem, WorkspaceWorkerDiscoveryPage, WorkspaceWorkerSubject,
|
||||
};
|
||||
|
||||
@@ -746,8 +746,7 @@ impl WorkspaceWorkerRemoveExecutor {
|
||||
));
|
||||
}
|
||||
|
||||
self.execute_target_removal(&runtime, &target, reason, false)
|
||||
.await
|
||||
self.execute_target_removal(&runtime, &target, reason).await
|
||||
}
|
||||
|
||||
async fn execute_target_removal(
|
||||
@@ -755,7 +754,6 @@ impl WorkspaceWorkerRemoveExecutor {
|
||||
runtime: &RuntimeRegistry,
|
||||
target: &RuntimeWorkerRef,
|
||||
reason: &str,
|
||||
allow_internal: bool,
|
||||
) -> std::result::Result<worker::WorkspaceResponse, String> {
|
||||
let remove_lock = {
|
||||
let mut locks = self
|
||||
@@ -857,7 +855,7 @@ impl WorkspaceWorkerRemoveExecutor {
|
||||
));
|
||||
}
|
||||
};
|
||||
if worker.singleton_key.is_some() && !allow_internal {
|
||||
if worker.singleton_key.is_some() {
|
||||
return Ok(worker_remove_error_response(
|
||||
StatusCode::CONFLICT,
|
||||
"internal_worker_forbidden",
|
||||
@@ -1009,6 +1007,7 @@ pub struct WorkspaceServerApi {
|
||||
store: Arc<dyn ControlPlaneStore>,
|
||||
catalog: WorkspaceCatalogService,
|
||||
routers: Arc<AsyncMutex<HashMap<String, Router>>>,
|
||||
apis: Arc<AsyncMutex<HashMap<String, WorkspaceApi>>>,
|
||||
}
|
||||
|
||||
impl WorkspaceServerApi {
|
||||
@@ -1018,9 +1017,71 @@ impl WorkspaceServerApi {
|
||||
catalog: WorkspaceCatalogService::new(store.clone()),
|
||||
store,
|
||||
routers: Arc::new(AsyncMutex::new(HashMap::new())),
|
||||
apis: Arc::new(AsyncMutex::new(HashMap::new())),
|
||||
}
|
||||
}
|
||||
|
||||
async fn api_for_workspace(&self, workspace_id: &str) -> Result<Option<WorkspaceApi>> {
|
||||
let mut apis = self.apis.lock().await;
|
||||
if let Some(api) = apis.get(workspace_id).cloned() {
|
||||
return Ok(Some(api));
|
||||
}
|
||||
let Some(workspace) = self.store.get_workspace(workspace_id).await? else {
|
||||
return Ok(None);
|
||||
};
|
||||
let repositories = self.store.list_repositories(workspace_id)?;
|
||||
let config = self
|
||||
.template
|
||||
.for_catalog_workspace(&workspace, repositories)?;
|
||||
let api = WorkspaceApi::new(config, self.store.clone()).await?;
|
||||
apis.insert(workspace_id.to_string(), api.clone());
|
||||
Ok(Some(api))
|
||||
}
|
||||
|
||||
async fn workspace_deletion_preflight(
|
||||
&self,
|
||||
actor_account_id: &str,
|
||||
workspace_id: &str,
|
||||
) -> Result<WorkspaceDeletionPreflightResponse> {
|
||||
let mut preflight = self
|
||||
.store
|
||||
.workspace_deletion_preflight(actor_account_id, workspace_id)?;
|
||||
let Some(api) = self.api_for_workspace(workspace_id).await? else {
|
||||
return Err(Error::InvalidInput("Workspace does not exist".to_string()));
|
||||
};
|
||||
for registry_worker in self.store.list_worker_registry(workspace_id, 10_000)? {
|
||||
let worker_key = registry_worker.display_name;
|
||||
match api.runtime.worker(®istry_worker.worker) {
|
||||
Ok(worker) if worker.state == "stopped" && worker.singleton_key.is_none() => {}
|
||||
Ok(worker) => preflight.blockers.push(WorkspaceDeletionBlocker {
|
||||
kind: if worker.singleton_key.is_some() {
|
||||
WorkspaceDeletionBlockerKind::RetentionHold
|
||||
} else {
|
||||
WorkspaceDeletionBlockerKind::WorkerRemovalBlocked
|
||||
},
|
||||
resource_kind: Some("worker".to_string()),
|
||||
resource_key: Some(worker_key),
|
||||
message: if worker.singleton_key.is_some() {
|
||||
"Internal or singleton Workers must be released by their owning service first."
|
||||
.to_string()
|
||||
} else {
|
||||
"Stop running, restoring, or otherwise active Workers before deleting the Workspace."
|
||||
.to_string()
|
||||
},
|
||||
}),
|
||||
Err(_) => preflight.blockers.push(WorkspaceDeletionBlocker {
|
||||
kind: WorkspaceDeletionBlockerKind::CleanupUnavailable,
|
||||
resource_kind: Some("worker".to_string()),
|
||||
resource_key: Some(worker_key),
|
||||
message: "Worker state is unavailable; retry after Runtime state is healthy."
|
||||
.to_string(),
|
||||
}),
|
||||
}
|
||||
}
|
||||
preflight.can_delete = preflight.blockers.is_empty();
|
||||
Ok(preflight)
|
||||
}
|
||||
|
||||
async fn execute_workspace_deletion(
|
||||
&self,
|
||||
operation_id: &str,
|
||||
@@ -1032,18 +1093,10 @@ impl WorkspaceServerApi {
|
||||
&[],
|
||||
None,
|
||||
)?;
|
||||
let workspace = self
|
||||
.store
|
||||
.get_workspace(&operation.workspace_id)
|
||||
let api = self
|
||||
.api_for_workspace(&operation.workspace_id)
|
||||
.await?
|
||||
.ok_or_else(|| Error::InvalidInput("Workspace no longer exists".to_string()))?;
|
||||
let repositories = self.store.list_repositories(&operation.workspace_id)?;
|
||||
let config = self
|
||||
.template
|
||||
.for_catalog_workspace(&workspace, repositories)?;
|
||||
let api = WorkspaceApi::new(config, self.store.clone()).await?;
|
||||
self.store
|
||||
.release_workspace_assignments_for_deletion(&operation.workspace_id)?;
|
||||
|
||||
let mut child_operation_ids = Vec::new();
|
||||
let mut blockers = Vec::new();
|
||||
@@ -1053,22 +1106,8 @@ impl WorkspaceServerApi {
|
||||
{
|
||||
let worker_key = worker.display_name.clone();
|
||||
let target = worker.worker;
|
||||
let lifecycle = WorkerLifecycleRequest {
|
||||
reason: Some("Workspace deletion".to_string()),
|
||||
ticket_assignment: None,
|
||||
};
|
||||
let _ = api.runtime.cancel_worker(&target, lifecycle.clone());
|
||||
if api.runtime.stop_worker(&target, lifecycle).is_err() {
|
||||
blockers.push(WorkspaceDeletionBlocker {
|
||||
kind: WorkspaceDeletionBlockerKind::WorkerRemovalBlocked,
|
||||
resource_kind: Some("worker".to_string()),
|
||||
resource_key: Some(worker_key.clone()),
|
||||
message: "Worker stop did not reach a retryable terminal state.".to_string(),
|
||||
});
|
||||
continue;
|
||||
}
|
||||
let response = WorkspaceWorkerRemoveExecutor::new(&api)
|
||||
.execute_target_removal(api.runtime.as_ref(), &target, "Workspace deletion", true)
|
||||
.execute_target_removal(api.runtime.as_ref(), &target, "Workspace deletion")
|
||||
.await
|
||||
.map_err(Error::Store)?;
|
||||
if let Some(child_operation_id) = self.store.latest_worker_removal_operation_id(
|
||||
@@ -1097,7 +1136,6 @@ impl WorkspaceServerApi {
|
||||
&api,
|
||||
&workdir.workdir_id,
|
||||
operation_id,
|
||||
operation.force_delete_dirty_workdirs,
|
||||
) {
|
||||
Ok(child) => {
|
||||
child_operation_ids.push(child.operation_id.clone());
|
||||
@@ -1115,7 +1153,7 @@ impl WorkspaceServerApi {
|
||||
resource_kind: Some("workdir".to_string()),
|
||||
resource_key: Some(workdir.workdir_id),
|
||||
message: if dirty {
|
||||
"Workdir is dirty or its cleanliness is unknown. Enable force deletion only after reviewing the impact."
|
||||
"Workdir is dirty or its cleanliness is unknown. Clean it and refresh status before retrying deletion."
|
||||
.to_string()
|
||||
} else {
|
||||
"Workdir removal did not complete; retry the Workspace deletion operation."
|
||||
@@ -1145,25 +1183,23 @@ impl WorkspaceServerApi {
|
||||
}
|
||||
let completed = self.store.finalize_workspace_deletion(operation_id)?;
|
||||
self.routers.lock().await.remove(&completed.workspace_id);
|
||||
self.apis.lock().await.remove(&completed.workspace_id);
|
||||
Ok(completed)
|
||||
}
|
||||
|
||||
async fn router_for_workspace(&self, workspace_id: &str) -> Result<Option<Router>> {
|
||||
let mut routers = self.routers.lock().await;
|
||||
if let Some(router) = routers.get(workspace_id) {
|
||||
return Ok(Some(router.clone()));
|
||||
if let Some(router) = self.routers.lock().await.get(workspace_id).cloned() {
|
||||
return Ok(Some(router));
|
||||
}
|
||||
let Some(workspace) = self.store.get_workspace(workspace_id).await? else {
|
||||
let Some(api) = self.api_for_workspace(workspace_id).await? else {
|
||||
return Ok(None);
|
||||
};
|
||||
let repositories = self.store.list_repositories(workspace_id)?;
|
||||
let config = self
|
||||
.template
|
||||
.for_catalog_workspace(&workspace, repositories)?;
|
||||
let api = WorkspaceApi::new(config, self.store.clone()).await?;
|
||||
tokio::spawn(run_orchestrator_turn_end_hook(api.clone()));
|
||||
let router = build_inner_router(api);
|
||||
routers.insert(workspace_id.to_string(), router.clone());
|
||||
self.routers
|
||||
.lock()
|
||||
.await
|
||||
.insert(workspace_id.to_string(), router.clone());
|
||||
Ok(Some(router))
|
||||
}
|
||||
|
||||
@@ -1269,8 +1305,8 @@ async fn preflight_server_workspace_deletion(
|
||||
Err(error) => return server_error_response(error),
|
||||
};
|
||||
match api
|
||||
.store
|
||||
.workspace_deletion_preflight(&actor_account_id, &workspace_id)
|
||||
.await
|
||||
{
|
||||
Ok(preflight) => Json(preflight).into_response(),
|
||||
Err(error) => server_error_response(error),
|
||||
@@ -1288,6 +1324,25 @@ async fn start_server_workspace_deletion(
|
||||
Ok(None) => return forbidden_server_response("Workspace deletion requires its owner"),
|
||||
Err(error) => return server_error_response(error),
|
||||
};
|
||||
let existing = match api
|
||||
.store
|
||||
.workspace_deletion_operation(&actor_account_id, &request.operation_id)
|
||||
{
|
||||
Ok(existing) => existing,
|
||||
Err(error) => return server_error_response(error),
|
||||
};
|
||||
if existing.is_none() {
|
||||
let preflight = match api
|
||||
.workspace_deletion_preflight(&actor_account_id, &workspace_id)
|
||||
.await
|
||||
{
|
||||
Ok(preflight) => preflight,
|
||||
Err(error) => return server_error_response(error),
|
||||
};
|
||||
if !preflight.can_delete {
|
||||
return (StatusCode::CONFLICT, Json(preflight)).into_response();
|
||||
}
|
||||
}
|
||||
let reservation =
|
||||
match api
|
||||
.store
|
||||
@@ -10520,27 +10575,7 @@ fn execute_reserved_workdir_removal(
|
||||
operation: WorkdirRemovalOperation,
|
||||
recovery: bool,
|
||||
) -> Result<WorkdirRemovalOperation> {
|
||||
execute_reserved_workdir_removal_with_provider(
|
||||
api,
|
||||
operation,
|
||||
recovery,
|
||||
api.runtime.as_ref(),
|
||||
false,
|
||||
)
|
||||
}
|
||||
|
||||
fn execute_reserved_workdir_removal_for_workspace_deletion(
|
||||
api: &WorkspaceApi,
|
||||
operation: WorkdirRemovalOperation,
|
||||
force_dirty: bool,
|
||||
) -> Result<WorkdirRemovalOperation> {
|
||||
execute_reserved_workdir_removal_with_provider(
|
||||
api,
|
||||
operation,
|
||||
false,
|
||||
api.runtime.as_ref(),
|
||||
force_dirty,
|
||||
)
|
||||
execute_reserved_workdir_removal_with_provider(api, operation, recovery, api.runtime.as_ref())
|
||||
}
|
||||
|
||||
fn execute_reserved_workdir_removal_with_provider(
|
||||
@@ -10548,7 +10583,6 @@ fn execute_reserved_workdir_removal_with_provider(
|
||||
operation: WorkdirRemovalOperation,
|
||||
recovery: bool,
|
||||
provider: &dyn WorkdirRemovalRuntimeProvider,
|
||||
force_dirty: bool,
|
||||
) -> Result<WorkdirRemovalOperation> {
|
||||
if operation.state == WorkdirRemovalOperationState::Completed {
|
||||
return Ok(operation);
|
||||
@@ -10623,9 +10657,8 @@ fn execute_reserved_workdir_removal_with_provider(
|
||||
true,
|
||||
);
|
||||
};
|
||||
if !force_dirty
|
||||
&& (status.summary.cleanliness.as_deref() != Some("clean")
|
||||
|| status.summary.status != WorkingDirectoryStatusKind::Active)
|
||||
if status.summary.cleanliness.as_deref() != Some("clean")
|
||||
|| status.summary.status != WorkingDirectoryStatusKind::Active
|
||||
{
|
||||
return api.config_store.complete_workdir_removal_retained(
|
||||
&operation,
|
||||
@@ -10723,7 +10756,6 @@ fn execute_workdir_removal_for_workspace_deletion(
|
||||
api: &WorkspaceApi,
|
||||
working_directory_id: &str,
|
||||
parent_operation_id: &str,
|
||||
force_dirty: bool,
|
||||
) -> Result<WorkdirRemovalOperation> {
|
||||
let source_actor = format!("workspace-deletion:{parent_operation_id}");
|
||||
let reason = "Workspace deletion";
|
||||
@@ -10750,7 +10782,7 @@ fn execute_workdir_removal_for_workspace_deletion(
|
||||
api.config_store
|
||||
.reserve_workdir_removal_operation(&intent)?
|
||||
};
|
||||
execute_reserved_workdir_removal_for_workspace_deletion(api, operation, force_dirty)
|
||||
execute_reserved_workdir_removal(api, operation, false)
|
||||
}
|
||||
|
||||
fn recover_workdir_removals(api: &WorkspaceApi) -> Result<()> {
|
||||
@@ -24371,7 +24403,6 @@ mod tests {
|
||||
clean_operation.clone(),
|
||||
false,
|
||||
&clean_provider,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
@@ -24384,7 +24415,6 @@ mod tests {
|
||||
removed.clone(),
|
||||
false,
|
||||
&clean_provider,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(replay, removed);
|
||||
@@ -24409,7 +24439,6 @@ mod tests {
|
||||
missing_operation,
|
||||
false,
|
||||
&missing_provider,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
@@ -24436,7 +24465,6 @@ mod tests {
|
||||
unknown_operation,
|
||||
false,
|
||||
&unknown_provider,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(unknown.state, WorkdirRemovalOperationState::Failed);
|
||||
@@ -24458,34 +24486,11 @@ mod tests {
|
||||
dirty_operation,
|
||||
false,
|
||||
&dirty_provider,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(dirty.disposition, Some(WorkdirRemovalDisposition::Retained));
|
||||
assert_eq!(dirty_provider.cleanup_calls(), 0);
|
||||
|
||||
let (forced_operation, mut forced_summary) =
|
||||
reserve_removal_fixture(&api, "provider-dirty-forced");
|
||||
forced_summary.cleanliness = Some("dirty".to_string());
|
||||
let forced_provider = FakeWorkdirRemovalProvider::new(
|
||||
workdir_removal_result(
|
||||
WorkerOperationState::Accepted,
|
||||
Some(forced_summary),
|
||||
Vec::new(),
|
||||
),
|
||||
workdir_removal_result(WorkerOperationState::Accepted, None, Vec::new()),
|
||||
);
|
||||
let forced = execute_reserved_workdir_removal_with_provider(
|
||||
&api,
|
||||
forced_operation,
|
||||
false,
|
||||
&forced_provider,
|
||||
true,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(forced.disposition, Some(WorkdirRemovalDisposition::Removed));
|
||||
assert_eq!(forced_provider.cleanup_calls(), 1);
|
||||
|
||||
let (unsupported_operation, unsupported_summary) =
|
||||
reserve_removal_fixture(&api, "provider-unsupported");
|
||||
let unsupported_provider = FakeWorkdirRemovalProvider::new(
|
||||
@@ -24509,7 +24514,6 @@ mod tests {
|
||||
unsupported_operation,
|
||||
false,
|
||||
&unsupported_provider,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(unsupported.state, WorkdirRemovalOperationState::Failed);
|
||||
@@ -24564,7 +24568,6 @@ mod tests {
|
||||
operation,
|
||||
true,
|
||||
provider.as_ref(),
|
||||
false,
|
||||
)
|
||||
}));
|
||||
}
|
||||
|
||||
@@ -6438,7 +6438,6 @@ fn migrate_workspace_deletion_v52_to_v53(conn: &Connection) -> Result<()> {
|
||||
workspace_revision TEXT NOT NULL,
|
||||
owner_account_id TEXT NOT NULL,
|
||||
actor_account_id TEXT NOT NULL,
|
||||
force_delete_dirty_workdirs INTEGER NOT NULL CHECK(force_delete_dirty_workdirs IN (0, 1)),
|
||||
state TEXT NOT NULL CHECK(state IN ('queued', 'running', 'blocked', 'failed', 'succeeded')),
|
||||
resource_counts_json TEXT NOT NULL,
|
||||
child_operation_ids_json TEXT NOT NULL,
|
||||
@@ -6475,7 +6474,6 @@ fn verify_workspace_deletion_schema(conn: &Connection) -> Result<()> {
|
||||
"workspace_revision",
|
||||
"owner_account_id",
|
||||
"actor_account_id",
|
||||
"force_delete_dirty_workdirs",
|
||||
"state",
|
||||
"resource_counts_json",
|
||||
"child_operation_ids_json",
|
||||
|
||||
@@ -12,7 +12,6 @@ use crate::store::{SqliteWorkspaceStore, WorkspaceRecord};
|
||||
use crate::{Error, Result};
|
||||
|
||||
const MAX_OPERATION_ID_BYTES: usize = 128;
|
||||
const CONFIRMATION_PREFIX: &str = "delete ";
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct WorkspaceDeletionReservation {
|
||||
@@ -40,8 +39,6 @@ pub trait WorkspaceDeletionStore: Send + Sync {
|
||||
operation_id: &str,
|
||||
) -> Result<Option<WorkspaceDeletionOperationResponse>>;
|
||||
|
||||
fn release_workspace_assignments_for_deletion(&self, workspace_id: &str) -> Result<u64>;
|
||||
|
||||
fn latest_worker_removal_operation_id(
|
||||
&self,
|
||||
workspace_id: &str,
|
||||
@@ -74,11 +71,11 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
let workspace = owner_workspace(conn, actor_account_id, workspace_id)?;
|
||||
let resources = resource_counts(conn, workspace_id)?;
|
||||
let accessible: u64 = conn.query_row(
|
||||
"SELECT COUNT(*) FROM workspaces WHERE owner_account_id = ?1 AND state = 'active'",
|
||||
"SELECT COUNT(*) FROM workspaces WHERE owner_account_id = ?1",
|
||||
params![actor_account_id],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
let mut blockers = Vec::new();
|
||||
let mut blockers = workspace_database_blockers(conn, workspace_id)?;
|
||||
if accessible <= 1 {
|
||||
blockers.push(WorkspaceDeletionBlocker {
|
||||
kind: WorkspaceDeletionBlockerKind::LastAccessibleWorkspace,
|
||||
@@ -92,7 +89,6 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
display_name: workspace.display_name,
|
||||
expected_revision: workspace.updated_at,
|
||||
can_delete: blockers.is_empty(),
|
||||
force_delete_dirty_workdirs_available: true,
|
||||
resources,
|
||||
blockers,
|
||||
})
|
||||
@@ -126,8 +122,7 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
}
|
||||
|
||||
let workspace = owner_workspace(tx, actor_account_id, workspace_id)?;
|
||||
let expected_confirmation = format!("{CONFIRMATION_PREFIX}{}", workspace.display_name);
|
||||
if request.confirmation != expected_confirmation {
|
||||
if request.confirmation != workspace.display_name {
|
||||
return Err(Error::InvalidInput(
|
||||
"confirmation must exactly match the displayed Workspace name".to_string(),
|
||||
));
|
||||
@@ -138,7 +133,7 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
));
|
||||
}
|
||||
let accessible: u64 = tx.query_row(
|
||||
"SELECT COUNT(*) FROM workspaces WHERE owner_account_id = ?1 AND state = 'active'",
|
||||
"SELECT COUNT(*) FROM workspaces WHERE owner_account_id = ?1",
|
||||
params![actor_account_id],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
@@ -148,6 +143,11 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
.to_string(),
|
||||
));
|
||||
}
|
||||
if !workspace_database_blockers(tx, workspace_id)?.is_empty() {
|
||||
return Err(Error::WorkspaceConfigConflict(
|
||||
"Workspace deletion preflight changed; reload current blockers".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let resources = resource_counts(tx, workspace_id)?;
|
||||
let now = Utc::now().to_rfc3339();
|
||||
@@ -156,10 +156,10 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
"INSERT INTO workspace_deletion_operations (
|
||||
operation_id, request_fingerprint, workspace_id, workspace_display_name,
|
||||
workspace_revision, owner_account_id, actor_account_id,
|
||||
force_delete_dirty_workdirs, state, resource_counts_json,
|
||||
state, resource_counts_json,
|
||||
child_operation_ids_json, blockers_json, failure_category,
|
||||
created_at, updated_at, completed_at
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'queued', ?9, '[]', '[]', NULL, ?10, ?10, NULL)",
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 'queued', ?8, '[]', '[]', NULL, ?9, ?9, NULL)",
|
||||
params![
|
||||
request.operation_id,
|
||||
fingerprint,
|
||||
@@ -168,7 +168,6 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
request.expected_revision,
|
||||
workspace.owner_account_id,
|
||||
actor_account_id,
|
||||
request.force_delete_dirty_workdirs as i64,
|
||||
serde_json::to_string(&resources).map_err(|error| Error::Store(error.to_string()))?,
|
||||
now,
|
||||
],
|
||||
@@ -210,17 +209,6 @@ impl WorkspaceDeletionStore for SqliteWorkspaceStore {
|
||||
})
|
||||
}
|
||||
|
||||
fn release_workspace_assignments_for_deletion(&self, workspace_id: &str) -> Result<u64> {
|
||||
self.with_conn(|conn| {
|
||||
let changed = conn.execute(
|
||||
"DELETE FROM ticket_current_worker_assignments WHERE workspace_id = ?1",
|
||||
params![workspace_id],
|
||||
)?;
|
||||
u64::try_from(changed)
|
||||
.map_err(|_| Error::Store("assignment deletion count overflow".to_string()))
|
||||
})
|
||||
}
|
||||
|
||||
fn latest_worker_removal_operation_id(
|
||||
&self,
|
||||
workspace_id: &str,
|
||||
@@ -373,30 +361,28 @@ fn read_operation(
|
||||
) -> Result<Option<StoredOperation>> {
|
||||
conn.query_row(
|
||||
"SELECT request_fingerprint, actor_account_id, workspace_id, workspace_display_name,
|
||||
state, force_delete_dirty_workdirs, resource_counts_json,
|
||||
child_operation_ids_json, blockers_json, failure_category,
|
||||
created_at, updated_at, completed_at
|
||||
state, resource_counts_json, child_operation_ids_json,
|
||||
blockers_json, failure_category, created_at, updated_at, completed_at
|
||||
FROM workspace_deletion_operations WHERE operation_id = ?1",
|
||||
params![operation_id],
|
||||
|row| {
|
||||
let state: String = row.get(4)?;
|
||||
let resource_counts_json: String = row.get(6)?;
|
||||
let child_operation_ids_json: String = row.get(7)?;
|
||||
let blockers_json: String = row.get(8)?;
|
||||
let resource_counts_json: String = row.get(5)?;
|
||||
let child_operation_ids_json: String = row.get(6)?;
|
||||
let blockers_json: String = row.get(7)?;
|
||||
Ok((
|
||||
row.get::<_, String>(0)?,
|
||||
row.get::<_, String>(1)?,
|
||||
row.get::<_, String>(2)?,
|
||||
row.get::<_, String>(3)?,
|
||||
state,
|
||||
row.get::<_, bool>(5)?,
|
||||
resource_counts_json,
|
||||
child_operation_ids_json,
|
||||
blockers_json,
|
||||
row.get::<_, Option<String>>(9)?,
|
||||
row.get::<_, Option<String>>(8)?,
|
||||
row.get::<_, String>(9)?,
|
||||
row.get::<_, String>(10)?,
|
||||
row.get::<_, String>(11)?,
|
||||
row.get::<_, Option<String>>(12)?,
|
||||
row.get::<_, Option<String>>(11)?,
|
||||
))
|
||||
},
|
||||
)
|
||||
@@ -408,7 +394,6 @@ fn read_operation(
|
||||
workspace_id,
|
||||
display_name,
|
||||
state,
|
||||
force,
|
||||
resources,
|
||||
children,
|
||||
blockers,
|
||||
@@ -425,7 +410,6 @@ fn read_operation(
|
||||
workspace_id,
|
||||
display_name,
|
||||
state: parse_deletion_state(&state)?,
|
||||
force_delete_dirty_workdirs: force,
|
||||
resources: serde_json::from_str(&resources)
|
||||
.map_err(|error| Error::Store(error.to_string()))?,
|
||||
child_operation_ids: serde_json::from_str(&children)
|
||||
@@ -474,6 +458,98 @@ fn owner_workspace(
|
||||
Ok(workspace)
|
||||
}
|
||||
|
||||
fn workspace_database_blockers(
|
||||
conn: &rusqlite::Connection,
|
||||
workspace_id: &str,
|
||||
) -> Result<Vec<WorkspaceDeletionBlocker>> {
|
||||
let mut blockers = Vec::new();
|
||||
for (sql, kind, resource_kind, message) in [
|
||||
(
|
||||
"SELECT workdir_id FROM worker_workdir_links WHERE workspace_id = ?1 AND unlinked_at IS NULL",
|
||||
WorkspaceDeletionBlockerKind::WorkdirRemovalBlocked,
|
||||
"workdir",
|
||||
"Release this active Worker–Workdir attachment before deleting the Workspace.",
|
||||
),
|
||||
(
|
||||
"SELECT workdir_id FROM worker_workdir_attachment_reservations WHERE workspace_id = ?1",
|
||||
WorkspaceDeletionBlockerKind::WorkdirRemovalBlocked,
|
||||
"workdir",
|
||||
"Wait for or cancel this pending Workdir attachment reservation.",
|
||||
),
|
||||
(
|
||||
"SELECT display_name FROM worker_registry WHERE workspace_id = ?1 AND retention_state = 'pinned'",
|
||||
WorkspaceDeletionBlockerKind::RetentionHold,
|
||||
"worker",
|
||||
"Remove this Worker retention pin before deleting the Workspace.",
|
||||
),
|
||||
(
|
||||
"SELECT workdir_id FROM workdir_registry WHERE workspace_id = ?1 AND COALESCE(cleanliness, '') != 'clean'",
|
||||
WorkspaceDeletionBlockerKind::DirtyWorkdir,
|
||||
"workdir",
|
||||
"Clean this Workdir and refresh unknown cleanliness before deleting the Workspace.",
|
||||
),
|
||||
] {
|
||||
for resource_key in query_resource_keys(conn, sql, workspace_id)? {
|
||||
blockers.push(WorkspaceDeletionBlocker {
|
||||
kind,
|
||||
resource_kind: Some(resource_kind.to_string()),
|
||||
resource_key: Some(resource_key),
|
||||
message: message.to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
for (sql, kind, resource_kind, message) in [
|
||||
(
|
||||
"SELECT COUNT(*) FROM ticket_current_worker_assignments WHERE workspace_id = ?1",
|
||||
WorkspaceDeletionBlockerKind::WorkerRemovalBlocked,
|
||||
"ticket",
|
||||
"Remove current Ticket assignments before deleting the Workspace.",
|
||||
),
|
||||
(
|
||||
"SELECT COUNT(*) FROM worker_removal_operations WHERE workspace_id = ?1 AND state IN ('planned', 'blocked', 'executing', 'failed', 'stale')",
|
||||
WorkspaceDeletionBlockerKind::CleanupUnavailable,
|
||||
"worker",
|
||||
"Resolve pending or failed Worker removal operations first.",
|
||||
),
|
||||
(
|
||||
"SELECT COUNT(*) FROM workdir_removal_operations WHERE workspace_id = ?1 AND state IN ('pending', 'failed')",
|
||||
WorkspaceDeletionBlockerKind::CleanupUnavailable,
|
||||
"workdir",
|
||||
"Resolve pending or failed Workdir removal operations first.",
|
||||
),
|
||||
(
|
||||
"SELECT COUNT(*) FROM workdir_create_operations WHERE workspace_id = ?1 AND state = 'pending'",
|
||||
WorkspaceDeletionBlockerKind::CleanupUnavailable,
|
||||
"workdir",
|
||||
"Wait for pending Workdir creation operations to finish.",
|
||||
),
|
||||
] {
|
||||
let count: u64 = conn.query_row(sql, params![workspace_id], |row| row.get(0))?;
|
||||
if count != 0 {
|
||||
blockers.push(WorkspaceDeletionBlocker {
|
||||
kind,
|
||||
resource_kind: Some(resource_kind.to_string()),
|
||||
resource_key: None,
|
||||
message: format!("{message} ({count})"),
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(blockers)
|
||||
}
|
||||
|
||||
fn query_resource_keys(
|
||||
conn: &rusqlite::Connection,
|
||||
sql: &str,
|
||||
workspace_id: &str,
|
||||
) -> Result<Vec<String>> {
|
||||
let mut statement = conn.prepare(sql)?;
|
||||
statement
|
||||
.query_map(params![workspace_id], |row| row.get(0))?
|
||||
.collect::<std::result::Result<Vec<_>, _>>()
|
||||
.map_err(Into::into)
|
||||
}
|
||||
|
||||
fn resource_counts(
|
||||
conn: &rusqlite::Connection,
|
||||
workspace_id: &str,
|
||||
@@ -515,8 +591,8 @@ fn request_fingerprint(
|
||||
request: &WorkspaceDeletionRequest,
|
||||
) -> String {
|
||||
let canonical = format!(
|
||||
"workspace-delete-v1\0{actor_account_id}\0{workspace_id}\0{}\0{}\0{}",
|
||||
request.expected_revision, request.confirmation, request.force_delete_dirty_workdirs
|
||||
"workspace-delete-v1\0{actor_account_id}\0{workspace_id}\0{}\0{}",
|
||||
request.expected_revision, request.confirmation
|
||||
);
|
||||
encode_hex(&Sha256::digest(canonical.as_bytes()))
|
||||
}
|
||||
@@ -590,8 +666,7 @@ mod tests {
|
||||
let request = WorkspaceDeletionRequest {
|
||||
operation_id: "delete-workspace-a".to_string(),
|
||||
expected_revision: preflight.expected_revision,
|
||||
confirmation: "delete Alpha".to_string(),
|
||||
force_delete_dirty_workdirs: false,
|
||||
confirmation: "Alpha".to_string(),
|
||||
};
|
||||
let first = store
|
||||
.reserve_workspace_deletion(&owner, &workspace_id, &request)
|
||||
@@ -644,21 +719,54 @@ mod tests {
|
||||
let mut request = WorkspaceDeletionRequest {
|
||||
operation_id: "delete-alpha-guarded".to_string(),
|
||||
expected_revision: "stale".to_string(),
|
||||
confirmation: "delete Alpha".to_string(),
|
||||
force_delete_dirty_workdirs: false,
|
||||
confirmation: "Alpha".to_string(),
|
||||
};
|
||||
assert!(matches!(
|
||||
store.reserve_workspace_deletion(&owner, &workspace_id, &request),
|
||||
Err(Error::WorkspaceConfigConflict(_))
|
||||
));
|
||||
request.expected_revision = preflight.expected_revision;
|
||||
request.confirmation = "Alpha".to_string();
|
||||
request.confirmation = "delete Alpha".to_string();
|
||||
assert!(matches!(
|
||||
store.reserve_workspace_deletion(&owner, &workspace_id, &request),
|
||||
Err(Error::InvalidInput(_))
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pinned_worker_blocks_preflight_before_operation_reservation() {
|
||||
let (store, owner, workspace_id) = setup();
|
||||
store
|
||||
.with_conn(|conn| {
|
||||
conn.execute(
|
||||
"INSERT INTO worker_registry (
|
||||
workspace_id, runtime_id, worker_id, display_name,
|
||||
created_at, updated_at, retention_state
|
||||
) VALUES (?1, 'runtime-a', 'worker-a', 'Pinned worker', '1', '1', 'pinned')",
|
||||
params![workspace_id],
|
||||
)?;
|
||||
Ok(())
|
||||
})
|
||||
.expect("worker");
|
||||
let preflight = store
|
||||
.workspace_deletion_preflight(&owner, &workspace_id)
|
||||
.expect("preflight");
|
||||
assert!(!preflight.can_delete);
|
||||
assert!(preflight.blockers.iter().any(|blocker| {
|
||||
blocker.kind == WorkspaceDeletionBlockerKind::RetentionHold
|
||||
&& blocker.resource_key.as_deref() == Some("Pinned worker")
|
||||
}));
|
||||
let request = WorkspaceDeletionRequest {
|
||||
operation_id: "delete-pinned".to_string(),
|
||||
expected_revision: preflight.expected_revision,
|
||||
confirmation: "Alpha".to_string(),
|
||||
};
|
||||
assert!(matches!(
|
||||
store.reserve_workspace_deletion(&owner, &workspace_id, &request),
|
||||
Err(Error::WorkspaceConfigConflict(_))
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn last_accessible_workspace_and_revision_conflicts_fail_closed() {
|
||||
let (store, owner, workspace_id) = setup();
|
||||
@@ -675,8 +783,7 @@ mod tests {
|
||||
&WorkspaceDeletionRequest {
|
||||
operation_id: "delete-beta".to_string(),
|
||||
expected_revision: preflight.expected_revision,
|
||||
confirmation: "delete Beta".to_string(),
|
||||
force_delete_dirty_workdirs: false,
|
||||
confirmation: "Beta".to_string(),
|
||||
},
|
||||
)
|
||||
.expect("reserve")
|
||||
|
||||
Reference in New Issue
Block a user