diff --git a/Cargo.lock b/Cargo.lock index 6d3627ed..d2650c4c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6618,6 +6618,7 @@ dependencies = [ "wasmtime", "wat", "workdir", + "workspace-api", "yoi-plugin-pdk", ] diff --git a/crates/manifest/src/config.rs b/crates/manifest/src/config.rs index 9671675d..1eb3e86c 100644 --- a/crates/manifest/src/config.rs +++ b/crates/manifest/src/config.rs @@ -92,6 +92,8 @@ pub struct FeatureConfigPartial { #[serde(default)] pub worker: Option, #[serde(default)] + pub workspace_worker_discovery: Option, + #[serde(default)] pub objective: Option, #[serde(default)] pub manage_workdir: Option, @@ -119,6 +121,11 @@ impl FeatureConfigPartial { ), flow: merge_option(self.flow, other.flow, FeatureFlagConfigPartial::merge), worker: merge_option(self.worker, other.worker, WorkerFeatureConfigPartial::merge), + workspace_worker_discovery: merge_option( + self.workspace_worker_discovery, + other.workspace_worker_discovery, + FeatureFlagConfigPartial::merge, + ), objective: merge_option( self.objective, other.objective, @@ -265,6 +272,10 @@ impl From for FeatureConfig { .worker .map(WorkerFeatureConfig::from) .unwrap_or_default(), + workspace_worker_discovery: value + .workspace_worker_discovery + .map(FeatureFlagConfig::from) + .unwrap_or_default(), objective: value .objective .map(FeatureFlagConfig::from) @@ -394,6 +405,7 @@ impl From for FeatureConfigPartial { sub_worker: Some(value.sub_worker.into()), flow: Some(value.flow.into()), worker: Some(value.worker.into()), + workspace_worker_discovery: Some(value.workspace_worker_discovery.into()), objective: Some(value.objective.into()), manage_workdir: Some(value.manage_workdir.into()), ticket: Some(value.ticket.into()), diff --git a/crates/manifest/src/lib.rs b/crates/manifest/src/lib.rs index 3d057c05..5cb31db9 100644 --- a/crates/manifest/src/lib.rs +++ b/crates/manifest/src/lib.rs @@ -118,6 +118,10 @@ pub struct FeatureConfig { pub flow: FeatureFlagConfig, #[serde(default)] pub worker: WorkerFeatureConfig, + /// Privileged read-only discovery of visible Workspace Workers. Backend + /// source proof remains required for every listing operation. + #[serde(default)] + pub workspace_worker_discovery: FeatureFlagConfig, #[serde(default)] pub objective: FeatureFlagConfig, #[serde(default)] @@ -142,6 +146,7 @@ impl Default for FeatureConfig { sub_worker: FeatureFlagConfig::disabled(), flow: FeatureFlagConfig::disabled(), worker: WorkerFeatureConfig::disabled(), + workspace_worker_discovery: FeatureFlagConfig::disabled(), objective: FeatureFlagConfig::disabled(), manage_workdir: FeatureFlagConfig::disabled(), ticket: TicketFeatureConfig::default(), diff --git a/crates/manifest/src/profile.rs b/crates/manifest/src/profile.rs index d84c195d..b6cf09d0 100644 --- a/crates/manifest/src/profile.rs +++ b/crates/manifest/src/profile.rs @@ -946,9 +946,11 @@ fn apply_role_profile( value["feature"]["sub_worker"] = serde_json::json!({ "enabled": sub_worker }); value["feature"]["flow"] = serde_json::json!({ "enabled": slug == "coder" }); value["feature"]["worker"] = serde_json::json!({ - "enabled": slug == "orchestrator", - "direct_spawn": slug != "orchestrator" + "enabled": matches!(slug, "companion" | "orchestrator"), + "direct_spawn": !matches!(slug, "companion" | "orchestrator") }); + value["feature"]["workspace_worker_discovery"] = + serde_json::json!({ "enabled": slug == "companion" }); value["feature"]["manage_workdir"] = serde_json::json!({ "enabled": matches!(slug, "companion" | "orchestrator") }); @@ -1423,7 +1425,7 @@ mod tests { } #[test] - fn builtin_companion_uses_sub_worker_control_without_worker_control() { + fn builtin_companion_combines_runtime_and_sub_worker_control_with_discovery() { let tmp = TempDir::new().unwrap(); let resolved = ProfileResolver::new() .with_workspace_base(tmp.path()) @@ -1435,7 +1437,9 @@ mod tests { assert!(resolved.manifest.feature.manage_workdir.enabled); assert!(resolved.manifest.feature.sub_worker.enabled); - assert!(!resolved.manifest.feature.worker.enabled); + assert!(resolved.manifest.feature.worker.enabled); + assert!(!resolved.manifest.feature.worker.direct_spawn); + assert!(resolved.manifest.feature.workspace_worker_discovery.enabled); } #[test] diff --git a/crates/worker-runtime/src/auth.rs b/crates/worker-runtime/src/auth.rs index bf02b35f..05279ea5 100644 --- a/crates/worker-runtime/src/auth.rs +++ b/crates/worker-runtime/src/auth.rs @@ -17,6 +17,7 @@ const WORKER_MUTATION_SOURCE_SIGNING_INPUT_PREFIX: &str = "yoi-worker-source-v1. pub const WORKER_REMOVE_PERMISSION: &str = "workspace:worker-remove"; pub const RUNTIME_REQUEST_SOURCE_PROOF_HEADER: &str = "x-yoi-runtime-request-proof"; pub const WORKSPACE_REQUEST_PERMISSION: &str = "workspace:request"; +pub const WORKSPACE_WORKER_DISCOVERY_PERMISSION: &str = "workspace:worker-discovery"; pub const BACKEND_RESOURCE_FETCH_PERMISSION: &str = "workspace:resource-fetch"; const RUNTIME_REQUEST_SOURCE_PROOF_PREFIX: &str = "yoi-runtime-request-v1"; const RUNTIME_REQUEST_SOURCE_SIGNING_INPUT_PREFIX: &str = "yoi-runtime-request-v1."; diff --git a/crates/worker-runtime/src/worker_source.rs b/crates/worker-runtime/src/worker_source.rs index 4c43e4a1..ed485ab0 100644 --- a/crates/worker-runtime/src/worker_source.rs +++ b/crates/worker-runtime/src/worker_source.rs @@ -9,8 +9,8 @@ use worker::{ use crate::auth::{ RUNTIME_REQUEST_SOURCE_PROOF_HEADER, RuntimeAuthError, RuntimeIdentityMaterial, RuntimeRequestSourceSigner, RuntimeWorkerMutationSourceSigner, WORKER_REMOVE_PERMISSION, - WORKSPACE_REQUEST_PERMISSION, WorkerMutationActorKind, WorkerMutationOperation, - WorkerMutationSourceClaims, new_token_id, + WORKSPACE_REQUEST_PERMISSION, WORKSPACE_WORKER_DISCOVERY_PERMISSION, WorkerMutationActorKind, + WorkerMutationOperation, WorkerMutationSourceClaims, new_token_id, }; use crate::runtime::RuntimeWorkspaceScope; use crate::worker_backend::WorkspacePromptProjectionCache; @@ -343,6 +343,51 @@ impl RuntimeOwnedWorkspaceClient { self.request_timeout = request_timeout; self } + + fn execute_with_permission( + &self, + request: WorkspaceRequest, + permission: &'static str, + ) -> Result { + let base_url = self.base_url.clone(); + let workspace_id = self.workspace_id.clone(); + let runtime_id = self.runtime_id.clone(); + let worker_id = self.worker_id.clone(); + let request_source_signer = self.request_source_signer.clone(); + let request_source_audience = self.request_source_audience.clone(); + let request_timeout = self.request_timeout; + if tokio::runtime::Handle::try_current().is_ok() { + std::thread::spawn(move || { + execute_runtime_owned_workspace_http( + &base_url, + &workspace_id, + &runtime_id, + &worker_id, + request_source_signer.as_ref(), + request_source_audience.as_deref(), + request_timeout, + permission, + request, + ) + }) + .join() + .map_err(|_| { + WorkspaceClientError::Request("workspace request thread panicked".to_string()) + })? + } else { + execute_runtime_owned_workspace_http( + &self.base_url, + &self.workspace_id, + &self.runtime_id, + &self.worker_id, + self.request_source_signer.as_ref(), + self.request_source_audience.as_deref(), + self.request_timeout, + permission, + request, + ) + } + } } impl std::fmt::Debug for RuntimeOwnedWorkspaceClient { @@ -377,42 +422,40 @@ impl WorkspaceClient for RuntimeOwnedWorkspaceClient { &self, request: WorkspaceRequest, ) -> Result { - let base_url = self.base_url.clone(); - let workspace_id = self.workspace_id.clone(); - let runtime_id = self.runtime_id.clone(); - let worker_id = self.worker_id.clone(); - let request_source_signer = self.request_source_signer.clone(); - let request_source_audience = self.request_source_audience.clone(); - let request_timeout = self.request_timeout; - if tokio::runtime::Handle::try_current().is_ok() { - std::thread::spawn(move || { - execute_runtime_owned_workspace_http( - &base_url, - &workspace_id, - &runtime_id, - &worker_id, - request_source_signer.as_ref(), - request_source_audience.as_deref(), - request_timeout, - request, - ) - }) - .join() - .map_err(|_| { - WorkspaceClientError::Request("workspace request thread panicked".to_string()) - })? - } else { - execute_runtime_owned_workspace_http( - &self.base_url, - &self.workspace_id, - &self.runtime_id, - &self.worker_id, - self.request_source_signer.as_ref(), - self.request_source_audience.as_deref(), - self.request_timeout, - request, - ) + self.execute_with_permission(request, WORKSPACE_REQUEST_PERMISSION) + } + + fn list_workspace_workers( + &self, + request: worker::WorkspaceWorkerDiscoveryRequest, + ) -> Result { + let mut path = format!( + "/api/w/{}/worker-discovery/workers?limit={}", + self.workspace_id, request.limit + ); + if let Some(cursor) = request.cursor.as_deref() { + path.push_str("&cursor="); + path.push_str(&percent_encode_query(cursor)); } + if let Some(query) = request.query.as_deref() { + path.push_str("&query="); + path.push_str(&percent_encode_query(query)); + } + let response = self.execute_with_permission( + WorkspaceRequest::get(path), + WORKSPACE_WORKER_DISCOVERY_PERMISSION, + )?; + if !(200..300).contains(&response.status) { + return Err(WorkspaceClientError::Request(format!( + "Workspace Worker discovery failed with HTTP {}: {}", + response.status, response.body + ))); + } + serde_json::from_str(&response.body).map_err(|error| { + WorkspaceClientError::Request(format!( + "invalid Workspace Worker discovery response: {error}" + )) + }) } fn current_prompt_projection( @@ -506,6 +549,19 @@ impl WorkspaceClient for RuntimeOwnedWorkspaceClient { } } +fn percent_encode_query(value: &str) -> String { + let mut encoded = String::with_capacity(value.len()); + for byte in value.bytes() { + if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'.' | b'_' | b'~') { + encoded.push(char::from(byte)); + } else { + use std::fmt::Write as _; + let _ = write!(encoded, "%{byte:02X}"); + } + } + encoded +} + fn execute_runtime_owned_workspace_http( base_url: &str, workspace_id: &str, @@ -514,6 +570,7 @@ fn execute_runtime_owned_workspace_http( request_source_signer: Option<&RuntimeRequestSourceSigner>, request_source_audience: Option<&str>, request_timeout: Option, + permission: &'static str, request: WorkspaceRequest, ) -> Result { if !request.path.starts_with('/') || request.path.starts_with("//") { @@ -553,7 +610,7 @@ fn execute_runtime_owned_workspace_http( audience, workspace_id, Some(worker_id), - WORKSPACE_REQUEST_PERMISSION, + permission, method.as_str(), &request.path, body.as_bytes(), @@ -903,6 +960,85 @@ mod tests { ); } + #[test] + fn workspace_worker_discovery_signs_dedicated_permission_and_encoded_query() { + use std::io::{Read, Write}; + use std::net::TcpListener; + use std::sync::Mutex; + + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let address = listener.local_addr().unwrap(); + let received = Arc::new(Mutex::new(String::new())); + let received_for_server = received.clone(); + let body = serde_json::json!({ + "workers": [{ + "subject": { + "kind": "runtime_worker", + "runtime_id": "runtime-b", + "worker_id": "worker-b" + }, + "resource_key": "W-2", + "display_name": "coder two", + "profile": "builtin:coder", + "status": "idle" + }], + "next_cursor": "v1:1" + }) + .to_string(); + let server = std::thread::spawn(move || { + let (mut stream, _) = listener.accept().unwrap(); + let mut bytes = [0_u8; 4096]; + let count = stream.read(&mut bytes).unwrap(); + *received_for_server.lock().unwrap() = + String::from_utf8_lossy(&bytes[..count]).into_owned(); + write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + body.len(), + body + ) + .unwrap(); + }); + + let identity = RuntimeIdentityMaterial::generate("runtime-a").unwrap(); + let client = RuntimeOwnedWorkspaceClient::new( + "workspace-a", + format!("http://{address}"), + "runtime-a", + "worker-a", + ) + .with_runtime_request_source(&identity, "server-a"); + let page = client + .list_workspace_workers(worker::WorkspaceWorkerDiscoveryRequest { + cursor: Some("v1:0".to_string()), + limit: 1, + query: Some("coder two".to_string()), + }) + .unwrap(); + assert_eq!(page.workers[0].resource_key, "W-2"); + server.join().unwrap(); + + let request = received.lock().unwrap().clone(); + assert!(request.contains( + "GET /api/w/workspace-a/worker-discovery/workers?limit=1&cursor=v1%3A0&query=coder%20two " + )); + let token = request + .lines() + .find_map(|line| { + line.split_once(':').and_then(|(name, value)| { + name.eq_ignore_ascii_case(RUNTIME_REQUEST_SOURCE_PROOF_HEADER) + .then(|| value.trim()) + }) + }) + .unwrap(); + let claims = decode_runtime_request_source_claims(token).unwrap(); + assert_eq!(claims.permission, WORKSPACE_WORKER_DISCOVERY_PERMISSION); + assert_eq!( + claims.path, + "/api/w/workspace-a/worker-discovery/workers?limit=1&cursor=v1%3A0&query=coder%20two" + ); + } + #[test] fn remote_authority_stamps_and_signs_worker_remove_without_caller_claim_choices() { let identity = RuntimeIdentityMaterial::generate("runtime-a").unwrap(); diff --git a/crates/worker/Cargo.toml b/crates/worker/Cargo.toml index bf951b6c..87437d17 100644 --- a/crates/worker/Cargo.toml +++ b/crates/worker/Cargo.toml @@ -27,6 +27,7 @@ toml = { workspace = true } tracing = { workspace = true } tools = { workspace = true } workdir = { workspace = true } +workspace-api = { workspace = true } minijinja = "2.19.0" chrono = "0.4" config-source = { path = "../config-source" } diff --git a/crates/worker/src/controller.rs b/crates/worker/src/controller.rs index ec4437db..ad7aa5b5 100644 --- a/crates/worker/src/controller.rs +++ b/crates/worker/src/controller.rs @@ -958,6 +958,23 @@ where crate::feature::builtin::manage_workdir::manage_workdir_feature(workspace_client), ); } + if feature_config.workspace_worker_discovery.enabled { + let workspace_client = worker.workspace_client_handle(); + let has_workspace_identity = workspace_client.workspace_id().is_some_and(|workspace_id| { + !workspace_id.is_empty() && !workspace_id.chars().any(char::is_control) + }); + if !workspace_client.is_available() || !has_workspace_identity { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "Workspace Worker discovery requires Backend Workspace API authority", + )); + } + feature_registry.add_module( + crate::feature::builtin::workspace_worker_discovery::workspace_worker_discovery_feature( + workspace_client, + ), + ); + } if feature_config.worker.enabled { let workspace_client = worker.workspace_client_handle(); let has_workspace_identity = workspace_client.workspace_id().is_some_and(|workspace_id| { diff --git a/crates/worker/src/feature/builtin.rs b/crates/worker/src/feature/builtin.rs index 010fcf96..2f1be0a7 100644 --- a/crates/worker/src/feature/builtin.rs +++ b/crates/worker/src/feature/builtin.rs @@ -17,6 +17,7 @@ pub mod session_explore; pub mod task; pub mod ticket; pub mod worker_observation; +pub mod workspace_worker_discovery; pub(crate) use memory_extract::{MemoryExtractFeature, MemoryExtractState, render_extract_input}; pub(crate) use session_explore::{SessionExploreFeature, SessionExploreState}; diff --git a/crates/worker/src/feature/builtin/workspace_worker_discovery.rs b/crates/worker/src/feature/builtin/workspace_worker_discovery.rs new file mode 100644 index 00000000..4d8ea268 --- /dev/null +++ b/crates/worker/src/feature/builtin/workspace_worker_discovery.rs @@ -0,0 +1,275 @@ +//! Privileged, read-only discovery of Workspace-visible Workers. +//! +//! This feature deliberately stays separate from the canonical `WorkerList` +//! control-grant surface. Discovery results carry the typed subject needed by a +//! later control operation, but discovery itself grants no control authority. + +use std::sync::Arc; + +use agen::tool::{Tool, ToolDefinition, ToolError, ToolExecutionContext, ToolMeta, ToolOutput}; +use async_trait::async_trait; +use serde::Deserialize; +use serde_json::json; + +use crate::feature::{ + FeatureDescriptor, FeatureInstallContext, FeatureInstallError, FeatureInstructionContribution, + FeatureInstructionDeclaration, FeatureInstructionId, FeatureModule, ToolContribution, + ToolDeclaration, +}; +use crate::worker::{WorkspaceClient, WorkspaceWorkerDiscoveryRequest}; + +const FEATURE_ID: &str = "workspace-worker-discovery"; +const TOOL_NAME: &str = "ListWorkspaceWorkers"; +const DEFAULT_LIMIT: usize = 50; +const MAX_LIMIT: usize = 100; +const INSTRUCTION_ID: &str = "workspace-worker-discovery.policy"; +const PROMPT_REF: &str = "common.workspace_worker_discovery"; +const DESCRIPTION: &str = "List or directly find Workspace-visible Workers through Backend authority. Results include each W-key and the typed runtime_worker subject needed by later Worker control calls, but do not grant control authority."; + +fn instruction() -> FeatureInstructionDeclaration { + FeatureInstructionDeclaration::new( + FeatureInstructionId::builtin(INSTRUCTION_ID), + PROMPT_REF, + "Workspace Worker discovery and control-authority separation", + ) + .expect("static Workspace Worker discovery instruction is valid") +} + +#[derive(Clone)] +pub struct WorkspaceWorkerDiscoveryFeature { + client: Arc, +} + +pub fn workspace_worker_discovery_feature( + client: Arc, +) -> WorkspaceWorkerDiscoveryFeature { + WorkspaceWorkerDiscoveryFeature { client } +} + +impl FeatureModule for WorkspaceWorkerDiscoveryFeature { + fn descriptor(&self) -> FeatureDescriptor { + FeatureDescriptor::builtin(FEATURE_ID, "Workspace Worker Discovery") + .with_description(DESCRIPTION) + .with_instruction(instruction()) + .with_tool(ToolDeclaration::new(TOOL_NAME, DESCRIPTION)) + } + + fn install(&self, context: &mut FeatureInstallContext<'_>) -> Result<(), FeatureInstallError> { + context + .instructions() + .register(FeatureInstructionContribution::new(instruction()))?; + let client = self.client.clone(); + let definition: ToolDefinition = Arc::new(move || { + ( + ToolMeta::new(TOOL_NAME) + .description(DESCRIPTION) + .input_schema(input_schema()), + Arc::new(ListWorkspaceWorkersTool { + client: client.clone(), + }) as Arc, + ) + }); + context + .tools() + .register(ToolContribution::new(TOOL_NAME, definition))?; + Ok(()) + } +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct ListWorkspaceWorkersInput { + #[serde(default)] + cursor: Option, + #[serde(default)] + limit: Option, + #[serde(default)] + query: Option, +} + +#[derive(Clone)] +struct ListWorkspaceWorkersTool { + client: Arc, +} + +#[async_trait] +impl Tool for ListWorkspaceWorkersTool { + async fn execute( + &self, + input_json: &str, + _ctx: ToolExecutionContext, + ) -> Result { + let input: ListWorkspaceWorkersInput = serde_json::from_str(input_json) + .map_err(|error| ToolError::InvalidArgument(error.to_string()))?; + let limit = input.limit.unwrap_or(DEFAULT_LIMIT); + if !(1..=MAX_LIMIT).contains(&limit) { + return Err(ToolError::InvalidArgument(format!( + "limit must be between 1 and {MAX_LIMIT}" + ))); + } + let query = input.query.map(|query| query.trim().to_string()); + if query.as_deref().is_some_and(str::is_empty) { + return Err(ToolError::InvalidArgument( + "query must not be empty".to_string(), + )); + } + if query.as_ref().is_some_and(|query| query.len() > 128) { + return Err(ToolError::InvalidArgument( + "query must not exceed 128 bytes".to_string(), + )); + } + + let page = self + .client + .list_workspace_workers(WorkspaceWorkerDiscoveryRequest { + cursor: input.cursor, + limit, + query, + }) + .map_err(|error| ToolError::ExecutionFailed(error.to_string()))?; + let count = page.workers.len(); + Ok(ToolOutput { + summary: format!("Listed {count} Workspace Worker(s)"), + content: Some(serde_json::to_string_pretty(&page).map_err(|error| { + ToolError::ExecutionFailed(format!( + "encode Workspace Worker discovery result: {error}" + )) + })?), + attachments: Vec::new(), + }) + } +} + +fn input_schema() -> serde_json::Value { + json!({ + "type": "object", + "additionalProperties": false, + "properties": { + "cursor": { + "type": "string", + "description": "Opaque cursor returned by a prior page." + }, + "limit": { + "type": "integer", + "minimum": 1, + "maximum": MAX_LIMIT, + "default": DEFAULT_LIMIT + }, + "query": { + "type": "string", + "maxLength": 128, + "description": "Exact W-key or Worker display name lookup." + } + } + }) +} + +#[cfg(test)] +mod tests { + use std::sync::Mutex; + + use workspace_api::{ + WorkspaceWorkerDiscoveryItem, WorkspaceWorkerDiscoveryPage, WorkspaceWorkerSubject, + }; + + use super::*; + use crate::worker::{WorkspaceClientError, WorkspaceRequest, WorkspaceResponse}; + + #[derive(Debug)] + struct RecordingClient { + requests: Mutex>, + result: WorkspaceWorkerDiscoveryPage, + unavailable: bool, + } + + impl WorkspaceClient for RecordingClient { + fn workspace_id(&self) -> Option<&str> { + Some("workspace-1") + } + + fn kind(&self) -> &str { + "recording" + } + + fn is_available(&self) -> bool { + !self.unavailable + } + + fn execute( + &self, + _request: WorkspaceRequest, + ) -> Result { + panic!("discovery must not use generic Workspace request authority") + } + + fn list_workspace_workers( + &self, + request: WorkspaceWorkerDiscoveryRequest, + ) -> Result { + self.requests.lock().unwrap().push(request); + if self.unavailable { + Err(WorkspaceClientError::Unavailable("denied".to_string())) + } else { + Ok(self.result.clone()) + } + } + } + + fn page() -> WorkspaceWorkerDiscoveryPage { + WorkspaceWorkerDiscoveryPage { + workers: vec![WorkspaceWorkerDiscoveryItem { + subject: WorkspaceWorkerSubject::RuntimeWorker { + runtime_id: "arcadia".to_string(), + worker_id: "worker-1".to_string(), + }, + resource_key: "W-12".to_string(), + display_name: "coder-one".to_string(), + profile: Some("builtin:coder".to_string()), + status: Some("idle".to_string()), + }], + next_cursor: Some("v1:1".to_string()), + } + } + + #[tokio::test] + async fn tool_preserves_typed_subject_and_forwards_lookup() { + let client = Arc::new(RecordingClient { + requests: Mutex::new(Vec::new()), + result: page(), + unavailable: false, + }); + let tool = ListWorkspaceWorkersTool { + client: client.clone(), + }; + let output = tool + .execute( + r#"{"query":" W-12 ","limit":1}"#, + ToolExecutionContext::default(), + ) + .await + .unwrap(); + let value: serde_json::Value = + serde_json::from_str(output.content.as_deref().unwrap()).unwrap(); + assert_eq!(value["workers"][0]["resource_key"], "W-12"); + assert_eq!(value["workers"][0]["subject"]["kind"], "runtime_worker"); + assert_eq!(value["workers"][0]["subject"]["runtime_id"], "arcadia"); + let requests = client.requests.lock().unwrap(); + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].query.as_deref(), Some("W-12")); + assert_eq!(requests[0].limit, 1); + } + + #[tokio::test] + async fn missing_backend_authority_fails_closed() { + let client = Arc::new(RecordingClient { + requests: Mutex::new(Vec::new()), + result: page(), + unavailable: true, + }); + let error = ListWorkspaceWorkersTool { client } + .execute("{}", ToolExecutionContext::default()) + .await + .unwrap_err(); + assert!(error.to_string().contains("denied")); + } +} diff --git a/crates/worker/src/lib.rs b/crates/worker/src/lib.rs index 8de12fbd..20042594 100644 --- a/crates/worker/src/lib.rs +++ b/crates/worker/src/lib.rs @@ -52,6 +52,6 @@ pub use worker::{ LocalWorkingDirectory, WORKER_INPUT_SUBMISSION_EXTENSION_DOMAIN, Worker, WorkerError, WorkerFilesystemAuthority, WorkerRunResult, WorkerWorkspaceContext, WorkspaceClient, WorkspaceClientError, WorkspaceId, WorkspaceIdError, WorkspacePromptCatalogResolution, - WorkspaceRequest, WorkspaceRequestMethod, WorkspaceResponse, apply_worker_manifest, - marker_workspace_client, unavailable_workspace_client, + WorkspaceRequest, WorkspaceRequestMethod, WorkspaceResponse, WorkspaceWorkerDiscoveryRequest, + apply_worker_manifest, marker_workspace_client, unavailable_workspace_client, }; diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 80dd85b2..bd7e5f22 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -251,6 +251,13 @@ impl WorkspacePromptCatalogResolution { } } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct WorkspaceWorkerDiscoveryRequest { + pub cursor: Option, + pub limit: usize, + pub query: Option, +} + /// Path-free Workspace operation authority injected by Runtime/host code. /// /// Workers receive this trait object rather than a Backend URL. The concrete @@ -263,6 +270,18 @@ pub trait WorkspaceClient: std::fmt::Debug + Send + Sync { fn execute(&self, request: WorkspaceRequest) -> Result; + /// Lists Workspace-visible Workers through dedicated Runtime-owned source + /// proof. Implementations must not fall back to generic Workspace request + /// authority. + fn list_workspace_workers( + &self, + _request: WorkspaceWorkerDiscoveryRequest, + ) -> Result { + Err(WorkspaceClientError::Unavailable( + "Workspace Worker discovery authority is unavailable".to_string(), + )) + } + /// Resolve the Workspace's current immutable Prompt projection for future /// operation boundaries. Creation and restore continue to use persisted /// launch/session state; this hook never reconstructs historical prompts. diff --git a/crates/worker/tests/controller_test.rs b/crates/worker/tests/controller_test.rs index 9a984657..364440b2 100644 --- a/crates/worker/tests/controller_test.rs +++ b/crates/worker/tests/controller_test.rs @@ -180,6 +180,30 @@ async fn make_worker_with_pwd( make_worker_with_pwd_and_manifest(client, MANIFEST_TOML).await } +#[derive(Debug)] +struct AvailableWorkspaceClient; + +impl WorkspaceClient for AvailableWorkspaceClient { + fn workspace_id(&self) -> Option<&str> { + Some("workspace-1") + } + + fn kind(&self) -> &str { + "test" + } + + fn is_available(&self) -> bool { + true + } + + fn execute( + &self, + _request: WorkspaceRequest, + ) -> Result { + Err(WorkspaceClientError::Unavailable("test".to_string())) + } +} + #[derive(Debug)] struct NoopWorkspaceClient; @@ -776,6 +800,45 @@ permission = "write" } } +#[tokio::test] +async fn workspace_worker_discovery_requires_workspace_authority_and_stays_separate_from_worker_list() + { + let manifest_toml = r#" +[worker] +name = "workspace-worker-discovery-test" +pwd = "./" + +[model] +scheme = "anthropic" +model_id = "test-model" + +[engine] +max_tokens = 100 + +[feature.workspace_worker_discovery] +enabled = true + +[[scope.allow]] +target = "./" +permission = "write" +"#; + let client = MockClient::new(simple_text_events()); + let client_for_assert = client.clone(); + let (worker, _pwd) = make_worker_with_pwd_manifest_and_workspace_context( + client, + manifest_toml, + WorkerWorkspaceContext::with_client(None, Arc::new(AvailableWorkspaceClient)), + ) + .await; + let handle = spawn_controller(worker).await; + handle.send(Method::run_text("Hello")).await.unwrap(); + wait_for_status(&handle, WorkerStatus::Idle).await; + let request = wait_for_captured_request(&client_for_assert).await; + let names = request_tool_names(&request); + assert!(names.iter().any(|name| name == "ListWorkspaceWorkers")); + assert!(!names.iter().any(|name| name == "WorkerList")); +} + #[tokio::test] async fn sub_worker_feature_exposure_does_not_require_delegation_scope() { let manifest = r#" diff --git a/crates/workspace-api/src/lib.rs b/crates/workspace-api/src/lib.rs index 99ce7440..79ffb565 100644 --- a/crates/workspace-api/src/lib.rs +++ b/crates/workspace-api/src/lib.rs @@ -283,6 +283,35 @@ pub struct WorkerCapabilitySummary { pub can_spawn_followup: bool, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum WorkspaceWorkerSubject { + RuntimeWorker { + runtime_id: String, + worker_id: String, + }, +} + +/// Bounded, model-safe projection used by privileged Workspace Worker discovery. +/// Runtime placement appears only in the typed subject required by Worker +/// control operations; provider and launch internals are intentionally omitted. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkspaceWorkerDiscoveryItem { + pub subject: WorkspaceWorkerSubject, + pub resource_key: String, + pub display_name: String, + pub profile: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub status: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkspaceWorkerDiscoveryPage { + pub workers: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub next_cursor: Option, +} + /// Workspace-authoritative Worker projection. /// /// `resource_key` is required here even though Runtime-internal Worker summaries diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 8fb9cd5d..bc7a73d2 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -4450,7 +4450,9 @@ mod tests { .unwrap(); assert!(companion.feature.manage_workdir.enabled); assert!(companion.feature.sub_worker.enabled); - assert!(!companion.feature.worker.enabled); + assert!(companion.feature.worker.enabled); + assert!(!companion.feature.worker.direct_spawn); + assert!(companion.feature.workspace_worker_discovery.enabled); let coder = archive .resolve_profile("builtin:coder", root.path(), "embedded-test-coder") .unwrap(); diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 38b198c5..231a472b 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -65,7 +65,8 @@ use workspace_api::{ ObjectiveLinkTicketRequest, ObjectiveStateRequest, PutRepositorySshHostTrustRequest, RepositoryAccessProjection, RepositorySshCredential, RepositorySshHostTrust, RotateRepositorySshCredentialRequest, TICKET_ORCHESTRATION_PLANS_QUERY_PATH, - TICKET_RELATIONS_QUERY_PATH, + TICKET_RELATIONS_QUERY_PATH, WorkspaceWorkerDiscoveryItem, WorkspaceWorkerDiscoveryPage, + WorkspaceWorkerSubject, }; use crate::auth::{ @@ -1063,7 +1064,13 @@ async fn authorize_scoped_workspace_request( .map_err(|_| StatusCode::BAD_REQUEST.into_response())?; let digest = worker_runtime::auth::request_body_digest(&body); *request.body_mut() = axum::body::Body::from(body); - let permission = if path.starts_with("/api/runtime/v1/workspaces/") + let permission = if path + .split('?') + .next() + .is_some_and(|path| path.ends_with("/worker-discovery/workers")) + { + worker_runtime::auth::WORKSPACE_WORKER_DISCOVERY_PERMISSION + } else if path.starts_with("/api/runtime/v1/workspaces/") || path.contains("/profile-source-archives/") { worker_runtime::auth::BACKEND_RESOURCE_FETCH_PERMISSION @@ -1135,7 +1142,13 @@ async fn authorize_workspace_api_request( }; let digest = worker_runtime::auth::request_body_digest(&body); *request.body_mut() = axum::body::Body::from(body); - let permission = if path.starts_with("/api/runtime/v1/workspaces/") + let permission = if path + .split('?') + .next() + .is_some_and(|path| path.ends_with("/worker-discovery/workers")) + { + worker_runtime::auth::WORKSPACE_WORKER_DISCOVERY_PERMISSION + } else if path.starts_with("/api/runtime/v1/workspaces/") || path.contains("/profile-source-archives/") { worker_runtime::auth::BACKEND_RESOURCE_FETCH_PERMISSION @@ -2468,6 +2481,10 @@ fn build_inner_router(api: WorkspaceApi) -> Router { "/api/w/{workspace_id}/worker-observation/session", post(scoped_capture_worker_observation_session), ) + .route( + "/api/w/{workspace_id}/worker-discovery/workers", + get(scoped_discover_workspace_workers), + ) .route( "/api/w/{workspace_id}/workers", get(scoped_list_workers).post(scoped_create_workspace_worker), @@ -8218,6 +8235,144 @@ async fn scoped_get_workspace_worker( }) } +#[derive(Debug, Deserialize)] +struct WorkspaceWorkerDiscoveryQuery { + cursor: Option, + limit: Option, + query: Option, +} + +async fn scoped_discover_workspace_workers( + State(api): State, + AxumPath(path): AxumPath, + Query(query): Query, + source: Option>, +) -> Response { + let Some(Extension(_source)) = source else { + return ( + StatusCode::FORBIDDEN, + Json(serde_json::json!({ + "error": "workspace_worker_discovery_source_required" + })), + ) + .into_response(); + }; + + if let Err(error) = validate_workspace_scope(&api, &path.workspace_id) { + return error.into_response(); + } + let limit = query.limit.unwrap_or(50).clamp(1, 100); + let cursor = query.cursor.as_deref(); + let query = match query.query.as_deref().map(str::trim) { + Some("") | None => None, + Some(query) if query.len() <= 128 => Some(query), + Some(_) => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({ + "error": "invalid_workspace_worker_discovery_query" + })), + ) + .into_response(); + } + }; + let filter_fingerprint = workspace_worker_discovery_filter_fingerprint(query); + let offset = match query_cursor_offset(cursor, filter_fingerprint) { + Ok(offset) => offset, + Err(()) => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({ + "error": "invalid_workspace_worker_discovery_cursor" + })), + ) + .into_response(); + } + }; + + let workers = match workers_response(api) { + Ok(workers) => workers, + Err(error) => return error.into_response(), + }; + Json(workspace_worker_discovery_page( + workers.items, + query, + offset, + limit, + filter_fingerprint, + )) + .into_response() +} + +fn workspace_worker_discovery_page( + workers: Vec, + query: Option<&str>, + offset: usize, + limit: usize, + filter_fingerprint: u64, +) -> WorkspaceWorkerDiscoveryPage { + let mut workers = workers + .into_iter() + .filter(|worker| worker.workspace.visibility == "workspace_scoped") + .filter(|worker| { + query.is_none_or(|query| { + worker.resource_key == query + || worker.display_name == query + || worker.label == query + }) + }) + .collect::>(); + workers.sort_by(|left, right| left.resource_key.cmp(&right.resource_key)); + let total = workers.len(); + let page = workers + .into_iter() + .skip(offset) + .take(limit) + .map(|worker| WorkspaceWorkerDiscoveryItem { + subject: WorkspaceWorkerSubject::RuntimeWorker { + runtime_id: worker.runtime_id, + worker_id: worker.worker_id, + }, + resource_key: worker.resource_key, + display_name: worker.display_name, + profile: worker.profile, + status: Some(worker.state), + }) + .collect::>(); + let next_offset = offset.saturating_add(page.len()); + WorkspaceWorkerDiscoveryPage { + workers: page, + next_cursor: (next_offset < total) + .then(|| format!("v1:{next_offset}:{filter_fingerprint:016x}")), + } +} + +fn query_cursor_offset( + cursor: Option<&str>, + filter_fingerprint: u64, +) -> std::result::Result { + let Some(cursor) = cursor else { + return Ok(0); + }; + let mut parts = cursor.split(':'); + match (parts.next(), parts.next(), parts.next(), parts.next()) { + (Some("v1"), Some(offset), Some(fingerprint), None) + if u64::from_str_radix(fingerprint, 16).ok() == Some(filter_fingerprint) => + { + offset.parse().map_err(|_| ()) + } + _ => Err(()), + } +} + +fn workspace_worker_discovery_filter_fingerprint(query: Option<&str>) -> u64 { + use std::hash::{Hash, Hasher}; + + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + query.hash(&mut hasher); + hasher.finish() +} + async fn scoped_list_workers( State(api): State, AxumPath(path): AxumPath, @@ -22316,6 +22471,171 @@ mod tests { assert_eq!(valid.status(), StatusCode::OK); } + #[test] + fn workspace_worker_discovery_filters_internal_workers_and_paginates_typed_subjects() { + fn summary( + key: &str, + worker_id: &str, + name: &str, + visibility: &str, + ) -> workspace_api::WorkerSummary { + workspace_api::WorkerSummary { + runtime_id: "runtime-test".to_string(), + worker_id: worker_id.to_string(), + resource_key: key.to_string(), + host_id: "host-test".to_string(), + display_name: name.to_string(), + label: name.to_string(), + profile: Some("builtin:coder".to_string()), + singleton_key: None, + tags: Vec::new(), + workspace: workspace_api::WorkerWorkspaceSummary { + visibility: visibility.to_string(), + identity: "workspace-test".to_string(), + workspace_id: Some(TEST_WORKSPACE_ID.to_string()), + }, + state: "idle".to_string(), + last_seen_at: None, + pinned: false, + retention_state: "normal".to_string(), + implementation: workspace_api::WorkerImplementationSummary { + kind: "remote".to_string(), + display_hint: "remote".to_string(), + }, + capabilities: workspace_api::WorkerCapabilitySummary { + can_stop: true, + can_spawn_followup: false, + }, + working_directory: None, + diagnostics: Vec::new(), + } + } + + let fingerprint = workspace_worker_discovery_filter_fingerprint(None); + let first = workspace_worker_discovery_page( + vec![ + summary("W-3", "internal", "service", "backend_internal"), + summary("W-2", "worker-b", "beta", "workspace_scoped"), + summary("W-1", "worker-a", "alpha", "workspace_scoped"), + ], + None, + 0, + 1, + fingerprint, + ); + assert_eq!(first.workers.len(), 1); + assert_eq!(first.workers[0].resource_key, "W-1"); + assert_eq!( + first.workers[0].subject, + WorkspaceWorkerSubject::RuntimeWorker { + runtime_id: "runtime-test".to_string(), + worker_id: "worker-a".to_string(), + } + ); + let cursor = first.next_cursor.as_deref().unwrap(); + let offset = query_cursor_offset(Some(cursor), fingerprint).unwrap(); + let second = workspace_worker_discovery_page( + vec![ + summary("W-3", "internal", "service", "backend_internal"), + summary("W-2", "worker-b", "beta", "workspace_scoped"), + summary("W-1", "worker-a", "alpha", "workspace_scoped"), + ], + None, + offset, + 1, + fingerprint, + ); + assert_eq!(second.workers[0].resource_key, "W-2"); + assert!(second.next_cursor.is_none()); + + let lookup_fingerprint = workspace_worker_discovery_filter_fingerprint(Some("beta")); + let lookup = workspace_worker_discovery_page( + vec![summary("W-2", "worker-b", "beta", "workspace_scoped")], + Some("beta"), + 0, + 10, + lookup_fingerprint, + ); + assert_eq!(lookup.workers[0].resource_key, "W-2"); + assert!(query_cursor_offset(Some(cursor), lookup_fingerprint).is_err()); + } + + #[tokio::test] + async fn workspace_worker_discovery_requires_dedicated_source_proof() { + let workspace = tempfile::tempdir().unwrap(); + let mut api = test_api(workspace.path()).await; + let identity = + worker_runtime::auth::RuntimeIdentityMaterial::generate("runtime-test").unwrap(); + configure_runtime_request_auth(&mut api, &identity, "runtime-test"); + seed_worker_source_member(&api, "runtime-test", "worker-caller"); + seed_worker_source_member(&api, "runtime-test", "worker-target"); + let target_key = api + .store + .resource_key( + TEST_WORKSPACE_ID, + WorkspaceResourceKind::Worker, + "worker-target", + ) + .unwrap() + .unwrap(); + let target = format!( + "/api/w/{TEST_WORKSPACE_ID}/worker-discovery/workers?limit=1&query={target_key}" + ); + let signer = worker_runtime::auth::RuntimeRequestSourceSigner::from_identity(&identity); + let issue = |permission| { + signer + .issue( + "server-test", + TEST_WORKSPACE_ID, + Some("worker-caller"), + permission, + "GET", + &target, + b"", + i64::try_from(worker_runtime::auth::unix_now_seconds()).unwrap_or(i64::MAX), + 30, + ) + .unwrap() + }; + let app = build_router(api); + + let generic = app + .clone() + .oneshot( + Request::builder() + .uri(&target) + .header( + worker_runtime::auth::RUNTIME_REQUEST_SOURCE_PROOF_HEADER, + issue(worker_runtime::auth::WORKSPACE_REQUEST_PERMISSION), + ) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(generic.status(), StatusCode::UNAUTHORIZED); + + let response = app + .oneshot( + Request::builder() + .uri(&target) + .header( + worker_runtime::auth::RUNTIME_REQUEST_SOURCE_PROOF_HEADER, + issue(worker_runtime::auth::WORKSPACE_WORKER_DISCOVERY_PERMISSION), + ) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body: WorkspaceWorkerDiscoveryPage = + serde_json::from_slice(&to_bytes(response.into_body(), usize::MAX).await.unwrap()) + .unwrap(); + assert!(body.workers.is_empty()); + assert!(body.next_cursor.is_none()); + } + #[tokio::test] async fn internal_resource_fetch_rest_returns_typed_missing_resource() { let workspace = tempfile::tempdir().unwrap(); diff --git a/resources/profiles/base.dcdl b/resources/profiles/base.dcdl index 2e960d15..c91afed8 100644 --- a/resources/profiles/base.dcdl +++ b/resources/profiles/base.dcdl @@ -28,6 +28,7 @@ feature = { image = { enabled = true; }; sub_worker = { enabled = false; }; worker = { enabled = false; }; + workspace_worker_discovery = { enabled = false; }; objective = { enabled = true; }; ticket = { enabled = true; authoring = true; thread = true; }; merge_request = { diff --git a/resources/profiles/companion.dcdl b/resources/profiles/companion.dcdl index d3720e3a..85f4fbee 100644 --- a/resources/profiles/companion.dcdl +++ b/resources/profiles/companion.dcdl @@ -8,6 +8,8 @@ import "./base.dcdl" // { memory = { enabled = true; }; web = { enabled = true; }; sub_worker = { enabled = true; }; + worker = { enabled = true; direct_spawn = false; }; + workspace_worker_discovery = { enabled = true; }; manage_workdir = { enabled = true; }; ticket = { enabled = true; authoring = true; thread = true; }; }; diff --git a/resources/prompts/catalog.dcdl b/resources/prompts/catalog.dcdl index caa82038..0f7a2a90 100644 --- a/resources/prompts/catalog.dcdl +++ b/resources/prompts/catalog.dcdl @@ -10,6 +10,7 @@ commonMergeRequest = import "./common/merge-request.md"; commonTickets = import "./common/tickets.md"; commonToolUsage = import "./common/tool-usage.md"; commonWorkerObservation = import "./common/worker-observation.md"; +commonWorkspaceWorkerDiscovery = import "./common/workspace-worker-discovery.md"; commonWorkerOrchestration = import "./common/worker-orchestration.md"; commonWorkspace = import "./common/workspace.md"; commonWriting = import "./common/writing.md"; @@ -40,6 +41,7 @@ in tickets = commonTickets.content; tool_usage = commonToolUsage.content; worker_observation = commonWorkerObservation.content; + workspace_worker_discovery = commonWorkspaceWorkerDiscovery.content; worker_orchestration = commonWorkerOrchestration.content; workspace = commonWorkspace.content; writing = commonWriting.content; diff --git a/resources/prompts/common/workspace-worker-discovery.md b/resources/prompts/common/workspace-worker-discovery.md new file mode 100644 index 00000000..e4f43743 --- /dev/null +++ b/resources/prompts/common/workspace-worker-discovery.md @@ -0,0 +1,6 @@ +Workspace Worker discovery is separate from Worker control authority. + +- Use `ListWorkspaceWorkers` to list accessible Workspace Workers or directly find one by its exact `W-*` key or display name. +- Reuse the returned typed `subject` unchanged when a later `Worker*` control tool requires a target. Do not guess `runtime_id` or `worker_id`. +- Discovery does not grant control. `WorkerList` remains the authoritative list of Workers this Worker may control, and a discovered Worker can still be rejected by every control operation. +- Results exclude service-private/Internal Workers and omit provider, launch, credential, and capability internals.