From 7f0d02531290defcc08159513df6863fc0ab6f28 Mon Sep 17 00:00:00 2001 From: Hare Date: Tue, 11 Aug 2026 05:41:48 +0900 Subject: [PATCH] server: scope repository identity by workspace --- crates/workspace-server/src/server.rs | 252 ++++++++++-- crates/workspace-server/src/store.rs | 526 +++++++++++++++++++++++++- 2 files changed, 743 insertions(+), 35 deletions(-) diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index fdc60a38..6529ad6d 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -396,6 +396,7 @@ impl WorkspaceApi { runtime_id: &str, mut request: WorkerSpawnRequest, ) -> ApiResult { + self.validate_worker_spawn_repository_scope(&request)?; let workspace_api = self.workspace_api_ref(runtime_id); request.resolved_workspace_api = Some(workspace_api.clone()); let attachment_reservation = @@ -573,6 +574,64 @@ impl WorkspaceApi { fn repository_reader(&self) -> RepositoryRegistryReader { RepositoryRegistryReader::new(self.config.repositories.clone()) } + + fn require_workspace_repository(&self, repository_id: &str) -> ApiResult { + self.store + .get_repository(&self.config.workspace_id, repository_id)? + .ok_or_else(|| ApiError::from(Error::UnknownRepository(repository_id.to_string()))) + } + + fn require_configured_workspace_repository( + &self, + repository_id: &str, + ) -> ApiResult { + self.require_workspace_repository(repository_id)?; + self.config + .repositories + .iter() + .find(|repository| repository.id == repository_id) + .cloned() + .ok_or_else(|| ApiError::from(Error::UnknownRepository(repository_id.to_string()))) + } + + fn validate_worker_spawn_repository_scope( + &self, + request: &WorkerSpawnRequest, + ) -> ApiResult<()> { + let selected_repository_id = + if let Some(working_directory) = request.resolved_working_directory_request.as_ref() { + let repository_id = working_directory.repository.id.as_str(); + self.require_workspace_repository(repository_id)?; + Some(repository_id.to_string()) + } else if let Some(claim) = request.resolved_working_directory.as_ref() { + let workdir = self + .store + .get_workdir_registry(&self.config.workspace_id, &claim.working_directory_id)? + .ok_or_else(|| { + ApiError::from(Error::Config(format!( + "unknown working directory `{}` in this Workspace", + claim.working_directory_id + ))) + })?; + self.require_workspace_repository(&workdir.repository_id)?; + Some(workdir.repository_id) + } else { + None + }; + + if let WorkerSpawnIntent::TicketRole { ticket_id, .. } = &request.intent { + let ticket = self.authority.ticket(ticket_id)?; + if let Some(repository_id) = ticket.repository_id.as_deref() { + self.require_workspace_repository(repository_id)?; + if selected_repository_id.as_deref() != Some(repository_id) { + return Err(ApiError::from(Error::Config(format!( + "Ticket `{ticket_id}` targets repository `{repository_id}`, but the Worker launch does not resolve that repository in this Workspace" + )))); + } + } + } + Ok(()) + } } fn import_configured_repositories( @@ -2401,11 +2460,10 @@ async fn scoped_edit_ticket_item( ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; if let Some(TicketTargetEdit::Set { repository_id, .. }) = request.target.as_ref() { - if !api + if api .store - .list_repositories(&api.config.workspace_id)? - .iter() - .any(|repository| repository.repository_id == *repository_id) + .get_repository(&api.config.workspace_id, repository_id)? + .is_none() { return Err(settings_bad_request( "unknown_ticket_repository", @@ -2563,6 +2621,7 @@ async fn execute_worker_ticket_rest_operation( let is_mutation = operation_kind != "read"; let target = ticket_mutation_target(&operation).cloned(); let source = authenticate_worker_mutation_source(api, workspace_id, &headers)?; + validate_ticket_repository_operation(api, &operation)?; let before = target.as_ref().and_then(|id| backend.show(id.clone()).ok()); let previous_state = before .as_ref() @@ -3106,6 +3165,24 @@ fn ticket_mutation_target(operation: &TicketBackendOperation) -> Option<&TicketI } } +fn validate_ticket_repository_operation( + api: &WorkspaceApi, + operation: &TicketBackendOperation, +) -> ApiResult<()> { + let repository_id = match operation { + TicketBackendOperation::Create { input } => input.repository_id.as_deref(), + TicketBackendOperation::EditItem { edit, .. } => match edit.target.as_ref() { + Some(TicketTargetEdit::Set { repository_id, .. }) => Some(repository_id.as_str()), + _ => None, + }, + _ => None, + }; + if let Some(repository_id) = repository_id { + api.require_workspace_repository(repository_id)?; + } + Ok(()) +} + fn bind_worker_ticket_operation_source( source: &WorkerMutationSource, operation: &mut TicketBackendOperation, @@ -6819,19 +6896,22 @@ fn working_directory_request_from_repository( } fn configured_working_directory_request( - config: &ServerConfig, + api: &WorkspaceApi, request: &WorkerSpawnWorkingDirectoryRequest, ) -> Result { - let repository = config + if api + .store + .get_repository(&api.config.workspace_id, &request.repository_id)? + .is_none() + { + return Err(Error::UnknownRepository(request.repository_id.clone())); + } + let repository = api + .config .repositories .iter() .find(|repository| repository.id == request.repository_id) - .ok_or_else(|| { - Error::Config(format!( - "unknown repository id `{}` for Worker working directory", - request.repository_id - )) - })?; + .ok_or_else(|| Error::UnknownRepository(request.repository_id.clone()))?; Ok(working_directory_request_from_repository( repository, request.selector.as_deref(), @@ -7635,9 +7715,7 @@ async fn create_runtime_worker( request.resolved_working_directory_request = request .working_directory_request .as_ref() - .map(|working_directory| { - configured_working_directory_request(&api.config, working_directory) - }) + .map(|working_directory| configured_working_directory_request(&api, working_directory)) .transpose()?; let prepared_workdir_id = if let Some(working_directory_request) = request.resolved_working_directory_request.as_mut() @@ -9568,12 +9646,7 @@ fn working_directory_request_for_browser( api: &WorkspaceApi, request: BrowserWorkingDirectoryCreateRequest, ) -> ApiResult { - let repository = api - .config - .repositories - .iter() - .find(|repository| repository.id == request.repository_id) - .ok_or_else(|| Error::UnknownRepository(request.repository_id.clone()))?; + let repository = api.require_configured_workspace_repository(&request.repository_id)?; let selector = request .selector .or_else(|| repository.default_selector.clone()) @@ -10336,6 +10409,111 @@ mod tests { ); } + #[tokio::test] + async fn repository_bound_ticket_flow_and_workdir_launches_fail_closed_across_workspaces() { + let dir = tempfile::tempdir().unwrap(); + let api = test_api(dir.path()).await; + let other_workspace = WorkspaceRecord { + workspace_id: "other-workspace".to_string(), + owner_account_id: None, + display_name: "Other Workspace".to_string(), + state: "active".to_string(), + created_at: "1".to_string(), + updated_at: "1".to_string(), + }; + api.store.upsert_workspace(&other_workspace).await.unwrap(); + api.store + .upsert_repository(&RepositoryRecord { + workspace_id: other_workspace.workspace_id.clone(), + repository_id: "foreign".to_string(), + name: "Foreign".to_string(), + kind: "git".to_string(), + provider: Some("git".to_string()), + uri: dir.path().join("foreign").display().to_string(), + default_ref: Some("HEAD".to_string()), + auth_ref_kind: None, + auth_ref_key: None, + created_at: "1".to_string(), + updated_at: "1".to_string(), + }) + .unwrap(); + + let mut create_input = ticket::NewTicket::new("Foreign repository target"); + create_input.repository_id = Some("foreign".to_string()); + assert!( + validate_ticket_repository_operation( + &api, + &TicketBackendOperation::Create { + input: create_input.clone(), + }, + ) + .is_err() + ); + + let ticket = browser_ticket_backend(&api) + .unwrap() + .create(create_input) + .unwrap(); + let flow_ticket_launch = WorkerSpawnRequest { + requested_worker_name: Some("cross-workspace-ticket".to_string()), + intent: WorkerSpawnIntent::TicketRole { + ticket_id: ticket.id, + role: TicketWorkerRole::Coder, + }, + acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { + expected_segments: 2, + }, + profile: ProfileSelector::Builtin("builtin:coder".to_string()), + ticket_assignment: None, + initial_submit: vec![ + Segment::Flow { + selector: "builtin:coder-review".to_string(), + }, + Segment::text("Implement the Ticket"), + ], + working_directory_request: None, + resolved_working_directory_request: None, + resolved_working_directory: None, + resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), + resolved_workspace_api: None, + }; + assert!( + api.validate_worker_spawn_repository_scope(&flow_ticket_launch) + .is_err() + ); + + let mut foreign_repository = api.config.repositories[0].clone(); + foreign_repository.id = "foreign".to_string(); + let workdir_flow_launch = WorkerSpawnRequest { + requested_worker_name: Some("cross-workspace-workdir".to_string()), + intent: WorkerSpawnIntent::WorkspaceCoding, + acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { + expected_segments: 1, + }, + profile: ProfileSelector::Builtin("builtin:coder".to_string()), + ticket_assignment: None, + initial_submit: vec![Segment::Flow { + selector: "builtin:coder-review".to_string(), + }], + working_directory_request: None, + resolved_working_directory_request: Some(working_directory_request_from_repository( + &foreign_repository, + Some("HEAD"), + )), + resolved_working_directory: None, + resolved_config_bundle: None, + resolved_worker_observation_enabled: false, + resolved_worker_observation_grants: Vec::new(), + resolved_workspace_api: None, + }; + assert!( + api.validate_worker_spawn_repository_scope(&workdir_flow_launch) + .is_err() + ); + } + #[test] fn worker_ticket_assignment_projects_coder_intent_and_run_acceptance() { let initial_submit = vec![ @@ -10678,6 +10856,7 @@ mod tests { async fn workspace_workdir_summaries_include_runtime_observed_rows() { let dir = tempfile::tempdir().unwrap(); let api = test_api(dir.path()).await; + seed_test_repository(&api, "repo"); api.store .upsert_workdir_registry(&WorkdirRegistryRecord { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -12563,7 +12742,34 @@ mod tests { runtime_worker_id.to_string() } + fn seed_test_repository(api: &WorkspaceApi, repository_id: &str) { + if api + .store + .get_repository(&api.config.workspace_id, repository_id) + .unwrap() + .is_some() + { + return; + } + api.store + .upsert_repository(&RepositoryRecord { + workspace_id: api.config.workspace_id.clone(), + repository_id: repository_id.to_string(), + name: repository_id.to_string(), + kind: "git".to_string(), + provider: Some("git".to_string()), + uri: api.config.workspace_root.display().to_string(), + default_ref: Some("HEAD".to_string()), + auth_ref_kind: None, + auth_ref_key: None, + created_at: "1".to_string(), + updated_at: "1".to_string(), + }) + .unwrap(); + } + fn seed_cleanup_workdir(api: &WorkspaceApi, workdir_id: &str, status: &str, cleanliness: &str) { + seed_test_repository(api, "repo-test"); let now = now_registry_timestamp(); api.store .upsert_workdir_registry(&WorkdirRegistryRecord { @@ -14449,9 +14655,7 @@ mod tests { "/api/runtimes/embedded-worker-runtime/workers", json!({ "intent": { - "kind": "ticket_role", - "ticket_id": "00001KVZSGT0Q", - "role": "coder" + "kind": "workspace_coding" }, "requested_worker_name": "api-friendly-name", "acceptance": { diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 0b3617a6..e98b77b2 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -151,6 +151,11 @@ const MIGRATIONS: &[Migration] = &[ name: "remove Backend-owned Flow runtime authority", apply: remove_backend_flow_runtime_authority, }, + Migration { + version: 27, + name: "scope Repository identity and references by Workspace", + apply: scope_repository_identity_by_workspace, + }, ]; struct Migration { @@ -452,6 +457,11 @@ pub trait ControlPlaneStore: Send + Sync { async fn get_workspace(&self, workspace_id: &str) -> Result>; fn list_workspaces(&self) -> Result>; fn upsert_repository(&self, record: &RepositoryRecord) -> Result<()>; + fn get_repository( + &self, + workspace_id: &str, + repository_id: &str, + ) -> Result>; fn list_repositories(&self, workspace_id: &str) -> Result>; fn put_flow_source_for_kind( @@ -755,6 +765,7 @@ impl SqliteWorkspaceStore { configure_sqlite(&conn)?; apply_migrations(&conn)?; ticket::migrate_sqlite_ticket_schema(&conn)?; + validate_workspace_repository_references(&conn)?; Ok(Self { conn: Arc::new(Mutex::new(conn)), }) @@ -888,8 +899,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore { workspace_id, repository_id, name, kind, provider, uri, default_ref, auth_ref_kind, auth_ref_key, created_at, updated_at ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11) - ON CONFLICT(repository_id) DO UPDATE SET - workspace_id = excluded.workspace_id, + ON CONFLICT(workspace_id, repository_id) DO UPDATE SET name = excluded.name, kind = excluded.kind, provider = excluded.provider, @@ -916,6 +926,25 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }) } + fn get_repository( + &self, + workspace_id: &str, + repository_id: &str, + ) -> Result> { + self.with_conn(|conn| { + conn.query_row( + r#"SELECT workspace_id, repository_id, name, kind, provider, uri, default_ref, + auth_ref_kind, auth_ref_key, created_at, updated_at + FROM repositories + WHERE workspace_id = ?1 AND repository_id = ?2"#, + params![workspace_id, repository_id], + read_repository_record, + ) + .optional() + .map_err(Error::from) + }) + } + fn list_repositories(&self, workspace_id: &str) -> Result> { self.with_conn(|conn| { let mut stmt = conn.prepare( @@ -4020,6 +4049,205 @@ DROP TABLE IF EXISTS flow_instances; Ok(()) } +fn scope_repository_identity_by_workspace(conn: &Connection) -> Result<()> { + validate_workspace_repository_references(conn)?; + conn.execute_batch( + r#" +CREATE TABLE repositories_v27 ( + workspace_id TEXT NOT NULL, + repository_id TEXT NOT NULL, + name TEXT NOT NULL, + kind TEXT NOT NULL, + provider TEXT, + uri TEXT NOT NULL, + default_ref TEXT, + auth_ref_kind TEXT, + auth_ref_key TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, repository_id), + FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE +); +INSERT INTO repositories_v27 ( + workspace_id, repository_id, name, kind, provider, uri, default_ref, + auth_ref_kind, auth_ref_key, created_at, updated_at +) +SELECT workspace_id, repository_id, name, kind, provider, uri, default_ref, + auth_ref_kind, auth_ref_key, created_at, updated_at +FROM repositories; + +CREATE TABLE artifacts_v27 ( + workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE, + artifact_id TEXT PRIMARY KEY, + kind TEXT NOT NULL, + uri TEXT NOT NULL, + media_type TEXT, + sha256 TEXT, + size_bytes INTEGER, + summary TEXT, + created_at TEXT NOT NULL, + created_by_kind TEXT NOT NULL, + created_by_key TEXT NOT NULL, + created_by_display TEXT NOT NULL, + created_by_source_kind TEXT, + created_by_source_key TEXT, + ticket_id TEXT, + objective_id TEXT, + event_id TEXT, + worker_ref_kind TEXT, + worker_ref_key TEXT, + worker_display TEXT, + repository_id TEXT, + source_kind TEXT, + source_revision TEXT, + FOREIGN KEY (workspace_id, repository_id) + REFERENCES repositories_v27(workspace_id, repository_id) +); +INSERT INTO artifacts_v27 ( + workspace_id, artifact_id, kind, uri, media_type, sha256, size_bytes, summary, + created_at, created_by_kind, created_by_key, created_by_display, + created_by_source_kind, created_by_source_key, ticket_id, objective_id, event_id, + worker_ref_kind, worker_ref_key, worker_display, repository_id, source_kind, + source_revision +) +SELECT workspace_id, artifact_id, kind, uri, media_type, sha256, size_bytes, summary, + created_at, created_by_kind, created_by_key, created_by_display, + created_by_source_kind, created_by_source_key, ticket_id, objective_id, event_id, + worker_ref_kind, worker_ref_key, worker_display, repository_id, source_kind, + source_revision +FROM artifacts; + +CREATE TABLE workdir_registry_v27 ( + workspace_id TEXT NOT NULL, + workdir_id TEXT NOT NULL, + runtime_id TEXT NOT NULL, + repository_id TEXT NOT NULL, + creation_selector TEXT, + creation_ref TEXT, + materialization_status TEXT NOT NULL CHECK (materialization_status IN ('pending', 'present', 'not_found', 'corrupted', 'unknown', 'failed')), + cleanliness TEXT NOT NULL CHECK (cleanliness IN ('clean', 'dirty', 'unknown')), + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + current_selector TEXT, + current_ref TEXT, + PRIMARY KEY (workspace_id, workdir_id), + FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE, + FOREIGN KEY (workspace_id, repository_id) + REFERENCES repositories_v27(workspace_id, repository_id) +); +INSERT INTO workdir_registry_v27 ( + workspace_id, workdir_id, runtime_id, repository_id, + creation_selector, creation_ref, current_selector, current_ref, + materialization_status, cleanliness, created_at, updated_at +) +SELECT workspace_id, workdir_id, runtime_id, repository_id, + creation_selector, creation_ref, current_selector, current_ref, + materialization_status, cleanliness, created_at, updated_at +FROM workdir_registry; + +CREATE TABLE worker_workdir_links_v27 ( + workspace_id TEXT NOT NULL, + runtime_id TEXT NOT NULL, + runtime_worker_id INTEGER NOT NULL, + workdir_id TEXT NOT NULL, + role TEXT NOT NULL, + linked_at TEXT NOT NULL, + unlinked_at TEXT, + PRIMARY KEY (workspace_id, runtime_id, runtime_worker_id, workdir_id, role), + FOREIGN KEY (workspace_id, runtime_id, runtime_worker_id) + REFERENCES worker_registry(workspace_id, runtime_id, runtime_worker_id) ON DELETE CASCADE, + FOREIGN KEY (workspace_id, workdir_id) + REFERENCES workdir_registry_v27(workspace_id, workdir_id) ON DELETE CASCADE +); +INSERT INTO worker_workdir_links_v27 ( + workspace_id, runtime_id, runtime_worker_id, workdir_id, role, linked_at, unlinked_at +) +SELECT workspace_id, runtime_id, runtime_worker_id, workdir_id, role, linked_at, unlinked_at +FROM worker_workdir_links; + +CREATE TABLE worker_workdir_attachment_reservations_v27 ( + workspace_id TEXT NOT NULL, + workdir_id TEXT NOT NULL, + reservation_id TEXT NOT NULL, + reserved_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, workdir_id), + FOREIGN KEY (workspace_id, workdir_id) + REFERENCES workdir_registry_v27(workspace_id, workdir_id) ON DELETE CASCADE +); +INSERT INTO worker_workdir_attachment_reservations_v27 ( + workspace_id, workdir_id, reservation_id, reserved_at +) +SELECT workspace_id, workdir_id, reservation_id, reserved_at +FROM worker_workdir_attachment_reservations; + +DROP TABLE worker_workdir_links; +DROP TABLE worker_workdir_attachment_reservations; +DROP TABLE workdir_registry; +DROP TABLE artifacts; +DROP TABLE repositories; +ALTER TABLE repositories_v27 RENAME TO repositories; +ALTER TABLE artifacts_v27 RENAME TO artifacts; +ALTER TABLE workdir_registry_v27 RENAME TO workdir_registry; +ALTER TABLE worker_workdir_links_v27 RENAME TO worker_workdir_links; +ALTER TABLE worker_workdir_attachment_reservations_v27 + RENAME TO worker_workdir_attachment_reservations; + +CREATE INDEX idx_workdir_registry_workspace_updated + ON workdir_registry(workspace_id, updated_at DESC); +CREATE INDEX idx_worker_workdir_links_worker + ON worker_workdir_links(workspace_id, runtime_id, runtime_worker_id, linked_at DESC); +CREATE UNIQUE INDEX ux_worker_workdir_links_active_worker + ON worker_workdir_links(workspace_id, runtime_id, runtime_worker_id) + WHERE unlinked_at IS NULL; +CREATE UNIQUE INDEX ux_worker_workdir_links_active_workdir + ON worker_workdir_links(workspace_id, workdir_id) + WHERE unlinked_at IS NULL; +CREATE UNIQUE INDEX ux_worker_workdir_attachment_reservation_id + ON worker_workdir_attachment_reservations(workspace_id, reservation_id); +"#, + )?; + Ok(()) +} + +fn validate_workspace_repository_references(conn: &Connection) -> Result<()> { + for (table, repository_nullable) in [ + ("workdir_registry", false), + ("artifacts", true), + // `typed_tickets` is owned and migrated by the Ticket component. The control-plane + // migration may reject an already-invalid integrated reference, but must not rebuild + // that component table or claim its schema authority. + ("typed_tickets", true), + ] { + if !table_exists(conn, table)? || !column_exists(conn, table, "repository_id")? { + continue; + } + let null_filter = if repository_nullable { + "child.repository_id IS NOT NULL AND" + } else { + "" + }; + let sql = format!( + "SELECT child.workspace_id, child.repository_id FROM {table} AS child \ + WHERE {null_filter} NOT EXISTS (\ + SELECT 1 FROM repositories AS repository \ + WHERE repository.workspace_id = child.workspace_id \ + AND repository.repository_id = child.repository_id\ + ) LIMIT 1" + ); + let invalid = conn + .query_row(&sql, [], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + }) + .optional()?; + if let Some((workspace_id, repository_id)) = invalid { + return Err(Error::Store(format!( + "invalid Workspace-owned repository reference: {table} contains repository `{repository_id}` outside Workspace `{workspace_id}`" + ))); + } + } + Ok(()) +} + fn current_schema_version(conn: &Connection) -> Result { conn.query_row( "SELECT COALESCE(MAX(version), 0) FROM __yoi_schema_migrations", @@ -4618,7 +4846,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 26); + assert_eq!(current_schema_version(&conn).unwrap(), 27); assert!(table_exists(&conn, "worker_workdir_attachment_reservations").unwrap()); } @@ -4651,7 +4879,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY); apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 26); + assert_eq!(current_schema_version(&conn).unwrap(), 27); assert!(table_exists(&conn, "flow_sources").unwrap()); assert!(table_exists(&conn, "flow_source_revisions").unwrap()); assert!(!table_exists(&conn, "flow_instances").unwrap()); @@ -4659,13 +4887,246 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY); assert!(!table_exists(&conn, "flow_events").unwrap()); } + #[test] + fn schema_v27_upgrades_repository_identity_without_losing_workspace_owned_references() { + let conn = Connection::open_in_memory().unwrap(); + configure_sqlite(&conn).unwrap(); + for migration in MIGRATIONS + .iter() + .filter(|migration| migration.version <= 26) + { + let tx = conn.unchecked_transaction().unwrap(); + (migration.apply)(&tx).unwrap(); + tx.execute( + "INSERT INTO __yoi_schema_migrations (version, name) VALUES (?1, ?2)", + params![migration.version, migration.name], + ) + .unwrap(); + tx.commit().unwrap(); + } + conn.execute_batch( + r#" +INSERT INTO workspaces ( + workspace_id, display_name, state, created_at, updated_at +) VALUES + ('workspace-a', 'Workspace A', 'active', '1', '1'), + ('workspace-b', 'Workspace B', 'active', '1', '1'); +INSERT INTO repositories ( + repository_id, workspace_id, name, kind, uri, created_at, updated_at +) VALUES ('main', 'workspace-a', 'Main', 'git', '/repo-a', '1', '1'); +INSERT INTO artifacts ( + workspace_id, artifact_id, kind, uri, created_at, + created_by_kind, created_by_key, created_by_display, repository_id +) VALUES ( + 'workspace-a', 'artifact-1', 'report', 'artifact://1', '1', + 'worker', 'worker-1', 'Worker 1', 'main' +); +INSERT INTO worker_registry ( + workspace_id, runtime_id, runtime_worker_id, display_name, + retention_state, created_at, updated_at +) VALUES ('workspace-a', 'runtime-a', 1, 'Worker 1', 'normal', '1', '1'); +INSERT INTO workdir_registry ( + workspace_id, workdir_id, runtime_id, repository_id, + creation_selector, creation_ref, materialization_status, + cleanliness, created_at, updated_at, current_selector, current_ref +) VALUES ( + 'workspace-a', 'workdir-1', 'runtime-a', 'main', + 'develop', 'abc', 'present', 'clean', '1', '1', 'develop', 'abc' +); +INSERT INTO worker_workdir_links ( + workspace_id, runtime_id, runtime_worker_id, workdir_id, + role, linked_at, unlinked_at +) VALUES ('workspace-a', 'runtime-a', 1, 'workdir-1', 'attachment', '1', NULL); +INSERT INTO worker_workdir_attachment_reservations ( + workspace_id, workdir_id, reservation_id, reserved_at +) VALUES ('workspace-a', 'workdir-1', 'reservation-1', '1'); +"#, + ) + .unwrap(); + + apply_migrations(&conn).unwrap(); + + assert_eq!(current_schema_version(&conn).unwrap(), 27); + let repositories_sql: String = conn + .query_row( + "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'repositories'", + [], + |row| row.get(0), + ) + .unwrap(); + assert!(repositories_sql.contains("PRIMARY KEY (workspace_id, repository_id)")); + let preserved: (i64, i64, i64) = ( + conn.query_row("SELECT COUNT(*) FROM artifacts", [], |row| row.get(0)) + .unwrap(), + conn.query_row("SELECT COUNT(*) FROM worker_workdir_links", [], |row| { + row.get(0) + }) + .unwrap(), + conn.query_row( + "SELECT COUNT(*) FROM worker_workdir_attachment_reservations", + [], + |row| row.get(0), + ) + .unwrap(), + ); + assert_eq!(preserved, (1, 1, 1)); + let foreign_key_violations: i64 = conn + .query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |row| { + row.get(0) + }) + .unwrap(); + assert_eq!(foreign_key_violations, 0); + conn.execute( + r#"INSERT INTO repositories ( + workspace_id, repository_id, name, kind, uri, created_at, updated_at + ) VALUES ('workspace-b', 'main', 'Other Main', 'git', '/repo-b', '2', '2')"#, + [], + ) + .unwrap(); + assert_eq!( + conn.query_row( + "SELECT COUNT(*) FROM repositories WHERE repository_id = 'main'", + [], + |row| row.get::<_, i64>(0), + ) + .unwrap(), + 2 + ); + assert!( + conn.execute( + r#"INSERT INTO workdir_registry ( + workspace_id, workdir_id, runtime_id, repository_id, + materialization_status, cleanliness, created_at, updated_at + ) VALUES ('workspace-b', 'invalid', 'runtime-b', 'missing', + 'present', 'clean', '2', '2')"#, + [], + ) + .is_err() + ); + assert!( + conn.execute( + r#"INSERT INTO artifacts ( + workspace_id, artifact_id, kind, uri, created_at, + created_by_kind, created_by_key, created_by_display, repository_id + ) VALUES ('workspace-b', 'invalid-artifact', 'report', 'artifact://invalid', '2', + 'worker', 'worker-2', 'Worker 2', 'missing')"#, + [], + ) + .is_err() + ); + } + + #[test] + fn schema_v27_rejects_cross_workspace_legacy_repository_references() { + let conn = Connection::open_in_memory().unwrap(); + configure_sqlite(&conn).unwrap(); + for migration in MIGRATIONS + .iter() + .filter(|migration| migration.version <= 26) + { + let tx = conn.unchecked_transaction().unwrap(); + (migration.apply)(&tx).unwrap(); + tx.execute( + "INSERT INTO __yoi_schema_migrations (version, name) VALUES (?1, ?2)", + params![migration.version, migration.name], + ) + .unwrap(); + tx.commit().unwrap(); + } + conn.execute_batch( + r#" +INSERT INTO workspaces ( + workspace_id, display_name, state, created_at, updated_at +) VALUES + ('workspace-a', 'Workspace A', 'active', '1', '1'), + ('workspace-b', 'Workspace B', 'active', '1', '1'); +INSERT INTO repositories ( + repository_id, workspace_id, name, kind, uri, created_at, updated_at +) VALUES ('main', 'workspace-a', 'Main', 'git', '/repo-a', '1', '1'); +INSERT INTO workdir_registry ( + workspace_id, workdir_id, runtime_id, repository_id, + materialization_status, cleanliness, created_at, updated_at +) VALUES ('workspace-b', 'foreign-workdir', 'runtime-b', 'main', + 'present', 'clean', '1', '1'); +"#, + ) + .unwrap(); + + let error = apply_migrations(&conn).unwrap_err(); + + assert!(error.to_string().contains("workdir_registry")); + assert!(error.to_string().contains("workspace-b")); + assert_eq!(current_schema_version(&conn).unwrap(), 26); + assert_eq!( + conn.query_row("SELECT COUNT(*) FROM workdir_registry", [], |row| { + row.get::<_, i64>(0) + }) + .unwrap(), + 1 + ); + } + + #[tokio::test] + async fn startup_rejects_cross_workspace_ticket_repository_reference_without_claiming_ticket_schema() + { + let dir = tempfile::tempdir().unwrap(); + let database_path = dir.path().join("workspace.sqlite"); + let store = SqliteWorkspaceStore::open(&database_path).unwrap(); + for workspace_id in ["workspace-a", "workspace-b"] { + store + .upsert_workspace(&WorkspaceRecord { + workspace_id: workspace_id.to_string(), + owner_account_id: None, + display_name: workspace_id.to_string(), + state: "active".to_string(), + created_at: "1".to_string(), + updated_at: "1".to_string(), + }) + .await + .unwrap(); + } + store + .upsert_repository(&RepositoryRecord { + workspace_id: "workspace-a".to_string(), + repository_id: "main".to_string(), + name: "Main".to_string(), + kind: "git".to_string(), + provider: Some("git".to_string()), + uri: "/repo-a".to_string(), + default_ref: Some("HEAD".to_string()), + auth_ref_kind: None, + auth_ref_key: None, + created_at: "1".to_string(), + updated_at: "1".to_string(), + }) + .unwrap(); + drop(store); + + let backend = ticket::SqliteTicketBackend::open_verified( + database_path.clone(), + "workspace-b".to_string(), + ) + .unwrap(); + let mut input = ticket::NewTicket::new("Foreign repository"); + input.repository_id = Some("main".to_string()); + ticket::TicketBackend::create(&backend, input).unwrap(); + drop(backend); + + let error = match SqliteWorkspaceStore::open(&database_path) { + Ok(_) => panic!("cross-Workspace Ticket repository reference must fail closed"), + Err(error) => error, + }; + assert!(error.to_string().contains("typed_tickets")); + assert!(error.to_string().contains("workspace-b")); + } + #[tokio::test] async fn migrates_sqlite_and_preserves_workspace_record() { let dir = tempfile::tempdir().unwrap(); let db = dir.path().join("control-plane.sqlite"); let store = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 26); + assert_eq!(store.schema_version().await.unwrap(), 27); assert!( !store .with_conn(|conn| table_exists(conn, "worker_workspace_credentials")) @@ -4682,7 +5143,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY); store.upsert_workspace(&record).await.unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(reopened.schema_version().await.unwrap(), 26); + assert_eq!(reopened.schema_version().await.unwrap(), 27); assert_eq!( reopened.get_workspace("local-dev").await.unwrap(), Some(record) @@ -5229,7 +5690,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY); .unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 26); + assert_eq!(store.schema_version().await.unwrap(), 27); store .with_conn(|conn| { @@ -5418,7 +5879,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn repository_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 26); + assert_eq!(store.schema_version().await.unwrap(), 27); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -5443,20 +5904,48 @@ CREATE TABLE ticket_assignment_operations ( updated_at: "2".to_string(), }; store.upsert_repository(&repository).unwrap(); + assert_eq!( + store.get_repository("local-dev", "main").unwrap(), + Some(repository.clone()) + ); assert_eq!( store.list_repositories("local-dev").unwrap(), - vec![repository] + vec![repository.clone()] ); assert_eq!( store.list_repositories("other-workspace").unwrap(), Vec::new() ); + + let other_workspace = WorkspaceRecord { + workspace_id: "other-workspace".to_string(), + owner_account_id: None, + display_name: "Other Workspace".to_string(), + state: "active".to_string(), + created_at: "3".to_string(), + updated_at: "3".to_string(), + }; + store.upsert_workspace(&other_workspace).await.unwrap(); + let mut other_repository = repository.clone(); + other_repository.workspace_id = other_workspace.workspace_id.clone(); + other_repository.name = "Other Yoi".to_string(); + other_repository.uri = "/other/yoi".to_string(); + store.upsert_repository(&other_repository).unwrap(); + + assert_eq!( + store.get_repository("local-dev", "main").unwrap(), + Some(repository) + ); + assert_eq!( + store.get_repository("other-workspace", "main").unwrap(), + Some(other_repository) + ); } #[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(), 26); + assert_eq!(store.schema_version().await.unwrap(), 27); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -5542,6 +6031,21 @@ CREATE TABLE ticket_assignment_operations ( updated_at: "1".to_string(), }; store.upsert_workspace(&workspace).await.unwrap(); + store + .upsert_repository(&RepositoryRecord { + workspace_id: workspace.workspace_id.clone(), + repository_id: "repo".to_string(), + name: "Repository".to_string(), + kind: "git".to_string(), + provider: Some("git".to_string()), + uri: ".".to_string(), + default_ref: Some("HEAD".to_string()), + auth_ref_kind: None, + auth_ref_key: None, + created_at: "1".to_string(), + updated_at: "1".to_string(), + }) + .unwrap(); let worker = WorkerRegistryRecord { workspace_id: "local-dev".to_string(), @@ -5704,7 +6208,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn account_and_login_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 26); + assert_eq!(store.schema_version().await.unwrap(), 27); let now = "2026-07-22T00:00:00Z".to_string(); let account = AccountRecord { account_id: "acct-user-alice".to_string(),