diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index f6f1db17..2ba5f157 100644 --- a/crates/worker-runtime/src/catalog.rs +++ b/crates/worker-runtime/src/catalog.rs @@ -202,6 +202,10 @@ impl std::fmt::Debug for WorkspaceApiRef { /// summarized without exposing raw host paths. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct CreateWorkerRequest { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub idempotency_key: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub idempotency_fingerprint: Option, pub profile: ProfileSelector, #[serde(default, skip_serializing_if = "Option::is_none")] pub display_name: Option, diff --git a/crates/worker-runtime/src/http_server.rs b/crates/worker-runtime/src/http_server.rs index 47c13cd2..4f0f4c95 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -1412,6 +1412,8 @@ mod tests { let profile = ProfileSelector::Builtin("builtin:coder".to_string()); let bundle = test_bundle(profile.clone()); CreateWorkerRequest { + idempotency_key: None, + idempotency_fingerprint: None, profile, display_name: None, profile_source: crate::catalog::ProfileSourceArchiveSource::Http { @@ -1812,6 +1814,8 @@ mod ws_tests { fn ws_create_request() -> CreateWorkerRequest { let bundle = ws_test_bundle(ProfileSelector::Builtin("builtin:companion".to_string())); CreateWorkerRequest { + idempotency_key: None, + idempotency_fingerprint: None, profile: ProfileSelector::Builtin("builtin:companion".to_string()), display_name: None, profile_source: crate::catalog::ProfileSourceArchiveSource::Http { diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index e6a71785..e7b69fdc 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -354,6 +354,11 @@ impl Runtime { request: CreateWorkerRequest, scope: Option<&RuntimeWorkspaceScope>, ) -> Result { + if request.idempotency_key.is_some() != request.idempotency_fingerprint.is_some() { + return Err(RuntimeError::InvalidRequest( + "idempotency_key and idempotency_fingerprint must be provided together".to_string(), + )); + } let (backend, worker_ref, spawn_request) = { let mut state = self.lock()?; state.ensure_running()?; @@ -375,6 +380,20 @@ impl Runtime { ))); } } + if let Some(idempotency_key) = request.idempotency_key.as_deref() { + let workspace_id = scope.map(|scope| scope.workspace_id.as_str()); + if let Some(existing) = state.workers.values().find(|record| { + record.workspace_id.as_deref() == workspace_id + && record.request.idempotency_key.as_deref() == Some(idempotency_key) + }) { + if existing.request.idempotency_fingerprint != request.idempotency_fingerprint { + return Err(RuntimeError::InvalidRequest(format!( + "worker creation idempotency key {idempotency_key} was already used with different input" + ))); + } + return Ok(existing.detail()); + } + } let backend = state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "worker creation requires an execution backend".to_string(), @@ -2126,6 +2145,8 @@ mod tests { let profile = ProfileSelector::Builtin("builtin:coder".to_string()); let bundle = test_bundle_for_profile(profile.clone()); CreateWorkerRequest { + idempotency_key: None, + idempotency_fingerprint: None, profile, display_name: None, profile_source: crate::catalog::ProfileSourceArchiveSource::Http { @@ -2614,6 +2635,24 @@ mod tests { )); } + #[test] + fn create_worker_idempotency_reuses_worker_and_rejects_different_input() { + let runtime = runtime_with_backend(); + let mut request = task_request("idempotent"); + request.idempotency_key = Some("operation-1".to_string()); + request.idempotency_fingerprint = Some("sha256:input-1".to_string()); + + let first = runtime.create_worker(request.clone()).unwrap(); + let replayed = runtime.create_worker(request.clone()).unwrap(); + assert_eq!(replayed.worker_ref, first.worker_ref); + assert_eq!(runtime.list_workers().unwrap().len(), 1); + + request.idempotency_fingerprint = Some("sha256:different".to_string()); + let error = runtime.create_worker(request).unwrap_err(); + assert!(matches!(error, RuntimeError::InvalidRequest(_))); + assert_eq!(runtime.list_workers().unwrap().len(), 1); + } + #[test] fn create_worker_rejects_system_initial_input_without_persisting_worker() { let runtime = runtime_with_backend(); diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 4916ce62..79f5fe27 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -1451,6 +1451,8 @@ mod tests { fn create_request(_name: &str) -> CreateWorkerRequest { let bundle = test_bundle(); CreateWorkerRequest { + idempotency_key: None, + idempotency_fingerprint: None, profile: ProfileSelector::Builtin("builtin:companion".to_string()), display_name: None, profile_source: crate::catalog::ProfileSourceArchiveSource::Embedded { diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index f0aa99b1..e44f828a 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -315,6 +315,20 @@ pub struct WorkerTicketAssignmentRequest { pub operation_id: String, } +pub(crate) fn worker_spawn_idempotency( + request: &WorkerSpawnRequest, +) -> Result, String> { + let Some(assignment) = request.ticket_assignment.as_ref() else { + return Ok(None); + }; + let encoded = serde_json::to_vec(request) + .map_err(|error| format!("serialize Worker spawn idempotency input: {error}"))?; + Ok(Some(( + assignment.operation_id.clone(), + format!("sha256:{}", digest_hex(&encoded, 64)), + ))) +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] pub struct WorkerSpawnRequest { @@ -1708,7 +1722,14 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { }; } }; + let (idempotency_key, idempotency_fingerprint) = worker_spawn_idempotency(&request) + .expect("WorkerSpawnRequest serialization is infallible") + .map_or((None, None), |(key, fingerprint)| { + (Some(key), Some(fingerprint)) + }); let create_request = CreateWorkerRequest { + idempotency_key, + idempotency_fingerprint, profile, display_name: request.requested_worker_name.clone(), config_bundle: None, @@ -2685,7 +2706,14 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { }; } }; + let (idempotency_key, idempotency_fingerprint) = worker_spawn_idempotency(&request) + .expect("WorkerSpawnRequest serialization is infallible") + .map_or((None, None), |(key, fingerprint)| { + (Some(key), Some(fingerprint)) + }); let create = CreateWorkerRequest { + idempotency_key, + idempotency_fingerprint, profile, display_name: request.requested_worker_name.clone(), config_bundle: None, diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index a5634da1..071d65f7 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -1774,11 +1774,14 @@ fn existing_lifecycle_assignment_worker( let Some(worker_id) = operation.worker_id else { return Ok(None); }; - Ok(Some( - api.runtime - .worker(runtime_id, &worker_id) - .map_err(|error| error.into_error())?, - )) + let worker = api + .runtime + .worker(runtime_id, &worker_id) + .map_err(|error| error.into_error())?; + if operation.assignment_id.is_none() && worker.state == "stopped" { + return Ok(None); + } + Ok(Some(worker)) } fn require_ticket_assignment_value(field: &str, value: String) -> Result { @@ -2258,6 +2261,16 @@ impl WorkerTicketSourceContext { } } +fn worker_source_actor_role(is_current_assignment: bool, is_orchestrator: bool) -> &'static str { + if is_current_assignment { + "coder" + } else if is_orchestrator { + "orchestrator" + } else { + "worker" + } +} + fn worker_ticket_source_context( api: &WorkspaceApi, workspace_id: &str, @@ -2271,17 +2284,13 @@ fn worker_ticket_source_context( .flatten() }); let orchestrator = find_workspace_orchestrator(api); - let actor_role = if assignment.as_ref().is_some_and(|assignment| { + let is_current_assignment = assignment.as_ref().is_some_and(|assignment| { assignment.runtime_id == source.runtime_id && assignment.worker_id == source.worker_id - }) { - "assigned" - } else if orchestrator.as_ref().is_some_and(|worker| { + }); + let is_orchestrator = orchestrator.as_ref().is_some_and(|worker| { worker.runtime_id == source.runtime_id && worker.worker_id == source.worker_id - }) { - "orchestrator" - } else { - "worker" - }; + }); + let actor_role = worker_source_actor_role(is_current_assignment, is_orchestrator); WorkerTicketSourceContext { workspace_id: workspace_id.to_string(), runtime_id: source.runtime_id.clone(), @@ -3987,6 +3996,25 @@ async fn scoped_restore_runtime_worker( } }; if let Some(assignment) = assignment_request.as_ref() { + let fingerprint = format!( + "sha256:{}", + Sha256::digest(format!( + "restore\0{}\0{}\0{}", + assignment.ticket_id, runtime_id, worker_id + )) + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::() + ); + api.store.reserve_ticket_assignment_operation( + &workspace_id, + &assignment.operation_id, + &assignment.ticket_id, + &runtime_id, + Some(&worker_id), + &fingerprint, + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + )?; if let Some(worker) = existing_lifecycle_assignment_worker(&api, assignment, &runtime_id)? { if worker.worker_id != worker_id { return Err(Error::TicketAssignmentConflict(format!( @@ -4015,14 +4043,6 @@ async fn scoped_restore_runtime_worker( ) .await?; if let Some(assignment) = assignment_request.as_ref() { - api.store.reserve_ticket_assignment_operation( - &workspace_id, - &assignment.operation_id, - &assignment.ticket_id, - &runtime_id, - &worker_id, - &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), - )?; assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker_id)?; } dispatch_pending_ticket_notifications(&api, &workspace_id); @@ -5760,6 +5780,21 @@ async fn create_runtime_worker( access_token: Some(credential.token), }); } + let spawn_idempotency = + crate::hosts::worker_spawn_idempotency(&request).map_err(Error::Config)?; + if let (Some(assignment), Some((_, fingerprint))) = + (lifecycle_assignment.as_ref(), spawn_idempotency.as_ref()) + { + api.store.reserve_ticket_assignment_operation( + &api.config.workspace_id, + &assignment.operation_id, + &assignment.ticket_id, + &runtime_id, + None, + fingerprint, + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + )?; + } let result = api .runtime .spawn_worker(&runtime_id, request) @@ -5778,13 +5813,10 @@ async fn create_runtime_worker( WorkerRegistryDisplayNamePolicy::UseProvided, )?; if let Some(assignment) = lifecycle_assignment.as_ref() { - api.store.reserve_ticket_assignment_operation( + api.store.bind_ticket_assignment_operation_worker( &api.config.workspace_id, &assignment.operation_id, - &assignment.ticket_id, - &runtime_id, &worker.worker_id, - &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), )?; assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker.worker_id)?; } @@ -8871,6 +8903,14 @@ mod tests { } } + #[test] + fn worker_source_actor_roles_use_canonical_vocabulary() { + assert_eq!(worker_source_actor_role(true, false), "coder"); + assert_eq!(worker_source_actor_role(false, true), "orchestrator"); + assert_eq!(worker_source_actor_role(false, false), "worker"); + assert_eq!(worker_source_actor_role(true, true), "coder"); + } + #[tokio::test] async fn ticket_assignment_endpoints_read_and_clear_current_assignment() { let dir = tempfile::tempdir().unwrap(); @@ -9169,7 +9209,7 @@ mod tests { .attributes .get("source_actor_role") .map(String::as_str), - Some("assigned") + Some("coder") ); let unauthorized = scoped_ticket_backend_operation( @@ -9455,16 +9495,39 @@ mod tests { TEST_CREATED_AT, ) .unwrap(); + let pending_request = WorkerSpawnRequest { + ticket_assignment: Some(crate::hosts::WorkerTicketAssignmentRequest { + ticket_id: second_ticket.id.clone(), + operation_id: "pending-spawn-operation".to_string(), + }), + ..request + }; + let (_, pending_fingerprint) = crate::hosts::worker_spawn_idempotency(&pending_request) + .unwrap() + .unwrap(); api.store .reserve_ticket_assignment_operation( TEST_WORKSPACE_ID, "pending-spawn-operation", &second_ticket.id, EMBEDDED_WORKER_RUNTIME_ID, - &first_worker.worker_id, + None, + &pending_fingerprint, TEST_CREATED_AT, ) .unwrap(); + let spawned_before_backend_failure = api + .runtime + .spawn_worker(EMBEDDED_WORKER_RUNTIME_ID, pending_request.clone()) + .unwrap() + .worker + .unwrap(); + assert!( + api.store + .get_ticket_assignment_operation(TEST_WORKSPACE_ID, "pending-spawn-operation") + .unwrap() + .is_some_and(|operation| operation.worker_id.is_none()) + ); let worker_count_before_retry = api.runtime.list_workers(100).items.len(); let Json(reconciled) = scoped_create_runtime_worker( State(api.clone()), @@ -9472,17 +9535,14 @@ mod tests { workspace_id: TEST_WORKSPACE_ID.to_string(), runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), }), - Json(WorkerSpawnRequest { - ticket_assignment: Some(crate::hosts::WorkerTicketAssignmentRequest { - ticket_id: second_ticket.id.clone(), - operation_id: "pending-spawn-operation".to_string(), - }), - ..request - }), + Json(pending_request), ) .await .unwrap(); - assert_eq!(reconciled.worker.unwrap().worker_id, first_worker.worker_id); + assert_eq!( + reconciled.worker.unwrap().worker_id, + spawned_before_backend_failure.worker_id + ); assert_eq!( api.runtime.list_workers(100).items.len(), worker_count_before_retry, @@ -10101,6 +10161,8 @@ mod tests { fn runtime_create_request() -> worker_runtime::catalog::CreateWorkerRequest { let bundle = runtime_test_bundle(); worker_runtime::catalog::CreateWorkerRequest { + idempotency_key: None, + idempotency_fingerprint: None, profile: worker_runtime::catalog::ProfileSelector::Builtin( "builtin:companion".to_string(), ), diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index b635a14c..da9e6636 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -107,6 +107,11 @@ const MIGRATIONS: &[Migration] = &[ name: "atomic Ticket notification identity credentials and cursors", apply: strengthen_ticket_notifications, }, + Migration { + version: 19, + name: "crash safe Worker lifecycle assignment reservations", + apply: strengthen_ticket_assignment_lifecycle_reservations, + }, ]; struct Migration { @@ -580,9 +585,16 @@ pub trait ControlPlaneStore: Send + Sync { operation_id: &str, ticket_id: &str, runtime_id: &str, - worker_id: &str, + worker_id: Option<&str>, + request_fingerprint: &str, created_at: &str, ) -> Result<()>; + fn bind_ticket_assignment_operation_worker( + &self, + workspace_id: &str, + operation_id: &str, + worker_id: &str, + ) -> Result<()>; fn get_current_ticket_worker_assignment( &self, workspace_id: &str, @@ -1799,15 +1811,16 @@ impl ControlPlaneStore for SqliteWorkspaceStore { operation_id: &str, ticket_id: &str, runtime_id: &str, - worker_id: &str, + worker_id: Option<&str>, + request_fingerprint: &str, created_at: &str, ) -> Result<()> { self.with_conn(|conn| { let inserted = conn.execute( r#"INSERT OR IGNORE INTO ticket_assignment_operations ( workspace_id, operation_id, action, ticket_id, runtime_id, worker_id, - assignment_id, expected_assignment_id, created_at - ) VALUES (?1, ?2, 'assign', ?3, ?4, ?5, NULL, NULL, ?6)"#, + assignment_id, expected_assignment_id, created_at, request_fingerprint + ) VALUES (?1, ?2, 'assign', ?3, ?4, ?5, NULL, NULL, ?6, ?7)"#, params![ workspace_id, operation_id, @@ -1815,6 +1828,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { runtime_id, worker_id, created_at, + request_fingerprint, ], )?; if inserted > 0 { @@ -1829,8 +1843,9 @@ impl ControlPlaneStore for SqliteWorkspaceStore { if existing.action == "assign" && existing.ticket_id == ticket_id && existing.runtime_id.as_deref() == Some(runtime_id) - && existing.worker_id.as_deref() == Some(worker_id) + && (worker_id.is_none() || existing.worker_id.as_deref() == worker_id) && existing.expected_assignment_id.is_none() + && existing.request_fingerprint.as_deref() == Some(request_fingerprint) { Ok(()) } else { @@ -1841,6 +1856,29 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }) } + fn bind_ticket_assignment_operation_worker( + &self, + workspace_id: &str, + operation_id: &str, + worker_id: &str, + ) -> Result<()> { + self.with_conn(|conn| { + let updated = conn.execute( + r#"UPDATE ticket_assignment_operations + SET worker_id = ?3 + WHERE workspace_id = ?1 AND operation_id = ?2 + AND assignment_id IS NULL AND (worker_id IS NULL OR worker_id = ?3)"#, + params![workspace_id, operation_id, worker_id], + )?; + if updated == 1 { + return Ok(()); + } + Err(Error::TicketAssignmentConflict(format!( + "assignment operation {operation_id} cannot bind Worker {worker_id}" + ))) + }) + } + fn get_current_ticket_worker_assignment( &self, workspace_id: &str, @@ -1890,11 +1928,30 @@ impl ControlPlaneStore for SqliteWorkspaceStore { params![record.workspace_id, record.ticket_id, assignment_id], read_ticket_worker_assignment_record, )?; + let previous = if existing.action == "reassign" { + existing + .expected_assignment_id + .as_deref() + .map(|previous_assignment_id| { + tx.query_row( + r#"SELECT workspace_id, ticket_id, assignment_id, runtime_id, worker_id, + assigned_by, assigned_at + FROM ticket_worker_assignments + WHERE workspace_id = ?1 AND ticket_id = ?2 AND assignment_id = ?3"#, + params![ + record.workspace_id, + record.ticket_id, + previous_assignment_id, + ], + read_ticket_worker_assignment_record, + ) + }) + .transpose()? + } else { + None + }; tx.commit()?; - return Ok(TicketWorkerAssignmentUpdate { - current, - previous: None, - }); + return Ok(TicketWorkerAssignmentUpdate { current, previous }); } reserved_operation = true; } @@ -3003,6 +3060,7 @@ pub struct TicketAssignmentOperationRecord { pub worker_id: Option, pub assignment_id: Option, pub expected_assignment_id: Option, + pub request_fingerprint: Option, } fn read_assignment_operation( @@ -3011,7 +3069,8 @@ fn read_assignment_operation( operation_id: &str, ) -> Result> { conn.query_row( - r#"SELECT action, ticket_id, runtime_id, worker_id, assignment_id, expected_assignment_id + r#"SELECT action, ticket_id, runtime_id, worker_id, assignment_id, expected_assignment_id, + request_fingerprint FROM ticket_assignment_operations WHERE workspace_id = ?1 AND operation_id = ?2"#, params![workspace_id, operation_id], @@ -3023,6 +3082,7 @@ fn read_assignment_operation( worker_id: row.get(3)?, assignment_id: row.get(4)?, expected_assignment_id: row.get(5)?, + request_fingerprint: row.get(6)?, }) }, ) @@ -3397,6 +3457,15 @@ CREATE TABLE ticket_notification_cursors ( Ok(()) } +fn strengthen_ticket_assignment_lifecycle_reservations(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" +ALTER TABLE ticket_assignment_operations ADD COLUMN request_fingerprint TEXT; +"#, + )?; + Ok(()) +} + fn create_objective_event_tables(conn: &Connection) -> Result<()> { conn.execute_batch( r#" @@ -4079,7 +4148,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); let db = dir.path().join("control-plane.sqlite"); let store = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 18); + assert_eq!(store.schema_version().await.unwrap(), 19); let record = WorkspaceRecord { workspace_id: "local-dev".to_string(), @@ -4092,7 +4161,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); store.upsert_workspace(&record).await.unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(reopened.schema_version().await.unwrap(), 18); + assert_eq!(reopened.schema_version().await.unwrap(), 19); assert_eq!( reopened.get_workspace("local-dev").await.unwrap(), Some(record) @@ -4102,7 +4171,8 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn ticket_worker_assignment_replaces_current_and_preserves_audit_history() { let dir = tempfile::tempdir().unwrap(); - let store = SqliteWorkspaceStore::open(dir.path().join("server.db")).unwrap(); + let db = dir.path().join("server.db"); + let store = SqliteWorkspaceStore::open(&db).unwrap(); store .upsert_workspace(&WorkspaceRecord { workspace_id: "workspace-a".to_string(), @@ -4204,6 +4274,16 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); .unwrap(); assert_eq!(replaced.current, second); assert_eq!(replaced.previous, Some(first.clone())); + let replayed_reassignment = store + .set_current_ticket_worker_assignment( + &second, + Some("assignment-1"), + "ignored-reassign-event", + "operation-2", + true, + ) + .unwrap(); + assert_eq!(replayed_reassignment, replaced); assert_eq!( store .get_current_ticket_worker_assignment("workspace-a", "ticket-1") @@ -4254,10 +4334,29 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); "reserved-operation", "ticket-3", "runtime-3", - "worker-3", + None, + "sha256:reserved", "2026-07-31T00:00:05Z", ) .unwrap(); + drop(store); + let store = SqliteWorkspaceStore::open(&db).unwrap(); + let pending = store + .get_ticket_assignment_operation("workspace-a", "reserved-operation") + .unwrap() + .unwrap(); + assert_eq!(pending.worker_id, None); + assert_eq!( + pending.request_fingerprint.as_deref(), + Some("sha256:reserved") + ); + store + .bind_ticket_assignment_operation_worker( + "workspace-a", + "reserved-operation", + "worker-3", + ) + .unwrap(); let reserved_assignment = TicketWorkerAssignmentRecord { workspace_id: "workspace-a".to_string(), ticket_id: "ticket-3".to_string(), @@ -4652,7 +4751,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); .unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 18); + assert_eq!(store.schema_version().await.unwrap(), 19); store .with_conn(|conn| { @@ -4755,7 +4854,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn repository_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 18); + assert_eq!(store.schema_version().await.unwrap(), 19); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -4793,7 +4892,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn memory_authority_records_round_trip_and_close_staging() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 18); + assert_eq!(store.schema_version().await.unwrap(), 19); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -4967,7 +5066,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn account_and_login_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 18); + assert_eq!(store.schema_version().await.unwrap(), 19); let now = "2026-07-22T00:00:00Z".to_string(); let account = AccountRecord { account_id: "acct-user-alice".to_string(),