diff --git a/crates/manifest/src/config.rs b/crates/manifest/src/config.rs index d488cdb3..b71520ac 100644 --- a/crates/manifest/src/config.rs +++ b/crates/manifest/src/config.rs @@ -87,6 +87,8 @@ pub struct FeatureConfigPartial { #[serde(default)] pub objective: Option, #[serde(default)] + pub manage_workdir: Option, + #[serde(default)] pub ticket: Option, #[serde(default)] pub plugins: Option, @@ -104,6 +106,11 @@ impl FeatureConfigPartial { other.objective, FeatureFlagConfigPartial::merge, ), + manage_workdir: merge_option( + self.manage_workdir, + other.manage_workdir, + FeatureFlagConfigPartial::merge, + ), ticket: merge_option(self.ticket, other.ticket, TicketFeatureConfigPartial::merge), plugins: merge_option(self.plugins, other.plugins, FeatureFlagConfigPartial::merge), } @@ -180,6 +187,10 @@ impl From for FeatureConfig { .objective .map(FeatureFlagConfig::from) .unwrap_or_default(), + manage_workdir: value + .manage_workdir + .map(FeatureFlagConfig::from) + .unwrap_or_default(), ticket: value .ticket .map(TicketFeatureConfig::from) @@ -258,6 +269,7 @@ impl From for FeatureConfigPartial { web: Some(value.web.into()), workers: Some(value.workers.into()), objective: Some(value.objective.into()), + manage_workdir: Some(value.manage_workdir.into()), ticket: Some(value.ticket.into()), plugins: Some(value.plugins.into()), } @@ -1804,6 +1816,7 @@ worker_max_turns = 7 assert!(!manifest.feature.web.enabled); assert!(!manifest.feature.workers.enabled); assert!(!manifest.feature.objective.enabled); + assert!(!manifest.feature.manage_workdir.enabled); assert!(!manifest.feature.ticket.enabled); } @@ -1814,6 +1827,9 @@ worker_max_turns = 7 [feature.task] enabled = true +[feature.manage_workdir] +enabled = true + [feature.ticket] enabled = true authoring = false @@ -1848,6 +1864,7 @@ orchestration_control = false .try_into() .unwrap(); assert!(manifest.feature.task.enabled); + assert!(manifest.feature.manage_workdir.enabled); assert!(manifest.feature.ticket.enabled); assert!(!manifest.feature.ticket.authoring); assert!(!manifest.feature.ticket.thread); @@ -1865,6 +1882,9 @@ orchestration_control = false [feature.memory] enabled = true +[feature.manage_workdir] +enabled = false + [feature.ticket] enabled = true authoring = false @@ -1883,6 +1903,9 @@ orchestration_control = true [feature.memory] staging = true +[feature.manage_workdir] +enabled = true + [feature.objective] enabled = true @@ -1918,6 +1941,7 @@ enabled = true .unwrap(); assert!(manifest.feature.memory.enabled); assert!(manifest.feature.memory.staging); + assert!(manifest.feature.manage_workdir.enabled); assert!(manifest.feature.ticket.enabled); assert!(!manifest.feature.ticket.authoring); assert!(manifest.feature.ticket.thread); diff --git a/crates/manifest/src/lib.rs b/crates/manifest/src/lib.rs index 1db252aa..54d16032 100644 --- a/crates/manifest/src/lib.rs +++ b/crates/manifest/src/lib.rs @@ -115,6 +115,8 @@ pub struct FeatureConfig { #[serde(default)] pub objective: FeatureFlagConfig, #[serde(default)] + pub manage_workdir: FeatureFlagConfig, + #[serde(default)] pub ticket: TicketFeatureConfig, #[serde(default)] pub plugins: FeatureFlagConfig, @@ -128,6 +130,7 @@ impl Default for FeatureConfig { web: FeatureFlagConfig::disabled(), workers: FeatureFlagConfig::disabled(), objective: FeatureFlagConfig::disabled(), + manage_workdir: FeatureFlagConfig::disabled(), ticket: TicketFeatureConfig::default(), plugins: FeatureFlagConfig::disabled(), } diff --git a/crates/manifest/src/profile.rs b/crates/manifest/src/profile.rs index f107a42a..9a30d6b2 100644 --- a/crates/manifest/src/profile.rs +++ b/crates/manifest/src/profile.rs @@ -998,6 +998,7 @@ fn apply_role_profile( value["feature"]["memory"] = serde_json::json!({ "enabled": memory }); value["feature"]["web"] = serde_json::json!({ "enabled": web }); value["feature"]["workers"] = serde_json::json!({ "enabled": workers }); + value["feature"]["manage_workdir"] = serde_json::json!({ "enabled": slug == "orchestrator" }); let ticket = match slug { "companion" => serde_json::json!({ "enabled": true, "authoring": true, "thread": true }), "intake" => { @@ -1492,6 +1493,7 @@ mod tests { assert!(intake.feature.ticket.authoring); assert!(intake.feature.ticket.thread); assert!(intake.feature.objective.enabled); + assert!(!intake.feature.manage_workdir.enabled); assert!(intake.feature.ticket.intake); assert!(!intake.feature.ticket.orchestration_control); assert!(intake.scope.allow.is_empty()); @@ -1508,6 +1510,7 @@ mod tests { assert!(!orchestrator.feature.ticket.authoring); assert!(orchestrator.feature.ticket.thread); assert!(orchestrator.feature.objective.enabled); + assert!(orchestrator.feature.manage_workdir.enabled); assert!(!orchestrator.feature.ticket.intake); assert!(orchestrator.feature.ticket.orchestration_control); assert!(orchestrator.scope.allow.is_empty()); @@ -1532,6 +1535,7 @@ mod tests { assert!(!coder.feature.ticket.authoring); assert!(coder.feature.ticket.thread); assert!(coder.feature.objective.enabled); + assert!(!coder.feature.manage_workdir.enabled); assert!(!coder.feature.ticket.intake); assert!(!coder.feature.ticket.orchestration_control); let reviewer = resolve("reviewer"); @@ -1542,6 +1546,7 @@ mod tests { assert!(!reviewer.feature.ticket.authoring); assert!(reviewer.feature.ticket.thread); assert!(reviewer.feature.objective.enabled); + assert!(!reviewer.feature.manage_workdir.enabled); assert!(!reviewer.feature.ticket.intake); assert!(!reviewer.feature.ticket.orchestration_control); assert!(reviewer.scope.allow.is_empty()); diff --git a/crates/worker/src/controller.rs b/crates/worker/src/controller.rs index 426ceafe..ca4e49f9 100644 --- a/crates/worker/src/controller.rs +++ b/crates/worker/src/controller.rs @@ -665,6 +665,24 @@ where ), ); } + if feature_config.manage_workdir.enabled { + // Workdir lifecycle is Workspace control-plane authority. The Worker + // receives only the injected WorkspaceClient and never Runtime URLs, + // repository paths, materializer handles, or cleanup sessions. + 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, + "manage Workdir tools require Backend Workspace API authority", + )); + } + feature_registry.add_module( + crate::feature::builtin::manage_workdir::manage_workdir_feature(workspace_client), + ); + } for module in crate::feature::plugin::plugin_tool_features_if_enabled( feature_config.plugins.enabled, &worker.manifest().plugins, diff --git a/crates/worker/src/feature/builtin.rs b/crates/worker/src/feature/builtin.rs index 1cc2f109..45057c0d 100644 --- a/crates/worker/src/feature/builtin.rs +++ b/crates/worker/src/feature/builtin.rs @@ -4,6 +4,7 @@ //! same descriptor-approved registry path used by feature modules. They are not //! an external plugin-loading surface. +pub mod manage_workdir; pub mod memory; pub mod objective; pub mod session_explore; diff --git a/crates/worker/src/feature/builtin/manage_workdir.rs b/crates/worker/src/feature/builtin/manage_workdir.rs new file mode 100644 index 00000000..07eb0b11 --- /dev/null +++ b/crates/worker/src/feature/builtin/manage_workdir.rs @@ -0,0 +1,670 @@ +//! Workspace-authority Workdir lifecycle tools for embedded Orchestrators. +//! +//! The model-visible contract contains only stable Workspace, Runtime, +//! repository, selector, and Workdir identities. Repository paths, Runtime +//! endpoints, credentials, materializer handles, and operation sessions stay +//! behind [`WorkspaceClient`]. + +use std::sync::Arc; + +use async_trait::async_trait; +use llm_engine::tool::{ + Tool, ToolDefinition, ToolError, ToolExecutionContext, ToolMeta, ToolOutput, +}; +use serde::{Deserialize, Serialize}; +use serde_json::json; + +use crate::feature::{ + FeatureDescriptor, FeatureInstallContext, FeatureInstallError, FeatureModule, ToolContribution, + ToolDeclaration, +}; +use crate::worker::{WorkspaceClient, WorkspaceRequest, WorkspaceRequestMethod, WorkspaceResponse}; + +const FEATURE_ID: &str = "manage-workdir"; +const FEATURE_NAME: &str = "Manage Workdir"; +const FEATURE_DESCRIPTION: &str = + "Workspace-authority tools for listing, materializing, and deleting persistent Workdirs."; + +const LIST_TOOL: &str = "WorkdirList"; +const CREATE_TOOL: &str = "WorkdirCreate"; +const DELETE_TOOL: &str = "WorkdirDelete"; + +const LIST_DESCRIPTION: &str = "List persistent Workdirs in the current Workspace through Backend Workspace API authority. The result contains safe summaries and diagnostics, never host paths or Runtime connection details."; +const CREATE_DESCRIPTION: &str = "Materialize a persistent Workdir on a selected Runtime from a Workspace repository and optional selector. Repository resolution and materialization remain Backend/Runtime authority."; +const DELETE_DESCRIPTION: &str = "Delete one persistent Workdir by id through Backend Workspace API authority. Occupied, blocked, or dirty Workdirs requiring confirmation are rejected."; + +#[derive(Clone, Debug)] +pub struct ManageWorkdirFeature { + client: Arc, +} + +impl ManageWorkdirFeature { + pub fn new(client: Arc) -> Self { + Self { client } + } +} + +pub fn manage_workdir_feature(client: Arc) -> ManageWorkdirFeature { + ManageWorkdirFeature::new(client) +} + +impl FeatureModule for ManageWorkdirFeature { + fn descriptor(&self) -> FeatureDescriptor { + FeatureDescriptor::builtin(FEATURE_ID, FEATURE_NAME) + .with_description(FEATURE_DESCRIPTION) + .with_tool(ToolDeclaration::new(LIST_TOOL, LIST_DESCRIPTION)) + .with_tool(ToolDeclaration::new(CREATE_TOOL, CREATE_DESCRIPTION)) + .with_tool(ToolDeclaration::new(DELETE_TOOL, DELETE_DESCRIPTION)) + } + + fn install(&self, context: &mut FeatureInstallContext<'_>) -> Result<(), FeatureInstallError> { + let backend = WorkspaceHttpWorkdirBackend::new(self.client.clone()); + for (name, definition) in [ + ( + LIST_TOOL, + workdir_tool( + LIST_TOOL, + LIST_DESCRIPTION, + list_schema(), + backend.clone(), + WorkdirOperation::List, + ), + ), + ( + CREATE_TOOL, + workdir_tool( + CREATE_TOOL, + CREATE_DESCRIPTION, + create_schema(), + backend.clone(), + WorkdirOperation::Create, + ), + ), + ( + DELETE_TOOL, + workdir_tool( + DELETE_TOOL, + DELETE_DESCRIPTION, + delete_schema(), + backend.clone(), + WorkdirOperation::Delete, + ), + ), + ] { + context + .tools() + .register(ToolContribution::new(name, definition))?; + } + Ok(()) + } +} + +#[derive(Clone, Debug)] +struct WorkspaceHttpWorkdirBackend { + client: Arc, +} + +impl WorkspaceHttpWorkdirBackend { + fn new(client: Arc) -> Self { + Self { client } + } + + fn workspace_id(&self) -> Result<&str, ToolError> { + let Some(workspace_id) = self.client.workspace_id() else { + return Err(ToolError::ExecutionFailed( + "manage Workdir tools require a Workspace identity".to_string(), + )); + }; + if workspace_id.is_empty() || workspace_id.chars().any(char::is_control) { + return Err(ToolError::ExecutionFailed( + "manage Workdir tools require a valid Workspace identity".to_string(), + )); + } + Ok(workspace_id) + } + + fn list(&self) -> Result { + let workspace_id = encode_path_segment(self.workspace_id()?); + let response = self.execute_json::(WorkspaceRequest::get(format!( + "/api/w/{workspace_id}/working-directories" + )))?; + let count = response.items.len(); + workdir_output(format!("Listed {count} Workdir(s)"), &response) + } + + fn create(&self, input: WorkdirCreateInput) -> Result { + let runtime_id = validate_identity(&input.runtime_id, CREATE_TOOL, "runtime_id")?; + let repository_id = validate_identity(&input.repository_id, CREATE_TOOL, "repository_id")?; + let selector = validate_optional_selector(input.selector)?; + let workspace_id = encode_path_segment(self.workspace_id()?); + let runtime_path = encode_path_segment(runtime_id); + let request = WorkdirCreateRequest { + runtime_id: runtime_id.to_string(), + repository_id: repository_id.to_string(), + selector, + }; + let response = self.execute_json::(WorkspaceRequest::json( + WorkspaceRequestMethod::Post, + format!("/api/w/{workspace_id}/runtimes/{runtime_path}/working-directories"), + serde_json::to_string(&request).map_err(decode_error)?, + ))?; + workdir_output( + format!("Created Workdir {}", response.item.working_directory_id), + &response, + ) + } + + fn delete(&self, input: WorkdirDeleteInput) -> Result { + let workdir_id = validate_identity( + &input.working_directory_id, + DELETE_TOOL, + "working_directory_id", + )?; + let workspace_id = encode_path_segment(self.workspace_id()?); + let workdir_path = encode_path_segment(workdir_id); + let response = self.execute_json::(WorkspaceRequest { + method: WorkspaceRequestMethod::Delete, + path: format!("/api/w/{workspace_id}/working-directories/{workdir_path}"), + body: None, + })?; + workdir_output(format!("Deleted Workdir {workdir_id}"), &response) + } + + fn execute_json Deserialize<'de>>( + &self, + request: WorkspaceRequest, + ) -> Result { + let response = self.client.execute(request).map_err(|error| { + ToolError::ExecutionFailed(format!("Workspace Workdir request failed: {error}")) + })?; + decode_response(response) + } +} + +#[derive(Clone, Copy)] +enum WorkdirOperation { + List, + Create, + Delete, +} + +fn workdir_tool( + name: &'static str, + description: &'static str, + schema: serde_json::Value, + backend: WorkspaceHttpWorkdirBackend, + operation: WorkdirOperation, +) -> ToolDefinition { + Arc::new(move || { + ( + ToolMeta::new(name) + .description(description) + .input_schema(schema.clone()), + Arc::new(WorkspaceHttpWorkdirTool { + backend: backend.clone(), + operation, + }) as Arc, + ) + }) +} + +#[derive(Clone)] +struct WorkspaceHttpWorkdirTool { + backend: WorkspaceHttpWorkdirBackend, + operation: WorkdirOperation, +} + +#[async_trait] +impl Tool for WorkspaceHttpWorkdirTool { + async fn execute( + &self, + input_json: &str, + _ctx: ToolExecutionContext, + ) -> Result { + match self.operation { + WorkdirOperation::List => { + let _input = parse_input::(input_json)?; + self.backend.list() + } + WorkdirOperation::Create => self + .backend + .create(parse_input::(input_json)?), + WorkdirOperation::Delete => self + .backend + .delete(parse_input::(input_json)?), + } + } +} + +fn parse_input Deserialize<'de>>(input: &str) -> Result { + serde_json::from_str(input).map_err(|error| ToolError::InvalidArgument(error.to_string())) +} + +fn decode_response Deserialize<'de>>( + response: WorkspaceResponse, +) -> Result { + if !response.is_success() { + return Err(ToolError::ExecutionFailed(format!( + "Workspace Workdir API returned HTTP {}: {}", + response.status, + bounded_error_body(&response.body) + ))); + } + serde_json::from_str(&response.body).map_err(decode_error) +} + +fn decode_error(error: serde_json::Error) -> ToolError { + ToolError::ExecutionFailed(format!("decode Workspace Workdir response: {error}")) +} + +fn workdir_output(summary: String, value: &T) -> Result { + Ok(ToolOutput { + summary, + content: Some(serde_json::to_string_pretty(value).map_err(decode_error)?), + }) +} + +fn validate_identity<'a>( + value: &'a str, + tool_name: &str, + field: &str, +) -> Result<&'a str, ToolError> { + let value = value.trim(); + if value.is_empty() || value.chars().any(char::is_control) { + return Err(ToolError::InvalidArgument(format!( + "{tool_name} requires non-empty {field} without control characters" + ))); + } + Ok(value) +} + +fn validate_optional_selector(selector: Option) -> Result, ToolError> { + let Some(selector) = selector else { + return Ok(None); + }; + let selector = selector.trim(); + if selector.is_empty() || selector.chars().any(char::is_control) { + return Err(ToolError::InvalidArgument( + "WorkdirCreate selector must be non-empty and contain no control characters" + .to_string(), + )); + } + Ok(Some(selector.to_string())) +} + +fn encode_path_segment(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(byte as char); + } else { + use std::fmt::Write as _; + let _ = write!(encoded, "%{byte:02X}"); + } + } + encoded +} + +fn bounded_error_body(body: &str) -> String { + const MAX_CHARS: usize = 4096; + let mut chars = body.chars(); + let bounded: String = chars.by_ref().take(MAX_CHARS).collect(); + if chars.next().is_some() { + format!("{bounded}…") + } else { + bounded + } +} + +fn list_schema() -> serde_json::Value { + json!({ + "type": "object", + "additionalProperties": false, + "properties": {} + }) +} + +fn create_schema() -> serde_json::Value { + json!({ + "type": "object", + "additionalProperties": false, + "required": ["runtime_id", "repository_id"], + "properties": { + "runtime_id": {"type": "string", "minLength": 1}, + "repository_id": {"type": "string", "minLength": 1}, + "selector": {"type": ["string", "null"], "minLength": 1} + } + }) +} + +fn delete_schema() -> serde_json::Value { + json!({ + "type": "object", + "additionalProperties": false, + "required": ["working_directory_id"], + "properties": { + "working_directory_id": {"type": "string", "minLength": 1} + } + }) +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct WorkdirListInput {} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct WorkdirCreateInput { + runtime_id: String, + repository_id: String, + #[serde(default)] + selector: Option, +} + +#[derive(Debug, Serialize)] +struct WorkdirCreateRequest { + runtime_id: String, + repository_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + selector: Option, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct WorkdirDeleteInput { + working_directory_id: String, +} + +#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] +struct WorkdirListResponse { + workspace_id: String, + items: Vec, + diagnostics: Vec, +} + +#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] +struct WorkdirDetailResponse { + workspace_id: String, + item: WorkdirSummary, + diagnostics: Vec, +} + +#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] +struct WorkdirSummary { + working_directory_id: String, + repository_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + creation_selector: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + creation_ref: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + current_selector: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + current_ref: Option, + materializer_kind: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + cleanup_target: Option, + status: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + cleanliness: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + primary_worker_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + occupied_by: Option, +} + +#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] +struct WorkdirCleanupTarget { + kind: String, + working_directory_id: String, + repository_id: String, +} + +#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] +struct WorkdirOccupancy { + runtime_id: String, + runtime_worker_id: u64, + worker_id: String, + display_name: String, + linked_at: String, +} + +#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] +struct WorkdirDiagnostic { + code: String, + severity: String, + message: String, +} + +#[cfg(test)] +mod tests { + use std::sync::Mutex; + + use super::*; + use crate::feature::{FeatureModule, FeatureRegistryBuilder}; + use crate::hook::HookRegistryBuilder; + use crate::worker::{WorkspaceClientError, WorkspaceResponse}; + + #[derive(Debug)] + struct RecordingWorkspaceClient { + requests: Mutex>, + responses: Mutex>, + } + + impl RecordingWorkspaceClient { + fn new(responses: Vec) -> Self { + Self { + requests: Mutex::new(Vec::new()), + responses: Mutex::new(responses.into_iter().rev().collect()), + } + } + + fn requests(&self) -> Vec { + self.requests.lock().unwrap().clone() + } + } + + impl WorkspaceClient for RecordingWorkspaceClient { + fn workspace_id(&self) -> Option<&str> { + Some("workspace/test") + } + + fn kind(&self) -> &str { + "recording" + } + + fn is_available(&self) -> bool { + true + } + + fn execute( + &self, + request: WorkspaceRequest, + ) -> Result { + self.requests.lock().unwrap().push(request); + self.responses + .lock() + .unwrap() + .pop() + .ok_or_else(|| WorkspaceClientError::Request("missing response".to_string())) + } + } + + fn response(body: serde_json::Value) -> WorkspaceResponse { + WorkspaceResponse { + status: 200, + body: serde_json::to_string(&body).unwrap(), + } + } + + fn workdir_json(id: &str) -> serde_json::Value { + json!({ + "working_directory_id": id, + "repository_id": "main", + "creation_selector": "refs/heads/main", + "creation_ref": "0123456789abcdef", + "materializer_kind": "local_git_worktree", + "cleanup_target": { + "kind": "git_worktree", + "working_directory_id": id, + "repository_id": "main" + }, + "status": "active", + "cleanliness": "clean" + }) + } + + #[test] + fn descriptor_declares_only_workdir_lifecycle_tools() { + let feature = manage_workdir_feature(Arc::new(RecordingWorkspaceClient::new(Vec::new()))); + let descriptor = feature.descriptor(); + let names = descriptor + .tools + .into_iter() + .map(|tool| tool.name) + .collect::>(); + + assert_eq!(descriptor.id.as_str(), "builtin:manage-workdir"); + assert_eq!(names, [LIST_TOOL, CREATE_TOOL, DELETE_TOOL]); + } + + #[test] + fn feature_installs_all_declared_tools_through_registry() { + let feature = manage_workdir_feature(Arc::new(RecordingWorkspaceClient::new(Vec::new()))); + let mut pending_tools = Vec::new(); + let mut hook_builder = HookRegistryBuilder::default(); + let report = FeatureRegistryBuilder::new() + .with_module(feature) + .install_into_pending(&mut pending_tools, &mut hook_builder); + let names = pending_tools + .iter() + .map(|definition| definition().0.name) + .collect::>(); + + assert!(report.reports[0].installed); + assert_eq!( + report.installed_tool_names(), + [LIST_TOOL, CREATE_TOOL, DELETE_TOOL] + ); + assert_eq!(names, [LIST_TOOL, CREATE_TOOL, DELETE_TOOL]); + } + + #[test] + fn missing_workspace_identity_fails_closed() { + #[derive(Debug)] + struct MissingWorkspaceClient; + + impl WorkspaceClient for MissingWorkspaceClient { + fn workspace_id(&self) -> Option<&str> { + None + } + + fn kind(&self) -> &str { + "missing" + } + + fn is_available(&self) -> bool { + false + } + + fn execute( + &self, + _request: WorkspaceRequest, + ) -> Result { + panic!("missing Workspace client must not execute a request") + } + } + + let backend = WorkspaceHttpWorkdirBackend::new(Arc::new(MissingWorkspaceClient)); + let error = backend.list().unwrap_err(); + assert!(matches!(error, ToolError::ExecutionFailed(_))); + assert!(error.to_string().contains("Workspace identity")); + } + + #[test] + fn schemas_expose_identities_without_paths_or_session_handles() { + let create = create_schema(); + assert_eq!(create["required"], json!(["runtime_id", "repository_id"])); + assert!(create["properties"].get("path").is_none()); + assert!(create["properties"].get("session_id").is_none()); + assert_eq!(delete_schema()["required"], json!(["working_directory_id"])); + } + + #[test] + fn list_create_and_delete_use_scoped_workspace_authority_paths() { + let client = Arc::new(RecordingWorkspaceClient::new(vec![ + response(json!({ + "workspace_id": "workspace/test", + "items": [workdir_json("wd-list")], + "diagnostics": [] + })), + response(json!({ + "workspace_id": "workspace/test", + "item": workdir_json("wd-created"), + "diagnostics": [] + })), + response(json!({ + "workspace_id": "workspace/test", + "item": { + "working_directory_id": "wd-created", + "repository_id": "main", + "materializer_kind": "local_git_worktree", + "status": "not_found" + }, + "diagnostics": [] + })), + ])); + let backend = WorkspaceHttpWorkdirBackend::new(client.clone()); + + backend.list().unwrap(); + backend + .create(WorkdirCreateInput { + runtime_id: "runtime/one".to_string(), + repository_id: "main".to_string(), + selector: Some("refs/heads/topic".to_string()), + }) + .unwrap(); + backend + .delete(WorkdirDeleteInput { + working_directory_id: "wd-created".to_string(), + }) + .unwrap(); + + let requests = client.requests(); + assert_eq!( + requests[0].path, + "/api/w/workspace%2Ftest/working-directories" + ); + assert_eq!(requests[0].method, WorkspaceRequestMethod::Get); + assert_eq!( + requests[1].path, + "/api/w/workspace%2Ftest/runtimes/runtime%2Fone/working-directories" + ); + assert_eq!(requests[1].method, WorkspaceRequestMethod::Post); + let body: serde_json::Value = + serde_json::from_str(requests[1].body.as_deref().unwrap()).unwrap(); + assert_eq!(body["repository_id"], "main"); + assert_eq!(body["selector"], "refs/heads/topic"); + assert_eq!( + requests[2].path, + "/api/w/workspace%2Ftest/working-directories/wd-created" + ); + assert_eq!(requests[2].method, WorkspaceRequestMethod::Delete); + } + + #[test] + fn invalid_or_extra_inputs_are_rejected_before_workspace_request() { + let client = Arc::new(RecordingWorkspaceClient::new(Vec::new())); + let backend = WorkspaceHttpWorkdirBackend::new(client.clone()); + let error = backend + .create(WorkdirCreateInput { + runtime_id: " ".to_string(), + repository_id: "main".to_string(), + selector: None, + }) + .unwrap_err(); + assert!(matches!(error, ToolError::InvalidArgument(_))); + assert!(client.requests().is_empty()); + assert!(parse_input::(r#"{"path":"/tmp"}"#).is_err()); + } +} diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 73148196..145ea6e7 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -3924,6 +3924,41 @@ mod tests { } } + #[test] + fn embedded_orchestrator_profile_enables_manage_workdir() { + let root = tempfile::tempdir().unwrap(); + let broker = BackendResourceBroker::default(); + let runtime_id = "runtime-test"; + let selector = ProfileSelector::Builtin("builtin:orchestrator".to_string()); + let bundle = builtin_profile_config_bundle( + &selector, + "workspace-test", + Some(runtime_id), + &broker, + ProfileSourceArchiveTransport::BackendResourceHandle, + ) + .unwrap(); + let handle = bundle.profile_source_archive_handle.as_ref().unwrap(); + let response = broker + .fetch_profile_source_archive(worker_runtime::resource::BackendResourceFetchRequest { + handle: handle.clone(), + runtime_id: runtime_id.to_string(), + worker_id: None, + audit_correlation_id: handle.audit_correlation_id.clone(), + }) + .unwrap(); + let archive = + worker_runtime::resource::profile_source_archive_from_response(handle, response) + .unwrap() + .verify() + .unwrap(); + let manifest = archive + .resolve_profile("builtin:orchestrator", root.path(), "embedded-orchestrator") + .unwrap(); + + assert!(manifest.feature.manage_workdir.enabled); + } + #[test] fn embedded_builtin_decodal_profiles_resolve_through_archive() { let root = tempfile::tempdir().unwrap(); diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 2d028889..16637c8f 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -3736,11 +3736,8 @@ async fn scoped_working_directory_detail( AxumPath(path): AxumPath, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - working_directory_detail_for_runtime( - api, - EMBEDDED_WORKER_RUNTIME_ID, - &path.working_directory_id, - ) + let runtime_id = registered_workdir_runtime_id(&api, &path.working_directory_id)?; + working_directory_detail_for_runtime(api, &runtime_id, &path.working_directory_id) } async fn scoped_cleanup_working_directory( @@ -3748,11 +3745,27 @@ async fn scoped_cleanup_working_directory( AxumPath(path): AxumPath, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - cleanup_working_directory_for_runtime( - api, - EMBEDDED_WORKER_RUNTIME_ID, - &path.working_directory_id, - ) + let runtime_id = registered_workdir_runtime_id(&api, &path.working_directory_id)?; + cleanup_working_directory_for_runtime(api, &runtime_id, &path.working_directory_id) +} + +fn registered_workdir_runtime_id( + api: &WorkspaceApi, + working_directory_id: &str, +) -> ApiResult { + api.store + .get_workdir_registry(&api.config.workspace_id, working_directory_id)? + .map(|record| record.runtime_id) + .ok_or_else(|| { + ApiError::with_diagnostics( + Error::RuntimeOperationFailed { + runtime_id: "workspace-backend".to_string(), + code: "working_directory_not_found".to_string(), + message: format!("Unknown Workdir `{working_directory_id}`"), + }, + Vec::new(), + ) + }) } fn create_working_directory_for_runtime( @@ -10525,6 +10538,47 @@ mod tests { .unwrap(); } + #[tokio::test] + async fn workspace_scoped_workdir_routes_resolve_runtime_owner_from_registry() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + seed_cleanup_workdir(&api, "remote-workdir", "present", "clean"); + + assert_eq!( + registered_workdir_runtime_id(&api, "remote-workdir").unwrap(), + "runtime-test" + ); + } + + #[tokio::test] + async fn simple_workdir_cleanup_rejects_dirty_and_blocked_candidates() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + seed_cleanup_workdir(&api, "dirty-workdir", "present", "dirty"); + let dirty = + cleanup_working_directory_for_runtime(api.clone(), "runtime-test", "dirty-workdir") + .unwrap_err(); + assert!(matches!( + dirty.error, + Error::RuntimeOperationFailed { ref code, .. } + if code == "workspace_cleanup_dirty_confirmation_required" + )); + + let pinned = seed_cleanup_worker(&api, 17, "pinned"); + seed_cleanup_workdir(&api, "blocked-workdir", "present", "clean"); + seed_cleanup_link(&api, pinned.as_str(), "blocked-workdir"); + let blocked = + cleanup_working_directory_for_runtime(api.clone(), "runtime-test", "blocked-workdir") + .unwrap_err(); + assert!(matches!( + blocked.error, + Error::RuntimeOperationFailed { ref code, .. } + if code == "workspace_cleanup_workdir_blocked" + )); + } + #[tokio::test] async fn synthetic_verified_clean_workdir_can_still_use_clean_cleanup_path() { let workspace = tempfile::tempdir().unwrap(); diff --git a/resources/profiles/orchestrator.dcdl b/resources/profiles/orchestrator.dcdl index 29ddd573..a2f0f157 100644 --- a/resources/profiles/orchestrator.dcdl +++ b/resources/profiles/orchestrator.dcdl @@ -8,6 +8,7 @@ import "./base.dcdl" // { memory = { enabled = true; }; web = { enabled = true; }; workers = { enabled = true; }; + manage_workdir = { enabled = true; }; ticket = { enabled = true; thread = true; orchestration_control = true; }; }; }