server: scope repository identity by workspace
This commit is contained in:
@@ -396,6 +396,7 @@ impl WorkspaceApi {
|
|||||||
runtime_id: &str,
|
runtime_id: &str,
|
||||||
mut request: WorkerSpawnRequest,
|
mut request: WorkerSpawnRequest,
|
||||||
) -> ApiResult<WorkerSpawnResult> {
|
) -> ApiResult<WorkerSpawnResult> {
|
||||||
|
self.validate_worker_spawn_repository_scope(&request)?;
|
||||||
let workspace_api = self.workspace_api_ref(runtime_id);
|
let workspace_api = self.workspace_api_ref(runtime_id);
|
||||||
request.resolved_workspace_api = Some(workspace_api.clone());
|
request.resolved_workspace_api = Some(workspace_api.clone());
|
||||||
let attachment_reservation =
|
let attachment_reservation =
|
||||||
@@ -573,6 +574,64 @@ impl WorkspaceApi {
|
|||||||
fn repository_reader(&self) -> RepositoryRegistryReader {
|
fn repository_reader(&self) -> RepositoryRegistryReader {
|
||||||
RepositoryRegistryReader::new(self.config.repositories.clone())
|
RepositoryRegistryReader::new(self.config.repositories.clone())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn require_workspace_repository(&self, repository_id: &str) -> ApiResult<RepositoryRecord> {
|
||||||
|
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<ConfiguredRepository> {
|
||||||
|
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(
|
fn import_configured_repositories(
|
||||||
@@ -2401,11 +2460,10 @@ async fn scoped_edit_ticket_item(
|
|||||||
) -> ApiResult<Json<TicketDetail>> {
|
) -> ApiResult<Json<TicketDetail>> {
|
||||||
validate_workspace_scope(&api, &path.workspace_id)?;
|
validate_workspace_scope(&api, &path.workspace_id)?;
|
||||||
if let Some(TicketTargetEdit::Set { repository_id, .. }) = request.target.as_ref() {
|
if let Some(TicketTargetEdit::Set { repository_id, .. }) = request.target.as_ref() {
|
||||||
if !api
|
if api
|
||||||
.store
|
.store
|
||||||
.list_repositories(&api.config.workspace_id)?
|
.get_repository(&api.config.workspace_id, repository_id)?
|
||||||
.iter()
|
.is_none()
|
||||||
.any(|repository| repository.repository_id == *repository_id)
|
|
||||||
{
|
{
|
||||||
return Err(settings_bad_request(
|
return Err(settings_bad_request(
|
||||||
"unknown_ticket_repository",
|
"unknown_ticket_repository",
|
||||||
@@ -2563,6 +2621,7 @@ async fn execute_worker_ticket_rest_operation(
|
|||||||
let is_mutation = operation_kind != "read";
|
let is_mutation = operation_kind != "read";
|
||||||
let target = ticket_mutation_target(&operation).cloned();
|
let target = ticket_mutation_target(&operation).cloned();
|
||||||
let source = authenticate_worker_mutation_source(api, workspace_id, &headers)?;
|
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 before = target.as_ref().and_then(|id| backend.show(id.clone()).ok());
|
||||||
let previous_state = before
|
let previous_state = before
|
||||||
.as_ref()
|
.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(
|
fn bind_worker_ticket_operation_source(
|
||||||
source: &WorkerMutationSource,
|
source: &WorkerMutationSource,
|
||||||
operation: &mut TicketBackendOperation,
|
operation: &mut TicketBackendOperation,
|
||||||
@@ -6819,19 +6896,22 @@ fn working_directory_request_from_repository(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn configured_working_directory_request(
|
fn configured_working_directory_request(
|
||||||
config: &ServerConfig,
|
api: &WorkspaceApi,
|
||||||
request: &WorkerSpawnWorkingDirectoryRequest,
|
request: &WorkerSpawnWorkingDirectoryRequest,
|
||||||
) -> Result<WorkingDirectoryRequest> {
|
) -> Result<WorkingDirectoryRequest> {
|
||||||
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
|
.repositories
|
||||||
.iter()
|
.iter()
|
||||||
.find(|repository| repository.id == request.repository_id)
|
.find(|repository| repository.id == request.repository_id)
|
||||||
.ok_or_else(|| {
|
.ok_or_else(|| Error::UnknownRepository(request.repository_id.clone()))?;
|
||||||
Error::Config(format!(
|
|
||||||
"unknown repository id `{}` for Worker working directory",
|
|
||||||
request.repository_id
|
|
||||||
))
|
|
||||||
})?;
|
|
||||||
Ok(working_directory_request_from_repository(
|
Ok(working_directory_request_from_repository(
|
||||||
repository,
|
repository,
|
||||||
request.selector.as_deref(),
|
request.selector.as_deref(),
|
||||||
@@ -7635,9 +7715,7 @@ async fn create_runtime_worker(
|
|||||||
request.resolved_working_directory_request = request
|
request.resolved_working_directory_request = request
|
||||||
.working_directory_request
|
.working_directory_request
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.map(|working_directory| {
|
.map(|working_directory| configured_working_directory_request(&api, working_directory))
|
||||||
configured_working_directory_request(&api.config, working_directory)
|
|
||||||
})
|
|
||||||
.transpose()?;
|
.transpose()?;
|
||||||
let prepared_workdir_id = if let Some(working_directory_request) =
|
let prepared_workdir_id = if let Some(working_directory_request) =
|
||||||
request.resolved_working_directory_request.as_mut()
|
request.resolved_working_directory_request.as_mut()
|
||||||
@@ -9568,12 +9646,7 @@ fn working_directory_request_for_browser(
|
|||||||
api: &WorkspaceApi,
|
api: &WorkspaceApi,
|
||||||
request: BrowserWorkingDirectoryCreateRequest,
|
request: BrowserWorkingDirectoryCreateRequest,
|
||||||
) -> ApiResult<WorkingDirectoryRequest> {
|
) -> ApiResult<WorkingDirectoryRequest> {
|
||||||
let repository = api
|
let repository = api.require_configured_workspace_repository(&request.repository_id)?;
|
||||||
.config
|
|
||||||
.repositories
|
|
||||||
.iter()
|
|
||||||
.find(|repository| repository.id == request.repository_id)
|
|
||||||
.ok_or_else(|| Error::UnknownRepository(request.repository_id.clone()))?;
|
|
||||||
let selector = request
|
let selector = request
|
||||||
.selector
|
.selector
|
||||||
.or_else(|| repository.default_selector.clone())
|
.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]
|
#[test]
|
||||||
fn worker_ticket_assignment_projects_coder_intent_and_run_acceptance() {
|
fn worker_ticket_assignment_projects_coder_intent_and_run_acceptance() {
|
||||||
let initial_submit = vec![
|
let initial_submit = vec![
|
||||||
@@ -10678,6 +10856,7 @@ mod tests {
|
|||||||
async fn workspace_workdir_summaries_include_runtime_observed_rows() {
|
async fn workspace_workdir_summaries_include_runtime_observed_rows() {
|
||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
let api = test_api(dir.path()).await;
|
let api = test_api(dir.path()).await;
|
||||||
|
seed_test_repository(&api, "repo");
|
||||||
api.store
|
api.store
|
||||||
.upsert_workdir_registry(&WorkdirRegistryRecord {
|
.upsert_workdir_registry(&WorkdirRegistryRecord {
|
||||||
workspace_id: TEST_WORKSPACE_ID.to_string(),
|
workspace_id: TEST_WORKSPACE_ID.to_string(),
|
||||||
@@ -12563,7 +12742,34 @@ mod tests {
|
|||||||
runtime_worker_id.to_string()
|
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) {
|
fn seed_cleanup_workdir(api: &WorkspaceApi, workdir_id: &str, status: &str, cleanliness: &str) {
|
||||||
|
seed_test_repository(api, "repo-test");
|
||||||
let now = now_registry_timestamp();
|
let now = now_registry_timestamp();
|
||||||
api.store
|
api.store
|
||||||
.upsert_workdir_registry(&WorkdirRegistryRecord {
|
.upsert_workdir_registry(&WorkdirRegistryRecord {
|
||||||
@@ -14449,9 +14655,7 @@ mod tests {
|
|||||||
"/api/runtimes/embedded-worker-runtime/workers",
|
"/api/runtimes/embedded-worker-runtime/workers",
|
||||||
json!({
|
json!({
|
||||||
"intent": {
|
"intent": {
|
||||||
"kind": "ticket_role",
|
"kind": "workspace_coding"
|
||||||
"ticket_id": "00001KVZSGT0Q",
|
|
||||||
"role": "coder"
|
|
||||||
},
|
},
|
||||||
"requested_worker_name": "api-friendly-name",
|
"requested_worker_name": "api-friendly-name",
|
||||||
"acceptance": {
|
"acceptance": {
|
||||||
|
|||||||
@@ -151,6 +151,11 @@ const MIGRATIONS: &[Migration] = &[
|
|||||||
name: "remove Backend-owned Flow runtime authority",
|
name: "remove Backend-owned Flow runtime authority",
|
||||||
apply: remove_backend_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 {
|
struct Migration {
|
||||||
@@ -452,6 +457,11 @@ pub trait ControlPlaneStore: Send + Sync {
|
|||||||
async fn get_workspace(&self, workspace_id: &str) -> Result<Option<WorkspaceRecord>>;
|
async fn get_workspace(&self, workspace_id: &str) -> Result<Option<WorkspaceRecord>>;
|
||||||
fn list_workspaces(&self) -> Result<Vec<WorkspaceRecord>>;
|
fn list_workspaces(&self) -> Result<Vec<WorkspaceRecord>>;
|
||||||
fn upsert_repository(&self, record: &RepositoryRecord) -> Result<()>;
|
fn upsert_repository(&self, record: &RepositoryRecord) -> Result<()>;
|
||||||
|
fn get_repository(
|
||||||
|
&self,
|
||||||
|
workspace_id: &str,
|
||||||
|
repository_id: &str,
|
||||||
|
) -> Result<Option<RepositoryRecord>>;
|
||||||
fn list_repositories(&self, workspace_id: &str) -> Result<Vec<RepositoryRecord>>;
|
fn list_repositories(&self, workspace_id: &str) -> Result<Vec<RepositoryRecord>>;
|
||||||
|
|
||||||
fn put_flow_source_for_kind(
|
fn put_flow_source_for_kind(
|
||||||
@@ -755,6 +765,7 @@ impl SqliteWorkspaceStore {
|
|||||||
configure_sqlite(&conn)?;
|
configure_sqlite(&conn)?;
|
||||||
apply_migrations(&conn)?;
|
apply_migrations(&conn)?;
|
||||||
ticket::migrate_sqlite_ticket_schema(&conn)?;
|
ticket::migrate_sqlite_ticket_schema(&conn)?;
|
||||||
|
validate_workspace_repository_references(&conn)?;
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
conn: Arc::new(Mutex::new(conn)),
|
conn: Arc::new(Mutex::new(conn)),
|
||||||
})
|
})
|
||||||
@@ -888,8 +899,7 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
|
|||||||
workspace_id, repository_id, name, kind, provider, uri, default_ref,
|
workspace_id, repository_id, name, kind, provider, uri, default_ref,
|
||||||
auth_ref_kind, auth_ref_key, created_at, updated_at
|
auth_ref_kind, auth_ref_key, created_at, updated_at
|
||||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)
|
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)
|
||||||
ON CONFLICT(repository_id) DO UPDATE SET
|
ON CONFLICT(workspace_id, repository_id) DO UPDATE SET
|
||||||
workspace_id = excluded.workspace_id,
|
|
||||||
name = excluded.name,
|
name = excluded.name,
|
||||||
kind = excluded.kind,
|
kind = excluded.kind,
|
||||||
provider = excluded.provider,
|
provider = excluded.provider,
|
||||||
@@ -916,6 +926,25 @@ impl ControlPlaneStore for SqliteWorkspaceStore {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn get_repository(
|
||||||
|
&self,
|
||||||
|
workspace_id: &str,
|
||||||
|
repository_id: &str,
|
||||||
|
) -> Result<Option<RepositoryRecord>> {
|
||||||
|
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<Vec<RepositoryRecord>> {
|
fn list_repositories(&self, workspace_id: &str) -> Result<Vec<RepositoryRecord>> {
|
||||||
self.with_conn(|conn| {
|
self.with_conn(|conn| {
|
||||||
let mut stmt = conn.prepare(
|
let mut stmt = conn.prepare(
|
||||||
@@ -4020,6 +4049,205 @@ DROP TABLE IF EXISTS flow_instances;
|
|||||||
Ok(())
|
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<i64> {
|
fn current_schema_version(conn: &Connection) -> Result<i64> {
|
||||||
conn.query_row(
|
conn.query_row(
|
||||||
"SELECT COALESCE(MAX(version), 0) FROM __yoi_schema_migrations",
|
"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();
|
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());
|
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();
|
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_sources").unwrap());
|
||||||
assert!(table_exists(&conn, "flow_source_revisions").unwrap());
|
assert!(table_exists(&conn, "flow_source_revisions").unwrap());
|
||||||
assert!(!table_exists(&conn, "flow_instances").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());
|
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]
|
#[tokio::test]
|
||||||
async fn migrates_sqlite_and_preserves_workspace_record() {
|
async fn migrates_sqlite_and_preserves_workspace_record() {
|
||||||
let dir = tempfile::tempdir().unwrap();
|
let dir = tempfile::tempdir().unwrap();
|
||||||
let db = dir.path().join("control-plane.sqlite");
|
let db = dir.path().join("control-plane.sqlite");
|
||||||
let store = SqliteWorkspaceStore::open(&db).unwrap();
|
let store = SqliteWorkspaceStore::open(&db).unwrap();
|
||||||
|
|
||||||
assert_eq!(store.schema_version().await.unwrap(), 26);
|
assert_eq!(store.schema_version().await.unwrap(), 27);
|
||||||
assert!(
|
assert!(
|
||||||
!store
|
!store
|
||||||
.with_conn(|conn| table_exists(conn, "worker_workspace_credentials"))
|
.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();
|
store.upsert_workspace(&record).await.unwrap();
|
||||||
|
|
||||||
let reopened = SqliteWorkspaceStore::open(&db).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!(
|
assert_eq!(
|
||||||
reopened.get_workspace("local-dev").await.unwrap(),
|
reopened.get_workspace("local-dev").await.unwrap(),
|
||||||
Some(record)
|
Some(record)
|
||||||
@@ -5229,7 +5690,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY);
|
|||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
let store = SqliteWorkspaceStore::from_connection(conn).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
|
store
|
||||||
.with_conn(|conn| {
|
.with_conn(|conn| {
|
||||||
@@ -5418,7 +5879,7 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn repository_records_round_trip() {
|
async fn repository_records_round_trip() {
|
||||||
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
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 {
|
let workspace = WorkspaceRecord {
|
||||||
workspace_id: "local-dev".to_string(),
|
workspace_id: "local-dev".to_string(),
|
||||||
owner_account_id: None,
|
owner_account_id: None,
|
||||||
@@ -5443,20 +5904,48 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
updated_at: "2".to_string(),
|
updated_at: "2".to_string(),
|
||||||
};
|
};
|
||||||
store.upsert_repository(&repository).unwrap();
|
store.upsert_repository(&repository).unwrap();
|
||||||
|
assert_eq!(
|
||||||
|
store.get_repository("local-dev", "main").unwrap(),
|
||||||
|
Some(repository.clone())
|
||||||
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
store.list_repositories("local-dev").unwrap(),
|
store.list_repositories("local-dev").unwrap(),
|
||||||
vec![repository]
|
vec![repository.clone()]
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
store.list_repositories("other-workspace").unwrap(),
|
store.list_repositories("other-workspace").unwrap(),
|
||||||
Vec::new()
|
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]
|
#[tokio::test]
|
||||||
async fn memory_authority_records_round_trip_and_close_staging() {
|
async fn memory_authority_records_round_trip_and_close_staging() {
|
||||||
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
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 {
|
let workspace = WorkspaceRecord {
|
||||||
workspace_id: "local-dev".to_string(),
|
workspace_id: "local-dev".to_string(),
|
||||||
owner_account_id: None,
|
owner_account_id: None,
|
||||||
@@ -5542,6 +6031,21 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
updated_at: "1".to_string(),
|
updated_at: "1".to_string(),
|
||||||
};
|
};
|
||||||
store.upsert_workspace(&workspace).await.unwrap();
|
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 {
|
let worker = WorkerRegistryRecord {
|
||||||
workspace_id: "local-dev".to_string(),
|
workspace_id: "local-dev".to_string(),
|
||||||
@@ -5704,7 +6208,7 @@ CREATE TABLE ticket_assignment_operations (
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn account_and_login_records_round_trip() {
|
async fn account_and_login_records_round_trip() {
|
||||||
let store = SqliteWorkspaceStore::in_memory().unwrap();
|
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 now = "2026-07-22T00:00:00Z".to_string();
|
||||||
let account = AccountRecord {
|
let account = AccountRecord {
|
||||||
account_id: "acct-user-alice".to_string(),
|
account_id: "acct-user-alice".to_string(),
|
||||||
|
|||||||
Reference in New Issue
Block a user