From 81e631e640d4dbf1e6a58a1057f2498bd9cbb7f6 Mon Sep 17 00:00:00 2001 From: Hare Date: Sat, 1 Aug 2026 17:23:59 +0900 Subject: [PATCH] auth: enforce workspace worker credentials over ticket REST --- crates/ticket/src/lib.rs | 11 - crates/worker-runtime/src/execution.rs | 21 + crates/worker-runtime/src/http_server.rs | 65 +- crates/worker-runtime/src/runtime.rs | 149 ++- crates/worker-runtime/src/worker_backend.rs | 28 + crates/worker/src/controller.rs | 16 +- crates/worker/src/feature/builtin/ticket.rs | 386 +++++- crates/worker/src/worker.rs | 31 +- crates/workspace-server/src/hosts.rs | 210 +++- crates/workspace-server/src/server.rs | 1186 ++++++++++++++++--- crates/workspace-server/src/store.rs | 33 + 11 files changed, 1853 insertions(+), 283 deletions(-) diff --git a/crates/ticket/src/lib.rs b/crates/ticket/src/lib.rs index 88f1811c..e613a841 100644 --- a/crates/ticket/src/lib.rs +++ b/crates/ticket/src/lib.rs @@ -1791,17 +1791,6 @@ where }) } -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(tag = "status", rename_all = "snake_case")] -pub enum TicketBackendHttpResponse { - Ok { - result: TicketBackendOperationResult, - }, - Error { - message: String, - }, -} - #[derive(Debug, Clone)] pub struct LocalTicketBackend { root: PathBuf, diff --git a/crates/worker-runtime/src/execution.rs b/crates/worker-runtime/src/execution.rs index 6f5144a8..7bcc7439 100644 --- a/crates/worker-runtime/src/execution.rs +++ b/crates/worker-runtime/src/execution.rs @@ -31,6 +31,7 @@ pub enum WorkerExecutionOperation { Restore, Input, ProtocolMethod, + ReplaceWorkspaceAccessToken, Stop, Cancel, } @@ -331,6 +332,17 @@ pub trait WorkerExecutionBackend: Send + Sync + 'static { Vec::new() } + fn replace_workspace_access_token( + &self, + _handle: &WorkerExecutionHandle, + _access_token: String, + ) -> WorkerExecutionResult { + WorkerExecutionResult::unsupported( + WorkerExecutionOperation::ReplaceWorkspaceAccessToken, + "execution backend does not support replacing Workspace access tokens", + ) + } + fn stop_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult { WorkerExecutionResult::unsupported( WorkerExecutionOperation::Stop, @@ -443,6 +455,15 @@ impl WorkerExecutionBackendRef { self.backend.worker_completions(handle, kind, prefix) } + pub(crate) fn replace_workspace_access_token( + &self, + handle: &WorkerExecutionHandle, + access_token: String, + ) -> WorkerExecutionResult { + self.backend + .replace_workspace_access_token(handle, access_token) + } + pub(crate) fn stop_worker(&self, handle: &WorkerExecutionHandle) -> WorkerExecutionResult { self.backend.stop_worker(handle) } diff --git a/crates/worker-runtime/src/http_server.rs b/crates/worker-runtime/src/http_server.rs index 4f0f4c95..f44f52a8 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -12,7 +12,7 @@ use crate::auth::{ }; use crate::catalog::{ ConfigBundleRef, CreateWorkerRequest, WorkerDetail, WorkerLifecycleAck, WorkerSummary, - WorkingDirectoryRequest, WorkingDirectoryStatus, + WorkingDirectoryRequest, WorkingDirectoryStatus, WorkspaceApiRef, }; use crate::config_bundle::{ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary}; use crate::error::RuntimeError; @@ -193,6 +193,10 @@ fn runtime_http_router_with_optional_auth( ) .route("/v1/workers/{worker_id}/input", post(send_worker_input)) .route("/v1/workers/{worker_id}/restore", post(restore_worker)) + .route( + "/v1/workers/{worker_id}/workspace-api", + post(replace_worker_workspace_api), + ) .route( "/v1/workers/{worker_id}/completions", post(worker_completions), @@ -282,6 +286,12 @@ pub struct RuntimeHttpWorkerResponse { pub worker: WorkerDetail, } +/// Replace the Workspace API binding for an existing Worker. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct RuntimeHttpWorkerWorkspaceApiRequest { + pub workspace_api: WorkspaceApiRef, +} + /// Worker delete response. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct RuntimeHttpWorkerDeleteResponse { @@ -509,6 +519,28 @@ async fn create_worker( Ok(Json(RuntimeHttpWorkerResponse { worker })) } +async fn replace_worker_workspace_api( + State(state): State, + auth: Option>, + Path(worker_id): Path, + body: Result, JsonRejection>, +) -> RestResult { + let Json(request) = body.map_err(RuntimeHttpRestError::json_rejection)?; + let worker_ref = worker_ref_for(&state.runtime, worker_id)?; + let worker = match auth_workspace_scope(&state, auth.as_ref())? { + Some(scope) => state.runtime.replace_worker_workspace_api_scoped( + &scope, + &worker_ref, + request.workspace_api, + ), + None => state + .runtime + .replace_worker_workspace_api(&worker_ref, request.workspace_api), + } + .map_err(RuntimeHttpRestError::runtime)?; + Ok(Json(RuntimeHttpWorkerResponse { worker })) +} + async fn restore_worker( State(state): State, auth: Option>, @@ -960,6 +992,9 @@ fn required_runtime_permission(method: &Method, path: &str) -> Option<&'static s if path.starts_with("/v1/config-bundles") || path.starts_with("/v1/working-directories") { return Some("workers:create"); } + if path.ends_with("/workspace-api") { + return Some("workers:create"); + } if path.ends_with("/input") || path.ends_with("/restore") { return Some("workers:input"); } @@ -1484,6 +1519,17 @@ mod tests { ) } + fn replace_workspace_access_token( + &self, + _handle: &WorkerExecutionHandle, + _access_token: String, + ) -> WorkerExecutionResult { + WorkerExecutionResult::accepted( + WorkerExecutionOperation::ReplaceWorkspaceAccessToken, + WorkerExecutionRunState::Idle, + ) + } + fn stop_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult { WorkerExecutionResult::accepted( WorkerExecutionOperation::Stop, @@ -1575,6 +1621,23 @@ mod tests { created.worker.worker_id ); + let response = authed_json_request( + app.clone(), + Method::POST, + &format!("/v1/workers/{}/workspace-api", created.worker.worker_id), + token, + &RuntimeHttpWorkerWorkspaceApiRequest { + workspace_api: WorkspaceApiRef { + workspace_id: "local".to_string(), + base_url: "http://127.0.0.1:8787".to_string(), + runtime_id: None, + access_token: Some("workspace-access-token".to_string()), + }, + }, + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let input = WorkerInput::user("hello from backend"); let response = authed_json_request( app.clone(), diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 678de3f5..5a885106 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -1,7 +1,7 @@ use crate::catalog::{ ConfigBundleRef, CreateWorkerRequest, WorkerDetail, WorkerLifecycleAck, WorkerStatus, WorkerSummary, WorkingDirectoryRequest, - WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, + WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, WorkspaceApiRef, }; use crate::config_bundle::{ ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary, validate_config_bundle, @@ -589,6 +589,94 @@ impl Runtime { Ok(()) } + /// Replace the Workspace API binding persisted for a Worker and update the + /// live execution when one is connected. + pub fn replace_worker_workspace_api_scoped( + &self, + scope: &RuntimeWorkspaceScope, + worker_ref: &WorkerRef, + workspace_api: WorkspaceApiRef, + ) -> Result { + self.ensure_worker_in_workspace(scope, worker_ref)?; + if workspace_api.workspace_id != scope.workspace_id { + return Err(RuntimeError::InvalidRequest(format!( + "Workspace API scope `{}` does not match authorized workspace `{}`", + workspace_api.workspace_id, scope.workspace_id + ))); + } + self.replace_worker_workspace_api(worker_ref, workspace_api) + } + + pub fn replace_worker_workspace_api( + &self, + worker_ref: &WorkerRef, + workspace_api: WorkspaceApiRef, + ) -> Result { + let access_token = workspace_api + .access_token + .as_ref() + .filter(|token| !token.trim().is_empty()) + .cloned() + .ok_or_else(|| { + RuntimeError::InvalidRequest( + "Workspace API replacement requires an access token".to_string(), + ) + })?; + let (previous_workspace_api, live_execution) = { + let state = self.lock()?; + let worker = state.worker(worker_ref)?; + if let Some(existing) = worker.request.workspace_api.as_ref() + && (existing.workspace_id != workspace_api.workspace_id + || existing.base_url.trim_end_matches('/') + != workspace_api.base_url.trim_end_matches('/') + || existing.runtime_id.as_ref().is_some_and(|runtime_id| { + workspace_api.runtime_id.as_ref() != Some(runtime_id) + })) + { + return Err(RuntimeError::InvalidRequest( + "Workspace API replacement cannot change Worker Workspace identity, Runtime identity, or base URL" + .to_string(), + )); + } + let live_execution = match ( + state.execution_backend.clone(), + worker.execution_handle.clone(), + ) { + (Some(backend), Some(handle)) => Some((backend, handle)), + _ => None, + }; + (worker.request.workspace_api.clone(), live_execution) + }; + + { + let mut state = self.lock()?; + state.worker_mut(worker_ref)?.request.workspace_api = Some(workspace_api); + if let Err(error) = state.persist_runtime_snapshot() { + state.worker_mut(worker_ref)?.request.workspace_api = previous_workspace_api; + return Err(error); + } + } + + if let Some((backend, handle)) = live_execution { + let result = backend.replace_workspace_access_token(&handle, access_token); + if !result.is_accepted() { + let mut state = self.lock()?; + state.worker_mut(worker_ref)?.request.workspace_api = previous_workspace_api; + state.persist_runtime_snapshot()?; + return Err(RuntimeError::WorkerExecutionRejected { + worker_id: worker_ref.worker_id.clone(), + operation: result.operation, + outcome: result.outcome, + message: result.message_or_default(), + result, + }); + } + } + + let state = self.lock()?; + Ok(state.worker(worker_ref)?.detail()) + } + /// Attach a live execution through a workspace-scoped Runtime authorization context. pub fn restore_worker_scoped( &self, @@ -953,7 +1041,8 @@ impl Runtime { WorkerExecutionOperation::Spawn | WorkerExecutionOperation::Restore | WorkerExecutionOperation::Input - | WorkerExecutionOperation::ProtocolMethod => return Ok(()), + | WorkerExecutionOperation::ProtocolMethod + | WorkerExecutionOperation::ReplaceWorkspaceAccessToken => return Ok(()), }; if result.is_accepted() { return Ok(()); @@ -2228,6 +2317,7 @@ mod tests { restore_result: Mutex>, restore_count: Mutex, contexts: Mutex>, + workspace_access_tokens: Mutex>, #[cfg(feature = "ws-server")] snapshots: Mutex>, } @@ -2316,6 +2406,21 @@ mod tests { }) } + fn replace_workspace_access_token( + &self, + handle: &WorkerExecutionHandle, + access_token: String, + ) -> WorkerExecutionResult { + self.workspace_access_tokens + .lock() + .unwrap() + .insert(handle.worker_ref().worker_id.clone(), access_token); + WorkerExecutionResult::accepted( + WorkerExecutionOperation::ReplaceWorkspaceAccessToken, + WorkerExecutionRunState::Idle, + ) + } + fn stop_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult { WorkerExecutionResult::accepted( WorkerExecutionOperation::Stop, @@ -2442,6 +2547,46 @@ mod tests { ); } + #[test] + fn workspace_api_replacement_updates_live_execution_and_persisted_request() { + let (runtime, backend) = runtime_and_backend(); + let scope = scope("workspace-a", "server-a"); + let worker = runtime + .create_worker_scoped( + &scope, + scoped_task_request("repair credential", "workspace-a"), + ) + .unwrap(); + let replacement = WorkspaceApiRef { + workspace_id: "workspace-a".to_string(), + base_url: "https://workspace.example/workspace-a/".to_string(), + runtime_id: Some("runtime-a".to_string()), + access_token: Some("replacement-token".to_string()), + }; + + runtime + .replace_worker_workspace_api_scoped(&scope, &worker.worker_ref, replacement.clone()) + .unwrap(); + + assert_eq!( + backend + .workspace_access_tokens + .lock() + .unwrap() + .get(&worker.worker_ref.worker_id), + Some(&"replacement-token".to_string()) + ); + let state = runtime.lock().unwrap(); + assert_eq!( + state + .worker(&worker.worker_ref) + .unwrap() + .request + .workspace_api, + Some(replacement) + ); + } + #[test] fn workspace_owner_binding_rejects_other_backend_and_forgets_after_last_worker_delete() { let runtime = runtime_with_backend(); diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 79f5fe27..2c166be4 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -1124,6 +1124,34 @@ where result } + fn replace_workspace_access_token( + &self, + handle: &WorkerExecutionHandle, + access_token: String, + ) -> WorkerExecutionResult { + let (worker, _busy) = match self.get_execution(handle) { + Ok(execution) => execution, + Err(mut result) => { + result.operation = WorkerExecutionOperation::ReplaceWorkspaceAccessToken; + return result; + } + }; + worker + .replace_workspace_access_token(access_token) + .map(|_| { + WorkerExecutionResult::accepted( + WorkerExecutionOperation::ReplaceWorkspaceAccessToken, + WorkerExecutionRunState::Idle, + ) + }) + .unwrap_or_else(|error| { + WorkerExecutionResult::errored( + WorkerExecutionOperation::ReplaceWorkspaceAccessToken, + error.to_string(), + ) + }) + } + fn stop_worker(&self, handle: &WorkerExecutionHandle) -> WorkerExecutionResult { if handle.backend_id() != self.backend_id() { return WorkerExecutionResult::rejected( diff --git a/crates/worker/src/controller.rs b/crates/worker/src/controller.rs index 48a43235..5869b340 100644 --- a/crates/worker/src/controller.rs +++ b/crates/worker/src/controller.rs @@ -26,7 +26,10 @@ use crate::shutdown_after_idle::{ use crate::spawn::comm_tools::{read_worker_output_tool, send_to_worker_tool, stop_worker_tool}; use crate::spawn::registry::SpawnedWorkerRegistry; use crate::spawn::tool::spawn_worker_tool; -use crate::worker::{SystemItemCommitter, Worker, WorkerError, WorkerRunResult}; +use crate::worker::{ + SystemItemCommitter, Worker, WorkerError, WorkerRunResult, WorkspaceClient, + WorkspaceClientError, +}; use protocol::{ AlertLevel, AlertSource, ErrorCode, Event, Method, RewindTargetId, RunResult, Segment, TurnResult, WorkerStatus, @@ -40,6 +43,7 @@ use protocol::{ pub struct WorkerHandle { method_tx: mpsc::Sender, event_tx: broadcast::Sender, + workspace_client: Arc, pub shared_state: Arc, pub runtime_dir: Arc, pub alerter: Alerter, @@ -113,6 +117,14 @@ impl WorkerHandle { pub fn alert(&self, level: AlertLevel, source: AlertSource, message: String) { self.alerter.alert(level, source, message); } + + /// Replace the Runtime-issued Workspace access token used by this live Worker. + pub fn replace_workspace_access_token( + &self, + access_token: String, + ) -> Result<(), WorkspaceClientError> { + self.workspace_client.replace_access_token(access_token) + } } async fn set_controller_status( @@ -234,6 +246,7 @@ impl WorkerController { let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>(); let (method_tx, method_rx) = mpsc::channel::(32); let (event_tx, _) = broadcast::channel::(256); + let workspace_client = worker.workspace_client_handle(); let alerter = Alerter::new(event_tx.clone()); let in_flight = InFlightEvents::new(event_tx.clone()); worker.attach_in_flight_events(in_flight.clone()); @@ -352,6 +365,7 @@ impl WorkerController { let handle = WorkerHandle { method_tx, event_tx: event_tx.clone(), + workspace_client, shared_state: shared_state.clone(), runtime_dir: runtime_dir.clone(), alerter: alerter.clone(), diff --git a/crates/worker/src/feature/builtin/ticket.rs b/crates/worker/src/feature/builtin/ticket.rs index b9e36510..2f9ca265 100644 --- a/crates/worker/src/feature/builtin/ticket.rs +++ b/crates/worker/src/feature/builtin/ticket.rs @@ -12,10 +12,10 @@ use std::{ use ticket::{ LocalTicketBackend, MarkdownText, NewOrchestrationPlanRecord, NewTicket, NewTicketEvent, NewTicketRelation, OrchestrationPlanKind, OrchestrationPlanRecord, Result as TicketResult, - Ticket, TicketBackend, TicketBackendHttpResponse, TicketBackendOperation, - TicketBackendOperationResult, TicketDoctorReport, TicketError, TicketIdOrSlug, - TicketIntakeSummary, TicketListQuery, TicketRef, TicketRelation, TicketRelationKind, - TicketRelationView, TicketReview, TicketStateChange, TicketSummary, + Ticket, TicketBackend, TicketBackendOperation, TicketBackendOperationResult, + TicketDoctorReport, TicketError, TicketIdOrSlug, TicketIntakeSummary, TicketListQuery, + TicketRef, TicketRelation, TicketRelationKind, TicketRelationView, TicketReview, + TicketStateChange, TicketSummary, config::{DEFAULT_TICKET_BACKEND_RELATIVE_PATH, TicketConfig}, tool::{TICKET_TOOL_NAMES, TicketToolBackend, ticket_tool_description, ticket_tools}, }; @@ -387,57 +387,310 @@ impl WorkspaceHttpTicketBackend { Self { client } } - fn endpoint(&self) -> String { - format!( - "/api/w/{}/tickets/backend", - self.client.workspace_id().unwrap_or_default() - ) - } - fn invoke( &self, operation: TicketBackendOperation, ) -> TicketResult { let client = self.client.clone(); - let endpoint = self.endpoint(); + let workspace_id = self.client.workspace_id().unwrap_or_default().to_string(); if tokio::runtime::Handle::try_current().is_ok() { - return std::thread::spawn(move || Self::invoke_client(client, endpoint, operation)) - .join() - .map_err(|_| { - TicketError::Conflict("ticket backend request thread panicked".to_string()) - })?; + return std::thread::spawn(move || { + Self::invoke_client(client, workspace_id, operation) + }) + .join() + .map_err(|_| { + TicketError::Conflict("ticket REST request thread panicked".to_string()) + })?; } - Self::invoke_client(client, endpoint, operation) + Self::invoke_client(client, workspace_id, operation) + } + + fn ticket_path(id: &TicketIdOrSlug) -> String { + let value = match id { + TicketIdOrSlug::Id(value) + | TicketIdOrSlug::Slug(value) + | TicketIdOrSlug::Query(value) => value, + }; + 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 request( + client: Arc, + method: WorkspaceRequestMethod, + endpoint: String, + body: Option, + ) -> TicketResult { + let request = match body { + Some(body) => WorkspaceRequest::json(method, endpoint, body.to_string()), + None if method == WorkspaceRequestMethod::Get => WorkspaceRequest::get(endpoint), + None => WorkspaceRequest { + method, + path: endpoint, + body: None, + }, + }; + let response = client.execute(request).map_err(|error| { + TicketError::Conflict(format!("ticket REST request failed: {error}")) + })?; + if !response.is_success() { + return Err(TicketError::Conflict(format!( + "ticket REST API returned HTTP {}: {}", + response.status, response.body + ))); + } + serde_json::from_str(&response.body) + .map_err(|error| TicketError::Conflict(format!("decode ticket REST response: {error}"))) + } + + fn request_unit( + client: Arc, + method: WorkspaceRequestMethod, + endpoint: String, + body: Option, + ) -> TicketResult { + let request = match body { + Some(body) => WorkspaceRequest::json(method, endpoint, body.to_string()), + None => WorkspaceRequest { + method, + path: endpoint, + body: None, + }, + }; + let response = client.execute(request).map_err(|error| { + TicketError::Conflict(format!("ticket REST request failed: {error}")) + })?; + if !response.is_success() { + return Err(TicketError::Conflict(format!( + "ticket REST API returned HTTP {}: {}", + response.status, response.body + ))); + } + Ok(TicketBackendOperationResult::Unit) } fn invoke_client( client: Arc, - endpoint: String, + workspace_id: String, operation: TicketBackendOperation, ) -> TicketResult { - let body = serde_json::to_string(&operation).map_err(|error| { - TicketError::Conflict(format!("serialize ticket operation: {error}")) - })?; - let response = client - .execute(WorkspaceRequest::json( + let base = format!("/api/w/{workspace_id}/tickets"); + match operation { + TicketBackendOperation::DefaultIntakeReadyStateChangeBody { from } => { + let value = Self::request::( + client, + WorkspaceRequestMethod::Post, + format!("{base}/default-intake-ready-body"), + Some(serde_json::json!({ "from": from })), + )?; + Ok(TicketBackendOperationResult::Text(value)) + } + TicketBackendOperation::List { filter } => { + let state = match filter.state { + ticket::TicketStateSelector::Active => "active".to_string(), + ticket::TicketStateSelector::All => "all".to_string(), + ticket::TicketStateSelector::States(states) => states + .into_iter() + .map(|state| state.as_str().to_string()) + .collect::>() + .join(","), + }; + let tickets = Self::request( + client, + WorkspaceRequestMethod::Get, + format!("{base}/search?state={state}"), + None, + )?; + Ok(TicketBackendOperationResult::Tickets(tickets)) + } + TicketBackendOperation::Show { id } => { + let ticket = Self::request( + client, + WorkspaceRequestMethod::Get, + format!("{base}/{}/record", Self::ticket_path(&id)), + None, + )?; + Ok(TicketBackendOperationResult::Ticket(ticket)) + } + TicketBackendOperation::Create { input } => { + let ticket = Self::request( + client, + WorkspaceRequestMethod::Post, + base, + Some(serde_json::to_value(input).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket create: {error}")) + })?), + )?; + Ok(TicketBackendOperationResult::TicketRef(ticket)) + } + TicketBackendOperation::EditItem { id, edit } => { + let ticket = Self::request( + client, + WorkspaceRequestMethod::Patch, + format!("{base}/{}/item", Self::ticket_path(&id)), + Some(serde_json::to_value(edit).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket edit: {error}")) + })?), + )?; + Ok(TicketBackendOperationResult::Ticket(ticket)) + } + TicketBackendOperation::DependencyCheck { id } => { + let check = Self::request( + client, + WorkspaceRequestMethod::Get, + format!("{base}/{}/dependency-check", Self::ticket_path(&id)), + None, + )?; + Ok(TicketBackendOperationResult::DependencyCheck(check)) + } + TicketBackendOperation::AddEvent { id, event } => Self::request_unit( + client, WorkspaceRequestMethod::Post, - endpoint, - body, - )) - .map_err(|error| { - TicketError::Conflict(format!("ticket backend request failed: {error}")) - })?; - if !response.is_success() { - return Err(TicketError::Conflict(format!( - "ticket backend returned HTTP {}: {}", - response.status, response.body - ))); - } - match serde_json::from_str::(&response.body).map_err( - |error| TicketError::Conflict(format!("decode ticket backend response: {error}")), - )? { - TicketBackendHttpResponse::Ok { result } => Ok(result), - TicketBackendHttpResponse::Error { message } => Err(TicketError::Conflict(message)), + format!("{base}/{}/thread-events", Self::ticket_path(&id)), + Some(serde_json::to_value(event).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket event: {error}")) + })?), + ), + TicketBackendOperation::AddStateChanged { id, change } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/state-changes", Self::ticket_path(&id)), + Some(serde_json::to_value(change).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket state change: {error}")) + })?), + ), + TicketBackendOperation::AddIntakeSummary { id, summary } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/intake-summaries", Self::ticket_path(&id)), + Some(serde_json::to_value(summary).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket intake summary: {error}")) + })?), + ), + TicketBackendOperation::SetStateField { id, field, change } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!( + "{base}/{}/state-fields/{}", + Self::ticket_path(&id), + Self::ticket_path(&TicketIdOrSlug::Query(field)) + ), + Some(serde_json::to_value(change).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket state field change: {error}")) + })?), + ), + TicketBackendOperation::SetWorkflowState { id, change } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/workflow-state", Self::ticket_path(&id)), + Some(serde_json::to_value(change).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket workflow change: {error}")) + })?), + ), + TicketBackendOperation::MarkIntakeReady { + id, + summary, + change, + } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/intake-ready", Self::ticket_path(&id)), + Some(serde_json::json!({ "summary": summary, "change": change })), + ), + TicketBackendOperation::QueueReady { id, .. } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/workflow/queue", Self::ticket_path(&id)), + None, + ), + TicketBackendOperation::Review { id, review } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/workflow/review", Self::ticket_path(&id)), + Some(serde_json::to_value(review).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket review: {error}")) + })?), + ), + TicketBackendOperation::Close { id, resolution } => Self::request_unit( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/workflow/close", Self::ticket_path(&id)), + Some(serde_json::to_value(resolution).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket close: {error}")) + })?), + ), + TicketBackendOperation::AddTicketRelation { id, relation } => { + let relation = Self::request( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/relations", Self::ticket_path(&id)), + Some(serde_json::to_value(relation).map_err(|error| { + TicketError::Conflict(format!("serialize Ticket relation: {error}")) + })?), + )?; + Ok(TicketBackendOperationResult::Relation(relation)) + } + TicketBackendOperation::QueryTicketRelations { ticket, kind } => { + let relations = Self::request( + client, + WorkspaceRequestMethod::Post, + format!("{base}/relations/search"), + Some(serde_json::json!({ "ticket": ticket, "kind": kind })), + )?; + Ok(TicketBackendOperationResult::Relations(relations)) + } + TicketBackendOperation::RelationView { id } => { + let view = Self::request( + client, + WorkspaceRequestMethod::Get, + format!("{base}/{}/relation-view", Self::ticket_path(&id)), + None, + )?; + Ok(TicketBackendOperationResult::RelationView(view)) + } + TicketBackendOperation::AddOrchestrationPlanRecord { id, record } => { + let record = Self::request( + client, + WorkspaceRequestMethod::Post, + format!("{base}/{}/orchestration-plans", Self::ticket_path(&id)), + Some(serde_json::to_value(record).map_err(|error| { + TicketError::Conflict(format!( + "serialize Ticket orchestration plan: {error}" + )) + })?), + )?; + Ok(TicketBackendOperationResult::OrchestrationPlanRecord( + record, + )) + } + TicketBackendOperation::QueryOrchestrationPlanRecords { ticket, kind } => { + let records = Self::request( + client, + WorkspaceRequestMethod::Post, + format!("{base}/orchestration-plans/search"), + Some(serde_json::json!({ "ticket": ticket, "kind": kind })), + )?; + Ok(TicketBackendOperationResult::OrchestrationPlanRecords( + records, + )) + } + TicketBackendOperation::Doctor => { + let report = Self::request( + client, + WorkspaceRequestMethod::Get, + format!("{base}/doctor"), + None, + )?; + Ok(TicketBackendOperationResult::DoctorReport(report)) + } } } } @@ -1129,7 +1382,42 @@ provider = "github" }) .unwrap_err(); - assert!(error.to_string().contains("ticket backend request failed")); + assert!(error.to_string().contains("ticket REST request failed")); + } + + #[test] + fn workspace_http_backend_posts_ticket_event_subresource() { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let address = listener.local_addr().unwrap(); + let server = std::thread::spawn(move || { + let (mut stream, _) = listener.accept().unwrap(); + let mut buffer = [0_u8; 8192]; + let size = stream.read(&mut buffer).unwrap(); + let request = String::from_utf8_lossy(&buffer[..size]); + assert!( + request + .starts_with("POST /api/w/workspace-a/tickets/01TEST/thread-events HTTP/1.1") + ); + assert!(!request.contains("\"operation\"")); + assert!(request.contains("\"kind\":\"comment\"")); + stream + .write_all(b"HTTP/1.1 204 No Content\r\nContent-Length: 0\r\n\r\n") + .unwrap(); + }); + let client = Arc::new(crate::worker::RuntimeWorkspaceHttpClient::new( + "workspace-a", + format!("http://{address}"), + "worker-a", + )); + let backend = WorkspaceHttpTicketBackend::new(client); + + backend + .add_event( + TicketIdOrSlug::Id("01TEST".to_string()), + NewTicketEvent::new(ticket::TicketEventKind::Comment, "REST comment"), + ) + .unwrap(); + server.join().unwrap(); } #[test] @@ -1141,15 +1429,13 @@ provider = "github" let mut buffer = [0_u8; 8192]; let len = stream.read(&mut buffer).unwrap(); let request = String::from_utf8_lossy(&buffer[..len]); - assert!(request.starts_with("POST /api/w/workspace-a/tickets/backend HTTP/1.1")); - assert!(request.contains("\"operation\":\"create\"")); + assert!(request.starts_with("POST /api/w/workspace-a/tickets HTTP/1.1")); + assert!(!request.contains("\"operation\"")); assert!(request.contains("\"title\":\"HTTP ticket\"")); - let response_body = serde_json::to_string(&TicketBackendHttpResponse::Ok { - result: TicketBackendOperationResult::TicketRef(TicketRef { - id: "01TEST".to_string(), - slug: "http-ticket".to_string(), - status: ticket::TicketStatus::Open, - }), + let response_body = serde_json::to_string(&TicketRef { + id: "01TEST".to_string(), + slug: "http-ticket".to_string(), + status: ticket::TicketStatus::Open, }) .unwrap(); write!( diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 82c10245..895d7d17 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -215,6 +215,13 @@ pub trait WorkspaceClient: std::fmt::Debug + Send + Sync { fn is_available(&self) -> bool; fn execute(&self, request: WorkspaceRequest) -> Result; + + /// Replace the Runtime-issued Workspace access token for this live client. + fn replace_access_token(&self, _access_token: String) -> Result<(), WorkspaceClientError> { + Err(WorkspaceClientError::Unavailable( + "Workspace client does not support access token replacement".to_string(), + )) + } } /// HTTP forwarding client created by Runtime for one concrete Worker execution. @@ -323,6 +330,13 @@ impl WorkspaceClient for RuntimeWorkspaceHttpClient { } Ok(result.0) } + + fn replace_access_token(&self, access_token: String) -> Result<(), WorkspaceClientError> { + *self.access_token.lock().map_err(|_| { + WorkspaceClientError::Request("workspace credential lock poisoned".to_string()) + })? = Some(access_token); + Ok(()) + } } fn execute_runtime_workspace_http_with_refresh( @@ -6318,13 +6332,28 @@ mod build_summary_prompt_tests { .with_access_token(Some("expired-token".to_string())); let response = client .execute(WorkspaceRequest::get( - "/api/w/workspace-refresh/tickets/backend", + "/api/w/workspace-refresh/tickets/search", )) .unwrap(); assert_eq!(response.status, 200); server.join().unwrap(); } + #[test] + fn runtime_workspace_client_can_install_missing_access_token() { + let client = + RuntimeWorkspaceHttpClient::new("workspace-a", "https://workspace.example", "worker-a"); + + client + .replace_access_token("replacement-token".to_string()) + .unwrap(); + + assert_eq!( + client.access_token.lock().unwrap().as_deref(), + Some("replacement-token") + ); + } + fn minimal_manifest() -> WorkerManifest { let toml_str = r#" [worker] diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index e44f828a..386b92b1 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -34,7 +34,8 @@ use worker_runtime::http_server::{ RuntimeHttpErrorResponse, RuntimeHttpSummaryResponse, RuntimeHttpWorkerCompletionsRequest, RuntimeHttpWorkerCompletionsResponse, RuntimeHttpWorkerDeleteResponse, RuntimeHttpWorkerInputResponse, RuntimeHttpWorkerLifecycleRequest, - RuntimeHttpWorkerLifecycleResponse, RuntimeHttpWorkerResponse, RuntimeHttpWorkersResponse, + RuntimeHttpWorkerLifecycleResponse, RuntimeHttpWorkerResponse, + RuntimeHttpWorkerWorkspaceApiRequest, RuntimeHttpWorkersResponse, RuntimeHttpWorkingDirectoriesResponse, RuntimeHttpWorkingDirectoryResponse, }; use worker_runtime::identity::{WorkerId as EmbeddedWorkerId, WorkerRef as EmbeddedWorkerRef}; @@ -260,6 +261,14 @@ pub struct WorkerRestoreResult { pub diagnostics: Vec, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkerWorkspaceApiResult { + pub state: WorkerOperationState, + #[serde(skip_serializing_if = "Option::is_none")] + pub worker: Option, + pub diagnostics: Vec, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct RuntimeList { pub items: Vec, @@ -415,6 +424,30 @@ pub struct ConfigBundleListResult { pub diagnostics: Vec, } +fn required_worker_workspace_api( + request: &WorkerSpawnRequest, +) -> Result { + let workspace_api = request.resolved_workspace_api.clone().ok_or_else(|| { + diagnostic( + "worker_workspace_credential_missing", + DiagnosticSeverity::Error, + "Workspace-bound Worker spawn requires a resolved Workspace API credential", + ) + })?; + if workspace_api + .access_token + .as_deref() + .is_none_or(|token| token.trim().is_empty()) + { + return Err(diagnostic( + "worker_workspace_credential_missing", + DiagnosticSeverity::Error, + "Workspace-bound Worker spawn requires a non-empty Workspace API access token", + )); + } + Ok(workspace_api) +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum WorkerOperationState { @@ -622,6 +655,24 @@ pub trait WorkspaceWorkerRuntime: Send + Sync { } } + fn replace_worker_workspace_api( + &self, + worker_id: &str, + _workspace_api: WorkspaceApiRef, + ) -> WorkerWorkspaceApiResult { + WorkerWorkspaceApiResult { + state: WorkerOperationState::Unsupported, + worker: None, + diagnostics: vec![diagnostic( + "worker_workspace_api_replace_unsupported", + DiagnosticSeverity::Info, + format!( + "runtime does not support replacing the Workspace API for worker `{worker_id}`" + ), + )], + } + } + fn create_working_directory( &self, _request: WorkingDirectoryRequest, @@ -1072,6 +1123,18 @@ impl RuntimeRegistry { Ok(runtime.restore_worker(worker_id)) } + pub fn replace_worker_workspace_api( + &self, + runtime_id: &str, + worker_id: &str, + workspace_api: WorkspaceApiRef, + ) -> Result { + validate_backend_identifier("runtime_id", runtime_id)?; + validate_backend_identifier("worker_id", worker_id)?; + let runtime = self.runtime(runtime_id)?; + Ok(runtime.replace_worker_workspace_api(worker_id, workspace_api)) + } + pub fn spawn_worker( &self, runtime_id: &str, @@ -1312,8 +1375,6 @@ impl RuntimeRegistry { pub struct EmbeddedWorkerRuntime { runtime_id: String, host_id: String, - workspace_id: String, - backend_base_url: Option, runtime: worker_runtime::Runtime, execution_enabled: bool, resource_broker: BackendResourceBroker, @@ -1366,18 +1427,11 @@ impl EmbeddedWorkerRuntime { self } - pub fn with_backend_base_url(mut self, backend_base_url: impl Into) -> Self { - self.backend_base_url = Some(backend_base_url.into().trim_end_matches('/').to_string()); - self - } - pub fn from_runtime(workspace_id: impl AsRef, runtime: worker_runtime::Runtime) -> Self { let workspace_id = workspace_id.as_ref().to_string(); Self { runtime_id: EMBEDDED_RUNTIME_ID.to_string(), host_id: host_id_for_embedded_workspace(&workspace_id), - workspace_id, - backend_base_url: None, runtime, execution_enabled: false, resource_broker: BackendResourceBroker::default(), @@ -1626,6 +1680,39 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { } } + fn replace_worker_workspace_api( + &self, + worker_id: &str, + workspace_api: WorkspaceApiRef, + ) -> WorkerWorkspaceApiResult { + let Some(worker_ref) = self.worker_ref(worker_id) else { + return WorkerWorkspaceApiResult { + state: WorkerOperationState::Rejected, + worker: None, + diagnostics: vec![diagnostic( + "embedded_worker_id_invalid", + DiagnosticSeverity::Warning, + "Worker id was empty and cannot receive Workspace access".to_string(), + )], + }; + }; + match self + .runtime + .replace_worker_workspace_api(&worker_ref, workspace_api) + { + Ok(detail) => WorkerWorkspaceApiResult { + state: WorkerOperationState::Accepted, + worker: Some(self.map_worker_detail(detail)), + diagnostics: Vec::new(), + }, + Err(err) => WorkerWorkspaceApiResult { + state: WorkerOperationState::Rejected, + worker: None, + diagnostics: vec![embedded_runtime_diagnostic(&err)], + }, + } + } + fn create_working_directory( &self, _request: WorkingDirectoryRequest, @@ -1727,6 +1814,18 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { .map_or((None, None), |(key, fingerprint)| { (Some(key), Some(fingerprint)) }); + let workspace_api = match required_worker_workspace_api(&request) { + Ok(workspace_api) => workspace_api, + Err(diagnostic) => { + diagnostics.push(diagnostic); + return WorkerSpawnResult { + state: WorkerOperationState::Rejected, + worker: None, + acceptance_evidence: Vec::new(), + diagnostics, + }; + } + }; let create_request = CreateWorkerRequest { idempotency_key, idempotency_fingerprint, @@ -1737,16 +1836,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { initial_input: request.initial_input.clone(), working_directory_request: request.resolved_working_directory_request.clone(), working_directory: request.resolved_working_directory.clone(), - workspace_api: request.resolved_workspace_api.clone().or_else(|| { - self.backend_base_url - .as_ref() - .map(|base_url| WorkspaceApiRef { - workspace_id: self.workspace_id.clone(), - base_url: base_url.clone(), - runtime_id: Some(self.runtime_id.clone()), - access_token: None, - }) - }), + workspace_api: Some(workspace_api), }; match self.runtime.create_worker(create_request) { Ok(detail) => WorkerSpawnResult { @@ -2602,6 +2692,28 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { } } + fn replace_worker_workspace_api( + &self, + worker_id: &str, + workspace_api: WorkspaceApiRef, + ) -> WorkerWorkspaceApiResult { + match self.post_json::<_, RuntimeHttpWorkerResponse>( + &format!("/v1/workers/{worker_id}/workspace-api"), + &RuntimeHttpWorkerWorkspaceApiRequest { workspace_api }, + ) { + Ok(response) => WorkerWorkspaceApiResult { + state: WorkerOperationState::Accepted, + worker: Some(self.map_worker_detail(response.worker)), + diagnostics: Vec::new(), + }, + Err(diagnostic) => WorkerWorkspaceApiResult { + state: WorkerOperationState::Rejected, + worker: None, + diagnostics: vec![diagnostic], + }, + } + } + fn create_working_directory( &self, request: WorkingDirectoryRequest, @@ -2711,6 +2823,17 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { .map_or((None, None), |(key, fingerprint)| { (Some(key), Some(fingerprint)) }); + let workspace_api = match required_worker_workspace_api(&request) { + Ok(workspace_api) => workspace_api, + Err(diagnostic) => { + return WorkerSpawnResult { + state: WorkerOperationState::Rejected, + worker: None, + acceptance_evidence: Vec::new(), + diagnostics: vec![diagnostic], + }; + } + }; let create = CreateWorkerRequest { idempotency_key, idempotency_fingerprint, @@ -2721,14 +2844,7 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { initial_input: request.initial_input.clone(), working_directory_request: request.resolved_working_directory_request.clone(), working_directory: request.resolved_working_directory.clone(), - workspace_api: request.resolved_workspace_api.clone().or_else(|| { - Some(WorkspaceApiRef { - workspace_id: self.workspace_id.clone(), - base_url: self.backend_base_url.clone(), - runtime_id: Some(self.runtime_id.clone()), - access_token: None, - }) - }), + workspace_api: Some(workspace_api), }; match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) { Ok(response) => WorkerSpawnResult { @@ -3700,6 +3816,15 @@ mod tests { use std::sync::{Arc, Mutex}; use std::thread; + fn test_workspace_api() -> WorkspaceApiRef { + WorkspaceApiRef { + workspace_id: "workspace-test".to_string(), + base_url: "http://127.0.0.1:8787".to_string(), + runtime_id: Some("runtime-test".to_string()), + access_token: Some("workspace-access-token".to_string()), + } + } + #[test] fn embedded_builtin_decodal_profiles_resolve_through_archive() { let root = tempfile::tempdir().unwrap(); @@ -4222,10 +4347,31 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_workspace_api()), } } + #[test] + fn embedded_runtime_rejects_tokenless_workspace_spawn() { + let runtime = EmbeddedWorkerRuntime::new_memory_with_execution_backend( + "local:test", + Arc::new(AcceptingExecutionBackend::default()), + ) + .expect("test backend should connect"); + let mut request = embedded_spawn_request(); + request.resolved_workspace_api = None; + + let spawned = runtime.spawn_worker(request); + + assert_eq!(spawned.state, WorkerOperationState::Rejected); + assert!( + spawned + .diagnostics + .iter() + .any(|diagnostic| { diagnostic.code == "worker_workspace_credential_missing" }) + ); + } + #[test] fn embedded_runtime_spawn_execution_failure_is_rejected_and_not_input_capable() { let runtime = EmbeddedWorkerRuntime::new_memory_with_execution_backend( @@ -4350,7 +4496,7 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_workspace_api()), }, ) .unwrap(); @@ -4448,7 +4594,7 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_workspace_api()), }, ) .unwrap(); @@ -4482,7 +4628,7 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_workspace_api()), }, ) .unwrap(); diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index bd0c7d37..1d853be3 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -8,7 +8,7 @@ use axum::extract::{Path as AxumPath, Query, State}; use axum::http::header::{CONTENT_TYPE, ETAG, IF_NONE_MATCH, LOCATION, ORIGIN, SET_COOKIE}; use axum::http::{HeaderMap, StatusCode, Uri}; use axum::response::{IntoResponse, Response}; -use axum::routing::{delete, get, post, put}; +use axum::routing::{delete, get, patch, post, put}; use axum::{Json, Router}; use chrono::{Duration, SecondsFormat, Utc}; use futures::{SinkExt, StreamExt}; @@ -26,7 +26,7 @@ use ticket::{ TicketTargetEdit, TicketWorkflowState, }; use ticket::{ - SqliteTicketBackend, TicketBackendHttpResponse, TicketBackendOperation, + SqliteTicketBackend, TicketBackendOperation, TicketBackendOperationResult, execute_ticket_backend_operation, }; use tokio::net::TcpListener; @@ -65,7 +65,8 @@ use crate::hosts::{ WorkerInputKind, WorkerInputRequest, WorkerInputResult, WorkerLifecycleRequest, WorkerLifecycleResult, WorkerOperationState, WorkerRestoreResult, WorkerSpawnAcceptanceRequirement, WorkerSpawnIntent, WorkerSpawnRequest, WorkerSpawnResult, - WorkerSpawnWorkingDirectoryRequest, WorkerSummary, WorkerWorkspaceSummary, + WorkerSpawnWorkingDirectoryRequest, WorkerSummary, WorkerWorkspaceApiResult, + WorkerWorkspaceSummary, }; use crate::identity::WorkspaceIdentity; use crate::memory_backend::execute_memory_backend_operation_with_authority; @@ -101,7 +102,7 @@ use worker_runtime::catalog::{ ConfigBundleRef, MaterializerKind, ProfileSelector, RepositorySelector as RuntimeRepositorySelector, WorkingDirectoryClaim, WorkingDirectoryOccupancy, WorkingDirectoryRepository, WorkingDirectoryRequest, - WorkingDirectoryStatusKind, WorkingDirectorySummary, + WorkingDirectoryStatusKind, WorkingDirectorySummary, WorkspaceApiRef, }; use worker_runtime::config_bundle::ConfigBundle; use worker_runtime::http_server::{ @@ -248,6 +249,7 @@ pub struct WorkspaceApi { companion: Arc, observation_proxy: BackendObservationProxy, resource_broker: BackendResourceBroker, + credential_operation_lock: Arc>, } impl WorkspaceApi { @@ -311,14 +313,6 @@ impl WorkspaceApi { execution_backend, ) .map(|runtime| runtime.with_resource_broker(resource_broker.clone())) - .map(|runtime| { - runtime.with_backend_base_url( - config - .backend_base_url - .clone() - .unwrap_or_else(|| "http://127.0.0.1:8787".to_string()), - ) - }) .map_err(|err| { crate::Error::Store(format!("invalid embedded Worker backend: {err}")) })?, @@ -351,6 +345,7 @@ impl WorkspaceApi { companion, observation_proxy, resource_broker, + credential_operation_lock: Arc::new(std::sync::Mutex::new(())), }) } @@ -358,6 +353,164 @@ impl WorkspaceApi { self.config.workspace_id.as_str() } + fn mint_worker_workspace_credential( + &self, + runtime_id: &str, + worker_id: Option<&str>, + ) -> ApiResult<(WorkerWorkspaceCredentialRecord, WorkspaceApiRef)> { + let now = Utc::now(); + let token = mint_secret("wac"); + let credential = WorkerWorkspaceCredentialRecord { + credential_id: new_id("wac"), + token: token.clone(), + workspace_id: self.config.workspace_id.clone(), + runtime_id: runtime_id.to_string(), + worker_id: worker_id.map(ToOwned::to_owned), + created_at: now.to_rfc3339_opts(SecondsFormat::Secs, true), + expires_at: (now + chrono::Duration::hours(1)) + .to_rfc3339_opts(SecondsFormat::Secs, true), + revoked_at: None, + }; + self.store.upsert_worker_workspace_credential(&credential)?; + Ok(( + credential, + WorkspaceApiRef { + workspace_id: self.config.workspace_id.clone(), + base_url: self + .config + .backend_base_url + .clone() + .unwrap_or_else(|| "http://127.0.0.1:8787".to_string()) + .trim_end_matches('/') + .to_string(), + runtime_id: Some(runtime_id.to_string()), + access_token: Some(token), + }, + )) + } + + fn revoke_credential_record( + &self, + credential: &mut WorkerWorkspaceCredentialRecord, + ) -> ApiResult<()> { + credential.revoked_at = Some(Utc::now().to_rfc3339()); + self.store.upsert_worker_workspace_credential(credential)?; + Ok(()) + } + + fn spawn_workspace_worker( + &self, + runtime_id: &str, + mut request: WorkerSpawnRequest, + ) -> ApiResult { + let (mut credential, workspace_api) = + self.mint_worker_workspace_credential(runtime_id, None)?; + request.resolved_workspace_api = Some(workspace_api.clone()); + let result = match self.runtime.spawn_worker(runtime_id, request) { + Ok(result) => result, + Err(error) => { + self.revoke_credential_record(&mut credential)?; + return Err(error.into_error().into()); + } + }; + let Some(worker) = result.worker.as_ref() else { + self.revoke_credential_record(&mut credential)?; + return Ok(result); + }; + let replacement = match self.runtime.replace_worker_workspace_api( + runtime_id, + &worker.worker_id, + workspace_api, + ) { + Ok(replacement) => replacement, + Err(error) => { + let _ = self.runtime.delete_worker(runtime_id, &worker.worker_id); + self.revoke_credential_record(&mut credential)?; + return Err(error.into_error().into()); + } + }; + if replacement.state != WorkerOperationState::Accepted { + let _ = self.runtime.delete_worker(runtime_id, &worker.worker_id); + self.revoke_credential_record(&mut credential)?; + return Err(Error::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: "worker_workspace_api_replace_failed".to_string(), + message: replacement + .diagnostics + .first() + .map(|diagnostic| diagnostic.message.clone()) + .unwrap_or_else(|| { + "Runtime rejected Workspace API replacement after spawn".to_string() + }), + } + .into()); + } + credential.worker_id = Some(worker.worker_id.clone()); + self.store.upsert_worker_workspace_credential(&credential)?; + self.store.revoke_worker_workspace_credentials_except( + &self.config.workspace_id, + runtime_id, + &worker.worker_id, + &credential.credential_id, + &Utc::now().to_rfc3339(), + )?; + Ok(result) + } + + fn rotate_worker_workspace_credential( + &self, + runtime_id: &str, + worker_id: &str, + ) -> ApiResult { + let _credential_guard = self.credential_operation_lock.lock().map_err(|_| { + Error::Config("Workspace credential operation lock poisoned".to_string()) + })?; + let (mut credential, workspace_api) = + self.mint_worker_workspace_credential(runtime_id, Some(worker_id))?; + let result = + match self + .runtime + .replace_worker_workspace_api(runtime_id, worker_id, workspace_api) + { + Ok(result) => result, + Err(error) => { + self.revoke_credential_record(&mut credential)?; + return Err(error.into_error().into()); + } + }; + if result.state != WorkerOperationState::Accepted { + self.revoke_credential_record(&mut credential)?; + return Ok(result); + } + self.store.revoke_worker_workspace_credentials_except( + &self.config.workspace_id, + runtime_id, + worker_id, + &credential.credential_id, + &Utc::now().to_rfc3339(), + )?; + Ok(result) + } + + fn restore_workspace_worker( + &self, + runtime_id: &str, + worker_id: &str, + ) -> ApiResult { + let rotation = self.rotate_worker_workspace_credential(runtime_id, worker_id)?; + if rotation.state != WorkerOperationState::Accepted { + return Ok(WorkerRestoreResult { + state: rotation.state, + worker: rotation.worker, + diagnostics: rotation.diagnostics, + }); + } + Ok(self + .runtime + .restore_worker(runtime_id, worker_id) + .map_err(|error| error.into_error())?) + } + fn repository_reader(&self) -> RepositoryRegistryReader { RepositoryRegistryReader::new(self.config.repositories.clone()) } @@ -489,15 +642,14 @@ pub fn build_router(api: WorkspaceApi) -> Router { .delete(scoped_delete_profile_source), ) .route("/api/tickets", get(list_tickets)) - .route("/api/w/{workspace_id}/tickets", get(scoped_list_tickets)) + .route( + "/api/w/{workspace_id}/tickets", + get(scoped_list_tickets).post(scoped_create_ticket_record), + ) .route( "/api/w/{workspace_id}/worker-credentials/refresh", post(scoped_refresh_worker_workspace_credential), ) - .route( - "/api/w/{workspace_id}/tickets/backend", - post(scoped_ticket_backend_operation), - ) .route( "/api/w/{workspace_id}/memory", get(scoped_get_memory_document), @@ -522,6 +674,86 @@ pub fn build_router(api: WorkspaceApi) -> Router { "/api/w/{workspace_id}/skills/{name}/activate", get(scoped_activate_skill), ) + .route( + "/api/w/{workspace_id}/tickets/default-intake-ready-body", + post(scoped_default_intake_ready_body), + ) + .route( + "/api/w/{workspace_id}/tickets/search", + get(scoped_list_ticket_summaries), + ) + .route( + "/api/w/{workspace_id}/tickets/doctor", + get(scoped_ticket_doctor), + ) + .route( + "/api/w/{workspace_id}/tickets/relations/search", + post(scoped_query_ticket_relations), + ) + .route( + "/api/w/{workspace_id}/tickets/orchestration-plans/search", + post(scoped_query_ticket_orchestration_plans), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/record", + get(scoped_get_ticket_record), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/item", + patch(scoped_edit_ticket_record_item), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/dependency-check", + get(scoped_ticket_dependency_check), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/thread-events", + post(scoped_add_ticket_thread_event), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/state-changes", + post(scoped_add_ticket_state_change), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/intake-summaries", + post(scoped_add_ticket_intake_summary), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/state-fields/{field}", + post(scoped_set_ticket_state_field), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/workflow-state", + post(scoped_set_ticket_workflow_state), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/intake-ready", + post(scoped_prepare_ticket_intake_ready), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/workflow/queue", + post(scoped_queue_ticket_record), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/workflow/review", + post(scoped_review_ticket_record), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/workflow/close", + post(scoped_close_ticket_record), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/relation-view", + get(scoped_ticket_relation_view), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/relations", + post(scoped_record_ticket_relation), + ) + .route( + "/api/w/{workspace_id}/tickets/{id}/orchestration-plans", + post(scoped_record_ticket_orchestration_plan), + ) .route( "/api/w/{workspace_id}/tickets/{id}", get(scoped_get_ticket).patch(scoped_edit_ticket_item), @@ -2022,6 +2254,10 @@ async fn scoped_refresh_worker_workspace_credential( headers: HeaderMap, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; + let _credential_guard = api + .credential_operation_lock + .lock() + .map_err(|_| Error::Config("Workspace credential operation lock poisoned".to_string()))?; let token = headers .get(axum::http::header::AUTHORIZATION) .and_then(|value| value.to_str().ok()) @@ -2063,13 +2299,19 @@ async fn scoped_refresh_worker_workspace_credential( })) } -async fn scoped_ticket_backend_operation( - State(api): State, - AxumPath(path): AxumPath, +async fn execute_worker_ticket_rest_operation( + api: &WorkspaceApi, + workspace_id: &str, headers: HeaderMap, - Json(mut operation): Json, -) -> ApiResult> { - validate_workspace_scope(&api, &path.workspace_id)?; + mut operation: TicketBackendOperation, +) -> ApiResult { + validate_workspace_scope(api, workspace_id)?; + if headers.get(axum::http::header::AUTHORIZATION).is_none() { + return Err(Error::WorkerWorkspaceAuthentication( + "missing Runtime Workspace credential".to_string(), + ) + .into()); + } let config = ticket::config::TicketConfig::load_workspace(&api.config.workspace_root) .map_err(|error| Error::Config(format!("load Ticket workspace settings: {error}")))?; let mut backend = SqliteTicketBackend::new( @@ -2081,64 +2323,535 @@ async fn scoped_ticket_backend_operation( let is_mutation = operation_kind != "read"; let target = ticket_mutation_target(&operation).cloned(); let read_target = ticket_read_target(&operation).cloned(); - let has_worker_credential = headers.contains_key(axum::http::header::AUTHORIZATION); - let source = if is_mutation || (read_target.is_some() && has_worker_credential) { - Some(authenticate_worker_mutation_source( - &api, - &path.workspace_id, - &headers, - )?) - } else { - None - }; + let source = authenticate_worker_mutation_source(api, workspace_id, &headers)?; let before = target.as_ref().and_then(|id| backend.show(id.clone()).ok()); - if let Some(source) = source.as_ref() { - bind_worker_ticket_operation_source(source, &mut operation); - let source_context = - worker_ticket_source_context(&api, &path.workspace_id, source, before.as_ref()); - backend = backend - .with_event_attributes(source_context.attributes(operation_kind)) - .with_mutation_hook(build_ticket_notification_hook( - &api, - source_context, - operation_kind, - before - .as_ref() - .map(|ticket| ticket.meta.workflow_state.as_str().to_string()) - .unwrap_or_else(|| ticket_operation_initial_state(&operation)), - )); + bind_worker_ticket_operation_source(&source, &mut operation); + let source_context = worker_ticket_source_context(api, workspace_id, &source, before.as_ref()); + backend = backend + .with_event_attributes(source_context.attributes(operation_kind)) + .with_mutation_hook(build_ticket_notification_hook( + api, + source_context, + operation_kind, + before + .as_ref() + .map(|ticket| ticket.meta.workflow_state.as_str().to_string()) + .unwrap_or_else(|| ticket_operation_initial_state(&operation)), + )); + + let result = execute_ticket_backend_operation(&backend, operation).map_err(Error::from)?; + if let Some(read_target) = read_target.as_ref() + && let Ok(ticket) = backend.show(read_target.clone()) + && let Some(event_index) = ticket.events.last().and_then(|event| { + event + .attributes + .get("event_sequence") + .and_then(|value| value.parse::().ok()) + }) + { + api.store.upsert_ticket_notification_cursor( + workspace_id, + &ticket.meta.id, + &source.runtime_id, + &source.worker_id, + event_index, + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + )?; } - let response = match execute_ticket_backend_operation(&backend, operation) { - Ok(result) => { - if let (Some(source), Some(read_target)) = (source.as_ref(), read_target.as_ref()) { - if let Ok(ticket) = backend.show(read_target.clone()) { - if let Some(event_index) = ticket.events.last().and_then(|event| { - event - .attributes - .get("event_sequence") - .and_then(|value| value.parse::().ok()) - }) { - api.store.upsert_ticket_notification_cursor( - &path.workspace_id, - &ticket.meta.id, - &source.runtime_id, - &source.worker_id, - event_index, - &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), - )?; - } - } + if is_mutation { + dispatch_pending_ticket_notifications(api, workspace_id); + } + Ok(result) +} + +#[cfg(test)] +async fn execute_worker_ticket_test_operation( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Json(operation): Json, +) -> ApiResult> { + execute_worker_ticket_rest_operation(&api, &path.workspace_id, headers, operation) + .await + .map(Json) +} + +fn ticket_rest_result( + result: TicketBackendOperationResult, + extract: impl FnOnce(TicketBackendOperationResult) -> Option, +) -> ApiResult> { + extract(result).map(Json).ok_or_else(|| { + Error::Config("Ticket REST handler received an unexpected backend result".to_string()) + .into() + }) +} + +fn ticket_rest_unit(result: TicketBackendOperationResult) -> ApiResult { + match result { + TicketBackendOperationResult::Unit => Ok(StatusCode::NO_CONTENT), + _ => Err(Error::Config( + "Ticket REST handler received an unexpected backend result".to_string(), + ) + .into()), + } +} + +#[derive(Debug, Deserialize)] +struct DefaultIntakeReadyBodyRequest { + from: String, +} + +async fn scoped_default_intake_ready_body( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Json(request): Json, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &path.workspace_id, + headers, + TicketBackendOperation::DefaultIntakeReadyStateChangeBody { from: request.from }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::Text(body) => Some(body), + _ => None, + }) +} + +#[derive(Debug, Default, Deserialize)] +struct TicketSummarySearchQuery { + state: Option, +} + +async fn scoped_list_ticket_summaries( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Query(query): Query, +) -> ApiResult>> { + let filter = match query.state.as_deref().unwrap_or("active") { + "active" => ticket::TicketListQuery::active(), + "all" => ticket::TicketListQuery::all(), + states => { + let mut selected = Vec::new(); + for state in states.split(',') { + selected.push(ticket::TicketListState::parse(state).ok_or_else(|| { + Error::Ticket(ticket::TicketError::InvalidPathComponent(state.to_string())) + })?); } - if is_mutation && source.is_some() { - dispatch_pending_ticket_notifications(&api, &path.workspace_id); - } - TicketBackendHttpResponse::Ok { result } + ticket::TicketListQuery::states(selected) } - Err(error) => TicketBackendHttpResponse::Error { - message: error.to_string(), - }, }; - Ok(Json(response)) + let result = execute_worker_ticket_rest_operation( + &api, + &path.workspace_id, + headers, + TicketBackendOperation::List { filter }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::Tickets(tickets) => Some(tickets), + _ => None, + }) +} + +async fn scoped_get_ticket_record( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::Show { + id: TicketIdOrSlug::Query(id), + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::Ticket(ticket) => Some(ticket), + _ => None, + }) +} + +async fn scoped_create_ticket_record( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Json(input): Json, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &path.workspace_id, + headers, + TicketBackendOperation::Create { input }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::TicketRef(ticket) => Some(ticket), + _ => None, + }) +} + +async fn scoped_edit_ticket_record_item( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(edit): Json, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::EditItem { + id: TicketIdOrSlug::Query(id), + edit, + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::Ticket(ticket) => Some(ticket), + _ => None, + }) +} + +async fn scoped_ticket_dependency_check( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::DependencyCheck { + id: TicketIdOrSlug::Query(id), + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::DependencyCheck(check) => Some(check), + _ => None, + }) +} + +async fn scoped_add_ticket_thread_event( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(event): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::AddEvent { + id: TicketIdOrSlug::Query(id), + event, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_add_ticket_state_change( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(change): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::AddStateChanged { + id: TicketIdOrSlug::Query(id), + change, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_add_ticket_intake_summary( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(summary): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::AddIntakeSummary { + id: TicketIdOrSlug::Query(id), + summary, + }, + ) + .await?; + ticket_rest_unit(result) +} + +#[derive(Debug, Deserialize)] +struct TicketIntakeReadyRequest { + summary: ticket::TicketIntakeSummary, + change: TicketStateChange, +} + +async fn scoped_set_ticket_state_field( + State(api): State, + AxumPath((workspace_id, id, field)): AxumPath<(String, String, String)>, + headers: HeaderMap, + Json(change): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::SetStateField { + id: TicketIdOrSlug::Query(id), + field, + change, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_set_ticket_workflow_state( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(change): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::SetWorkflowState { + id: TicketIdOrSlug::Query(id), + change, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_prepare_ticket_intake_ready( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(request): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::MarkIntakeReady { + id: TicketIdOrSlug::Query(id), + summary: request.summary, + change: request.change, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_queue_ticket_record( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::QueueReady { + id: TicketIdOrSlug::Query(id), + queued_by: String::new(), + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_review_ticket_record( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(review): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::Review { + id: TicketIdOrSlug::Query(id), + review, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_close_ticket_record( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(resolution): Json, +) -> ApiResult { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::Close { + id: TicketIdOrSlug::Query(id), + resolution, + }, + ) + .await?; + ticket_rest_unit(result) +} + +async fn scoped_record_ticket_relation( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(relation): Json, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::AddTicketRelation { + id: TicketIdOrSlug::Query(id), + relation, + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::Relation(relation) => Some(relation), + _ => None, + }) +} + +#[derive(Debug, Deserialize)] +struct TicketRelationSearchRequest { + ticket: Option, + kind: Option, +} + +async fn scoped_query_ticket_relations( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Json(query): Json, +) -> ApiResult>> { + let result = execute_worker_ticket_rest_operation( + &api, + &path.workspace_id, + headers, + TicketBackendOperation::QueryTicketRelations { + ticket: query.ticket, + kind: query.kind, + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::Relations(relations) => Some(relations), + _ => None, + }) +} + +async fn scoped_ticket_relation_view( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::RelationView { + id: TicketIdOrSlug::Query(id), + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::RelationView(view) => Some(view), + _ => None, + }) +} + +async fn scoped_record_ticket_orchestration_plan( + State(api): State, + AxumPath((workspace_id, id)): AxumPath<(String, String)>, + headers: HeaderMap, + Json(record): Json, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &workspace_id, + headers, + TicketBackendOperation::AddOrchestrationPlanRecord { + id: TicketIdOrSlug::Query(id), + record, + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::OrchestrationPlanRecord(record) => Some(record), + _ => None, + }) +} + +#[derive(Debug, Deserialize)] +struct TicketOrchestrationPlanSearchRequest { + ticket: Option, + kind: Option, +} + +async fn scoped_query_ticket_orchestration_plans( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, + Json(query): Json, +) -> ApiResult>> { + let result = execute_worker_ticket_rest_operation( + &api, + &path.workspace_id, + headers, + TicketBackendOperation::QueryOrchestrationPlanRecords { + ticket: query.ticket, + kind: query.kind, + }, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::OrchestrationPlanRecords(records) => Some(records), + _ => None, + }) +} + +async fn scoped_ticket_doctor( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, +) -> ApiResult> { + let result = execute_worker_ticket_rest_operation( + &api, + &path.workspace_id, + headers, + TicketBackendOperation::Doctor, + ) + .await?; + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::DoctorReport(report) => Some(report), + _ => None, + }) } #[derive(Debug, Clone, PartialEq, Eq)] @@ -2697,27 +3410,24 @@ fn start_memory_staging_consolidation( content: input_content, segments: None, }; - let result = api - .runtime - .spawn_worker( - &runtime_id, - WorkerSpawnRequest { - requested_worker_name: Some(MEMORY_CONSOLIDATION_PROFILE.to_string()), - intent: WorkerSpawnIntent::WorkspaceOrchestrator, - acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { - expected_segments: 1, - }, - profile: profile_selector, - ticket_assignment: None, - initial_input: Some(input), - working_directory_request: None, - resolved_working_directory_request: None, - resolved_working_directory: None, - resolved_config_bundle, - resolved_workspace_api: None, + let result = api.spawn_workspace_worker( + &runtime_id, + WorkerSpawnRequest { + requested_worker_name: Some(MEMORY_CONSOLIDATION_PROFILE.to_string()), + intent: WorkerSpawnIntent::WorkspaceOrchestrator, + acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { + expected_segments: 1, }, - ) - .map_err(|err| err.into_error())?; + profile: profile_selector, + ticket_assignment: None, + initial_input: Some(input), + working_directory_request: None, + resolved_working_directory_request: None, + resolved_working_directory: None, + resolved_config_bundle, + resolved_workspace_api: None, + }, + )?; if result.state != WorkerOperationState::Accepted { return Ok(MemoryConsolidationOutput { status: "skipped_spawn_rejected".to_string(), @@ -5410,27 +6120,24 @@ async fn create_workspace_worker( if resolved_working_directory.is_none() { reject_no_workdir_for_non_embedded_runtime(&request.runtime_id)?; } - let result = api - .runtime - .spawn_worker( - &request.runtime_id, - WorkerSpawnRequest { - requested_worker_name: Some(display_name.clone()), - intent: WorkerSpawnIntent::WorkspaceCoding, - acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { - expected_segments: if initial_input.is_some() { 1 } else { 0 }, - }, - profile: profile_selector, - ticket_assignment: None, - initial_input, - working_directory_request: None, - resolved_working_directory_request: None, - resolved_working_directory, - resolved_config_bundle, - resolved_workspace_api: None, + let result = api.spawn_workspace_worker( + &request.runtime_id, + WorkerSpawnRequest { + requested_worker_name: Some(display_name.clone()), + intent: WorkerSpawnIntent::WorkspaceCoding, + acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { + expected_segments: if initial_input.is_some() { 1 } else { 0 }, }, - ) - .map_err(|err| err.into_error())?; + profile: profile_selector, + ticket_assignment: None, + initial_input, + working_directory_request: None, + resolved_working_directory_request: None, + resolved_working_directory, + resolved_config_bundle, + resolved_workspace_api: None, + }, + )?; Ok(Json(record_browser_worker_spawn( &api, request.runtime_id, @@ -5604,10 +6311,7 @@ async fn restore_runtime_worker( State(api): State, AxumPath((runtime_id, worker_id)): AxumPath<(String, String)>, ) -> ApiResult> { - let mut result = api - .runtime - .restore_worker(&runtime_id, &worker_id) - .map_err(|err| err.into_error())?; + let mut result = api.restore_workspace_worker(&runtime_id, &worker_id)?; if let Some(worker) = result.worker.as_ref() { let record = sync_worker_observation(&api, worker)?; let links = api.store.list_worker_workdir_links( @@ -5762,26 +6466,6 @@ async fn create_runtime_worker( .map(|claim| claim.working_directory_id.clone()) }; let requested_worker_name = request.requested_worker_name.clone(); - if let Some(base_url) = api.config.backend_base_url.clone() { - let credential = WorkerWorkspaceCredentialRecord { - credential_id: new_id("wac"), - token: mint_secret("wac"), - workspace_id: api.config.workspace_id.clone(), - runtime_id: runtime_id.clone(), - worker_id: None, - created_at: Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), - expires_at: (Utc::now() + chrono::Duration::hours(1)) - .to_rfc3339_opts(SecondsFormat::Secs, true), - revoked_at: None, - }; - api.store.upsert_worker_workspace_credential(&credential)?; - request.resolved_workspace_api = Some(worker_runtime::catalog::WorkspaceApiRef { - workspace_id: api.config.workspace_id.clone(), - base_url, - runtime_id: Some(runtime_id.clone()), - access_token: Some(credential.token), - }); - } let spawn_idempotency = crate::hosts::worker_spawn_idempotency(&request).map_err(Error::Config)?; if let (Some(assignment), Some((_, fingerprint))) = @@ -5797,10 +6481,7 @@ async fn create_runtime_worker( &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), )?; } - let result = api - .runtime - .spawn_worker(&runtime_id, request) - .map_err(|err| err.into_error())?; + let result = api.spawn_workspace_worker(&runtime_id, request)?; if let Some(worker) = result.worker.as_ref() { let display_name = requested_worker_name .as_deref() @@ -8042,6 +8723,15 @@ mod tests { const TEST_REPOSITORY_ID: &str = "main"; const TEST_CREATED_AT: &str = "2026-06-23T06:43:28Z"; + fn test_worker_workspace_api(runtime_id: &str) -> WorkspaceApiRef { + WorkspaceApiRef { + workspace_id: TEST_WORKSPACE_ID.to_string(), + base_url: "http://127.0.0.1:8787".to_string(), + runtime_id: Some(runtime_id.to_string()), + access_token: Some("workspace-access-token".to_string()), + } + } + #[test] fn ticket_api_errors_preserve_http_status() { let not_found = ApiError::from(Error::Ticket(ticket::TicketError::NotFound( @@ -8693,6 +9383,17 @@ mod tests { } } + fn replace_workspace_access_token( + &self, + _handle: &worker_runtime::execution::WorkerExecutionHandle, + _access_token: String, + ) -> worker_runtime::execution::WorkerExecutionResult { + worker_runtime::execution::WorkerExecutionResult::accepted( + worker_runtime::execution::WorkerExecutionOperation::ReplaceWorkspaceAccessToken, + worker_runtime::execution::WorkerExecutionRunState::Idle, + ) + } + fn dispatch_input( &self, handle: &worker_runtime::execution::WorkerExecutionHandle, @@ -8846,7 +9547,9 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_worker_workspace_api( + EMBEDDED_WORKER_RUNTIME_ID, + )), }, ) .unwrap(); @@ -9022,7 +9725,7 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_worker_workspace_api(EMBEDDED_WORKER_RUNTIME_ID)), }; let source_worker = api .runtime @@ -9078,23 +9781,29 @@ mod tests { "x-yoi-worker-id", axum::http::HeaderValue::from_str(&source_worker.worker_id).unwrap(), ); - let Json(response) = scoped_ticket_backend_operation( - State(api.clone()), - AxumPath(ScopedWorkspacePath { - workspace_id: TEST_WORKSPACE_ID.to_string(), - }), - headers.clone(), - Json(TicketBackendOperation::AddEvent { - id: ticket_ref.id.clone().into(), - event: NewTicketEvent::new(TicketEventKind::Comment, "implementation update"), - }), - ) - .await - .unwrap(); - assert!( - matches!(response, TicketBackendHttpResponse::Ok { .. }), - "unexpected response: {response:?}" - ); + let response = build_router(api.clone()) + .oneshot( + Request::builder() + .method("POST") + .uri(format!( + "/api/w/{TEST_WORKSPACE_ID}/tickets/{}/thread-events", + ticket_ref.id + )) + .header("content-type", "application/json") + .header("authorization", "Bearer source-secret") + .header("x-yoi-worker-id", &source_worker.worker_id) + .body(Body::from( + serde_json::to_vec(&NewTicketEvent::new( + TicketEventKind::Comment, + "implementation update", + )) + .unwrap(), + )) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::NO_CONTENT); let committed = backend.show(ticket_ref.id.clone().into()).unwrap(); let committed_event = committed.events.last().unwrap(); assert_eq!( @@ -9126,18 +9835,22 @@ mod tests { .unwrap() .parse::() .unwrap(); - let _ = scoped_ticket_backend_operation( - State(api.clone()), - AxumPath(ScopedWorkspacePath { - workspace_id: TEST_WORKSPACE_ID.to_string(), - }), - headers.clone(), - Json(TicketBackendOperation::Show { - id: ticket_ref.id.clone().into(), - }), - ) - .await - .unwrap(); + let response = build_router(api.clone()) + .oneshot( + Request::builder() + .method("GET") + .uri(format!( + "/api/w/{TEST_WORKSPACE_ID}/tickets/{}/record", + ticket_ref.id + )) + .header("authorization", "Bearer source-secret") + .header("x-yoi-worker-id", &source_worker.worker_id) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); assert_eq!( api.store .get_ticket_notification_cursor( @@ -9157,7 +9870,7 @@ mod tests { "accepted Runtime system input must complete the outbox delivery" ); - let Json(stale_report) = scoped_ticket_backend_operation( + let Json(stale_report) = execute_worker_ticket_test_operation( State(api.clone()), AxumPath(ScopedWorkspacePath { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -9173,7 +9886,7 @@ mod tests { ) .await .unwrap(); - assert!(matches!(stale_report, TicketBackendHttpResponse::Ok { .. })); + assert!(matches!(stale_report, TicketBackendOperationResult::Unit)); let non_assigned_report = backend.show(ticket_ref.id.clone().into()).unwrap(); assert!( !non_assigned_report @@ -9201,7 +9914,7 @@ mod tests { true, ) .unwrap(); - let _ = scoped_ticket_backend_operation( + let _ = execute_worker_ticket_test_operation( State(api.clone()), AxumPath(ScopedWorkspacePath { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -9234,7 +9947,7 @@ mod tests { Some("coder") ); - let unauthorized = scoped_ticket_backend_operation( + let unauthorized = execute_worker_ticket_test_operation( State(api.clone()), AxumPath(ScopedWorkspacePath { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -9275,7 +9988,9 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_worker_workspace_api( + EMBEDDED_WORKER_RUNTIME_ID, + )), }, ) .unwrap() @@ -9298,7 +10013,9 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_worker_workspace_api( + EMBEDDED_WORKER_RUNTIME_ID, + )), }, ) .unwrap() @@ -9339,7 +10056,7 @@ mod tests { "x-yoi-worker-id", axum::http::HeaderValue::from_str(&source.worker_id).unwrap(), ); - let _ = scoped_ticket_backend_operation( + let _ = execute_worker_ticket_test_operation( State(api.clone()), AxumPath(ScopedWorkspacePath { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -9365,6 +10082,66 @@ mod tests { ); } + #[tokio::test] + async fn workspace_credential_rotation_revokes_prior_bound_credential() { + let dir = tempfile::tempdir().unwrap(); + let api = test_api(dir.path()).await; + let result = api + .spawn_workspace_worker( + EMBEDDED_WORKER_RUNTIME_ID, + WorkerSpawnRequest { + requested_worker_name: Some("credential-boundary".to_string()), + intent: WorkerSpawnIntent::WorkspaceCoding, + acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { + expected_segments: 0, + }, + profile: ProfileSelector::Builtin("builtin:coder".to_string()), + ticket_assignment: None, + initial_input: None, + working_directory_request: None, + resolved_working_directory_request: None, + resolved_working_directory: None, + resolved_config_bundle: None, + resolved_workspace_api: Some(test_worker_workspace_api( + "embedded-worker-runtime", + )), + }, + ) + .unwrap(); + let worker = result.worker.unwrap(); + let old_token = "legacy-worker-token"; + api.store + .upsert_worker_workspace_credential(&WorkerWorkspaceCredentialRecord { + credential_id: new_id("wac"), + token: old_token.to_string(), + workspace_id: TEST_WORKSPACE_ID.to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + worker_id: Some(worker.worker_id.clone()), + created_at: Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + expires_at: (Utc::now() + chrono::Duration::hours(1)) + .to_rfc3339_opts(SecondsFormat::Secs, true), + revoked_at: None, + }) + .unwrap(); + + let rotation = api + .rotate_worker_workspace_credential(EMBEDDED_WORKER_RUNTIME_ID, &worker.worker_id) + .unwrap(); + + assert_eq!(rotation.state, WorkerOperationState::Accepted); + assert!( + api.store + .authenticate_worker_workspace_credential( + old_token, + TEST_WORKSPACE_ID, + &worker.worker_id, + ) + .unwrap() + .is_none(), + "repair must revoke the prior bound credential" + ); + } + #[tokio::test] async fn worker_spawn_and_restore_assignment_operations_are_idempotent() { let dir = tempfile::tempdir().unwrap(); @@ -9517,13 +10294,15 @@ mod tests { TEST_CREATED_AT, ) .unwrap(); - let pending_request = WorkerSpawnRequest { + let mut pending_request = WorkerSpawnRequest { ticket_assignment: Some(crate::hosts::WorkerTicketAssignmentRequest { ticket_id: second_ticket.id.clone(), operation_id: "pending-spawn-operation".to_string(), }), ..request }; + pending_request.resolved_workspace_api = + Some(test_worker_workspace_api(EMBEDDED_WORKER_RUNTIME_ID)); let (_, pending_fingerprint) = crate::hosts::worker_spawn_idempotency(&pending_request) .unwrap() .unwrap(); @@ -9710,7 +10489,7 @@ mod tests { } #[tokio::test] - async fn ticket_backend_endpoint_uses_workspace_sqlite_backend() { + async fn ticket_rest_operations_use_workspace_sqlite_backend() { let dir = tempfile::tempdir().unwrap(); fs::create_dir_all(dir.path().join(".yoi")).unwrap(); fs::write( @@ -10568,6 +11347,41 @@ mod tests { ); } + #[tokio::test] + async fn ticket_rest_search_requires_worker_credential_and_rpc_route_is_removed() { + let dir = tempfile::tempdir().unwrap(); + let api = test_api(dir.path()).await; + let app = build_router(api); + + let response = app + .clone() + .oneshot( + Request::builder() + .method("GET") + .uri(format!( + "/api/w/{TEST_WORKSPACE_ID}/tickets/search?state=active" + )) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/w/{TEST_WORKSPACE_ID}/tickets/backend")) + .header("content-type", "application/json") + .body(Body::from("{}")) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::METHOD_NOT_ALLOWED); + } + #[tokio::test] async fn browser_worker_create_uses_workspace_default_and_preserves_unsupported_diagnostics() { let dir = tempfile::tempdir().unwrap(); @@ -11314,7 +12128,9 @@ mod tests { resolved_working_directory_request: None, resolved_working_directory: None, resolved_config_bundle: None, - resolved_workspace_api: None, + resolved_workspace_api: Some(test_worker_workspace_api( + "embedded-worker-runtime", + )), }, ) .expect("spawn worker"); diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 459c9123..77287faa 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -650,6 +650,14 @@ pub trait ControlPlaneStore: Send + Sync { new_token: &str, new_expires_at: &str, ) -> Result>; + fn revoke_worker_workspace_credentials_except( + &self, + workspace_id: &str, + runtime_id: &str, + worker_id: &str, + active_credential_id: &str, + revoked_at: &str, + ) -> Result<()>; fn revoke_worker_workspace_credentials( &self, workspace_id: &str, @@ -2338,6 +2346,31 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }) } + fn revoke_worker_workspace_credentials_except( + &self, + workspace_id: &str, + runtime_id: &str, + worker_id: &str, + active_credential_id: &str, + revoked_at: &str, + ) -> Result<()> { + self.with_conn(|conn| { + conn.execute( + r#"UPDATE worker_workspace_credentials SET revoked_at = ?5 + WHERE workspace_id = ?1 AND runtime_id = ?2 AND worker_id = ?3 + AND credential_id <> ?4 AND revoked_at IS NULL"#, + params![ + workspace_id, + runtime_id, + worker_id, + active_credential_id, + revoked_at + ], + )?; + Ok(()) + }) + } + fn revoke_worker_workspace_credentials( &self, workspace_id: &str,