diff --git a/crates/ticket/src/lib.rs b/crates/ticket/src/lib.rs index 999be609..88f1811c 100644 --- a/crates/ticket/src/lib.rs +++ b/crates/ticket/src/lib.rs @@ -9,6 +9,7 @@ use std::fmt; use std::fs::{self, File, OpenOptions}; use std::io::{self, Write}; use std::path::{Component, Path, PathBuf}; +use std::sync::Arc; use chrono::Utc; use fs4::fs_std::FileExt; @@ -2274,11 +2275,40 @@ impl LocalTicketBackend { } } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SqliteTicketMutationEvent { + pub workspace_id: String, + pub ticket_id: String, + pub event_index: i64, + pub event_kind: TicketEventKind, +} + +pub type SqliteTicketMutationHook = + dyn Fn(&Connection, &SqliteTicketMutationEvent) -> Result<()> + Send + Sync; + +#[derive(Clone)] pub struct SqliteTicketBackend { db_path: PathBuf, workspace_id: String, record_language: Option, + event_attributes: BTreeMap, + mutation_hook: Option>, +} + +impl fmt::Debug for SqliteTicketBackend { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("SqliteTicketBackend") + .field("db_path", &self.db_path) + .field("workspace_id", &self.workspace_id) + .field("record_language", &self.record_language) + .field("event_attributes", &self.event_attributes) + .field( + "mutation_hook", + &self.mutation_hook.as_ref().map(|_| "configured"), + ) + .finish() + } } impl SqliteTicketBackend { @@ -2287,9 +2317,21 @@ impl SqliteTicketBackend { db_path: db_path.into(), workspace_id: workspace_id.into(), record_language: None, + event_attributes: BTreeMap::new(), + mutation_hook: None, } } + pub fn with_event_attributes(mut self, attributes: BTreeMap) -> Self { + self.event_attributes = attributes; + self + } + + pub fn with_mutation_hook(mut self, hook: Arc) -> Self { + self.mutation_hook = Some(hook); + self + } + pub fn with_record_language(mut self, language: Option<&str>) -> Self { self.record_language = language.and_then(normalized_record_language); self @@ -2506,10 +2548,25 @@ CREATE TABLE IF NOT EXISTS typed_ticket_artifacts ( conn.execute("INSERT INTO typed_ticket_event_references (workspace_id, ticket_id, event_index, ordinal, kind, target) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", params![self.workspace_id, ticket_id, next_index, ordinal as i64, reference.kind, reference.target]).map_err(sqlite_err)?; } - for (key, value) in &event.attributes { + let mut attributes = event.attributes.clone(); + for (key, value) in &self.event_attributes { + attributes.insert(key.clone(), value.clone()); + } + for (key, value) in &attributes { conn.execute("INSERT INTO typed_ticket_event_attributes (workspace_id, ticket_id, event_index, key, value) VALUES (?1, ?2, ?3, ?4, ?5)", params![self.workspace_id, ticket_id, next_index, key, value]).map_err(sqlite_err)?; } + if let Some(hook) = &self.mutation_hook { + hook( + conn, + &SqliteTicketMutationEvent { + workspace_id: self.workspace_id.clone(), + ticket_id: ticket_id.to_string(), + event_index: next_index, + event_kind: event.kind.clone(), + }, + )?; + } Ok(()) } @@ -2718,6 +2775,9 @@ CREATE TABLE IF NOT EXISTS typed_ticket_artifacts ( for row in rows { let (index, kind, author, at, status, from, to, reason, state_field, heading, body) = row.map_err(sqlite_err)?; + let mut attributes = self.load_event_attributes(conn, ticket_id, index)?; + attributes.insert("event_id".to_string(), format!("{ticket_id}:{index}")); + attributes.insert("event_sequence".to_string(), index.to_string()); events.push(TicketEvent { kind: TicketEventKind::from(kind.as_str()), author, @@ -2730,7 +2790,7 @@ CREATE TABLE IF NOT EXISTS typed_ticket_artifacts ( heading, body: MarkdownText::new(body), references: self.load_event_references(conn, ticket_id, index)?, - attributes: self.load_event_attributes(conn, ticket_id, index)?, + attributes, }); } Ok(events) @@ -6393,6 +6453,41 @@ state: planning assert_partial_body_replacement_semantics(&backend); } + #[test] + fn sqlite_mutation_hook_failure_rolls_back_ticket_event() { + let tmp = TempDir::new().unwrap(); + let db_path = tmp.path().join("workspace.db"); + let backend = SqliteTicketBackend::new(&db_path, "workspace-test"); + let created = backend.create(NewTicket::new("Atomic mutation")).unwrap(); + let before = backend + .show(TicketIdOrSlug::Id(created.id.clone())) + .unwrap() + .events + .len(); + let failing = backend.clone().with_mutation_hook(Arc::new(|_, event| { + Err(TicketError::Conflict(format!( + "reject outbox for {}:{}", + event.ticket_id, event.event_index + ))) + })); + assert!( + failing + .add_event( + TicketIdOrSlug::Id(created.id.clone()), + NewTicketEvent::new(TicketEventKind::Comment, "must roll back"), + ) + .is_err() + ); + let after = backend.show(TicketIdOrSlug::Id(created.id)).unwrap(); + assert_eq!(after.events.len(), before); + assert!( + after + .events + .iter() + .all(|event| event.body.as_str() != "must roll back") + ); + } + #[test] fn sqlite_backend_persists_core_ticket_operations() { let tmp = TempDir::new().unwrap(); diff --git a/crates/ticket/src/tool.rs b/crates/ticket/src/tool.rs index fb88153e..0095f9b0 100644 --- a/crates/ticket/src/tool.rs +++ b/crates/ticket/src/tool.rs @@ -34,12 +34,15 @@ const MAX_BODY_MAX_BYTES: usize = 64 * 1024; const DEFAULT_DIAGNOSTIC_LIMIT: usize = 100; const MAX_DIAGNOSTIC_LIMIT: usize = 500; -pub const TICKET_BASE_TOOL_NAMES: [&str; 12] = [ +pub const TICKET_BASE_TOOL_NAMES: [&str; 15] = [ "TicketCreate", "TicketEditItem", "TicketList", "TicketShow", "TicketComment", + "TicketPlan", + "TicketDecision", + "TicketImplementationReport", "TicketReview", "TicketIntakeReady", "TicketQueue", @@ -66,12 +69,15 @@ pub const TICKET_ORCHESTRATION_TOOL_NAMES: [&str; 4] = [ pub const TICKET_ORCHESTRATION_READ_ONLY_TOOL_NAMES: [&str; 2] = ["TicketRelationQuery", "TicketOrchestrationPlanQuery"]; -pub const TICKET_TOOL_NAMES: [&str; 16] = [ +pub const TICKET_TOOL_NAMES: [&str; 19] = [ "TicketCreate", "TicketEditItem", "TicketList", "TicketShow", "TicketComment", + "TicketPlan", + "TicketDecision", + "TicketImplementationReport", "TicketReview", "TicketIntakeReady", "TicketQueue", @@ -94,10 +100,13 @@ pub const TICKET_READ_ONLY_TOOL_NAMES: [&str; 6] = [ "TicketOrchestrationPlanQuery", ]; -pub const TICKET_MUTATING_TOOL_NAMES: [&str; 10] = [ +pub const TICKET_MUTATING_TOOL_NAMES: [&str; 13] = [ "TicketCreate", "TicketEditItem", "TicketComment", + "TicketPlan", + "TicketDecision", + "TicketImplementationReport", "TicketReview", "TicketIntakeReady", "TicketQueue", @@ -120,9 +129,11 @@ routing, closing, planning, or implementation decisions."; const SHOW_DESCRIPTION: &str = "Show one Ticket by id or exact query through the configured \ typed Ticket backend. Output includes bounded Markdown body, recent thread events, resolution, and \ artifact metadata."; -const COMMENT_DESCRIPTION: &str = "Append a typed Ticket thread event. `role` must be `comment`, \ -`plan`, `decision`, or `implementation_report`; `body` is Markdown. Writes stay inside the \ -configured Ticket backend root."; +const COMMENT_DESCRIPTION: &str = "Append a typed Ticket comment event. `body` is Markdown."; +const PLAN_DESCRIPTION: &str = "Append a typed Ticket plan event. `body` is Markdown."; +const DECISION_DESCRIPTION: &str = "Append a typed Ticket decision event. `body` is Markdown."; +const IMPLEMENTATION_REPORT_DESCRIPTION: &str = + "Append a typed Ticket implementation_report event. `body` is Markdown."; const REVIEW_DESCRIPTION: &str = "Append a Ticket review event. `result` must be `approve` or \ `request_changes`; `body` is Markdown. Writes stay inside the configured Ticket backend root."; const INTAKE_READY_DESCRIPTION: &str = "Mark an existing Ticket planning lane ready through the typed \ @@ -161,6 +172,9 @@ fn base_tool_description(name: &str) -> &'static str { "TicketList" => LIST_DESCRIPTION, "TicketShow" => SHOW_DESCRIPTION, "TicketComment" => COMMENT_DESCRIPTION, + "TicketPlan" => PLAN_DESCRIPTION, + "TicketDecision" => DECISION_DESCRIPTION, + "TicketImplementationReport" => IMPLEMENTATION_REPORT_DESCRIPTION, "TicketReview" => REVIEW_DESCRIPTION, "TicketIntakeReady" => INTAKE_READY_DESCRIPTION, "TicketQueue" => QUEUE_DESCRIPTION, @@ -361,9 +375,6 @@ struct TicketCreateParams { /// Markdown body for item.md. If omitted, a small default body is used. #[serde(default)] body: Option, - /// Optional thread author for the create event. - #[serde(default)] - author: Option, /// Optional assignee frontmatter value. #[serde(default)] assignee: Option, @@ -376,9 +387,6 @@ struct TicketCreateParams { /// Optional state frontmatter value. Defaults to `planning`. #[serde(default)] state: Option, - /// Optional queued_by frontmatter value. - #[serde(default)] - queued_by: Option, /// Optional queued_at frontmatter value. #[serde(default)] queued_at: Option, @@ -412,9 +420,6 @@ struct TicketEditItemParams { /// Optional target repository/ref update. #[serde(default)] target: Option, - /// Optional thread author for the audited item_edit event. - #[serde(default)] - author: Option, } #[derive(Debug, Clone, Copy, Deserialize, schemars::JsonSchema)] @@ -542,25 +547,11 @@ struct TicketShowParams { } #[derive(Debug, Deserialize, schemars::JsonSchema)] -#[serde(rename_all = "snake_case")] -enum TicketCommentRoleParam { - Comment, - Plan, - Decision, - ImplementationReport, -} - -#[derive(Debug, Deserialize, schemars::JsonSchema)] -struct TicketCommentParams { +struct TicketThreadEventParams { /// Ticket id. ticket: String, - /// Thread event role: `comment`, `plan`, `decision`, or `implementation_report`. - role: TicketCommentRoleParam, /// Markdown event body. body: String, - /// Optional thread author. - #[serde(default)] - author: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -578,9 +569,6 @@ struct TicketReviewParams { result: TicketReviewResultParam, /// Markdown review body. body: String, - /// Optional thread author. - #[serde(default)] - author: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -589,9 +577,6 @@ struct TicketIntakeReadyParams { ticket: String, /// Concise bounded intake summary to append as a typed intake_summary event. intake_summary: String, - /// Optional author for both intake_summary and state_changed events. - #[serde(default)] - author: Option, /// Reason attached to the state_changed event. Defaults to `planning_ready`. #[serde(default)] reason: Option, @@ -604,9 +589,6 @@ struct TicketIntakeReadyParams { struct TicketQueueParams { /// Ticket id. ticket: String, - /// Optional queued_by frontmatter value. Defaults to the backend/user default. - #[serde(default)] - queued_by: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -621,9 +603,6 @@ struct TicketWorkflowStateParams { reason: String, /// Markdown body for the typed state_changed event. body: String, - /// Optional thread author. - #[serde(default)] - author: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -673,9 +652,6 @@ struct TicketRelationRecordParams { /// Optional bounded rationale/note. #[serde(default)] note: Option, - /// Optional record author. - #[serde(default)] - author: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -757,9 +733,6 @@ struct TicketOrchestrationPlanRecordParams { /// Accepted plan fields. Required for accepted_plan and invalid for other kinds. #[serde(default)] accepted_plan: Option, - /// Optional record author. - #[serde(default)] - author: Option, } #[derive(Debug, Deserialize, schemars::JsonSchema)] @@ -851,6 +824,21 @@ struct TicketCommentTool { backend: TicketToolBackend, } +#[derive(Clone)] +struct TicketPlanTool { + backend: TicketToolBackend, +} + +#[derive(Clone)] +struct TicketDecisionTool { + backend: TicketToolBackend, +} + +#[derive(Clone)] +struct TicketImplementationReportTool { + backend: TicketToolBackend, +} + #[derive(Clone)] struct TicketReviewTool { backend: TicketToolBackend, @@ -918,12 +906,12 @@ impl Tool for TicketCreateTool { if let Some(body) = params.body { input.body = MarkdownText::new(body); } - input.author = params.author; + input.author = None; input.assignee = params.assignee; input.readiness = params.readiness; input.risk_flags = params.risk_flags; input.workflow_state = params.state.map(TicketWorkflowStateParam::into_state); - input.queued_by = params.queued_by; + input.queued_by = None; input.queued_at = params.queued_at; input.repository_id = params.repository_id; input.ref_selector = params.ref_selector; @@ -971,7 +959,7 @@ impl Tool for TicketEditItemTool { body: params.body.map(MarkdownText::new), body_replacement, target: params.target, - author: params.author, + author: None, }; let ticket = self .backend @@ -1066,6 +1054,26 @@ impl Tool for TicketShowTool { } } +fn execute_ticket_thread_event( + backend: &TicketToolBackend, + tool_name: &str, + kind: TicketEventKind, + input_json: &str, +) -> Result { + let params: TicketThreadEventParams = parse_input(tool_name, input_json)?; + let role = kind.as_str().to_string(); + backend + .add_event( + TicketIdOrSlug::Query(params.ticket.clone()), + NewTicketEvent::new(kind, params.body), + ) + .map_err(|error| backend_error(tool_name, error))?; + Ok(json_output( + format!("Appended {role} event to ticket {}", params.ticket), + json!({ "ticket": params.ticket, "event": role, "ok": true }), + )) +} + #[async_trait] impl Tool for TicketCommentTool { async fn execute( @@ -1073,26 +1081,42 @@ impl Tool for TicketCommentTool { input_json: &str, _ctx: llm_engine::tool::ToolExecutionContext, ) -> Result { - let params: TicketCommentParams = parse_input("TicketComment", input_json)?; - let kind = match params.role { - TicketCommentRoleParam::Comment => TicketEventKind::Comment, - TicketCommentRoleParam::Plan => TicketEventKind::Plan, - TicketCommentRoleParam::Decision => TicketEventKind::Decision, - TicketCommentRoleParam::ImplementationReport => TicketEventKind::ImplementationReport, - }; - let role = kind.as_str().to_string(); - let mut event = NewTicketEvent::new(kind, params.body); - event.author = params.author; - self.backend - .add_event(TicketIdOrSlug::Query(params.ticket.clone()), event) - .map_err(|error| backend_error("TicketComment", error))?; - Ok(json_output( - format!("Appended {role} event to ticket {}", params.ticket), - json!({ "ticket": params.ticket, "event": role, "ok": true }), - )) + execute_ticket_thread_event( + &self.backend, + "TicketComment", + TicketEventKind::Comment, + input_json, + ) } } +macro_rules! impl_ticket_thread_event_tool { + ($tool:ty, $name:literal, $kind:expr) => { + #[async_trait] + impl Tool for $tool { + async fn execute( + &self, + input_json: &str, + _ctx: llm_engine::tool::ToolExecutionContext, + ) -> Result { + execute_ticket_thread_event(&self.backend, $name, $kind, input_json) + } + } + }; +} + +impl_ticket_thread_event_tool!(TicketPlanTool, "TicketPlan", TicketEventKind::Plan); +impl_ticket_thread_event_tool!( + TicketDecisionTool, + "TicketDecision", + TicketEventKind::Decision +); +impl_ticket_thread_event_tool!( + TicketImplementationReportTool, + "TicketImplementationReport", + TicketEventKind::ImplementationReport +); + #[async_trait] impl Tool for TicketReviewTool { async fn execute( @@ -1108,7 +1132,7 @@ impl Tool for TicketReviewTool { let result_str = result.as_str().to_string(); let review = TicketReview { result, - author: params.author, + author: None, body: MarkdownText::new(params.body), }; self.backend @@ -1138,14 +1162,14 @@ impl Tool for TicketIntakeReadyTool { .default_intake_ready_state_change_body(from.as_str()) }); let mut summary = TicketIntakeSummary::new(params.intake_summary); - summary.author = params.author.clone(); + summary.author = None; let mut change = TicketStateChange::new( from.as_str(), TicketWorkflowState::Ready.as_str(), reason, body, ); - change.author = params.author; + change.author = None; self.backend .mark_intake_ready( TicketIdOrSlug::Query(params.ticket.clone()), @@ -1168,7 +1192,7 @@ impl Tool for TicketQueueTool { _ctx: llm_engine::tool::ToolExecutionContext, ) -> Result { let params: TicketQueueParams = parse_input("TicketQueue", input_json)?; - let queued_by = params.queued_by.unwrap_or_else(default_author); + let queued_by = default_author(); self.backend .queue_ready(TicketIdOrSlug::Query(params.ticket.clone()), &queued_by) .map_err(|error| backend_error("TicketQueue", error))?; @@ -1196,7 +1220,7 @@ impl Tool for TicketWorkflowStateTool { } let mut change = TicketStateChange::new(from.as_str(), to.as_str(), params.reason, params.body); - change.author = params.author; + change.author = None; self.backend .set_workflow_state(TicketIdOrSlug::Query(params.ticket.clone()), change) .map_err(|error| backend_error("TicketWorkflowState", error))?; @@ -1251,7 +1275,7 @@ impl Tool for TicketRelationRecordTool { kind: params.kind.into_kind(), target: params.target.clone(), note: params.note, - author: params.author, + author: None, }; let output = self .backend @@ -1325,7 +1349,7 @@ impl Tool for TicketOrchestrationPlanRecordTool { related_ticket: params.related_ticket, note: params.note, accepted_plan, - author: params.author, + author: None, }; let output = self .backend @@ -1703,7 +1727,9 @@ fn input_schema(name: &str) -> Value { "TicketEditItem" => serde_json::to_value(schemars::schema_for!(TicketEditItemParams)), "TicketList" => serde_json::to_value(schemars::schema_for!(TicketListParams)), "TicketShow" => serde_json::to_value(schemars::schema_for!(TicketShowParams)), - "TicketComment" => serde_json::to_value(schemars::schema_for!(TicketCommentParams)), + "TicketComment" | "TicketPlan" | "TicketDecision" | "TicketImplementationReport" => { + serde_json::to_value(schemars::schema_for!(TicketThreadEventParams)) + } "TicketReview" => serde_json::to_value(schemars::schema_for!(TicketReviewParams)), "TicketIntakeReady" => serde_json::to_value(schemars::schema_for!(TicketIntakeReadyParams)), "TicketQueue" => serde_json::to_value(schemars::schema_for!(TicketQueueParams)), @@ -1747,6 +1773,9 @@ impl_from_backend!(TicketEditItemTool); impl_from_backend!(TicketListTool); impl_from_backend!(TicketShowTool); impl_from_backend!(TicketCommentTool); +impl_from_backend!(TicketPlanTool); +impl_from_backend!(TicketDecisionTool); +impl_from_backend!(TicketImplementationReportTool); impl_from_backend!(TicketReviewTool); impl_from_backend!(TicketIntakeReadyTool); impl_from_backend!(TicketQueueTool); @@ -1768,6 +1797,12 @@ pub fn ticket_tools(backend: impl Into) -> Vec("TicketList", backend.clone()), tool_definition::("TicketShow", backend.clone()), tool_definition::("TicketComment", backend.clone()), + tool_definition::("TicketPlan", backend.clone()), + tool_definition::("TicketDecision", backend.clone()), + tool_definition::( + "TicketImplementationReport", + backend.clone(), + ), tool_definition::("TicketReview", backend.clone()), tool_definition::("TicketIntakeReady", backend.clone()), tool_definition::("TicketQueue", backend.clone()), @@ -1841,6 +1876,9 @@ mod tests { "TicketCreate", "TicketEditItem", "TicketComment", + "TicketPlan", + "TicketDecision", + "TicketImplementationReport", "TicketReview", "TicketIntakeReady", "TicketQueue", @@ -2338,16 +2376,15 @@ mod tests { let temp = TempDir::new().unwrap(); let backend = backend(&temp); let created = backend.create(NewTicket::new("Flow Tool")).unwrap(); - let comment = tool_by_name(backend.clone(), "TicketComment"); + let report = tool_by_name(backend.clone(), "TicketImplementationReport"); let review = tool_by_name(backend.clone(), "TicketReview"); let close = tool_by_name(backend.clone(), "TicketClose"); let doctor = tool_by_name(backend.clone(), "TicketDoctor"); - comment + report .execute( &json!({ "ticket": created.id.clone(), - "role": "implementation_report", "body": "Implemented." }) .to_string(), @@ -2807,6 +2844,38 @@ mod tests { assert!(edit_schema.contains("old_string")); assert!(edit_schema.contains("new_string")); assert!(edit_schema.contains("replace_all")); + for name in [ + "TicketCreate", + "TicketEditItem", + "TicketComment", + "TicketPlan", + "TicketDecision", + "TicketImplementationReport", + "TicketReview", + "TicketIntakeReady", + "TicketQueue", + "TicketRelationRecord", + "TicketOrchestrationPlanRecord", + ] { + let schema = tools + .iter() + .map(|definition| definition().0) + .find(|meta| meta.name == name) + .unwrap() + .input_schema; + let properties = schema["properties"].as_object().unwrap(); + assert!(!properties.contains_key("author"), "{name} exposes author"); + assert!( + !properties.contains_key("queued_by"), + "{name} exposes queued_by" + ); + if matches!( + name, + "TicketComment" | "TicketPlan" | "TicketDecision" | "TicketImplementationReport" + ) { + assert!(!properties.contains_key("role"), "{name} exposes role"); + } + } let names = tools .into_iter() .map(|definition| definition().0) diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 20e05921..82c10245 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -226,7 +226,7 @@ pub struct RuntimeWorkspaceHttpClient { workspace_id: String, base_url: String, worker_id: String, - access_token: Option, + access_token: Mutex>, } impl std::fmt::Debug for RuntimeWorkspaceHttpClient { @@ -238,7 +238,11 @@ impl std::fmt::Debug for RuntimeWorkspaceHttpClient { .field("worker_id", &self.worker_id) .field( "access_token", - &self.access_token.as_ref().map(|_| "[redacted]"), + &self + .access_token + .lock() + .ok() + .and_then(|token| token.as_ref().map(|_| "[redacted]")), ) .finish() } @@ -254,12 +258,12 @@ impl RuntimeWorkspaceHttpClient { workspace_id: workspace_id.into(), base_url: base_url.into().trim_end_matches('/').to_string(), worker_id: worker_id.into(), - access_token: None, + access_token: Mutex::new(None), } } - pub fn with_access_token(mut self, access_token: Option) -> Self { - self.access_token = access_token; + pub fn with_access_token(self, access_token: Option) -> Self { + *self.access_token.lock().expect("new credential mutex") = access_token; self } } @@ -283,25 +287,93 @@ impl WorkspaceClient for RuntimeWorkspaceHttpClient { ) -> Result { let base_url = self.base_url.clone(); let worker_id = self.worker_id.clone(); - let access_token = self.access_token.clone(); - if tokio::runtime::Handle::try_current().is_ok() { - return std::thread::spawn(move || { - execute_runtime_workspace_http( + let access_token = self + .access_token + .lock() + .map_err(|_| { + WorkspaceClientError::Request("workspace credential lock poisoned".to_string()) + })? + .clone(); + let request_copy = request.clone(); + let result = if tokio::runtime::Handle::try_current().is_ok() { + std::thread::spawn(move || { + execute_runtime_workspace_http_with_refresh( &base_url, &worker_id, - access_token.as_deref(), - request, + access_token, + request_copy, ) }) .join() .map_err(|_| { WorkspaceClientError::Request("workspace request thread panicked".to_string()) - })?; + })? + } else { + execute_runtime_workspace_http_with_refresh( + &base_url, + &worker_id, + access_token, + request, + ) + }?; + if let Some(new_token) = result.1 { + *self.access_token.lock().map_err(|_| { + WorkspaceClientError::Request("workspace credential lock poisoned".to_string()) + })? = Some(new_token); } - execute_runtime_workspace_http(&base_url, &worker_id, access_token.as_deref(), request) + Ok(result.0) } } +fn execute_runtime_workspace_http_with_refresh( + base_url: &str, + worker_id: &str, + access_token: Option, + request: WorkspaceRequest, +) -> Result<(WorkspaceResponse, Option), WorkspaceClientError> { + let response = execute_runtime_workspace_http( + base_url, + worker_id, + access_token.as_deref(), + request.clone(), + )?; + if response.status != 401 { + return Ok((response, None)); + } + let Some(expired_token) = access_token else { + return Ok((response, None)); + }; + let workspace_id = request + .path + .strip_prefix("/api/w/") + .and_then(|path| path.split('/').next()) + .ok_or_else(|| WorkspaceClientError::InvalidPath(request.path.clone()))?; + let refresh_url = format!("{base_url}/api/w/{workspace_id}/worker-credentials/refresh"); + let refresh = reqwest::blocking::Client::new() + .post(refresh_url) + .bearer_auth(expired_token) + .header("x-yoi-worker-id", worker_id) + .send() + .map_err(|error| WorkspaceClientError::Request(error.to_string()))?; + if !refresh.status().is_success() { + return Ok((response, None)); + } + let body: serde_json::Value = refresh + .json() + .map_err(|error| WorkspaceClientError::Request(error.to_string()))?; + let new_token = body + .get("access_token") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| { + WorkspaceClientError::Request( + "Workspace credential refresh response omitted access_token".to_string(), + ) + })? + .to_string(); + let retried = execute_runtime_workspace_http(base_url, worker_id, Some(&new_token), request)?; + Ok((retried, Some(new_token))) +} + fn execute_runtime_workspace_http( base_url: &str, worker_id: &str, @@ -6185,6 +6257,74 @@ mod build_summary_prompt_tests { })); } + #[test] + fn runtime_workspace_client_refreshes_expired_credential_and_retries() { + use std::io::{BufRead, BufReader, Write}; + use std::net::TcpListener; + + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let address = listener.local_addr().unwrap(); + let server = std::thread::spawn(move || { + for step in 0..3 { + let (mut stream, _) = listener.accept().unwrap(); + let mut reader = BufReader::new(stream.try_clone().unwrap()); + let mut first_line = String::new(); + reader.read_line(&mut first_line).unwrap(); + let mut authorization = String::new(); + loop { + let mut line = String::new(); + reader.read_line(&mut line).unwrap(); + if let Some(value) = line.strip_prefix("authorization: ") { + authorization = value.trim().to_string(); + } + if line == "\r\n" || line.is_empty() { + break; + } + } + match step { + 0 => { + assert_eq!(authorization, "Bearer expired-token"); + stream + .write_all(b"HTTP/1.1 401 Unauthorized\r\nContent-Length: 0\r\n\r\n") + .unwrap(); + } + 1 => { + assert!(first_line.contains("/worker-credentials/refresh")); + assert_eq!(authorization, "Bearer expired-token"); + let body = + r#"{"access_token":"fresh-token","expires_at":"2099-01-01T00:00:00Z"}"#; + write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}", + body.len(), + body + ) + .unwrap(); + } + _ => { + assert_eq!(authorization, "Bearer fresh-token"); + stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\n{}") + .unwrap(); + } + } + } + }); + let client = RuntimeWorkspaceHttpClient::new( + "workspace-refresh", + format!("http://{address}"), + "worker-refresh", + ) + .with_access_token(Some("expired-token".to_string())); + let response = client + .execute(WorkspaceRequest::get( + "/api/w/workspace-refresh/tickets/backend", + )) + .unwrap(); + assert_eq!(response.status, 200); + server.join().unwrap(); + } + 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 31946869..f0aa99b1 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -308,6 +308,13 @@ pub struct WorkerSpawnWorkingDirectoryRequest { pub selector: Option, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct WorkerTicketAssignmentRequest { + pub ticket_id: String, + pub operation_id: String, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] pub struct WorkerSpawnRequest { @@ -317,6 +324,8 @@ pub struct WorkerSpawnRequest { pub acceptance: WorkerSpawnAcceptanceRequirement, pub profile: ProfileSelector, #[serde(default, skip_serializing_if = "Option::is_none")] + pub ticket_assignment: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] pub initial_input: Option, /// Optional safe working-directory creation request. The Workspace server resolves /// this into a runtime-internal `WorkingDirectoryRequest` from configured @@ -429,6 +438,8 @@ pub struct WorkerStopResult { pub struct WorkerLifecycleRequest { #[serde(default, skip_serializing_if = "Option::is_none")] pub reason: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub ticket_assignment: Option, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -4177,6 +4188,7 @@ mod tests { 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, @@ -4304,6 +4316,7 @@ mod tests { 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, @@ -4401,6 +4414,7 @@ mod tests { 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, @@ -4434,6 +4448,7 @@ mod tests { requested_worker_name: None, acceptance: WorkerSpawnAcceptanceRequirement::SocketReady, profile: ProfileSelector::Builtin("builtin:companion".to_string()), + ticket_assignment: None, initial_input: None, working_directory_request: None, resolved_working_directory_request: None, diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 3399aaab..a5634da1 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -1,7 +1,7 @@ -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeMap, HashMap, HashSet}; use std::path::{Component, Path, PathBuf}; use std::sync::Arc; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use axum::extract::ws::{Message as WsMessage, WebSocket, WebSocketUpgrade}; use axum::extract::{Path as AxumPath, Query, State}; @@ -17,6 +17,7 @@ use memory::backend::{ MemoryConsolidationOutput, }; use protocol::stream::{decode_method, encode_event}; +use rusqlite::OptionalExtension; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use ticket::{ @@ -91,9 +92,9 @@ use crate::resource_broker::BackendResourceBroker; use crate::skills; use crate::store::{ AccountRecord, ApiTokenRecord, AuthChallengeRecord, BrowserSessionRecord, ControlPlaneStore, - DeviceLoginFlowRecord, PasskeyCredentialRecord, RepositoryRecord, TicketNotificationRecipient, - TicketWorkerAssignmentRecord, UserRecord, WorkdirRegistryRecord, WorkerRegistryRecord, - WorkerWorkdirLinkRecord, WorkerWorkspaceCredentialRecord, WorkspaceRecord, + DeviceLoginFlowRecord, PasskeyCredentialRecord, RepositoryRecord, TicketWorkerAssignmentRecord, + UserRecord, WorkdirRegistryRecord, WorkerRegistryRecord, WorkerWorkdirLinkRecord, + WorkerWorkspaceCredentialRecord, WorkspaceRecord, }; use crate::{Error, Result}; use worker_runtime::catalog::{ @@ -489,6 +490,10 @@ pub fn build_router(api: WorkspaceApi) -> Router { ) .route("/api/tickets", get(list_tickets)) .route("/api/w/{workspace_id}/tickets", get(scoped_list_tickets)) + .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), @@ -527,6 +532,10 @@ pub fn build_router(api: WorkspaceApi) -> Router { .put(scoped_set_ticket_worker_assignment) .delete(scoped_clear_ticket_worker_assignment), ) + .route( + "/api/w/{workspace_id}/tickets/{id}/assignment/reassign", + post(scoped_reassign_ticket_worker_assignment), + ) .route( "/api/w/{workspace_id}/tickets/{id}/state", post(scoped_transition_ticket_state), @@ -1557,6 +1566,7 @@ struct TicketWorkerAssignmentResponse { workspace_id: String, ticket_id: String, assignment: Option, + worker: Option, } #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] @@ -1570,6 +1580,7 @@ struct TicketWorkerAssignmentMutationResponse { #[derive(Debug, Deserialize)] #[serde(deny_unknown_fields)] struct SetTicketWorkerAssignmentRequest { + operation_id: String, runtime_id: String, worker_id: String, expected_assignment_id: Option, @@ -1579,6 +1590,7 @@ struct SetTicketWorkerAssignmentRequest { #[derive(Debug, Default, Deserialize)] #[serde(deny_unknown_fields)] struct ClearTicketWorkerAssignmentQuery { + operation_id: Option, expected_assignment_id: Option, actor: Option, } @@ -1592,10 +1604,16 @@ async fn scoped_get_ticket_worker_assignment( let assignment = api .store .get_current_ticket_worker_assignment(&path.workspace_id, &ticket.id)?; + let worker = assignment.as_ref().and_then(|assignment| { + api.runtime + .worker(&assignment.runtime_id, &assignment.worker_id) + .ok() + }); Ok(Json(TicketWorkerAssignmentResponse { workspace_id: path.workspace_id, ticket_id: ticket.id, assignment, + worker, })) } @@ -1603,9 +1621,27 @@ async fn scoped_set_ticket_worker_assignment( State(api): State, AxumPath(path): AxumPath, Json(request): Json, +) -> ApiResult> { + set_ticket_worker_assignment(api, path, request, false).await +} + +async fn scoped_reassign_ticket_worker_assignment( + State(api): State, + AxumPath(path): AxumPath, + Json(request): Json, +) -> ApiResult> { + set_ticket_worker_assignment(api, path, request, true).await +} + +async fn set_ticket_worker_assignment( + api: WorkspaceApi, + path: ScopedRecordPath, + request: SetTicketWorkerAssignmentRequest, + allow_reassign: bool, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; let ticket = api.authority.ticket(&path.id)?; + let operation_id = require_ticket_assignment_value("operation_id", request.operation_id)?; let runtime_id = require_ticket_assignment_value("runtime_id", request.runtime_id)?; let worker_id = require_ticket_assignment_value("worker_id", request.worker_id)?; let expected_assignment_id = request @@ -1635,6 +1671,8 @@ async fn scoped_set_ticket_worker_assignment( &record, expected_assignment_id.as_deref(), &new_id("tasev"), + &operation_id, + allow_reassign, )?; Ok(Json(TicketWorkerAssignmentMutationResponse { workspace_id: path.workspace_id, @@ -1651,6 +1689,13 @@ async fn scoped_clear_ticket_worker_assignment( ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; let ticket = api.authority.ticket(&path.id)?; + let operation_id = query + .operation_id + .map(|value| require_ticket_assignment_value("operation_id", value)) + .transpose()? + .ok_or_else(|| { + Error::TicketAssignmentConflict("unassign requires operation_id".to_string()) + })?; let expected_assignment_id = query .expected_assignment_id .map(|value| require_ticket_assignment_value("expected_assignment_id", value)) @@ -1664,6 +1709,7 @@ async fn scoped_clear_ticket_worker_assignment( &path.workspace_id, &ticket.id, expected_assignment_id.as_deref(), + &operation_id, &new_id("tasev"), &actor, &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), @@ -1676,6 +1722,65 @@ async fn scoped_clear_ticket_worker_assignment( })) } +fn assign_ticket_worker_from_lifecycle( + api: &WorkspaceApi, + assignment: &crate::hosts::WorkerTicketAssignmentRequest, + runtime_id: &str, + worker_id: &str, +) -> Result { + let ticket = api.authority.ticket(&assignment.ticket_id)?; + let assigned_at = Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true); + let record = TicketWorkerAssignmentRecord { + workspace_id: api.config.workspace_id.clone(), + ticket_id: ticket.id, + assignment_id: new_id("tasg"), + runtime_id: runtime_id.to_string(), + worker_id: worker_id.to_string(), + assigned_by: "worker-lifecycle".to_string(), + assigned_at, + }; + Ok(api + .store + .set_current_ticket_worker_assignment( + &record, + None, + &new_id("tasev"), + &assignment.operation_id, + false, + )? + .current) +} + +fn existing_lifecycle_assignment_worker( + api: &WorkspaceApi, + assignment: &crate::hosts::WorkerTicketAssignmentRequest, + runtime_id: &str, +) -> Result> { + let Some(operation) = api + .store + .get_ticket_assignment_operation(&api.config.workspace_id, &assignment.operation_id)? + else { + return Ok(None); + }; + if operation.action != "assign" + || operation.ticket_id != assignment.ticket_id + || operation.runtime_id.as_deref() != Some(runtime_id) + { + return Err(Error::TicketAssignmentConflict(format!( + "assignment operation {} was already used with different lifecycle input", + assignment.operation_id + ))); + } + let Some(worker_id) = operation.worker_id else { + return Ok(None); + }; + Ok(Some( + api.runtime + .worker(runtime_id, &worker_id) + .map_err(|error| error.into_error())?, + )) +} + fn require_ticket_assignment_value(field: &str, value: String) -> Result { let value = value.trim(); if value.is_empty() { @@ -1902,6 +2007,59 @@ async fn scoped_close_ticket( browser_ticket_detail(&api, &path.id) } +#[derive(Debug, Serialize)] +struct WorkerWorkspaceCredentialRefreshResponse { + access_token: String, + expires_at: String, +} + +async fn scoped_refresh_worker_workspace_credential( + State(api): State, + AxumPath(path): AxumPath, + headers: HeaderMap, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + let token = headers + .get(axum::http::header::AUTHORIZATION) + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.strip_prefix("Bearer ")) + .filter(|value| !value.trim().is_empty()) + .ok_or_else(|| { + Error::WorkerWorkspaceAuthentication("missing expired credential".to_string()) + })?; + let worker_id = headers + .get("x-yoi-worker-id") + .and_then(|value| value.to_str().ok()) + .filter(|value| !value.trim().is_empty()) + .ok_or_else(|| { + Error::WorkerWorkspaceAuthentication("missing Runtime-bound Worker id".to_string()) + })?; + let new_token = mint_secret("wac"); + let expires_at = + (Utc::now() + chrono::Duration::hours(1)).to_rfc3339_opts(SecondsFormat::Secs, true); + let credential = api + .store + .refresh_worker_workspace_credential( + token, + &path.workspace_id, + worker_id, + &new_token, + &expires_at, + )? + .ok_or_else(|| { + Error::WorkerWorkspaceAuthentication("credential cannot be refreshed".to_string()) + })?; + api.runtime + .worker(&credential.runtime_id, worker_id) + .map_err(|_| { + Error::WorkerWorkspaceAuthentication("credential Worker no longer exists".to_string()) + })?; + Ok(Json(WorkerWorkspaceCredentialRefreshResponse { + access_token: new_token, + expires_at, + })) +} + async fn scoped_ticket_backend_operation( State(api): State, AxumPath(path): AxumPath, @@ -1911,37 +2069,64 @@ async fn scoped_ticket_backend_operation( validate_workspace_scope(&api, &path.workspace_id)?; let config = ticket::config::TicketConfig::load_workspace(&api.config.workspace_root) .map_err(|error| Error::Config(format!("load Ticket workspace settings: {error}")))?; - let backend = SqliteTicketBackend::new( + let mut backend = SqliteTicketBackend::new( api.config.database_path.clone(), api.config.workspace_id.clone(), ) .with_record_language(config.ticket_record_language()); + let operation_kind = ticket_mutation_operation_kind(&operation); + let is_mutation = operation_kind != "read"; let target = ticket_mutation_target(&operation).cloned(); - let source = target - .as_ref() - .map(|_| authenticate_worker_mutation_source(&api, &path.workspace_id, &headers)) - .transpose()?; - if let (Some(source), Some(target)) = (source.as_ref(), target.as_ref()) { - bind_worker_ticket_operation_source( + 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, - source, - target, - &mut operation, - )?; - } + &headers, + )?) + } else { + None + }; 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)), + )); + } let response = match execute_ticket_backend_operation(&backend, operation) { Ok(result) => { - if let (Some(source), Some(target)) = (source, target) { - let after = backend.show(target).map_err(Error::from)?; - enqueue_worker_ticket_notification( - &api, - &path.workspace_id, - &source, - before.as_ref(), - &after, - )?; + 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 && source.is_some() { dispatch_pending_ticket_notifications(&api, &path.workspace_id); } TicketBackendHttpResponse::Ok { result } @@ -1978,46 +2163,24 @@ fn ticket_mutation_target(operation: &TicketBackendOperation) -> Option<&TicketI } fn bind_worker_ticket_operation_source( - api: &WorkspaceApi, - workspace_id: &str, source: &WorkerMutationSource, - target: &TicketIdOrSlug, operation: &mut TicketBackendOperation, -) -> Result<()> { - let mut author = format!("worker:{}/{}", source.runtime_id, source.worker_id); - if let TicketBackendOperation::AddEvent { event, .. } = operation { - if event.kind == TicketEventKind::ImplementationReport { - let target_query = match target { - TicketIdOrSlug::Id(value) - | TicketIdOrSlug::Slug(value) - | TicketIdOrSlug::Query(value) => value.as_str(), - }; - let ticket = api.authority.ticket(target_query)?; - let assignment = api - .store - .get_current_ticket_worker_assignment(workspace_id, &ticket.id)? - .ok_or_else(|| { - Error::TicketAssignmentConflict(format!( - "Ticket {} has no current Worker assignment", - ticket.id - )) - })?; - if assignment.runtime_id != source.runtime_id - || assignment.worker_id != source.worker_id - { - return Err(Error::TicketAssignmentConflict(format!( - "Worker {}/{} does not hold current assignment {} for Ticket {}", - source.runtime_id, source.worker_id, assignment.assignment_id, ticket.id - ))); - } - author.push_str(&format!(" assignment:{}", assignment.assignment_id)); - } - } +) { + let author = format!("worker:{}/{}", source.runtime_id, source.worker_id); match operation { + TicketBackendOperation::Create { input } => input.author = Some(author), TicketBackendOperation::EditItem { edit, .. } => edit.author = Some(author), TicketBackendOperation::AddEvent { event, .. } => event.author = Some(author), - TicketBackendOperation::AddStateChanged { change, .. } => change.author = Some(author), + TicketBackendOperation::AddStateChanged { change, .. } + | TicketBackendOperation::SetStateField { change, .. } + | TicketBackendOperation::SetWorkflowState { change, .. } => change.author = Some(author), TicketBackendOperation::AddIntakeSummary { summary, .. } => summary.author = Some(author), + TicketBackendOperation::MarkIntakeReady { + summary, change, .. + } => { + summary.author = Some(author.clone()); + change.author = Some(author); + } TicketBackendOperation::QueueReady { queued_by, .. } => *queued_by = author, TicketBackendOperation::Review { review, .. } => review.author = Some(author), TicketBackendOperation::AddTicketRelation { relation, .. } => { @@ -2028,7 +2191,194 @@ fn bind_worker_ticket_operation_source( } _ => {} } - Ok(()) +} + +fn ticket_read_target(operation: &TicketBackendOperation) -> Option<&TicketIdOrSlug> { + match operation { + TicketBackendOperation::Show { id } => Some(id), + _ => None, + } +} + +fn ticket_mutation_operation_kind(operation: &TicketBackendOperation) -> &'static str { + match operation { + TicketBackendOperation::Create { .. } => "create", + TicketBackendOperation::EditItem { .. } => "edit_item", + TicketBackendOperation::AddEvent { .. } => "add_event", + TicketBackendOperation::AddStateChanged { .. } => "add_state_changed", + TicketBackendOperation::AddIntakeSummary { .. } => "add_intake_summary", + TicketBackendOperation::SetStateField { .. } => "set_state_field", + TicketBackendOperation::SetWorkflowState { .. } => "set_workflow_state", + TicketBackendOperation::MarkIntakeReady { .. } => "mark_intake_ready", + TicketBackendOperation::QueueReady { .. } => "queue_ready", + TicketBackendOperation::Review { .. } => "review", + TicketBackendOperation::Close { .. } => "close", + TicketBackendOperation::AddTicketRelation { .. } => "add_relation", + TicketBackendOperation::AddOrchestrationPlanRecord { .. } => "add_plan_record", + _ => "read", + } +} + +fn ticket_operation_initial_state(operation: &TicketBackendOperation) -> String { + match operation { + TicketBackendOperation::Create { input } => input + .workflow_state + .as_ref() + .map(|state| state.as_str().to_string()) + .unwrap_or_else(|| TicketWorkflowState::Planning.as_str().to_string()), + _ => TicketWorkflowState::Planning.as_str().to_string(), + } +} + +#[derive(Debug, Clone)] +struct WorkerTicketSourceContext { + workspace_id: String, + runtime_id: String, + worker_id: String, + actor_role: String, + assignment_id: Option, + orchestrator: Option<(String, String)>, +} + +impl WorkerTicketSourceContext { + fn attributes(&self, operation_kind: &str) -> BTreeMap { + let mut attributes = BTreeMap::from([ + ("source_runtime_id".to_string(), self.runtime_id.clone()), + ("source_worker_id".to_string(), self.worker_id.clone()), + ("source_actor_role".to_string(), self.actor_role.clone()), + ( + "source_operation_kind".to_string(), + operation_kind.to_string(), + ), + ]); + if let Some(assignment_id) = &self.assignment_id { + attributes.insert("source_assignment_id".to_string(), assignment_id.clone()); + } + attributes + } +} + +fn worker_ticket_source_context( + api: &WorkspaceApi, + workspace_id: &str, + source: &WorkerMutationSource, + ticket: Option<&ticket::Ticket>, +) -> WorkerTicketSourceContext { + let assignment = ticket.and_then(|ticket| { + api.store + .get_current_ticket_worker_assignment(workspace_id, &ticket.meta.id) + .ok() + .flatten() + }); + let orchestrator = find_workspace_orchestrator(api); + let actor_role = if assignment.as_ref().is_some_and(|assignment| { + assignment.runtime_id == source.runtime_id && assignment.worker_id == source.worker_id + }) { + "assigned" + } else if orchestrator.as_ref().is_some_and(|worker| { + worker.runtime_id == source.runtime_id && worker.worker_id == source.worker_id + }) { + "orchestrator" + } else { + "worker" + }; + WorkerTicketSourceContext { + workspace_id: workspace_id.to_string(), + runtime_id: source.runtime_id.clone(), + worker_id: source.worker_id.clone(), + actor_role: actor_role.to_string(), + assignment_id: assignment.and_then(|assignment| { + (assignment.runtime_id == source.runtime_id && assignment.worker_id == source.worker_id) + .then_some(assignment.assignment_id) + }), + orchestrator: orchestrator.map(|worker| (worker.runtime_id, worker.worker_id)), + } +} + +fn build_ticket_notification_hook( + _api: &WorkspaceApi, + source: WorkerTicketSourceContext, + operation_kind: &'static str, + previous_state: String, +) -> Arc { + let invoked = AtomicBool::new(false); + let notification_id = new_id("tnfy"); + Arc::new(move |conn, event| { + if invoked.swap(true, Ordering::SeqCst) { + return Ok(()); + } + let current_state: String = conn + .query_row( + "SELECT workflow_state FROM typed_tickets WHERE workspace_id = ?1 AND ticket_id = ?2", + rusqlite::params![source.workspace_id, event.ticket_id], + |row| row.get(0), + ) + .map_err(|error| ticket::TicketError::Conflict(format!("read committed Ticket state for outbox: {error}")))?; + conn.execute( + r#"INSERT INTO ticket_notification_outbox ( + notification_id, workspace_id, ticket_id, event_sequence, + source_runtime_id, source_worker_id, previous_state, current_state, created_at, + event_kind, source_operation_kind, source_actor_role, source_assignment_id + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)"#, + rusqlite::params![ + notification_id, + source.workspace_id, + event.ticket_id, + event.event_index, + source.runtime_id, + source.worker_id, + previous_state, + current_state, + Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + event.event_kind.as_str(), + operation_kind, + source.actor_role, + source.assignment_id, + ], + ) + .map_err(|error| { + ticket::TicketError::Conflict(format!("insert Ticket notification outbox: {error}")) + })?; + + let assigned: Option<(String, String)> = conn + .query_row( + r#"SELECT runtime_id, worker_id FROM ticket_current_worker_assignments + WHERE workspace_id = ?1 AND ticket_id = ?2"#, + rusqlite::params![source.workspace_id, event.ticket_id], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .optional() + .map_err(|error| { + ticket::TicketError::Conflict(format!( + "resolve assigned notification recipient: {error}" + )) + })?; + let mut recipients = Vec::new(); + if let Some((runtime_id, worker_id)) = assigned { + if runtime_id != source.runtime_id || worker_id != source.worker_id { + recipients.push((runtime_id, worker_id, "assigned")); + } + } + if (matches!(previous_state.as_str(), "queued" | "inprogress") + || matches!(current_state.as_str(), "queued" | "inprogress")) + && let Some((runtime_id, worker_id)) = &source.orchestrator + && (*runtime_id != source.runtime_id || *worker_id != source.worker_id) + { + recipients.push((runtime_id.clone(), worker_id.clone(), "orchestrator")); + } + recipients.sort(); + recipients.dedup_by(|left, right| left.0 == right.0 && left.1 == right.1); + for (runtime_id, worker_id, recipient_kind) in recipients { + conn.execute( + r#"INSERT OR IGNORE INTO ticket_notification_deliveries ( + notification_id, recipient_runtime_id, recipient_worker_id, recipient_kind, attempts + ) VALUES (?1, ?2, ?3, ?4, 0)"#, + rusqlite::params![notification_id, runtime_id, worker_id, recipient_kind], + ) + .map_err(|error| ticket::TicketError::Conflict(format!("insert Ticket notification delivery: {error}")))?; + } + Ok(()) + }) } fn authenticate_worker_mutation_source( @@ -2070,65 +2420,6 @@ fn authenticate_worker_mutation_source( }) } -fn enqueue_worker_ticket_notification( - api: &WorkspaceApi, - workspace_id: &str, - source: &WorkerMutationSource, - before: Option<&ticket::Ticket>, - after: &ticket::Ticket, -) -> Result<()> { - let mut recipients = Vec::new(); - if let Some(assignment) = api - .store - .get_current_ticket_worker_assignment(workspace_id, &after.meta.id)? - { - if assignment.runtime_id != source.runtime_id || assignment.worker_id != source.worker_id { - recipients.push(TicketNotificationRecipient { - runtime_id: assignment.runtime_id, - worker_id: assignment.worker_id, - recipient_kind: "assigned".to_string(), - }); - } - } - let previous_state = before - .map(|ticket| ticket.meta.workflow_state.as_str().to_string()) - .unwrap_or_else(|| after.meta.workflow_state.as_str().to_string()); - let current_state = after.meta.workflow_state.as_str().to_string(); - if matches!(previous_state.as_str(), "queued" | "inprogress") - || matches!(current_state.as_str(), "queued" | "inprogress") - { - if let Some(orchestrator) = find_workspace_orchestrator(api) { - if orchestrator.runtime_id != source.runtime_id - || orchestrator.worker_id != source.worker_id - { - recipients.push(TicketNotificationRecipient { - runtime_id: orchestrator.runtime_id, - worker_id: orchestrator.worker_id, - recipient_kind: "orchestrator".to_string(), - }); - } - } - } - recipients.sort_by(|left, right| { - (&left.runtime_id, &left.worker_id).cmp(&(&right.runtime_id, &right.worker_id)) - }); - recipients.dedup_by(|left, right| { - left.runtime_id == right.runtime_id && left.worker_id == right.worker_id - }); - api.store.enqueue_ticket_notification( - &new_id("tnfy"), - workspace_id, - &after.meta.id, - after.events.len().saturating_sub(1) as i64, - &source.runtime_id, - &source.worker_id, - &previous_state, - ¤t_state, - &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), - &recipients, - ) -} - fn dispatch_pending_ticket_notifications(api: &WorkspaceApi, workspace_id: &str) { let Ok(deliveries) = api .store @@ -2137,7 +2428,46 @@ fn dispatch_pending_ticket_notifications(api: &WorkspaceApi, workspace_id: &str) return; }; for delivery in deliveries { - if !ticket_notification_recipient_is_current(api, &delivery) { + let Some((current_runtime_id, current_worker_id)) = + current_ticket_notification_recipient(api, &delivery) + else { + continue; + }; + if current_runtime_id == delivery.source_runtime_id + && current_worker_id == delivery.source_worker_id + { + let _ = api.store.mark_ticket_notification_delivered( + &delivery.notification_id, + &delivery.recipient_runtime_id, + &delivery.recipient_worker_id, + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + ); + continue; + } + if current_runtime_id != delivery.recipient_runtime_id + || current_worker_id != delivery.recipient_worker_id + { + let _ = api.store.reroute_ticket_notification_delivery( + &delivery.notification_id, + &delivery.recipient_runtime_id, + &delivery.recipient_worker_id, + ¤t_runtime_id, + ¤t_worker_id, + ); + continue; + } + if api + .store + .get_ticket_notification_cursor( + &delivery.workspace_id, + &delivery.ticket_id, + ¤t_runtime_id, + ¤t_worker_id, + ) + .ok() + .flatten() + .is_some_and(|cursor| cursor >= delivery.event_sequence) + { let _ = api.store.mark_ticket_notification_delivered( &delivery.notification_id, &delivery.recipient_runtime_id, @@ -2152,8 +2482,14 @@ fn dispatch_pending_ticket_notifications(api: &WorkspaceApi, workspace_id: &str) WorkerInputRequest { kind: WorkerInputKind::System, content: format!( - "Ticket notification: workspace_id={} ticket_id={} event_sequence={}. Reread the Ticket before acting.", - delivery.workspace_id, delivery.ticket_id, delivery.event_sequence + "Ticket notification: workspace_id={} ticket_id={} event_sequence={} event_kind={} source_operation_kind={} source_runtime_id={} source_worker_id={}. Reread the Ticket before acting.", + delivery.workspace_id, + delivery.ticket_id, + delivery.event_sequence, + delivery.event_kind, + delivery.source_operation_kind, + delivery.source_runtime_id, + delivery.source_worker_id ), segments: None, }, @@ -2213,25 +2549,21 @@ fn find_workspace_orchestrator(api: &WorkspaceApi) -> Option { None } -fn ticket_notification_recipient_is_current( +fn current_ticket_notification_recipient( api: &WorkspaceApi, delivery: &crate::store::TicketNotificationDeliveryRecord, -) -> bool { +) -> Option<(String, String)> { match delivery.recipient_kind.as_str() { "assigned" => api .store .get_current_ticket_worker_assignment(&delivery.workspace_id, &delivery.ticket_id) .ok() .flatten() - .is_some_and(|assignment| { - assignment.runtime_id == delivery.recipient_runtime_id - && assignment.worker_id == delivery.recipient_worker_id - }), - "orchestrator" => find_workspace_orchestrator(api).is_some_and(|worker| { - worker.runtime_id == delivery.recipient_runtime_id - && worker.worker_id == delivery.recipient_worker_id - }), - _ => false, + .map(|assignment| (assignment.runtime_id, assignment.worker_id)), + "orchestrator" => { + find_workspace_orchestrator(api).map(|worker| (worker.runtime_id, worker.worker_id)) + } + _ => None, } } @@ -2367,6 +2699,7 @@ fn start_memory_staging_consolidation( expected_segments: 1, }, profile: profile_selector, + ticket_assignment: None, initial_input: Some(input), working_directory_request: None, resolved_working_directory_request: None, @@ -3427,6 +3760,7 @@ fn cleanup_runtime_worker_for_execution( candidate.runtime_worker_id.as_str(), WorkerLifecycleRequest { reason: Some("cleanup worker before deletion".to_string()), + ticket_assignment: None, }, ) { Ok(result) if result.state == WorkerOperationState::Accepted => {} @@ -3448,7 +3782,15 @@ fn cleanup_runtime_worker_for_execution( .runtime .delete_worker(runtime_id, candidate.runtime_worker_id.as_str()) { - Ok(result) if result.deleted && result.state == WorkerOperationState::Accepted => Ok(()), + Ok(result) if result.deleted && result.state == WorkerOperationState::Accepted => { + api.store.revoke_worker_workspace_credentials( + &api.config.workspace_id, + runtime_id, + candidate.runtime_worker_id.as_str(), + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + )?; + Ok(()) + } Ok(result) => Err(ApiError::with_diagnostics( Error::RuntimeOperationFailed { runtime_id: runtime_id.to_string(), @@ -3610,17 +3952,79 @@ async fn scoped_get_runtime_worker( get_runtime_worker(State(api), AxumPath((path.runtime_id, path.worker_id))).await } +#[derive(Debug, Default, Deserialize)] +struct RestoreTicketAssignmentQuery { + ticket_id: Option, + assignment_operation_id: Option, +} + async fn scoped_restore_runtime_worker( State(api): State, AxumPath(path): AxumPath, + Query(query): Query, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; let workspace_id = path.workspace_id.clone(); + let runtime_id = path.runtime_id.clone(); + let worker_id = path.worker_id.clone(); + let assignment_request = match ( + query.ticket_id.clone(), + query.assignment_operation_id.clone(), + ) { + (Some(ticket_id), Some(operation_id)) => { + Some(crate::hosts::WorkerTicketAssignmentRequest { + ticket_id, + operation_id, + }) + } + (None, None) => None, + _ => { + return Err(Error::TicketAssignmentConflict( + "restore assignment requires both ticket_id and assignment_operation_id" + .to_string(), + ) + .into()); + } + }; + if let Some(assignment) = assignment_request.as_ref() { + if let Some(worker) = existing_lifecycle_assignment_worker(&api, assignment, &runtime_id)? { + if worker.worker_id != worker_id { + return Err(Error::TicketAssignmentConflict(format!( + "assignment operation {} belongs to worker {}, not {}", + assignment.operation_id, worker.worker_id, worker_id + )) + .into()); + } + assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker_id)?; + dispatch_pending_ticket_notifications(&api, &workspace_id); + return Ok(Json(WorkerRestoreResponse { + workspace_id, + runtime_id, + worker_id, + result: crate::hosts::WorkerRestoreResult { + state: WorkerOperationState::Accepted, + worker: Some(worker), + diagnostics: Vec::new(), + }, + })); + } + } let response = restore_runtime_worker( State(api.clone()), - AxumPath((path.runtime_id, path.worker_id)), + AxumPath((runtime_id.clone(), worker_id.clone())), ) .await?; + if let Some(assignment) = assignment_request.as_ref() { + api.store.reserve_ticket_assignment_operation( + &workspace_id, + &assignment.operation_id, + &assignment.ticket_id, + &runtime_id, + &worker_id, + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + )?; + assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker_id)?; + } dispatch_pending_ticket_notifications(&api, &workspace_id); Ok(response) } @@ -4995,6 +5399,7 @@ async fn create_workspace_worker( 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, @@ -5293,6 +5698,19 @@ async fn create_runtime_worker( AxumPath(runtime_id): AxumPath, Json(mut request): Json, ) -> ApiResult> { + if let Some(assignment) = request.ticket_assignment.as_ref() { + if let Some(worker) = existing_lifecycle_assignment_worker(&api, assignment, &runtime_id)? { + assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker.worker_id)?; + dispatch_pending_ticket_notifications(&api, api.workspace_id()); + return Ok(Json(WorkerSpawnResult { + state: WorkerOperationState::Accepted, + worker: Some(worker), + acceptance_evidence: Vec::new(), + diagnostics: Vec::new(), + })); + } + } + let lifecycle_assignment = request.ticket_assignment.clone(); reject_workdir_for_embedded_runtime( &runtime_id, request.working_directory_request.is_some() || request.resolved_working_directory.is_some(), @@ -5330,6 +5748,9 @@ async fn create_runtime_worker( 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 { @@ -5356,6 +5777,17 @@ async fn create_runtime_worker( worker.profile.clone(), WorkerRegistryDisplayNamePolicy::UseProvided, )?; + if let Some(assignment) = lifecycle_assignment.as_ref() { + api.store.reserve_ticket_assignment_operation( + &api.config.workspace_id, + &assignment.operation_id, + &assignment.ticket_id, + &runtime_id, + &worker.worker_id, + &Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true), + )?; + assign_ticket_worker_from_lifecycle(&api, assignment, &runtime_id, &worker.worker_id)?; + } if worker.working_directory.is_none() { if let Some(workdir_id) = prepared_workdir_id.as_deref() { if api @@ -8354,6 +8786,7 @@ mod tests { expected_segments: 0, }, profile: ProfileSelector::Builtin(MEMORY_CONSOLIDATION_PROFILE.to_string()), + ticket_assignment: None, initial_input: None, working_directory_request: None, resolved_working_directory_request: None, @@ -8442,24 +8875,11 @@ mod tests { async fn ticket_assignment_endpoints_read_and_clear_current_assignment() { let dir = tempfile::tempdir().unwrap(); let api = test_api(dir.path()).await; - let Json(created) = scoped_ticket_backend_operation( - State(api.clone()), - AxumPath(ScopedWorkspacePath { - workspace_id: TEST_WORKSPACE_ID.to_string(), - }), - HeaderMap::new(), - Json(TicketBackendOperation::Create { - input: ticket::NewTicket::new("Assigned Ticket"), - }), - ) - .await - .unwrap(); - let ticket_id = match created { - TicketBackendHttpResponse::Ok { - result: ticket::TicketBackendOperationResult::TicketRef(ticket_ref), - } => ticket_ref.id, - other => panic!("unexpected create response: {other:?}"), - }; + let created = browser_ticket_backend(&api) + .unwrap() + .create(ticket::NewTicket::new("Assigned Ticket")) + .unwrap(); + let ticket_id = created.id; let assignment = TicketWorkerAssignmentRecord { workspace_id: TEST_WORKSPACE_ID.to_string(), ticket_id: ticket_id.clone(), @@ -8470,7 +8890,13 @@ mod tests { assigned_at: TEST_CREATED_AT.to_string(), }; api.store - .set_current_ticket_worker_assignment(&assignment, None, "event-api-1") + .set_current_ticket_worker_assignment( + &assignment, + None, + "event-api-1", + "operation-api-1", + false, + ) .unwrap(); let path = || ScopedRecordPath { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -8486,6 +8912,7 @@ mod tests { State(api.clone()), AxumPath(path()), Query(ClearTicketWorkerAssignmentQuery { + operation_id: Some("clear-stale".to_string()), expected_assignment_id: Some("stale-assignment".to_string()), actor: Some("test-user".to_string()), }), @@ -8499,6 +8926,7 @@ mod tests { State(api.clone()), AxumPath(path()), Query(ClearTicketWorkerAssignmentQuery { + operation_id: Some("clear-current".to_string()), expected_assignment_id: Some("assignment-api-1".to_string()), actor: Some("test-user".to_string()), }), @@ -8526,6 +8954,7 @@ mod tests { 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, @@ -8562,6 +8991,8 @@ mod tests { }, None, "notify-assignment-event", + "notify-assignment-operation", + false, ) .unwrap(); api.store @@ -8572,6 +9003,8 @@ mod tests { runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), worker_id: Some(source_worker.worker_id.clone()), created_at: TEST_CREATED_AT.to_string(), + expires_at: "2099-01-01T00:00:00Z".to_string(), + revoked_at: None, }) .unwrap(); let mut headers = HeaderMap::new(); @@ -8596,7 +9029,64 @@ mod tests { ) .await .unwrap(); - assert!(matches!(response, TicketBackendHttpResponse::Ok { .. })); + assert!( + matches!(response, TicketBackendHttpResponse::Ok { .. }), + "unexpected response: {response:?}" + ); + let committed = backend.show(ticket_ref.id.clone().into()).unwrap(); + let committed_event = committed.events.last().unwrap(); + assert_eq!( + committed_event + .attributes + .get("source_runtime_id") + .map(String::as_str), + Some(EMBEDDED_WORKER_RUNTIME_ID) + ); + assert_eq!( + committed_event + .attributes + .get("source_worker_id") + .map(String::as_str), + Some(source_worker.worker_id.as_str()) + ); + assert_eq!( + committed_event + .attributes + .get("source_operation_kind") + .map(String::as_str), + Some("add_event") + ); + assert!(committed_event.attributes.contains_key("event_id")); + assert!(committed_event.attributes.contains_key("event_sequence")); + let event_sequence = committed_event + .attributes + .get("event_sequence") + .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(); + assert_eq!( + api.store + .get_ticket_notification_cursor( + TEST_WORKSPACE_ID, + &ticket_ref.id, + EMBEDDED_WORKER_RUNTIME_ID, + &source_worker.worker_id, + ) + .unwrap(), + Some(event_sequence) + ); assert!( api.store .list_pending_ticket_notification_deliveries(TEST_WORKSPACE_ID, 10) @@ -8605,7 +9095,7 @@ mod tests { "accepted Runtime system input must complete the outbox delivery" ); - let stale_report = scoped_ticket_backend_operation( + let Json(stale_report) = scoped_ticket_backend_operation( State(api.clone()), AxumPath(ScopedWorkspacePath { workspace_id: TEST_WORKSPACE_ID.to_string(), @@ -8615,14 +9105,22 @@ mod tests { id: ticket_ref.id.clone().into(), event: NewTicketEvent::new( TicketEventKind::ImplementationReport, - "stale assignment report", + "non-assigned Worker report", ), }), ) .await - .unwrap_err() - .into_response(); - assert_eq!(stale_report.status(), StatusCode::CONFLICT); + .unwrap(); + assert!(matches!(stale_report, TicketBackendHttpResponse::Ok { .. })); + let non_assigned_report = backend.show(ticket_ref.id.clone().into()).unwrap(); + assert!( + !non_assigned_report + .events + .last() + .unwrap() + .attributes + .contains_key("source_assignment_id") + ); api.store .set_current_ticket_worker_assignment( @@ -8637,6 +9135,8 @@ mod tests { }, Some("notify-assignment"), "source-assignment-event", + "source-assignment-operation", + true, ) .unwrap(); let _ = scoped_ticket_backend_operation( @@ -8656,17 +9156,20 @@ mod tests { .await .unwrap(); let reported = backend.show(ticket_ref.id.clone().into()).unwrap(); - assert!( - reported - .events - .last() - .unwrap() - .author - .as_deref() - .is_some_and(|author| { - author.contains(&source_worker.worker_id) - && author.contains("assignment:source-assignment") - }) + let report_event = reported.events.last().unwrap(); + assert_eq!( + report_event + .attributes + .get("source_assignment_id") + .map(String::as_str), + Some("source-assignment") + ); + assert_eq!( + report_event + .attributes + .get("source_actor_role") + .map(String::as_str), + Some("assigned") ); let unauthorized = scoped_ticket_backend_operation( @@ -8704,6 +9207,7 @@ mod tests { 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, @@ -8726,6 +9230,7 @@ mod tests { expected_segments: 0, }, profile: ProfileSelector::Builtin("builtin:orchestrator".to_string()), + ticket_assignment: None, initial_input: None, working_directory_request: None, resolved_working_directory_request: None, @@ -8743,6 +9248,7 @@ mod tests { &orchestrator.worker_id, WorkerLifecycleRequest { reason: Some("test pending delivery".to_string()), + ticket_assignment: None, }, ) .unwrap(); @@ -8754,6 +9260,8 @@ mod tests { runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), worker_id: Some(source.worker_id.clone()), created_at: TEST_CREATED_AT.to_string(), + expires_at: "2099-01-01T00:00:00Z".to_string(), + revoked_at: None, }) .unwrap(); let backend = browser_ticket_backend(&api).unwrap(); @@ -8796,49 +9304,217 @@ mod tests { } #[tokio::test] - async fn ticket_browser_endpoints_mutate_typed_backend_and_return_thread() { + async fn worker_spawn_and_restore_assignment_operations_are_idempotent() { let dir = tempfile::tempdir().unwrap(); let api = test_api(dir.path()).await; - let Json(created) = scoped_ticket_backend_operation( - State(api.clone()), - AxumPath(ScopedWorkspacePath { - workspace_id: TEST_WORKSPACE_ID.to_string(), + let backend = browser_ticket_backend(&api).unwrap(); + let first_ticket = backend + .create(ticket::NewTicket::new("Spawn assignment")) + .unwrap(); + let request = WorkerSpawnRequest { + requested_worker_name: Some("assigned-spawn".to_string()), + intent: WorkerSpawnIntent::TicketRole { + ticket_id: first_ticket.id.clone(), + role: TicketWorkerRole::Coder, + }, + acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { + expected_segments: 0, + }, + profile: ProfileSelector::Builtin("builtin:coder".to_string()), + ticket_assignment: Some(crate::hosts::WorkerTicketAssignmentRequest { + ticket_id: first_ticket.id.clone(), + operation_id: "spawn-assignment-operation".to_string(), }), - HeaderMap::new(), - Json(TicketBackendOperation::Create { - input: ticket::NewTicket::new("Browser Ticket API"), + initial_input: None, + working_directory_request: None, + resolved_working_directory_request: None, + resolved_working_directory: None, + resolved_config_bundle: None, + resolved_workspace_api: None, + }; + let Json(first) = scoped_create_runtime_worker( + State(api.clone()), + AxumPath(ScopedRuntimePath { + workspace_id: TEST_WORKSPACE_ID.to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + }), + Json(request.clone()), + ) + .await + .unwrap(); + let first_worker = first.worker.unwrap(); + let Json(projected) = scoped_get_ticket_worker_assignment( + State(api.clone()), + AxumPath(ScopedRecordPath { + workspace_id: TEST_WORKSPACE_ID.to_string(), + id: first_ticket.id.clone(), }), ) .await .unwrap(); - let ticket_id = match created { - TicketBackendHttpResponse::Ok { - result: ticket::TicketBackendOperationResult::TicketRef(ticket_ref), - } => ticket_ref.id, - other => panic!("unexpected create response: {other:?}"), - }; + assert_eq!( + projected + .worker + .as_ref() + .map(|worker| worker.worker_id.as_str()), + Some(first_worker.worker_id.as_str()) + ); + let Json(retried) = scoped_create_runtime_worker( + State(api.clone()), + AxumPath(ScopedRuntimePath { + workspace_id: TEST_WORKSPACE_ID.to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + }), + Json(request.clone()), + ) + .await + .unwrap(); + assert_eq!(retried.worker.unwrap().worker_id, first_worker.worker_id); + assert_eq!( + api.store + .list_ticket_worker_assignment_events(TEST_WORKSPACE_ID, &first_ticket.id, 10,) + .unwrap() + .len(), + 1 + ); + + let current = api + .store + .get_current_ticket_worker_assignment(TEST_WORKSPACE_ID, &first_ticket.id) + .unwrap() + .unwrap(); + api.store + .clear_current_ticket_worker_assignment( + TEST_WORKSPACE_ID, + &first_ticket.id, + Some(¤t.assignment_id), + "spawn-unassign-operation", + "spawn-unassign-event", + "test-user", + TEST_CREATED_AT, + ) + .unwrap(); + api.runtime + .stop_worker( + EMBEDDED_WORKER_RUNTIME_ID, + &first_worker.worker_id, + WorkerLifecycleRequest { + reason: Some("restore assignment test".to_string()), + ticket_assignment: None, + }, + ) + .unwrap(); + let second_ticket = backend + .create(ticket::NewTicket::new("Restore assignment")) + .unwrap(); + let _ = scoped_restore_runtime_worker( + State(api.clone()), + AxumPath(ScopedRuntimeWorkerPath { + workspace_id: TEST_WORKSPACE_ID.to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + worker_id: first_worker.worker_id.clone(), + }), + Query(RestoreTicketAssignmentQuery { + ticket_id: Some(second_ticket.id.clone()), + assignment_operation_id: Some("restore-assignment-operation".to_string()), + }), + ) + .await + .unwrap(); + let Json(retried_restore) = scoped_restore_runtime_worker( + State(api.clone()), + AxumPath(ScopedRuntimeWorkerPath { + workspace_id: TEST_WORKSPACE_ID.to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + worker_id: first_worker.worker_id.clone(), + }), + Query(RestoreTicketAssignmentQuery { + ticket_id: Some(second_ticket.id.clone()), + assignment_operation_id: Some("restore-assignment-operation".to_string()), + }), + ) + .await + .unwrap(); + assert_eq!(retried_restore.worker_id, first_worker.worker_id); + assert_eq!(retried_restore.result.state, WorkerOperationState::Accepted); + let restored_assignment = api + .store + .get_current_ticket_worker_assignment(TEST_WORKSPACE_ID, &second_ticket.id) + .unwrap() + .unwrap(); + assert_eq!(restored_assignment.worker_id, first_worker.worker_id); + + api.store + .clear_current_ticket_worker_assignment( + TEST_WORKSPACE_ID, + &second_ticket.id, + Some(&restored_assignment.assignment_id), + "restore-clear-operation", + "restore-clear-event", + "test", + TEST_CREATED_AT, + ) + .unwrap(); + api.store + .reserve_ticket_assignment_operation( + TEST_WORKSPACE_ID, + "pending-spawn-operation", + &second_ticket.id, + EMBEDDED_WORKER_RUNTIME_ID, + &first_worker.worker_id, + TEST_CREATED_AT, + ) + .unwrap(); + let worker_count_before_retry = api.runtime.list_workers(100).items.len(); + let Json(reconciled) = scoped_create_runtime_worker( + State(api.clone()), + AxumPath(ScopedRuntimePath { + workspace_id: TEST_WORKSPACE_ID.to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + }), + Json(WorkerSpawnRequest { + ticket_assignment: Some(crate::hosts::WorkerTicketAssignmentRequest { + ticket_id: second_ticket.id.clone(), + operation_id: "pending-spawn-operation".to_string(), + }), + ..request + }), + ) + .await + .unwrap(); + assert_eq!(reconciled.worker.unwrap().worker_id, first_worker.worker_id); + assert_eq!( + api.runtime.list_workers(100).items.len(), + worker_count_before_retry, + "retrying a reserved lifecycle operation must not spawn another Worker" + ); + assert!( + api.store + .get_ticket_assignment_operation(TEST_WORKSPACE_ID, "pending-spawn-operation") + .unwrap() + .and_then(|operation| operation.assignment_id) + .is_some() + ); + } + + #[tokio::test] + async fn ticket_browser_endpoints_mutate_typed_backend_and_return_thread() { + let dir = tempfile::tempdir().unwrap(); + let api = test_api(dir.path()).await; + let ticket_ref = browser_ticket_backend(&api) + .unwrap() + .create(ticket::NewTicket::new("Browser Ticket API")) + .unwrap(); + let ticket_id = ticket_ref.id; let path = || ScopedRecordPath { workspace_id: TEST_WORKSPACE_ID.to_string(), id: ticket_id.clone(), }; - let Json(related) = scoped_ticket_backend_operation( - State(api.clone()), - AxumPath(ScopedWorkspacePath { - workspace_id: TEST_WORKSPACE_ID.to_string(), - }), - HeaderMap::new(), - Json(TicketBackendOperation::Create { - input: ticket::NewTicket::new("Related Browser Ticket"), - }), - ) - .await - .unwrap(); - let related_ticket_id = match related { - TicketBackendHttpResponse::Ok { - result: ticket::TicketBackendOperationResult::TicketRef(ticket_ref), - } => ticket_ref.id, - other => panic!("unexpected related create response: {other:?}"), - }; + let related_ticket_id = browser_ticket_backend(&api) + .unwrap() + .create(ticket::NewTicket::new("Related Browser Ticket")) + .unwrap() + .id; browser_ticket_backend(&api) .unwrap() .add_ticket_relation( @@ -8964,25 +9640,10 @@ mod tests { .unwrap(); let api = test_api(dir.path()).await; - let Json(response) = scoped_ticket_backend_operation( - State(api.clone()), - AxumPath(ScopedWorkspacePath { - workspace_id: TEST_WORKSPACE_ID.to_string(), - }), - HeaderMap::new(), - Json(TicketBackendOperation::Create { - input: ticket::NewTicket::new("Endpoint configured root"), - }), - ) - .await - .unwrap_or_else(|error| panic!("ticket backend operation failed: {}", error.error)); - - let ticket_ref = match response { - TicketBackendHttpResponse::Ok { - result: ticket::TicketBackendOperationResult::TicketRef(ticket_ref), - } => ticket_ref, - other => panic!("unexpected ticket backend response: {other:?}"), - }; + let ticket_ref = browser_ticket_backend(&api) + .unwrap() + .create(ticket::NewTicket::new("Endpoint configured root")) + .unwrap(); assert!(api.config.database_path.is_file()); assert!( !dir.path() @@ -9046,9 +9707,10 @@ mod tests { } async fn test_api(workspace_root: impl Into) -> WorkspaceApi { - let store = SqliteWorkspaceStore::in_memory().unwrap(); + let config = test_server_config(workspace_root); + let store = SqliteWorkspaceStore::open(config.database_path.clone()).unwrap(); WorkspaceApi::new_with_execution_backend( - test_server_config(workspace_root), + config, Arc::new(store), Arc::new(DeterministicExecutionBackend::default()), ) @@ -10560,6 +11222,7 @@ mod tests { 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, diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index ea3ea892..b635a14c 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -97,6 +97,16 @@ const MIGRATIONS: &[Migration] = &[ name: "worker workspace credentials and Ticket notification outbox", apply: create_ticket_notification_tables, }, + Migration { + version: 17, + name: "bidirectional idempotent Ticket Worker assignments", + apply: strengthen_ticket_worker_assignments, + }, + Migration { + version: 18, + name: "atomic Ticket notification identity credentials and cursors", + apply: strengthen_ticket_notifications, + }, ]; struct Migration { @@ -279,6 +289,8 @@ pub struct WorkerWorkspaceCredentialRecord { pub runtime_id: String, pub worker_id: Option, pub created_at: String, + pub expires_at: String, + pub revoked_at: Option, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -294,6 +306,10 @@ pub struct TicketNotificationDeliveryRecord { pub workspace_id: String, pub ticket_id: String, pub event_sequence: i64, + pub event_kind: String, + pub source_operation_kind: String, + pub source_actor_role: String, + pub source_assignment_id: Option, pub source_runtime_id: String, pub source_worker_id: String, pub recipient_runtime_id: String, @@ -553,6 +569,20 @@ pub trait ControlPlaneStore: Send + Sync { runtime_worker_id: u64, ) -> Result; + fn get_ticket_assignment_operation( + &self, + workspace_id: &str, + operation_id: &str, + ) -> Result>; + fn reserve_ticket_assignment_operation( + &self, + workspace_id: &str, + operation_id: &str, + ticket_id: &str, + runtime_id: &str, + worker_id: &str, + created_at: &str, + ) -> Result<()>; fn get_current_ticket_worker_assignment( &self, workspace_id: &str, @@ -563,12 +593,15 @@ pub trait ControlPlaneStore: Send + Sync { record: &TicketWorkerAssignmentRecord, expected_assignment_id: Option<&str>, event_id: &str, + operation_id: &str, + allow_reassign: bool, ) -> Result; fn clear_current_ticket_worker_assignment( &self, workspace_id: &str, ticket_id: &str, expected_assignment_id: Option<&str>, + operation_id: &str, event_id: &str, actor: &str, created_at: &str, @@ -590,6 +623,21 @@ pub trait ControlPlaneStore: Send + Sync { workspace_id: &str, worker_id: &str, ) -> Result>; + fn refresh_worker_workspace_credential( + &self, + token: &str, + workspace_id: &str, + worker_id: &str, + new_token: &str, + new_expires_at: &str, + ) -> Result>; + fn revoke_worker_workspace_credentials( + &self, + workspace_id: &str, + runtime_id: &str, + worker_id: &str, + revoked_at: &str, + ) -> Result<()>; fn enqueue_ticket_notification( &self, notification_id: &str, @@ -629,6 +677,30 @@ pub trait ControlPlaneStore: Send + Sync { recipient_worker_id: &str, error: &str, ) -> Result<()>; + fn reroute_ticket_notification_delivery( + &self, + notification_id: &str, + old_runtime_id: &str, + old_worker_id: &str, + new_runtime_id: &str, + new_worker_id: &str, + ) -> Result<()>; + fn upsert_ticket_notification_cursor( + &self, + workspace_id: &str, + ticket_id: &str, + runtime_id: &str, + worker_id: &str, + event_index: i64, + updated_at: &str, + ) -> Result<()>; + fn get_ticket_notification_cursor( + &self, + workspace_id: &str, + ticket_id: &str, + runtime_id: &str, + worker_id: &str, + ) -> Result>; fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()>; fn get_workdir_registry( @@ -1713,6 +1785,62 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }) } + fn get_ticket_assignment_operation( + &self, + workspace_id: &str, + operation_id: &str, + ) -> Result> { + self.with_conn(|conn| read_assignment_operation(conn, workspace_id, operation_id)) + } + + fn reserve_ticket_assignment_operation( + &self, + workspace_id: &str, + operation_id: &str, + ticket_id: &str, + runtime_id: &str, + worker_id: &str, + created_at: &str, + ) -> Result<()> { + self.with_conn(|conn| { + let inserted = conn.execute( + r#"INSERT OR IGNORE INTO ticket_assignment_operations ( + workspace_id, operation_id, action, ticket_id, runtime_id, worker_id, + assignment_id, expected_assignment_id, created_at + ) VALUES (?1, ?2, 'assign', ?3, ?4, ?5, NULL, NULL, ?6)"#, + params![ + workspace_id, + operation_id, + ticket_id, + runtime_id, + worker_id, + created_at, + ], + )?; + if inserted > 0 { + return Ok(()); + } + let existing = read_assignment_operation(conn, workspace_id, operation_id)? + .ok_or_else(|| { + Error::TicketAssignmentConflict(format!( + "assignment operation {operation_id} could not be reserved" + )) + })?; + if existing.action == "assign" + && existing.ticket_id == ticket_id + && existing.runtime_id.as_deref() == Some(runtime_id) + && existing.worker_id.as_deref() == Some(worker_id) + && existing.expected_assignment_id.is_none() + { + Ok(()) + } else { + Err(Error::TicketAssignmentConflict(format!( + "assignment operation {operation_id} was already used with different input" + ))) + } + }) + } + fn get_current_ticket_worker_assignment( &self, workspace_id: &str, @@ -1734,9 +1862,42 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record: &TicketWorkerAssignmentRecord, expected_assignment_id: Option<&str>, event_id: &str, + operation_id: &str, + allow_reassign: bool, ) -> Result { self.with_conn(|conn| { let tx = conn.unchecked_transaction()?; + let mut reserved_operation = false; + if let Some(existing) = + read_assignment_operation(&tx, &record.workspace_id, operation_id)? + { + if existing.action != if allow_reassign { "reassign" } else { "assign" } + || existing.ticket_id != record.ticket_id + || existing.runtime_id.as_deref() != Some(record.runtime_id.as_str()) + || existing.worker_id.as_deref() != Some(record.worker_id.as_str()) + || existing.expected_assignment_id.as_deref() != expected_assignment_id + { + return Err(Error::TicketAssignmentConflict(format!( + "assignment operation {operation_id} was already used with different input" + ))); + } + if let Some(assignment_id) = existing.assignment_id { + let current = tx.query_row( + r#"SELECT workspace_id, ticket_id, assignment_id, runtime_id, worker_id, + assigned_by, assigned_at + FROM ticket_worker_assignments + WHERE workspace_id = ?1 AND ticket_id = ?2 AND assignment_id = ?3"#, + params![record.workspace_id, record.ticket_id, assignment_id], + read_ticket_worker_assignment_record, + )?; + tx.commit()?; + return Ok(TicketWorkerAssignmentUpdate { + current, + previous: None, + }); + } + reserved_operation = true; + } let previous = tx .query_row( current_ticket_worker_assignment_select_sql().as_str(), @@ -1744,11 +1905,28 @@ impl ControlPlaneStore for SqliteWorkspaceStore { read_ticket_worker_assignment_record, ) .optional()?; - require_expected_ticket_assignment( - record.ticket_id.as_str(), - previous.as_ref(), - expected_assignment_id, - )?; + if previous.is_some() && !allow_reassign { + return Err(Error::TicketAssignmentConflict(format!( + "Ticket {} is already assigned; use the explicit reassign operation", + record.ticket_id + ))); + } + if allow_reassign { + let expected_assignment_id = expected_assignment_id.ok_or_else(|| { + Error::TicketAssignmentConflict( + "reassign requires expected_assignment_id".to_string(), + ) + })?; + require_expected_ticket_assignment( + record.ticket_id.as_str(), + previous.as_ref(), + Some(expected_assignment_id), + )?; + } else if expected_assignment_id.is_some() { + return Err(Error::TicketAssignmentConflict( + "assign does not accept expected_assignment_id".to_string(), + )); + } tx.execute( r#"INSERT INTO ticket_worker_assignments ( workspace_id, ticket_id, assignment_id, runtime_id, worker_id, @@ -1764,20 +1942,42 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.assigned_at, ], )?; - tx.execute( - r#"INSERT INTO ticket_current_worker_assignments ( - workspace_id, ticket_id, assignment_id, updated_at - ) VALUES (?1, ?2, ?3, ?4) - ON CONFLICT(workspace_id, ticket_id) DO UPDATE SET - assignment_id = excluded.assignment_id, - updated_at = excluded.updated_at"#, - params![ - record.workspace_id, - record.ticket_id, - record.assignment_id, - record.assigned_at, - ], - )?; + let current_write = if allow_reassign { + tx.execute( + r#"UPDATE ticket_current_worker_assignments + SET assignment_id = ?3, runtime_id = ?4, worker_id = ?5, updated_at = ?6 + WHERE workspace_id = ?1 AND ticket_id = ?2"#, + params![ + record.workspace_id, + record.ticket_id, + record.assignment_id, + record.runtime_id, + record.worker_id, + record.assigned_at, + ], + ) + } else { + tx.execute( + r#"INSERT INTO ticket_current_worker_assignments ( + workspace_id, ticket_id, assignment_id, runtime_id, worker_id, updated_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)"#, + params![ + record.workspace_id, + record.ticket_id, + record.assignment_id, + record.runtime_id, + record.worker_id, + record.assigned_at, + ], + ) + }; + if let Err(error) = current_write { + return Err(map_assignment_constraint( + error, + &record.ticket_id, + &record.worker_id, + )); + } tx.execute( r#"INSERT INTO ticket_worker_assignment_events ( workspace_id, ticket_id, event_id, action, assignment_id, @@ -1800,6 +2000,37 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.assigned_at, ], )?; + if reserved_operation { + let updated = tx.execute( + r#"UPDATE ticket_assignment_operations + SET assignment_id = ?3 + WHERE workspace_id = ?1 AND operation_id = ?2 AND assignment_id IS NULL"#, + params![record.workspace_id, operation_id, record.assignment_id], + )?; + if updated != 1 { + return Err(Error::TicketAssignmentConflict(format!( + "assignment operation {operation_id} reservation was not current" + ))); + } + } else { + tx.execute( + r#"INSERT INTO ticket_assignment_operations ( + workspace_id, operation_id, action, ticket_id, runtime_id, worker_id, + assignment_id, expected_assignment_id, created_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)"#, + params![ + record.workspace_id, + operation_id, + if allow_reassign { "reassign" } else { "assign" }, + record.ticket_id, + record.runtime_id, + record.worker_id, + record.assignment_id, + expected_assignment_id, + record.assigned_at, + ], + )?; + } tx.commit()?; Ok(TicketWorkerAssignmentUpdate { current: record.clone(), @@ -1813,12 +2044,38 @@ impl ControlPlaneStore for SqliteWorkspaceStore { workspace_id: &str, ticket_id: &str, expected_assignment_id: Option<&str>, + operation_id: &str, event_id: &str, actor: &str, created_at: &str, ) -> Result> { self.with_conn(|conn| { let tx = conn.unchecked_transaction()?; + if let Some(existing) = read_assignment_operation(&tx, workspace_id, operation_id)? { + if existing.action != "unassign" + || existing.ticket_id != ticket_id + || existing.expected_assignment_id.as_deref() != expected_assignment_id + { + return Err(Error::TicketAssignmentConflict(format!( + "assignment operation {operation_id} was already used with different input" + ))); + } + let assignment = existing + .assignment_id + .map(|assignment_id| { + tx.query_row( + r#"SELECT workspace_id, ticket_id, assignment_id, runtime_id, worker_id, + assigned_by, assigned_at + FROM ticket_worker_assignments + WHERE workspace_id = ?1 AND ticket_id = ?2 AND assignment_id = ?3"#, + params![workspace_id, ticket_id, assignment_id], + read_ticket_worker_assignment_record, + ) + }) + .transpose()?; + tx.commit()?; + return Ok(assignment); + } let previous = tx .query_row( current_ticket_worker_assignment_select_sql().as_str(), @@ -1849,6 +2106,22 @@ impl ControlPlaneStore for SqliteWorkspaceStore { created_at, ], )?; + tx.execute( + r#"INSERT INTO ticket_assignment_operations ( + workspace_id, operation_id, action, ticket_id, runtime_id, worker_id, + assignment_id, expected_assignment_id, created_at + ) VALUES (?1, ?2, 'unassign', ?3, ?4, ?5, ?6, ?7, ?8)"#, + params![ + workspace_id, + operation_id, + ticket_id, + previous.runtime_id, + previous.worker_id, + previous.assignment_id, + expected_assignment_id, + created_at, + ], + )?; tx.commit()?; Ok(Some(previous)) }) @@ -1885,14 +2158,17 @@ impl ControlPlaneStore for SqliteWorkspaceStore { self.with_conn(|conn| { conn.execute( r#"INSERT INTO worker_workspace_credentials ( - credential_id, token, workspace_id, runtime_id, worker_id, created_at - ) VALUES (?1, ?2, ?3, ?4, ?5, ?6) + credential_id, token, workspace_id, runtime_id, worker_id, created_at, + expires_at, revoked_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) ON CONFLICT(credential_id) DO UPDATE SET token = excluded.token, workspace_id = excluded.workspace_id, runtime_id = excluded.runtime_id, worker_id = excluded.worker_id, - created_at = excluded.created_at"#, + created_at = excluded.created_at, + expires_at = excluded.expires_at, + revoked_at = excluded.revoked_at"#, params![ record.credential_id, record.token, @@ -1900,6 +2176,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.runtime_id, record.worker_id, record.created_at, + record.expires_at, + record.revoked_at, ], )?; Ok(()) @@ -1916,9 +2194,11 @@ impl ControlPlaneStore for SqliteWorkspaceStore { let tx = conn.unchecked_transaction()?; let record = tx .query_row( - r#"SELECT credential_id, token, workspace_id, runtime_id, worker_id, created_at + r#"SELECT credential_id, token, workspace_id, runtime_id, worker_id, created_at, + expires_at, revoked_at FROM worker_workspace_credentials - WHERE token = ?1 AND workspace_id = ?2"#, + WHERE token = ?1 AND workspace_id = ?2 + AND revoked_at IS NULL AND datetime(expires_at) > datetime('now')"#, params![token, workspace_id], |row| { Ok(WorkerWorkspaceCredentialRecord { @@ -1928,6 +2208,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { runtime_id: row.get(3)?, worker_id: row.get(4)?, created_at: row.get(5)?, + expires_at: row.get(6)?, + revoked_at: row.get(7)?, }) }, ) @@ -1952,6 +2234,63 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }) } + fn refresh_worker_workspace_credential( + &self, + token: &str, + workspace_id: &str, + worker_id: &str, + new_token: &str, + new_expires_at: &str, + ) -> Result> { + self.with_conn(|conn| { + let tx = conn.unchecked_transaction()?; + let record = tx + .query_row( + r#"SELECT credential_id, runtime_id, created_at FROM worker_workspace_credentials + WHERE token = ?1 AND workspace_id = ?2 AND worker_id = ?3 AND revoked_at IS NULL"#, + params![token, workspace_id, worker_id], + |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?, row.get::<_, String>(2)?)), + ) + .optional()?; + let Some((credential_id, runtime_id, created_at)) = record else { + tx.commit()?; + return Ok(None); + }; + tx.execute( + "UPDATE worker_workspace_credentials SET token = ?1, expires_at = ?2 WHERE credential_id = ?3", + params![new_token, new_expires_at, credential_id], + )?; + tx.commit()?; + Ok(Some(WorkerWorkspaceCredentialRecord { + credential_id, + token: new_token.to_string(), + workspace_id: workspace_id.to_string(), + runtime_id, + worker_id: Some(worker_id.to_string()), + created_at, + expires_at: new_expires_at.to_string(), + revoked_at: None, + })) + }) + } + + fn revoke_worker_workspace_credentials( + &self, + workspace_id: &str, + runtime_id: &str, + worker_id: &str, + revoked_at: &str, + ) -> Result<()> { + self.with_conn(|conn| { + conn.execute( + r#"UPDATE worker_workspace_credentials SET revoked_at = ?4 + WHERE workspace_id = ?1 AND runtime_id = ?2 AND worker_id = ?3 AND revoked_at IS NULL"#, + params![workspace_id, runtime_id, worker_id, revoked_at], + )?; + Ok(()) + }) + } + fn enqueue_ticket_notification( &self, notification_id: &str, @@ -2011,7 +2350,8 @@ impl ControlPlaneStore for SqliteWorkspaceStore { self.with_conn(|conn| { let mut stmt = conn.prepare( r#"SELECT o.notification_id, o.workspace_id, o.ticket_id, o.event_sequence, - o.source_runtime_id, o.source_worker_id, + o.event_kind, o.source_operation_kind, o.source_actor_role, + o.source_assignment_id, o.source_runtime_id, o.source_worker_id, d.recipient_runtime_id, d.recipient_worker_id, d.recipient_kind, d.attempts FROM ticket_notification_deliveries AS d JOIN ticket_notification_outbox AS o ON o.notification_id = d.notification_id @@ -2025,12 +2365,16 @@ impl ControlPlaneStore for SqliteWorkspaceStore { workspace_id: row.get(1)?, ticket_id: row.get(2)?, event_sequence: row.get(3)?, - source_runtime_id: row.get(4)?, - source_worker_id: row.get(5)?, - recipient_runtime_id: row.get(6)?, - recipient_worker_id: row.get(7)?, - recipient_kind: row.get(8)?, - attempts: row.get(9)?, + event_kind: row.get(4)?, + source_operation_kind: row.get(5)?, + source_actor_role: row.get(6)?, + source_assignment_id: row.get(7)?, + source_runtime_id: row.get(8)?, + source_worker_id: row.get(9)?, + recipient_runtime_id: row.get(10)?, + recipient_worker_id: row.get(11)?, + recipient_kind: row.get(12)?, + attempts: row.get(13)?, }) })?; rows.collect::, _>>() @@ -2095,6 +2439,91 @@ impl ControlPlaneStore for SqliteWorkspaceStore { }) } + fn reroute_ticket_notification_delivery( + &self, + notification_id: &str, + old_runtime_id: &str, + old_worker_id: &str, + new_runtime_id: &str, + new_worker_id: &str, + ) -> Result<()> { + self.with_conn(|conn| { + let tx = conn.unchecked_transaction()?; + let recipient_kind: Option = tx + .query_row( + r#"SELECT recipient_kind FROM ticket_notification_deliveries + WHERE notification_id = ?1 AND recipient_runtime_id = ?2 AND recipient_worker_id = ?3"#, + params![notification_id, old_runtime_id, old_worker_id], + |row| row.get(0), + ) + .optional()?; + if let Some(recipient_kind) = recipient_kind { + tx.execute( + r#"INSERT OR IGNORE INTO ticket_notification_deliveries ( + notification_id, recipient_runtime_id, recipient_worker_id, recipient_kind, attempts + ) VALUES (?1, ?2, ?3, ?4, 0)"#, + params![notification_id, new_runtime_id, new_worker_id, recipient_kind], + )?; + tx.execute( + r#"DELETE FROM ticket_notification_deliveries + WHERE notification_id = ?1 AND recipient_runtime_id = ?2 AND recipient_worker_id = ?3"#, + params![notification_id, old_runtime_id, old_worker_id], + )?; + } + tx.commit()?; + Ok(()) + }) + } + + fn upsert_ticket_notification_cursor( + &self, + workspace_id: &str, + ticket_id: &str, + runtime_id: &str, + worker_id: &str, + event_index: i64, + updated_at: &str, + ) -> Result<()> { + self.with_conn(|conn| { + conn.execute( + r#"INSERT INTO ticket_notification_cursors ( + workspace_id, ticket_id, runtime_id, worker_id, last_event_index, updated_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6) + ON CONFLICT(workspace_id, ticket_id, runtime_id, worker_id) DO UPDATE SET + last_event_index = MAX(last_event_index, excluded.last_event_index), + updated_at = excluded.updated_at"#, + params![ + workspace_id, + ticket_id, + runtime_id, + worker_id, + event_index, + updated_at + ], + )?; + Ok(()) + }) + } + + fn get_ticket_notification_cursor( + &self, + workspace_id: &str, + ticket_id: &str, + runtime_id: &str, + worker_id: &str, + ) -> Result> { + self.with_conn(|conn| { + conn.query_row( + r#"SELECT last_event_index FROM ticket_notification_cursors + WHERE workspace_id = ?1 AND ticket_id = ?2 AND runtime_id = ?3 AND worker_id = ?4"#, + params![workspace_id, ticket_id, runtime_id, worker_id], + |row| row.get(0), + ) + .optional() + .map_err(Error::from) + }) + } + fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()> { self.with_conn(|conn| { conn.execute( @@ -2566,6 +2995,51 @@ fn read_ticket_worker_assignment_event_record( }) } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TicketAssignmentOperationRecord { + pub action: String, + pub ticket_id: String, + pub runtime_id: Option, + pub worker_id: Option, + pub assignment_id: Option, + pub expected_assignment_id: Option, +} + +fn read_assignment_operation( + conn: &Connection, + workspace_id: &str, + operation_id: &str, +) -> Result> { + conn.query_row( + r#"SELECT action, ticket_id, runtime_id, worker_id, assignment_id, expected_assignment_id + FROM ticket_assignment_operations + WHERE workspace_id = ?1 AND operation_id = ?2"#, + params![workspace_id, operation_id], + |row| { + Ok(TicketAssignmentOperationRecord { + action: row.get(0)?, + ticket_id: row.get(1)?, + runtime_id: row.get(2)?, + worker_id: row.get(3)?, + assignment_id: row.get(4)?, + expected_assignment_id: row.get(5)?, + }) + }, + ) + .optional() + .map_err(Error::from) +} + +fn map_assignment_constraint(error: rusqlite::Error, ticket_id: &str, worker_id: &str) -> Error { + if matches!(error, rusqlite::Error::SqliteFailure(_, _)) { + Error::TicketAssignmentConflict(format!( + "Ticket {ticket_id} or Worker {worker_id} already has a current assignment" + )) + } else { + Error::Sqlite(error) + } +} + fn require_expected_ticket_assignment( ticket_id: &str, current: Option<&TicketWorkerAssignmentRecord>, @@ -2850,6 +3324,79 @@ CREATE INDEX IF NOT EXISTS idx_ticket_notification_pending Ok(()) } +fn strengthen_ticket_worker_assignments(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" +ALTER TABLE ticket_current_worker_assignments RENAME TO ticket_current_worker_assignments_v16; + +CREATE TABLE ticket_current_worker_assignments ( + workspace_id TEXT NOT NULL, + ticket_id TEXT NOT NULL, + assignment_id TEXT NOT NULL, + runtime_id TEXT NOT NULL, + worker_id TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id), + UNIQUE (workspace_id, runtime_id, worker_id), + FOREIGN KEY (workspace_id, ticket_id, assignment_id) + REFERENCES ticket_worker_assignments(workspace_id, ticket_id, assignment_id) + ON DELETE CASCADE +); + +INSERT INTO ticket_current_worker_assignments ( + workspace_id, ticket_id, assignment_id, runtime_id, worker_id, updated_at +) +SELECT current.workspace_id, current.ticket_id, current.assignment_id, + assignment.runtime_id, assignment.worker_id, current.updated_at +FROM ticket_current_worker_assignments_v16 AS current +JOIN ticket_worker_assignments AS assignment + ON assignment.workspace_id = current.workspace_id + AND assignment.ticket_id = current.ticket_id + AND assignment.assignment_id = current.assignment_id; + +DROP TABLE ticket_current_worker_assignments_v16; + +CREATE TABLE ticket_assignment_operations ( + workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE, + operation_id TEXT NOT NULL, + action TEXT NOT NULL CHECK (action IN ('assign', 'reassign', 'unassign')), + ticket_id TEXT NOT NULL, + runtime_id TEXT, + worker_id TEXT, + assignment_id TEXT, + expected_assignment_id TEXT, + created_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, operation_id) +); +"#, + )?; + Ok(()) +} + +fn strengthen_ticket_notifications(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" +ALTER TABLE worker_workspace_credentials ADD COLUMN expires_at TEXT; +ALTER TABLE worker_workspace_credentials ADD COLUMN revoked_at TEXT; +ALTER TABLE ticket_notification_outbox ADD COLUMN event_kind TEXT NOT NULL DEFAULT 'comment'; +ALTER TABLE ticket_notification_outbox ADD COLUMN source_operation_kind TEXT NOT NULL DEFAULT 'unknown'; +ALTER TABLE ticket_notification_outbox ADD COLUMN source_actor_role TEXT NOT NULL DEFAULT 'worker'; +ALTER TABLE ticket_notification_outbox ADD COLUMN source_assignment_id TEXT; + +CREATE TABLE ticket_notification_cursors ( + workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE, + ticket_id TEXT NOT NULL, + runtime_id TEXT NOT NULL, + worker_id TEXT NOT NULL, + last_event_index INTEGER NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, runtime_id, worker_id) +); +"#, + )?; + Ok(()) +} + fn create_objective_event_tables(conn: &Connection) -> Result<()> { conn.execute_batch( r#" @@ -3532,7 +4079,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); let db = dir.path().join("control-plane.sqlite"); let store = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 16); + assert_eq!(store.schema_version().await.unwrap(), 18); let record = WorkspaceRecord { workspace_id: "local-dev".to_string(), @@ -3545,7 +4092,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); store.upsert_workspace(&record).await.unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(reopened.schema_version().await.unwrap(), 16); + assert_eq!(reopened.schema_version().await.unwrap(), 18); assert_eq!( reopened.get_workspace("local-dev").await.unwrap(), Some(record) @@ -3578,10 +4125,65 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); assigned_at: "2026-07-31T00:00:01Z".to_string(), }; let created = store - .set_current_ticket_worker_assignment(&first, None, "event-1") + .set_current_ticket_worker_assignment(&first, None, "event-1", "operation-1", false) .unwrap(); assert_eq!(created.current, first); assert_eq!(created.previous, None); + let retried = store + .set_current_ticket_worker_assignment( + &TicketWorkerAssignmentRecord { + assignment_id: "ignored-retry-assignment".to_string(), + ..first.clone() + }, + None, + "ignored-retry-event", + "operation-1", + false, + ) + .unwrap(); + assert_eq!(retried.current, first); + assert_eq!( + store + .list_ticket_worker_assignment_events("workspace-a", "ticket-1", 10) + .unwrap() + .len(), + 1, + "idempotent retry must not append another assignment event" + ); + let implicit_reassign = store + .set_current_ticket_worker_assignment( + &TicketWorkerAssignmentRecord { + assignment_id: "implicit-reassign".to_string(), + worker_id: "worker-other".to_string(), + ..first.clone() + }, + None, + "implicit-event", + "implicit-operation", + false, + ) + .unwrap_err(); + assert!(matches!( + implicit_reassign, + Error::TicketAssignmentConflict(_) + )); + let worker_conflict = store + .set_current_ticket_worker_assignment( + &TicketWorkerAssignmentRecord { + ticket_id: "ticket-2".to_string(), + assignment_id: "worker-conflict".to_string(), + ..first.clone() + }, + None, + "worker-conflict-event", + "worker-conflict-operation", + false, + ) + .unwrap_err(); + assert!(matches!( + worker_conflict, + Error::TicketAssignmentConflict(_) + )); let second = TicketWorkerAssignmentRecord { assignment_id: "assignment-2".to_string(), @@ -3592,7 +4194,13 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); ..first.clone() }; let replaced = store - .set_current_ticket_worker_assignment(&second, Some("assignment-1"), "event-2") + .set_current_ticket_worker_assignment( + &second, + Some("assignment-1"), + "event-2", + "operation-2", + true, + ) .unwrap(); assert_eq!(replaced.current, second); assert_eq!(replaced.previous, Some(first.clone())); @@ -3608,6 +4216,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); "workspace-a", "ticket-1", Some("assignment-1"), + "unassign-operation-stale", "event-stale", "user-1", "2026-07-31T00:00:03Z", @@ -3620,12 +4229,61 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); "workspace-a", "ticket-1", Some("assignment-2"), + "unassign-operation-2", "event-3", "user-2", "2026-07-31T00:00:03Z", ) .unwrap(); - assert_eq!(cleared, Some(second)); + assert_eq!(cleared, Some(second.clone())); + let retried_clear = store + .clear_current_ticket_worker_assignment( + "workspace-a", + "ticket-1", + Some("assignment-2"), + "unassign-operation-2", + "ignored-clear-event", + "user-2", + "2026-07-31T00:00:04Z", + ) + .unwrap(); + assert_eq!(retried_clear, Some(second)); + store + .reserve_ticket_assignment_operation( + "workspace-a", + "reserved-operation", + "ticket-3", + "runtime-3", + "worker-3", + "2026-07-31T00:00:05Z", + ) + .unwrap(); + let reserved_assignment = TicketWorkerAssignmentRecord { + workspace_id: "workspace-a".to_string(), + ticket_id: "ticket-3".to_string(), + assignment_id: "assignment-3".to_string(), + runtime_id: "runtime-3".to_string(), + worker_id: "worker-3".to_string(), + assigned_by: "runtime".to_string(), + assigned_at: "2026-07-31T00:00:06Z".to_string(), + }; + let completed_reservation = store + .set_current_ticket_worker_assignment( + &reserved_assignment, + None, + "reserved-event", + "reserved-operation", + false, + ) + .unwrap(); + assert_eq!(completed_reservation.current, reserved_assignment); + assert_eq!( + store + .get_ticket_assignment_operation("workspace-a", "reserved-operation") + .unwrap() + .and_then(|operation| operation.assignment_id), + Some("assignment-3".to_string()) + ); assert_eq!( store .get_current_ticket_worker_assignment("workspace-a", "ticket-1") @@ -3673,6 +4331,8 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); runtime_id: "runtime-1".to_string(), worker_id: None, created_at: "2026-07-31T00:00:01Z".to_string(), + expires_at: "2099-01-01T00:00:00Z".to_string(), + revoked_at: None, }) .unwrap(); let bound = store @@ -3680,6 +4340,53 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); .unwrap() .unwrap(); assert_eq!(bound.worker_id.as_deref(), Some("worker-1")); + let refreshed = store + .refresh_worker_workspace_credential( + "secret-token", + "workspace-a", + "worker-1", + "refreshed-token", + "2099-02-01T00:00:00Z", + ) + .unwrap() + .unwrap(); + assert_eq!(refreshed.token, "refreshed-token"); + assert!(store + .authenticate_worker_workspace_credential( + "secret-token", + "workspace-a", + "worker-1", + ) + .unwrap() + .is_none()); + assert!( + store + .authenticate_worker_workspace_credential( + "refreshed-token", + "workspace-a", + "worker-1", + ) + .unwrap() + .is_some() + ); + store + .revoke_worker_workspace_credentials( + "workspace-a", + "runtime-1", + "worker-1", + "2026-08-01T00:00:00Z", + ) + .unwrap(); + assert!( + store + .authenticate_worker_workspace_credential( + "refreshed-token", + "workspace-a", + "worker-1", + ) + .unwrap() + .is_none() + ); assert!( store .authenticate_worker_workspace_credential( @@ -3715,10 +4422,41 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); assert_eq!(pending.len(), 1); assert_eq!(pending[0].event_sequence, 4); store - .mark_ticket_notification_delivered( + .reroute_ticket_notification_delivery( "notification-1", "runtime-1", "worker-2", + "runtime-2", + "worker-3", + ) + .unwrap(); + assert_eq!( + store + .count_ticket_notification_deliveries_for_recipient( + "workspace-a", + "ticket-1", + "runtime-1", + "worker-2", + ) + .unwrap(), + 0 + ); + assert_eq!( + store + .count_ticket_notification_deliveries_for_recipient( + "workspace-a", + "ticket-1", + "runtime-2", + "worker-3", + ) + .unwrap(), + 1 + ); + store + .mark_ticket_notification_delivered( + "notification-1", + "runtime-2", + "worker-3", "2026-07-31T00:00:03Z", ) .unwrap(); @@ -3914,7 +4652,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); .unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 16); + assert_eq!(store.schema_version().await.unwrap(), 18); store .with_conn(|conn| { @@ -4017,7 +4755,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn repository_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 16); + assert_eq!(store.schema_version().await.unwrap(), 18); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -4055,7 +4793,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn memory_authority_records_round_trip_and_close_staging() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 16); + assert_eq!(store.schema_version().await.unwrap(), 18); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -4229,7 +4967,7 @@ CREATE TABLE ticket_worker_links (ticket_id TEXT, worker_ref_key TEXT); #[tokio::test] async fn account_and_login_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 16); + assert_eq!(store.schema_version().await.unwrap(), 18); let now = "2026-07-22T00:00:00Z".to_string(); let account = AccountRecord { account_id: "acct-user-alice".to_string(),