diff --git a/crates/workspace-server/src/latest_schema.sql b/crates/workspace-server/src/latest_schema.sql index e2978774..688d5224 100644 --- a/crates/workspace-server/src/latest_schema.sql +++ b/crates/workspace-server/src/latest_schema.sql @@ -628,6 +628,42 @@ CREATE TABLE workdir_create_operations ( PRIMARY KEY (workspace_id, operation_id), UNIQUE (workspace_id, working_directory_id) ); +CREATE TABLE workdir_create_credential_candidates ( + workspace_id TEXT NOT NULL, + operation_id TEXT NOT NULL, + ordinal INTEGER NOT NULL CHECK (ordinal >= 0 AND ordinal < 2), + role TEXT NOT NULL CHECK (role IN ('primary', 'workspace_default_fallback')), + credential_id TEXT NOT NULL CHECK (length(credential_id) BETWEEN 1 AND 128), + credential_revision INTEGER NOT NULL CHECK (credential_revision > 0), + PRIMARY KEY (workspace_id, operation_id, ordinal), + UNIQUE (workspace_id, operation_id, role), + UNIQUE (workspace_id, operation_id, credential_id), + FOREIGN KEY (workspace_id, operation_id) + REFERENCES workdir_create_operations(workspace_id, operation_id) + ON DELETE CASCADE +); +CREATE INDEX idx_workdir_create_credential_candidates_revision + ON workdir_create_credential_candidates( + workspace_id, credential_id, credential_revision + ); +CREATE TABLE workdir_create_credential_revision_retentions ( + workspace_id TEXT NOT NULL, + operation_id TEXT NOT NULL, + ordinal INTEGER NOT NULL, + credential_id TEXT NOT NULL, + credential_revision INTEGER NOT NULL, + PRIMARY KEY (workspace_id, operation_id, ordinal), + FOREIGN KEY (workspace_id, operation_id, ordinal) + REFERENCES workdir_create_credential_candidates( + workspace_id, operation_id, ordinal + ) + ON DELETE CASCADE, + FOREIGN KEY (workspace_id, credential_id, credential_revision) + REFERENCES repository_ssh_credential_revisions( + workspace_id, credential_id, revision + ) + ON DELETE RESTRICT +); CREATE TABLE "workdir_registry" ( workspace_id TEXT NOT NULL, workdir_id TEXT NOT NULL, diff --git a/crates/workspace-server/src/repository_access.rs b/crates/workspace-server/src/repository_access.rs index 7bea3352..1efc9ec5 100644 --- a/crates/workspace-server/src/repository_access.rs +++ b/crates/workspace-server/src/repository_access.rs @@ -762,6 +762,20 @@ impl RepositorySecretService { "credential `{credential_id}` revision changed" ))); } + let retained_by_workdir_create: bool = tx.query_row( + r#"SELECT EXISTS( + SELECT 1 + FROM workdir_create_credential_revision_retentions + WHERE workspace_id = ?1 AND credential_id = ?2 + )"#, + params![workspace_id, credential_id], + |row| row.get(0), + )?; + if retained_by_workdir_create { + return Err(Error::RepositoryConflict(format!( + "credential `{credential_id}` is retained by a retryable Workdir create operation" + ))); + } insert_audit(&tx, workspace_id, "credential_deleted", &credential_id, current.current_revision, actor_account_id, &now)?; let deleted = tx.execute( "DELETE FROM repository_ssh_credentials WHERE workspace_id = ?1 AND credential_id = ?2 AND current_revision = ?3", @@ -2442,6 +2456,114 @@ mod tests { ); } + #[test] + fn retryable_workdir_create_retains_candidate_revision_until_success() { + let (_dir, store, service) = test_service(); + let (private_key, _) = test_private_key(13); + service + .create_credential( + "workspace-a", + CreateRepositorySshCredentialRequest { + operation_id: "create-retained".to_string(), + credential_id: "retained-deploy".to_string(), + name: "Retained deploy".to_string(), + private_key, + passphrase: None, + }, + "owner-a", + ) + .unwrap(); + let operation = crate::store::WorkdirCreateOperationRecord { + workspace_id: "workspace-a".to_string(), + operation_id: "create-workdir-retained".to_string(), + request_fingerprint: "sha256:request".to_string(), + repository_id: "repo-a".to_string(), + selector: Some("develop".to_string()), + requested_runtime_id: Some("runtime-a".to_string()), + resolved_runtime_id: "runtime-a".to_string(), + config_revision: 1, + config_projection_digest: "sha256:projection".to_string(), + source_kind: Some("ssh".to_string()), + source_uri: Some("ssh://git@example.test/org/main.git".to_string()), + source_revision: Some(1), + source_fingerprint: Some("sha256:source".to_string()), + credential_id: None, + credential_revision: None, + host_trust_id: None, + host_trust_revision: None, + repository_access_mode: None, + credential_candidates: Vec::new(), + working_directory_id: "workdir-retained".to_string(), + state: "pending".to_string(), + failure: None, + created_at: "2026-08-24T00:00:00Z".to_string(), + updated_at: "2026-08-24T00:00:00Z".to_string(), + }; + store.reserve_workdir_create_operation(&operation).unwrap(); + let candidates = vec![crate::store::WorkdirCreateCredentialCandidate { + role: crate::store::WorkdirCreateCredentialCandidateRole::Primary, + credential_id: "retained-deploy".to_string(), + credential_revision: 1, + }]; + store + .bind_workdir_create_repository_access( + "workspace-a", + "create-workdir-retained", + "sha256:request", + "retained-deploy", + 1, + "host-a", + 1, + "read_only", + &candidates, + "2026-08-24T00:00:01Z", + ) + .unwrap(); + let projection = RepositoryAccessProjection { + workspace_id: "workspace-a".to_string(), + config_revision: 1, + projection_digest: "sha256:empty".to_string(), + bindings: Vec::new(), + }; + + let retained = service + .delete_credential( + "workspace-a", + "retained-deploy", + DeleteRepositorySshCredentialRequest { + operation_id: "delete-retained".to_string(), + expected_revision: 1, + }, + "owner-a", + &projection, + ) + .unwrap_err(); + assert!(matches!(retained, Error::RepositoryConflict(_))); + + store + .finish_workdir_create_operation( + "workspace-a", + "create-workdir-retained", + "sha256:request", + true, + None, + "2026-08-24T00:00:02Z", + ) + .unwrap(); + service + .delete_credential( + "workspace-a", + "retained-deploy", + DeleteRepositorySshCredentialRequest { + operation_id: "delete-released".to_string(), + expected_revision: 1, + }, + "owner-a", + &projection, + ) + .unwrap(); + } + #[test] fn referenced_resources_cannot_be_deleted() { let (_dir, _store, service) = test_service(); diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 09635831..2d02304a 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -172,7 +172,8 @@ use crate::store::{ AccountRecord, ApiTokenRecord, AuthChallengeRecord, BrowserSessionRecord, ControlPlaneStore, DeviceLoginFlowRecord, FlowSourceRecord, PasskeyCredentialRecord, RepositoryInsertOutcome, RepositoryRecord, TicketAssignmentPrincipal, TicketAssignmentRole, TicketCoderAssignmentRecord, - TicketRoleAssignmentRecord, UserRecord, WorkdirCreateOperationRecord, WorkdirRegistryRecord, + TicketRoleAssignmentRecord, UserRecord, WorkdirCreateCredentialCandidate, + WorkdirCreateCredentialCandidateRole, WorkdirCreateOperationRecord, WorkdirRegistryRecord, WorkerControlGrantRecord, WorkerRegistryRecord, WorkerWorkdirLinkRecord, WorkspaceRecord, WorkspaceResourceKind, WorkspaceRuntimeAuthenticationMode as StoredRuntimeAuthenticationMode, WorkspaceRuntimeBinding, WorkspaceRuntimeBindingAuditRecord, WorkspaceRuntimeBindingMutation, @@ -10917,6 +10918,7 @@ async fn create_workspace_working_directory( host_trust_id: None, host_trust_revision: None, repository_access_mode: None, + credential_candidates: Vec::new(), working_directory_id: next_backend_workdir_id(&request.repository_key), state: "pending".to_string(), failure: None, @@ -17624,10 +17626,10 @@ fn authorize_repository_materialization_operation( ) -> ApiResult<()> { let context = if request.repository.source.kind == workspace_api::RepositorySourceKind::Ssh { if let ( - Some(credential_id), - Some(credential_revision), - Some(host_trust_id), - Some(host_trust_revision), + Some(_credential_id), + Some(_credential_revision), + Some(_host_trust_id), + Some(_host_trust_revision), Some(access_mode), ) = ( operation.credential_id.as_deref(), @@ -17636,16 +17638,7 @@ fn authorize_repository_materialization_operation( operation.host_trust_revision, operation.repository_access_mode.as_deref(), ) { - let lease = api - .repository_secrets - .lease_ssh_materialization_access_revision( - &api.config.workspace_id, - credential_id, - credential_revision, - host_trust_id, - host_trust_revision, - )?; - let leases = repository_ssh_lease_candidates(api, lease)?; + let leases = repository_ssh_lease_candidates_from_operation(api, operation, request)?; let primary_lease = leases.first().ok_or_else(|| { settings_bad_request( "working_directory_repository_access_invalid", @@ -17766,6 +17759,20 @@ fn authorize_repository_materialization_operation( "Repository SSH access has no credential candidates", ) })?; + let credential_candidates = ssh + .credential_candidates + .iter() + .enumerate() + .map(|(index, candidate)| WorkdirCreateCredentialCandidate { + role: if index == 0 { + WorkdirCreateCredentialCandidateRole::Primary + } else { + WorkdirCreateCredentialCandidateRole::WorkspaceDefaultFallback + }, + credential_id: candidate.credential_id.clone(), + credential_revision: candidate.credential_revision, + }) + .collect::>(); api.config_store.bind_workdir_create_repository_access( &api.config.workspace_id, &operation.operation_id, @@ -17778,6 +17785,7 @@ fn authorize_repository_materialization_operation( workspace_api::RepositoryAccessMode::ReadOnly => "read_only", workspace_api::RepositoryAccessMode::ReadWrite => "read_write", }, + &credential_candidates, &now_registry_timestamp(), )?; context @@ -17814,6 +17822,112 @@ fn authorize_worker_spawn_workdir_materialization( authorize_repository_materialization(api, runtime_id, operation_id, &projection, request) } +fn repository_ssh_lease_candidates_from_operation( + api: &WorkspaceApi, + operation: &WorkdirCreateOperationRecord, + request: &WorkingDirectoryRequest, +) -> ApiResult> { + if operation.repository_id != request.repository.id + || operation.source_kind.as_deref() != Some(request.repository.source.kind.as_str()) + || operation.source_uri.as_deref() != Some(request.repository.source.uri.as_str()) + || operation.source_revision != Some(request.repository.source_revision) + || operation.source_fingerprint.as_deref() + != Some(request.repository.source_fingerprint.as_str()) + { + return Err(settings_bad_request( + "working_directory_repository_access_snapshot_mismatch", + "persisted Workdir Repository access snapshot does not match the create request", + )); + } + let host_trust_id = operation.host_trust_id.as_deref().ok_or_else(|| { + settings_bad_request( + "working_directory_repository_access_snapshot_missing", + "persisted Workdir Repository access snapshot is incomplete", + ) + })?; + let host_trust_revision = operation.host_trust_revision.ok_or_else(|| { + settings_bad_request( + "working_directory_repository_access_snapshot_missing", + "persisted Workdir Repository access snapshot is incomplete", + ) + })?; + if operation.credential_candidates.is_empty() { + return Err(settings_bad_request( + "working_directory_repository_access_snapshot_missing", + "persisted Workdir credential candidate snapshot is unavailable", + )); + } + let primary_credential_id = operation.credential_id.as_deref().ok_or_else(|| { + settings_bad_request( + "working_directory_repository_access_snapshot_missing", + "persisted Workdir primary credential evidence is unavailable", + ) + })?; + let primary_credential_revision = operation.credential_revision.ok_or_else(|| { + settings_bad_request( + "working_directory_repository_access_snapshot_missing", + "persisted Workdir primary credential evidence is unavailable", + ) + })?; + crate::workdir_create_operations::validate_workdir_create_credential_candidates( + primary_credential_id, + primary_credential_revision, + &operation.credential_candidates, + ) + .map_err(|_error| { + settings_bad_request( + "working_directory_repository_access_snapshot_invalid", + "persisted Workdir credential candidate snapshot is invalid", + ) + })?; + let _access_mode = operation + .repository_access_mode + .as_deref() + .map(|access_mode| match access_mode { + "read_only" => Ok(workspace_api::RepositoryAccessMode::ReadOnly), + "read_write" => Ok(workspace_api::RepositoryAccessMode::ReadWrite), + _other => Err(settings_bad_request( + "working_directory_repository_access_snapshot_invalid", + "persisted Repository access mode is invalid", + )), + }) + .transpose()? + .ok_or_else(|| { + settings_bad_request( + "working_directory_repository_access_snapshot_missing", + "persisted Workdir Repository access mode is unavailable", + ) + })?; + + let mut leases = Vec::with_capacity(operation.credential_candidates.len()); + for candidate in &operation.credential_candidates { + let lease = api + .repository_secrets + .lease_ssh_materialization_access_revision( + &api.config.workspace_id, + &candidate.credential_id, + candidate.credential_revision, + host_trust_id, + host_trust_revision, + ) + .map_err(|_| { + settings_bad_request( + "working_directory_repository_access_snapshot_unavailable", + "persisted Workdir credential or host-trust revision is unavailable", + ) + })?; + if lease.host_trust_id != host_trust_id || lease.host_trust_revision != host_trust_revision + { + return Err(settings_bad_request( + "working_directory_repository_access_snapshot_mismatch", + "persisted Workdir host-trust snapshot does not match the credential candidate", + )); + } + leases.push(lease); + } + Ok(leases) +} + fn authorize_repository_materialization( api: &WorkspaceApi, runtime_id: &str, @@ -19673,6 +19787,89 @@ mod tests { candidates[1].known_hosts_entry ); + let source = workspace_api::RepositorySource { + kind: workspace_api::RepositorySourceKind::Ssh, + uri: "ssh://git@example.test/org/repository.git".to_string(), + }; + let source_fingerprint = repository_source_fingerprint(&source); + let request = WorkingDirectoryRequest { + repository: worker_runtime::catalog::WorkingDirectoryRepository { + id: "repository-a".to_string(), + provider: "git".to_string(), + source: source.clone(), + source_revision: 1, + source_fingerprint: source_fingerprint.clone(), + selector: None, + }, + materializer: Default::default(), + backend_workdir_id: Some("workdir-a".to_string()), + materialization: None, + }; + let operation = WorkdirCreateOperationRecord { + workspace_id: api.config.workspace_id.clone(), + operation_id: "retry-workdir-a".to_string(), + request_fingerprint: "sha256:request".to_string(), + repository_id: "repository-a".to_string(), + selector: None, + requested_runtime_id: Some("runtime-1".to_string()), + resolved_runtime_id: "runtime-1".to_string(), + config_revision: 1, + config_projection_digest: "sha256:projection".to_string(), + source_kind: Some(source.kind.as_str().to_string()), + source_uri: Some(source.uri.clone()), + source_revision: Some(1), + source_fingerprint: Some(source_fingerprint), + credential_id: Some(candidates[0].credential_id.clone()), + credential_revision: Some(candidates[0].credential_revision), + host_trust_id: Some(candidates[0].host_trust_id.clone()), + host_trust_revision: Some(candidates[0].host_trust_revision), + repository_access_mode: Some("read_only".to_string()), + credential_candidates: candidates + .iter() + .enumerate() + .map(|(index, candidate)| WorkdirCreateCredentialCandidate { + role: if index == 0 { + WorkdirCreateCredentialCandidateRole::Primary + } else { + WorkdirCreateCredentialCandidateRole::WorkspaceDefaultFallback + }, + credential_id: candidate.credential_id.clone(), + credential_revision: candidate.credential_revision, + }) + .collect(), + working_directory_id: "workdir-a".to_string(), + state: "failed".to_string(), + failure: Some("runtime unavailable".to_string()), + created_at: "2026-08-24T00:00:00Z".to_string(), + updated_at: "2026-08-24T00:00:01Z".to_string(), + }; + let retry_candidates = + repository_ssh_lease_candidates_from_operation(&api, &operation, &request).unwrap(); + assert_eq!(retry_candidates.len(), 2); + assert_eq!( + retry_candidates + .iter() + .map(|candidate| ( + candidate.credential_id.as_str(), + candidate.credential_revision + )) + .collect::>(), + operation + .credential_candidates + .iter() + .map(|candidate| ( + candidate.credential_id.as_str(), + candidate.credential_revision + )) + .collect::>() + ); + let mut missing_snapshot = operation.clone(); + missing_snapshot.credential_candidates.clear(); + assert!( + repository_ssh_lease_candidates_from_operation(&api, &missing_snapshot, &request) + .is_err() + ); + let default_binding = workspace_api::RepositorySshAccessBinding { credential_id: crate::repository_access::WORKSPACE_DEFAULT_REPOSITORY_SSH_CREDENTIAL_ID .to_string(), @@ -28363,6 +28560,7 @@ mod tests { host_trust_id: None, host_trust_revision: None, repository_access_mode: None, + credential_candidates: Vec::new(), working_directory_id: "workdir-provider-rejection".to_string(), state: "pending".to_string(), failure: None, diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 76e283cb..daba7aea 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -18,7 +18,7 @@ use crate::workspace_deletion::WorkspaceDeletionStore; use crate::{Error, Result}; const OLDEST_SCHEMA_VERSION: i64 = 50; -const LATEST_SCHEMA_VERSION: i64 = 58; +const LATEST_SCHEMA_VERSION: i64 = 59; const SCHEMA_BASELINE_NAME: &str = "workspace schema baseline"; const WORKSPACE_RUNTIME_BINDINGS_MIGRATION_NAME: &str = "workspace runtime bindings"; const RUNTIME_BINDING_AUDIT_MIGRATION_NAME: &str = "workspace Runtime binding revision and audit"; @@ -32,6 +32,8 @@ const LEGACY_EXTERNAL_RUNTIME_BINDING_CUTOVER_MIGRATION_NAME: &str = "convert legacy Server-issued Runtime bindings to Workspace identity"; const REMOVE_WORKDIR_CACHE_GENERATION_MIGRATION_NAME: &str = "remove obsolete Workdir Repository cache generation"; +const WORKDIR_CREDENTIAL_CANDIDATE_SNAPSHOT_MIGRATION_NAME: &str = + "Workdir create credential candidate snapshots"; const MIGRATIONS: &[Migration] = &[ Migration { @@ -74,6 +76,11 @@ const MIGRATIONS: &[Migration] = &[ name: REMOVE_WORKDIR_CACHE_GENERATION_MIGRATION_NAME, apply: migrate_workdir_cache_generation_v57_to_v58, }, + Migration { + version: 59, + name: WORKDIR_CREDENTIAL_CANDIDATE_SNAPSHOT_MIGRATION_NAME, + apply: migrate_workdir_credential_candidate_snapshots_v58_to_v59, + }, ]; #[derive(Clone, Copy)] @@ -553,6 +560,39 @@ pub struct TicketWorkerAssignmentUpdate { pub previous: Option, } +#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum WorkdirCreateCredentialCandidateRole { + Primary, + WorkspaceDefaultFallback, +} + +impl WorkdirCreateCredentialCandidateRole { + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::Primary => "primary", + Self::WorkspaceDefaultFallback => "workspace_default_fallback", + } + } + + pub(crate) fn parse(value: &str) -> Result { + match value { + "primary" => Ok(Self::Primary), + "workspace_default_fallback" => Ok(Self::WorkspaceDefaultFallback), + other => Err(Error::Store(format!( + "invalid Workdir create credential candidate role `{other}`" + ))), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkdirCreateCredentialCandidate { + pub role: WorkdirCreateCredentialCandidateRole, + pub credential_id: String, + pub credential_revision: u64, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct WorkdirCreateOperationRecord { pub workspace_id: String, @@ -573,6 +613,8 @@ pub struct WorkdirCreateOperationRecord { pub host_trust_id: Option, pub host_trust_revision: Option, pub repository_access_mode: Option, + #[serde(default)] + pub credential_candidates: Vec, pub working_directory_id: String, pub state: String, pub failure: Option, @@ -9583,6 +9625,56 @@ fn table_columns(conn: &Connection, table_name: &str) -> Result> { .map_err(Error::from) } +fn migrate_workdir_credential_candidate_snapshots_v58_to_v59(conn: &Connection) -> Result<()> { + let tx = rusqlite::Transaction::new_unchecked(conn, TransactionBehavior::Exclusive)?; + tx.execute_batch( + r#" + CREATE TABLE workdir_create_credential_candidates ( + workspace_id TEXT NOT NULL, + operation_id TEXT NOT NULL, + ordinal INTEGER NOT NULL CHECK (ordinal >= 0 AND ordinal < 2), + role TEXT NOT NULL CHECK (role IN ('primary', 'workspace_default_fallback')), + credential_id TEXT NOT NULL CHECK (length(credential_id) BETWEEN 1 AND 128), + credential_revision INTEGER NOT NULL CHECK (credential_revision > 0), + PRIMARY KEY (workspace_id, operation_id, ordinal), + UNIQUE (workspace_id, operation_id, role), + UNIQUE (workspace_id, operation_id, credential_id), + FOREIGN KEY (workspace_id, operation_id) + REFERENCES workdir_create_operations(workspace_id, operation_id) + ON DELETE CASCADE + ); + CREATE INDEX idx_workdir_create_credential_candidates_revision + ON workdir_create_credential_candidates( + workspace_id, credential_id, credential_revision + ); + CREATE TABLE workdir_create_credential_revision_retentions ( + workspace_id TEXT NOT NULL, + operation_id TEXT NOT NULL, + ordinal INTEGER NOT NULL, + credential_id TEXT NOT NULL, + credential_revision INTEGER NOT NULL, + PRIMARY KEY (workspace_id, operation_id, ordinal), + FOREIGN KEY (workspace_id, operation_id, ordinal) + REFERENCES workdir_create_credential_candidates( + workspace_id, operation_id, ordinal + ) + ON DELETE CASCADE, + FOREIGN KEY (workspace_id, credential_id, credential_revision) + REFERENCES repository_ssh_credential_revisions( + workspace_id, credential_id, revision + ) + ON DELETE RESTRICT + ); + "#, + )?; + tx.execute( + "INSERT INTO __yoi_schema_migrations (version, name) VALUES (?1, ?2)", + params![59_i64, WORKDIR_CREDENTIAL_CANDIDATE_SNAPSHOT_MIGRATION_NAME], + )?; + tx.commit()?; + Ok(()) +} + #[cfg(test)] mod tests { use super::*; @@ -9651,6 +9743,8 @@ mod tests { create_latest_workspace_schema(&conn).unwrap(); conn.execute_batch( r#" + DROP TABLE workdir_create_credential_revision_retentions; + DROP TABLE workdir_create_credential_candidates; DROP INDEX workspace_signing_identity_audit_workspace_idx; DROP TABLE workspace_signing_identity_audit; DROP TABLE workspace_signing_identity_provisioning_operations; @@ -9792,6 +9886,10 @@ mod tests { version: 58, name: REMOVE_WORKDIR_CACHE_GENERATION_MIGRATION_NAME.to_string(), }, + WorkspaceSchemaMigrationStep { + version: 59, + name: WORKDIR_CREDENTIAL_CANDIDATE_SNAPSHOT_MIGRATION_NAME.to_string(), + }, ] ); @@ -9830,6 +9928,10 @@ mod tests { 58, REMOVE_WORKDIR_CACHE_GENERATION_MIGRATION_NAME.to_string(), ), + ( + 59, + WORKDIR_CREDENTIAL_CANDIDATE_SNAPSHOT_MIGRATION_NAME.to_string(), + ), ] ); assert!(!table_exists(conn, "trusted_runtime_records")?); @@ -9900,7 +10002,7 @@ mod tests { .iter() .map(|migration| migration.version) .collect::>(), - vec![52, 53, 54, 55, 56, 57, 58] + vec![52, 53, 54, 55, 56, 57, 58, 59] ); SqliteWorkspaceStore::migrate_database(&path).unwrap(); let conn = Connection::open(&path).unwrap(); @@ -9908,7 +10010,7 @@ mod tests { current_schema_version(&conn).unwrap(), LATEST_SCHEMA_VERSION ); - assert_eq!(workspace_schema_migration_history(&conn).unwrap().len(), 9); + assert_eq!(workspace_schema_migration_history(&conn).unwrap().len(), 10); } #[test] @@ -10696,6 +10798,33 @@ mod tests { ); } + #[test] + fn schema_v59_adds_workdir_create_credential_candidate_snapshots() { + let temp = tempfile::tempdir().unwrap(); + let path = temp.path().join("server.db"); + prepare_schema_v50(&path, Some("workspace-a")); + let conn = Connection::open(&path).unwrap(); + configure_sqlite(&conn).unwrap(); + for migration in MIGRATIONS + .iter() + .filter(|migration| migration.version <= 58) + { + (migration.apply)(&conn).unwrap(); + } + migrate_workdir_credential_candidate_snapshots_v58_to_v59(&conn).unwrap(); + + assert_eq!(current_schema_version(&conn).unwrap(), 59); + assert!(table_exists(&conn, "workdir_create_credential_candidates").unwrap()); + assert!(table_exists(&conn, "workdir_create_credential_revision_retentions").unwrap()); + assert_eq!( + conn.query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |row| { + row.get::<_, i64>(0) + }) + .unwrap(), + 0 + ); + } + #[test] fn runtime_binding_key_mutations_are_revisioned_idempotent_and_audited() { let store = SqliteWorkspaceStore::in_memory().unwrap(); diff --git a/crates/workspace-server/src/workdir_create_operations.rs b/crates/workspace-server/src/workdir_create_operations.rs index d255ec0a..22a71cbe 100644 --- a/crates/workspace-server/src/workdir_create_operations.rs +++ b/crates/workspace-server/src/workdir_create_operations.rs @@ -1,9 +1,15 @@ use rusqlite::{OptionalExtension, TransactionBehavior, params}; use sha2::{Digest, Sha256}; -use crate::store::WorkdirCreateOperationRecord; +use crate::store::{ + WorkdirCreateCredentialCandidate, WorkdirCreateCredentialCandidateRole, + WorkdirCreateOperationRecord, +}; use crate::{Error, Result, SqliteWorkspaceStore}; +const MAX_WORKDIR_CREATE_CREDENTIAL_CANDIDATES: usize = 2; +const MAX_CREDENTIAL_ID_BYTES: usize = 128; + pub fn selector_for_retry( explicit_selector: Option<&str>, persisted_selector: Option<&str>, @@ -172,10 +178,17 @@ impl SqliteWorkspaceStore { host_trust_id: &str, host_trust_revision: u64, repository_access_mode: &str, + credential_candidates: &[WorkdirCreateCredentialCandidate], now: &str, ) -> Result { + validate_workdir_create_credential_candidates( + credential_id, + credential_revision, + credential_candidates, + )?; self.with_conn_mut(|conn| { - let operation = read_workdir_create_operation(conn, workspace_id, operation_id)? + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + let operation = read_workdir_create_operation(&tx, workspace_id, operation_id)? .ok_or_else(|| { Error::RegistryInconsistency(format!( "Workdir create operation `{operation_id}` disappeared before Repository access binding" @@ -193,6 +206,7 @@ impl SqliteWorkspaceStore { || operation.host_trust_revision != Some(host_trust_revision) || operation.repository_access_mode.as_deref() != Some(repository_access_mode) + || operation.credential_candidates != credential_candidates { return Err(Error::InvalidInput(format!( "Workdir create operation `{operation_id}` Repository access evidence changed" @@ -200,7 +214,7 @@ impl SqliteWorkspaceStore { } return Ok(operation); } - conn.execute( + let updated = tx.execute( r#"UPDATE workdir_create_operations SET credential_id = ?4, credential_revision = ?5, host_trust_id = ?6, host_trust_revision = ?7, @@ -223,11 +237,60 @@ impl SqliteWorkspaceStore { now, ], )?; - read_workdir_create_operation(conn, workspace_id, operation_id)?.ok_or_else(|| { - Error::RegistryInconsistency(format!( - "Workdir create operation `{operation_id}` disappeared after Repository access binding" - )) - }) + if updated != 1 { + return Err(Error::RegistryInconsistency(format!( + "Workdir create operation `{operation_id}` changed before Repository access binding" + ))); + } + for (ordinal, candidate) in credential_candidates.iter().enumerate() { + tx.execute( + r#"INSERT INTO workdir_create_credential_candidates ( + workspace_id, operation_id, ordinal, role, + credential_id, credential_revision + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)"#, + params![ + workspace_id, + operation_id, + i64::try_from(ordinal).map_err(|_| Error::InvalidInput( + "credential candidate ordinal is out of range".to_string() + ))?, + candidate.role.as_str(), + candidate.credential_id, + i64::try_from(candidate.credential_revision).map_err(|_| { + Error::InvalidInput( + "credential candidate revision is out of range".to_string(), + ) + })?, + ], + )?; + tx.execute( + r#"INSERT INTO workdir_create_credential_revision_retentions ( + workspace_id, operation_id, ordinal, + credential_id, credential_revision + ) VALUES (?1, ?2, ?3, ?4, ?5)"#, + params![ + workspace_id, + operation_id, + i64::try_from(ordinal).map_err(|_| Error::InvalidInput( + "credential candidate ordinal is out of range".to_string() + ))?, + candidate.credential_id, + i64::try_from(candidate.credential_revision).map_err(|_| { + Error::InvalidInput( + "credential candidate revision is out of range".to_string(), + ) + })?, + ], + )?; + } + let bound = read_workdir_create_operation(&tx, workspace_id, operation_id)? + .ok_or_else(|| { + Error::RegistryInconsistency(format!( + "Workdir create operation `{operation_id}` disappeared after Repository access binding" + )) + })?; + tx.commit()?; + Ok(bound) }) } @@ -241,7 +304,8 @@ impl SqliteWorkspaceStore { updated_at: &str, ) -> Result { self.with_conn_mut(|conn| { - let changed = conn.execute( + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + let changed = tx.execute( r#"UPDATE workdir_create_operations SET state = ?1, failure = ?2, updated_at = ?3 WHERE workspace_id = ?4 AND operation_id = ?5 @@ -260,11 +324,21 @@ impl SqliteWorkspaceStore { "Workdir create operation `{operation_id}` could not be finalized" ))); } - read_workdir_create_operation(conn, workspace_id, operation_id)?.ok_or_else(|| { - Error::RegistryInconsistency(format!( - "Workdir create operation `{operation_id}` disappeared" - )) - }) + if succeeded { + tx.execute( + r#"DELETE FROM workdir_create_credential_revision_retentions + WHERE workspace_id = ?1 AND operation_id = ?2"#, + params![workspace_id, operation_id], + )?; + } + let finished = read_workdir_create_operation(&tx, workspace_id, operation_id)? + .ok_or_else(|| { + Error::RegistryInconsistency(format!( + "Workdir create operation `{operation_id}` disappeared" + )) + })?; + tx.commit()?; + Ok(finished) }) } @@ -282,8 +356,9 @@ fn read_workdir_create_operation( workspace_id: &str, operation_id: &str, ) -> Result> { - conn.query_row( - r#"SELECT workspace_id, operation_id, request_fingerprint, repository_id, selector, + let mut operation = conn + .query_row( + r#"SELECT workspace_id, operation_id, request_fingerprint, repository_id, selector, requested_runtime_id, resolved_runtime_id, config_revision, config_projection_digest, source_kind, source_uri, source_revision, source_fingerprint, credential_id, credential_revision, @@ -292,37 +367,146 @@ fn read_workdir_create_operation( created_at, updated_at FROM workdir_create_operations WHERE workspace_id = ?1 AND operation_id = ?2"#, - params![workspace_id, operation_id], - |row| { - Ok(WorkdirCreateOperationRecord { - workspace_id: row.get(0)?, - operation_id: row.get(1)?, - request_fingerprint: row.get(2)?, - repository_id: row.get(3)?, - selector: row.get(4)?, - requested_runtime_id: row.get(5)?, - resolved_runtime_id: row.get(6)?, - config_revision: row.get::<_, i64>(7)? as u64, - config_projection_digest: row.get(8)?, - source_kind: row.get(9)?, - source_uri: row.get(10)?, - source_revision: row.get::<_, Option>(11)?.map(|value| value as u64), - source_fingerprint: row.get(12)?, - credential_id: row.get(13)?, - credential_revision: row.get::<_, Option>(14)?.map(|value| value as u64), - host_trust_id: row.get(15)?, - host_trust_revision: row.get::<_, Option>(16)?.map(|value| value as u64), - repository_access_mode: row.get(17)?, - working_directory_id: row.get(18)?, - state: row.get(19)?, - failure: row.get(20)?, - created_at: row.get(21)?, - updated_at: row.get(22)?, - }) - }, - ) - .optional() - .map_err(Error::from) + params![workspace_id, operation_id], + |row| { + Ok(WorkdirCreateOperationRecord { + workspace_id: row.get(0)?, + operation_id: row.get(1)?, + request_fingerprint: row.get(2)?, + repository_id: row.get(3)?, + selector: row.get(4)?, + requested_runtime_id: row.get(5)?, + resolved_runtime_id: row.get(6)?, + config_revision: row.get::<_, i64>(7)? as u64, + config_projection_digest: row.get(8)?, + source_kind: row.get(9)?, + source_uri: row.get(10)?, + source_revision: row.get::<_, Option>(11)?.map(|value| value as u64), + source_fingerprint: row.get(12)?, + credential_id: row.get(13)?, + credential_revision: row.get::<_, Option>(14)?.map(|value| value as u64), + host_trust_id: row.get(15)?, + host_trust_revision: row.get::<_, Option>(16)?.map(|value| value as u64), + repository_access_mode: row.get(17)?, + credential_candidates: Vec::new(), + working_directory_id: row.get(18)?, + state: row.get(19)?, + failure: row.get(20)?, + created_at: row.get(21)?, + updated_at: row.get(22)?, + }) + }, + ) + .optional()?; + + if let Some(operation) = operation.as_mut() { + let mut statement = conn.prepare( + r#"SELECT ordinal, role, credential_id, credential_revision + FROM workdir_create_credential_candidates + WHERE workspace_id = ?1 AND operation_id = ?2 + ORDER BY ordinal ASC"#, + )?; + let rows = statement.query_map(params![workspace_id, operation_id], |row| { + Ok(( + row.get::<_, i64>(0)?, + row.get::<_, String>(1)?, + row.get::<_, String>(2)?, + row.get::<_, i64>(3)?, + )) + })?; + for (expected_ordinal, row) in rows.enumerate() { + let (ordinal, role, credential_id, credential_revision) = row?; + if ordinal + != i64::try_from(expected_ordinal).map_err(|_| { + Error::Store( + "Workdir create credential candidate ordinal is out of range".to_string(), + ) + })? + { + return Err(Error::Store( + "Workdir create credential candidate ordinals are not contiguous".to_string(), + )); + } + operation + .credential_candidates + .push(WorkdirCreateCredentialCandidate { + role: WorkdirCreateCredentialCandidateRole::parse(&role)?, + credential_id, + credential_revision: u64::try_from(credential_revision).map_err(|_| { + Error::Store(format!( + "invalid Workdir create credential candidate revision `{credential_revision}`" + )) + })?, + }); + } + if !operation.credential_candidates.is_empty() { + validate_workdir_create_credential_candidates( + operation.credential_id.as_deref().unwrap_or_default(), + operation.credential_revision.unwrap_or_default(), + &operation.credential_candidates, + ) + .map_err(|error| { + Error::Store(format!( + "invalid persisted Workdir create credential snapshot: {error}" + )) + })?; + } + } + Ok(operation) +} + +pub(crate) fn validate_workdir_create_credential_candidates( + credential_id: &str, + credential_revision: u64, + candidates: &[WorkdirCreateCredentialCandidate], +) -> Result<()> { + if candidates.is_empty() || candidates.len() > MAX_WORKDIR_CREATE_CREDENTIAL_CANDIDATES { + return Err(Error::InvalidInput(format!( + "Workdir create credential candidates must contain 1..={MAX_WORKDIR_CREATE_CREDENTIAL_CANDIDATES} entries" + ))); + } + if credential_id.is_empty() || credential_id.len() > MAX_CREDENTIAL_ID_BYTES { + return Err(Error::InvalidInput( + "Workdir create credential id is invalid".to_string(), + )); + } + if credential_revision == 0 { + return Err(Error::InvalidInput( + "Workdir create credential revision must be greater than zero".to_string(), + )); + } + let primary = &candidates[0]; + if primary.role != WorkdirCreateCredentialCandidateRole::Primary + || primary.credential_id != credential_id + || primary.credential_revision != credential_revision + { + return Err(Error::InvalidInput( + "Workdir create primary credential evidence does not match the ordered candidate snapshot" + .to_string(), + )); + } + if candidates.len() == 2 + && candidates[1].role != WorkdirCreateCredentialCandidateRole::WorkspaceDefaultFallback + { + return Err(Error::InvalidInput( + "Workdir create fallback credential role is invalid".to_string(), + )); + } + if candidates.iter().any(|candidate| { + candidate.credential_id.is_empty() + || candidate.credential_id.len() > MAX_CREDENTIAL_ID_BYTES + || candidate.credential_revision == 0 + }) { + return Err(Error::InvalidInput( + "Workdir create credential candidate identity is invalid".to_string(), + )); + } + if candidates.len() == 2 && candidates[0].credential_id == candidates[1].credential_id { + return Err(Error::InvalidInput( + "Workdir create credential candidates contain a duplicate credential id".to_string(), + )); + } + Ok(()) } #[cfg(test)] @@ -347,7 +531,7 @@ mod tests { } #[test] - fn retry_keeps_resolved_config_evidence_and_rejects_changed_input() { + fn retry_keeps_resolved_config_and_credential_candidates_after_fallback_moves() { let store = SqliteWorkspaceStore::in_memory().unwrap(); futures::executor::block_on(store.upsert_workspace(&WorkspaceRecord { workspace_id: "workspace".to_string(), @@ -403,6 +587,7 @@ mod tests { host_trust_id: None, host_trust_revision: None, repository_access_mode: None, + credential_candidates: Vec::new(), working_directory_id: "wd-1".to_string(), state: "pending".to_string(), failure: None, @@ -413,6 +598,44 @@ mod tests { store.reserve_workdir_create_operation(&record).unwrap(), record ); + store + .with_conn_mut(|conn| { + for (credential_id, revision) in + [("credential-1", 3_i64), ("workspace-default-ssh", 7_i64)] + { + conn.execute( + r#"INSERT INTO repository_ssh_credentials ( + workspace_id, credential_id, name, + public_key_algorithm, public_key_fingerprint, + current_revision, status, created_at + ) VALUES ('workspace', ?1, ?1, 'ssh-ed25519', ?1, ?2, + 'active', '2026-08-24T00:00:00Z')"#, + params![credential_id, revision], + )?; + conn.execute( + r#"INSERT INTO repository_ssh_credential_revisions ( + workspace_id, credential_id, revision, + public_key_algorithm, public_key_fingerprint, created_at + ) VALUES ('workspace', ?1, ?2, 'ssh-ed25519', ?1, + '2026-08-24T00:00:00Z')"#, + params![credential_id, revision], + )?; + } + Ok(()) + }) + .unwrap(); + let candidates = vec![ + WorkdirCreateCredentialCandidate { + role: WorkdirCreateCredentialCandidateRole::Primary, + credential_id: "credential-1".to_string(), + credential_revision: 3, + }, + WorkdirCreateCredentialCandidate { + role: WorkdirCreateCredentialCandidateRole::WorkspaceDefaultFallback, + credential_id: "workspace-default-ssh".to_string(), + credential_revision: 7, + }, + ]; let bound = store .bind_workdir_create_repository_access( "workspace", @@ -423,12 +646,22 @@ mod tests { "trust-1", 5, "read_only", + &candidates, "2026-08-24T00:00:01Z", ) .unwrap(); assert_eq!(bound.credential_id.as_deref(), Some("credential-1")); assert_eq!(bound.credential_revision, Some(3)); assert_eq!(bound.host_trust_revision, Some(5)); + assert_eq!(bound.credential_candidates, candidates); + let serialized = serde_json::to_string(&bound).unwrap(); + assert!(serialized.contains("workspace_default_fallback")); + assert!(!serialized.contains("private_key")); + assert!(!serialized.contains("known_hosts")); + // A concurrent Workspace-default rotation must not replace the fallback + // revision already bound to this operation. + let mut changed_candidates = candidates.clone(); + changed_candidates[1].credential_revision = 8; assert!( store .bind_workdir_create_repository_access( @@ -436,10 +669,11 @@ mod tests { "call-1", &record.request_fingerprint, "credential-1", - 4, + 3, "trust-1", 5, "read_only", + &changed_candidates, "2026-08-24T00:00:02Z", ) .is_err() @@ -475,6 +709,7 @@ mod tests { .unwrap(); assert_eq!(retry.state, "pending"); assert_eq!(retry.failure, None); + assert_eq!(retry.credential_candidates, candidates); assert_eq!( store .load_workdir_create_operation("workspace", "call-1") diff --git a/crates/workspace-server/src/workdir_removal.rs b/crates/workspace-server/src/workdir_removal.rs index 5edea488..87c2e9b1 100644 --- a/crates/workspace-server/src/workdir_removal.rs +++ b/crates/workspace-server/src/workdir_removal.rs @@ -1098,6 +1098,7 @@ mod tests { host_trust_id: None, host_trust_revision: None, repository_access_mode: None, + credential_candidates: Vec::new(), working_directory_id: "workdir-a".to_string(), state: "pending".to_string(), failure: None, diff --git a/crates/workspace-server/src/workspace_deletion.rs b/crates/workspace-server/src/workspace_deletion.rs index 5dda2f84..83af4ad7 100644 --- a/crates/workspace-server/src/workspace_deletion.rs +++ b/crates/workspace-server/src/workspace_deletion.rs @@ -55,6 +55,8 @@ const WORKSPACE_DELETION_PURGE_TABLES: &[&str] = &[ "typed_ticket_relations", "typed_ticket_risk_flags", "typed_tickets", + "workdir_create_credential_candidates", + "workdir_create_credential_revision_retentions", "workdir_create_operations", "workdir_registry", "workdir_removal_operations",