From 4b8dc302ee2ae49b1885b9b45d37697cee2a33ed Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 20 Aug 2026 17:52:23 +0900 Subject: [PATCH] fix: fence workspace bootstrap and resource routing --- crates/workspace-server/src/server.rs | 133 +++++++++++++----- crates/workspace-server/src/store.rs | 17 +++ .../workspace-server/src/workspace_catalog.rs | 75 +++++++++- 3 files changed, 187 insertions(+), 38 deletions(-) diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 3642531f..b63decb3 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -842,25 +842,19 @@ async fn create_server_workspace( headers: HeaderMap, Json(request): Json, ) -> Response { - let owner_account_id = match resolve_server_actor(&api, &headers).await { - Ok(Some(actor)) => Some(actor.account_id), - Ok(None) if api.template.allow_local_workspace_bootstrap => { - match api.store.list_workspaces() { - Ok(workspaces) if workspaces.is_empty() => None, - Ok(_) => { - return forbidden_server_response( - "Workspace creation requires an authenticated owner", - ); - } - Err(error) => return server_error_response(error), - } - } + let (owner_account_id, local_bootstrap) = match resolve_server_actor(&api, &headers).await { + Ok(Some(actor)) => (Some(actor.account_id), false), + Ok(None) if api.template.allow_local_workspace_bootstrap => (None, true), Ok(None) => { return forbidden_server_response("Workspace creation requires an authenticated owner"); } Err(error) => return server_error_response(error), }; - let created = match api.catalog.create(request, owner_account_id) { + let created = match if local_bootstrap { + api.catalog.create_first_ownerless(request) + } else { + api.catalog.create(request, owner_account_id) + } { Ok(created) => created, Err(error) => return server_error_response(error), }; @@ -935,7 +929,10 @@ fn is_server_global_forward(path: &str) -> bool { fn scoped_workspace_id(path: &str) -> Option<&str> { let mut segments = path.trim_start_matches('/').split('/'); match (segments.next(), segments.next(), segments.next()) { - (Some("api"), Some("w"), Some(workspace_id)) if !workspace_id.is_empty() => { + (Some("api"), Some("w"), Some(workspace_id)) + | (Some("internal"), Some("w"), Some(workspace_id)) + if !workspace_id.is_empty() => + { Some(workspace_id) } (Some("w"), Some(workspace_id), _) if !workspace_id.is_empty() => Some(workspace_id), @@ -1963,8 +1960,8 @@ pub fn build_router(api: WorkspaceApi) -> Router { post(scoped_test_remote_runtime_connection), ) .route( - "/internal/runtime/resources/fetch", - post(post_internal_runtime_resource_fetch), + "/internal/w/{workspace_id}/runtime/resources/fetch", + post(scoped_post_internal_runtime_resource_fetch), ) .route("/api/companion/status", get(get_companion_status)) .route( @@ -9792,13 +9789,20 @@ fn browser_worker_response_from_summary( }) } -async fn post_internal_runtime_resource_fetch( +async fn scoped_post_internal_runtime_resource_fetch( State(api): State, + AxumPath(workspace_id): AxumPath, Json(request): Json, ) -> std::result::Result< Json, (StatusCode, Json), > { + if workspace_id != api.workspace_id() { + return Err(( + StatusCode::NOT_FOUND, + Json(BackendResourceError::MissingResource), + )); + } api.resource_broker .fetch_profile_source_archive(request) .map(Json) @@ -14466,6 +14470,37 @@ mod tests { assert_eq!(b["workspace_id"], workspace_b.workspace.workspace_id); assert_eq!(b["display_name"], "Workspace B"); + let handle = missing_resource_handle(); + let resource_response = app + .clone() + .oneshot( + Request::post(format!( + "/internal/w/{}/runtime/resources/fetch", + workspace_b.workspace.workspace_id + )) + .header(axum::http::header::CONTENT_TYPE, "application/json") + .body(Body::from( + serde_json::to_vec(&BackendResourceFetchRequest { + audit_correlation_id: handle.audit_correlation_id.clone(), + runtime_id: "runtime-test".to_string(), + worker_id: None, + handle, + }) + .unwrap(), + )) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(resource_response.status(), StatusCode::NOT_FOUND); + let resource_error: BackendResourceError = serde_json::from_slice( + &to_bytes(resource_response.into_body(), usize::MAX) + .await + .unwrap(), + ) + .unwrap(); + assert_eq!(resource_error, BackendResourceError::MissingResource); + let missing = app .oneshot( Request::builder() @@ -14519,6 +14554,7 @@ mod tests { assert_eq!(workspace["display_name"], "Created Workspace"); let replayed = app + .clone() .oneshot( Request::builder() .method(Method::POST) @@ -14529,10 +14565,29 @@ mod tests { ) .await .unwrap(); - // Local bootstrap authority is consumed after the first Workspace; even - // an exact HTTP retry must authenticate rather than creating another - // ownerless Workspace accidentally. - assert_eq!(replayed.status(), StatusCode::FORBIDDEN); + assert_eq!(replayed.status(), StatusCode::OK); + + let second_payload = json!({ + "operation_key": "bootstrap-2", + "display_name": "Second Ownerless Workspace", + "repository": { + "uri": repository, + "display_name": "Repository", + "default_ref": "HEAD" + } + }); + let second = app + .oneshot( + Request::builder() + .method(Method::POST) + .uri("/api/workspaces") + .header(axum::http::header::CONTENT_TYPE, "application/json") + .body(Body::from(second_payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(second.status(), StatusCode::CONFLICT); } #[test] @@ -14545,6 +14600,10 @@ mod tests { scoped_workspace_id("/w/workspace-b/workers"), Some("workspace-b") ); + assert_eq!( + scoped_workspace_id("/internal/w/workspace-c/runtime/resources/fetch"), + Some("workspace-c") + ); assert_eq!(scoped_workspace_id("/api/workspaces"), None); assert_eq!(scoped_workspace_id("/api/workspace"), None); } @@ -17146,20 +17205,20 @@ mod tests { let handle = missing_resource_handle(); let response = app .oneshot( - Request::post("/internal/runtime/resources/fetch") - .header("content-type", "application/json") - .body(Body::from( - serde_json::to_vec( - &worker_runtime::resource::BackendResourceFetchRequest { - audit_correlation_id: handle.audit_correlation_id.clone(), - runtime_id: "runtime-test".to_string(), - worker_id: None, - handle, - }, - ) - .unwrap(), - )) + Request::post(format!( + "/internal/w/{TEST_WORKSPACE_ID}/runtime/resources/fetch" + )) + .header("content-type", "application/json") + .body(Body::from( + serde_json::to_vec(&worker_runtime::resource::BackendResourceFetchRequest { + audit_correlation_id: handle.audit_correlation_id.clone(), + runtime_id: "runtime-test".to_string(), + worker_id: None, + handle, + }) .unwrap(), + )) + .unwrap(), ) .await .unwrap(); @@ -17182,7 +17241,7 @@ mod tests { let archive = test_profile_archive(); let runtime_id = "runtime-test"; let handle = broker.issue_profile_source_archive_handle( - "workspace-test", + TEST_WORKSPACE_ID, crate::resource_broker::BackendResourceTarget::Runtime(runtime_id), archive, ); @@ -17191,7 +17250,7 @@ mod tests { let addr = listener.local_addr().unwrap(); let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let client = worker_runtime::resource::HttpBackendResourceClient::new( - format!("http://{addr}/internal/runtime/resources/fetch"), + format!("http://{addr}/internal/w/{TEST_WORKSPACE_ID}/runtime/resources/fetch"), None, ); diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 01a8f1b0..f1cf1a68 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -275,6 +275,9 @@ pub struct RepositoryRecord { pub struct WorkspaceBootstrapRecord { pub operation_key: String, pub request_fingerprint: String, + /// When true, the transaction must prove that no Workspace exists before + /// it inserts this ownerless local-bootstrap Workspace. + pub require_empty_catalog: bool, pub workspace: WorkspaceRecord, pub repository: RepositoryRecord, } @@ -1444,6 +1447,20 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }); } + if record.require_empty_catalog { + let workspace_exists = tx.query_row( + "SELECT EXISTS(SELECT 1 FROM workspaces LIMIT 1)", + [], + |row| row.get::<_, bool>(0), + )?; + if workspace_exists { + return Err(Error::WorkspaceConfigConflict( + "ownerless local bootstrap is available only while the Workspace catalog is empty" + .to_string(), + )); + } + } + if let Some(existing) = tx .query_row( r#"SELECT workspace_id, owner_account_id, display_name, state, created_at, updated_at diff --git a/crates/workspace-server/src/workspace_catalog.rs b/crates/workspace-server/src/workspace_catalog.rs index d0d3ed66..adf67b6f 100644 --- a/crates/workspace-server/src/workspace_catalog.rs +++ b/crates/workspace-server/src/workspace_catalog.rs @@ -76,7 +76,14 @@ impl WorkspaceCatalogService { request: WorkspaceCreateRequest, owner_account_id: Option, ) -> Result { - self.create_with_workspace_id(request, owner_account_id, None) + self.create_internal(request, owner_account_id, None, false) + } + + pub fn create_first_ownerless( + &self, + request: WorkspaceCreateRequest, + ) -> Result { + self.create_internal(request, None, None, true) } pub fn create_with_workspace_id( @@ -84,6 +91,16 @@ impl WorkspaceCatalogService { request: WorkspaceCreateRequest, owner_account_id: Option, requested_workspace_id: Option, + ) -> Result { + self.create_internal(request, owner_account_id, requested_workspace_id, false) + } + + fn create_internal( + &self, + request: WorkspaceCreateRequest, + owner_account_id: Option, + requested_workspace_id: Option, + require_empty_catalog: bool, ) -> Result { let operation_key = normalize_required( "operation_key", @@ -134,6 +151,7 @@ impl WorkspaceCatalogService { .create_workspace_bootstrap(&WorkspaceBootstrapRecord { operation_key, request_fingerprint: fingerprint.clone(), + require_empty_catalog, workspace: WorkspaceRecord { workspace_id: workspace_id.clone(), owner_account_id, @@ -288,6 +306,61 @@ mod tests { ); } + #[test] + fn concurrent_ownerless_bootstrap_commits_exactly_one_workspace() { + let store = Arc::new(SqliteWorkspaceStore::in_memory().unwrap()); + let service = WorkspaceCatalogService::new(store.clone()); + let repository_a = git_repository(); + let repository_b = git_repository(); + let requests = [ + WorkspaceCreateRequest { + operation_key: "bootstrap-a".to_string(), + display_name: "Workspace A".to_string(), + repository: InitialRepositoryIntent { + uri: repository_a.path().display().to_string(), + display_name: None, + default_ref: None, + }, + }, + WorkspaceCreateRequest { + operation_key: "bootstrap-b".to_string(), + display_name: "Workspace B".to_string(), + repository: InitialRepositoryIntent { + uri: repository_b.path().display().to_string(), + display_name: None, + default_ref: None, + }, + }, + ]; + let barrier = Arc::new(std::sync::Barrier::new(2)); + let results = std::thread::scope(|scope| { + requests + .into_iter() + .map(|request| { + let service = service.clone(); + let barrier = barrier.clone(); + scope.spawn(move || { + barrier.wait(); + service.create_first_ownerless(request) + }) + }) + .collect::>() + .into_iter() + .map(|handle| handle.join().unwrap()) + .collect::>() + }); + + assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1); + assert_eq!(results.iter().filter(|result| result.is_err()).count(), 1); + assert_eq!(store.list_workspaces().unwrap().len(), 1); + let error = results + .into_iter() + .find_map(Result::err) + .unwrap() + .to_string(); + assert!(error.contains("catalog is empty"), "{error}"); + } + #[tokio::test] async fn idempotency_key_reuse_with_different_payload_is_rejected() { let store = Arc::new(SqliteWorkspaceStore::in_memory().unwrap());