From f079479160b88255ab10191a6d8fb2238e9e0097 Mon Sep 17 00:00:00 2001 From: Hare Date: Tue, 25 Aug 2026 10:53:31 +0900 Subject: [PATCH] feat: queue ready dependency closures atomically --- crates/client/src/workspace_product.rs | 8 +- crates/ticket/src/lib.rs | 597 +++++++++++++++--- crates/ticket/src/tool.rs | 26 +- crates/tui/src/workspace_panel.rs | 10 +- crates/worker/src/feature/builtin/ticket.rs | 27 +- crates/workspace-server/src/authority.rs | 5 +- crates/workspace-server/src/server.rs | 198 +++++- .../tickets/[ticketId]/+page.svelte | 4 +- 8 files changed, 733 insertions(+), 142 deletions(-) diff --git a/crates/client/src/workspace_product.rs b/crates/client/src/workspace_product.rs index b892c210..46207451 100644 --- a/crates/client/src/workspace_product.rs +++ b/crates/client/src/workspace_product.rs @@ -473,8 +473,12 @@ impl TicketBackend for BackendWorkspaceProductClient { .map_err(ticket_client_error) } - fn queue_ready(&self, id: TicketIdOrSlug, _queued_by: &str) -> ticket::Result<()> { - self.send_unit::<()>( + fn queue_ready( + &self, + id: TicketIdOrSlug, + _queued_by: &str, + ) -> ticket::Result { + self.send_json::<(), _>( Method::POST, &format!( "/tickets/{}/workflow/queue", diff --git a/crates/ticket/src/lib.rs b/crates/ticket/src/lib.rs index 5e5a6a2f..21cf9b5a 100644 --- a/crates/ticket/src/lib.rs +++ b/crates/ticket/src/lib.rs @@ -807,6 +807,12 @@ pub struct TicketDependencyCheck { pub recommended_action: TicketWorkspaceNextAction, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct TicketQueueOutcome { + pub requested_ticket: String, + pub queued_tickets: Vec, +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] #[serde(rename_all = "lowercase")] pub enum TicketListState { @@ -1117,7 +1123,7 @@ pub fn project_ticket_workspace_item( pub fn ticket_queue_guard( summary: &TicketSummary, - _relation_blockers: &[TicketRelationBlocker], + relation_blockers: &[TicketRelationBlocker], orchestration_overlay: Option<&TicketWorkspaceStateOverlay>, ) -> TicketQueueGuard { if orchestration_overlay.is_some() { @@ -1140,9 +1146,34 @@ pub fn ticket_queue_guard( blocked_reason: None, }; } + if relation_blockers + .iter() + .any(|blocker| blocker.blocking_state == TicketWorkflowState::Planning) + { + let blocked_reason = relation_blockers + .iter() + .filter(|blocker| blocker.blocking_state == TicketWorkflowState::Planning) + .map(|blocker| { + format!( + "{} ({})", + blocker.blocking_ticket, + blocker.blocking_state.as_str() + ) + }) + .collect::>() + .join(", "); + return TicketQueueGuard { + can_queue_for_orchestrator: false, + reason: Some("Dependencies must leave planning before Queue can proceed".to_string()), + blocked_reason: Some(blocked_reason), + }; + } TicketQueueGuard { can_queue_for_orchestrator: true, - reason: None, + reason: (!relation_blockers.is_empty()).then(|| { + "Ready dependencies will be queued atomically; active dependencies remain unchanged" + .to_string() + }), blocked_reason: None, } } @@ -1156,7 +1187,7 @@ fn derive_ticket_workspace_projection( .iter() .filter(|blocker| !relation_blocker_allows_ready_queue(blocker)) .collect::>(); - if summary.workflow_state != TicketWorkflowState::Ready { + if summary.workflow_state != TicketWorkflowState::Ready || !active_blockers.is_empty() { let blockers_to_report = if active_blockers.is_empty() { relation_blockers.iter().collect::>() } else { @@ -1201,7 +1232,7 @@ fn derive_ticket_workspace_projection( visible_overlay: None, disabled_reason: None, key_hint: Some(format!( - "Queue records orchestration demand; dependency relations remain scheduling context ({blockers})." + "Queue records orchestration demand; dependency relations remain orchestration context ({blockers})." )), blocked_reason: Some(blockers), queue_guard: TicketQueueGuard { @@ -1409,7 +1440,7 @@ fn compact_ticket_state_label(state: TicketWorkflowState) -> &'static str { fn relation_blocker_allows_ready_queue(blocker: &TicketRelationBlocker) -> bool { matches!( blocker.blocking_state, - TicketWorkflowState::Queued | TicketWorkflowState::InProgress + TicketWorkflowState::Ready | TicketWorkflowState::Queued | TicketWorkflowState::InProgress ) } @@ -1733,7 +1764,7 @@ pub trait TicketBackend { ) -> Result<()>; fn set_workflow_state(&self, id: TicketIdOrSlug, change: TicketStateChange) -> Result<()>; fn mark_ready(&self, id: TicketIdOrSlug, request: TicketMarkReady) -> Result; - fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> Result<()>; + fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> Result; fn close(&self, id: TicketIdOrSlug, resolution: MarkdownText) -> Result<()>; fn add_ticket_relation( &self, @@ -1856,6 +1887,7 @@ pub enum TicketBackendOperationResult { Ticket(Ticket), TicketRef(TicketRef), DependencyCheck(TicketDependencyCheck), + QueueOutcome(TicketQueueOutcome), Relation(TicketRelation), Relations(Vec), RelationView(TicketRelationView), @@ -1916,8 +1948,7 @@ where TicketBackendOperationResult::Ticket(backend.mark_ready(id, request)?) } TicketBackendOperation::QueueReady { id, queued_by } => { - backend.queue_ready(id, &queued_by)?; - TicketBackendOperationResult::Unit + TicketBackendOperationResult::QueueOutcome(backend.queue_ready(id, &queued_by)?) } TicketBackendOperation::Close { id, resolution } => { backend.close(id, resolution)?; @@ -3893,25 +3924,90 @@ impl TicketBackend for SqliteTicketBackend { }) } - fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> Result<()> { + fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> Result { validate_required_event_value("queued_by", queued_by)?; self.with_write(|conn| { - let ticket_id = self.resolve_ticket_id(conn, id)?; - let ticket = self.load_ticket(conn, &ticket_id)?; - if ticket.meta.workflow_state != TicketWorkflowState::Ready { - return Err(TicketError::StaleWorkflowState { - expected: TicketWorkflowState::Ready.as_str().to_owned(), - actual: ticket.meta.workflow_state.as_str().to_owned(), - }); + let requested_ticket = self.resolve_ticket_id(conn, id)?; + let states = self.state_index(conn)?; + let relations = self.all_relations(conn)?; + let queued_tickets = dependency_queue_plan(&requested_ticket, &states, &relations)?; + if let Some(json) = self + .event_attributes + .get("queue_orchestrator_assignments") + { + let assignments = serde_json::from_str::>(json).map_err( + |error| { + TicketError::Conflict(format!( + "invalid Queue assignment fence: {error}" + )) + }, + )?; + let planned = assignments.keys().cloned().collect::>(); + let actual = queued_tickets.iter().cloned().collect::>(); + if planned != actual || assignments.values().any(|value| value.trim().is_empty()) { + return Err(TicketError::Conflict( + "Queue dependency plan changed after assignment validation".to_string(), + )); + } } - let target = resolve_ready_target( - self.target_authority.as_ref(), - &self.workspace_id, - &ticket, - )?; + + let mut targets = Vec::with_capacity(queued_tickets.len()); + for ticket_id in &queued_tickets { + let ticket = self.load_ticket(conn, ticket_id)?; + let target = resolve_ready_target( + self.target_authority.as_ref(), + &self.workspace_id, + &ticket, + )?; + targets.push((ticket_id.clone(), target)); + } + let at = now_utc(); - conn.execute("UPDATE typed_tickets SET workflow_state = 'queued', workflow_state_explicit = 1, queued_by = ?3, queued_at = ?4, repository_id = ?5, ref_selector = ?6, updated_at = ?4 WHERE workspace_id = ?1 AND ticket_id = ?2 AND workflow_state = 'ready'", params![self.workspace_id, ticket_id, queued_by, at, target.repository_id, target.ref_selector]).map_err(sqlite_err)?; - self.insert_event(conn, &ticket_id, &TicketEvent { kind: TicketEventKind::StateChanged, author: Some(queued_by.to_string()), at: Some(at.clone()), status: None, from: Some("ready".to_string()), to: Some("queued".to_string()), reason: Some("queued".to_string()), state_field: Some("state".to_string()), heading: Some(TicketEventKind::StateChanged.heading()), body: MarkdownText::new(format!("Queued for Orchestrator by {queued_by}.")), references: Vec::new(), attributes: BTreeMap::from([("queued_by".to_owned(), queued_by.to_owned()), ("queued_at".to_owned(), at), ("repository_id".to_owned(), target.repository_id), ("ref_selector".to_owned(), target.ref_selector)]) }) + for (ticket_id, target) in targets { + let updated = conn.execute( + "UPDATE typed_tickets SET workflow_state = 'queued', workflow_state_explicit = 1, queued_by = ?3, queued_at = ?4, repository_id = ?5, ref_selector = ?6, updated_at = ?4 WHERE workspace_id = ?1 AND ticket_id = ?2 AND workflow_state = 'ready'", + params![self.workspace_id, ticket_id, queued_by, at, target.repository_id, target.ref_selector], + ).map_err(sqlite_err)?; + if updated != 1 { + return Err(TicketError::Conflict(format!( + "Ticket {ticket_id} changed while the dependency queue plan was being applied" + ))); + } + let mut attributes = BTreeMap::from([ + ("queued_by".to_owned(), queued_by.to_owned()), + ("queued_at".to_owned(), at.clone()), + ("repository_id".to_owned(), target.repository_id), + ("ref_selector".to_owned(), target.ref_selector), + ("queue_root_ticket".to_owned(), requested_ticket.clone()), + ]); + if let Some(assignment_id) = self + .event_attributes + .get("queue_orchestrator_assignments") + .and_then(|json| serde_json::from_str::>(json).ok()) + .and_then(|assignments| assignments.get(&ticket_id).cloned()) + { + attributes.insert("orchestrator_assignment_id".to_owned(), assignment_id); + } + self.insert_event(conn, &ticket_id, &TicketEvent { + kind: TicketEventKind::StateChanged, + author: Some(queued_by.to_string()), + at: Some(at.clone()), + status: None, + from: Some("ready".to_string()), + to: Some("queued".to_string()), + reason: Some("queued".to_string()), + state_field: Some("state".to_string()), + heading: Some(TicketEventKind::StateChanged.heading()), + body: MarkdownText::new(format!("Queued for Orchestrator by {queued_by}.")), + references: Vec::new(), + attributes, + })?; + } + + Ok(TicketQueueOutcome { + requested_ticket, + queued_tickets, + }) }) } @@ -4530,40 +4626,56 @@ impl TicketBackend for LocalTicketBackend { self.ticket_from_dir(&dir) } - fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> Result<()> { + fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> Result { validate_required_event_value("queued_by", queued_by)?; let _lock = self.acquire_lock()?; - let dir = self.find_ticket_dir(&id)?; - let item = dir.join("item.md"); - let meta = ticket_meta_for_dir(&dir, read_item_file(&item)?.frontmatter)?; - if meta.workflow_state != TicketWorkflowState::Ready { - return Err(TicketError::StaleWorkflowState { - expected: TicketWorkflowState::Ready.as_str().to_owned(), - actual: meta.workflow_state.as_str().to_owned(), - }); + let requested_dir = self.find_ticket_dir(&id)?; + let requested_ticket = ticket_id_from_dir(&requested_dir)?; + let mut states = HashMap::new(); + for dir in self.iter_ticket_dirs(TicketListQuery::all())? { + let item = dir.join("item.md"); + let meta = ticket_meta_for_dir(&dir, read_item_file(&item)?.frontmatter)?; + states.insert(meta.id, meta.workflow_state); } - let ticket = self.ticket_from_dir(&dir)?; - let target = resolve_ready_target(self.target_authority.as_ref(), "local", &ticket)?; - let at = now_utc(); - let mut change = TicketStateChange::new( - TicketWorkflowState::Ready.as_str(), - TicketWorkflowState::Queued.as_str(), - "queued", - self.queued_ready_body(queued_by), - ); - change.author = Some(queued_by.to_string()); - self.apply_workflow_state_change( - &dir, - TicketWorkflowState::Ready, - TicketWorkflowState::Queued, - change, - &[ - ("queued_by", queued_by), - ("queued_at", at.as_str()), - ("repository_id", target.repository_id.as_str()), - ("ref_selector", target.ref_selector.as_str()), - ], - ) + let relations = self.all_ticket_relation_records()?; + let queued_tickets = dependency_queue_plan(&requested_ticket, &states, &relations)?; + + let mut planned = Vec::with_capacity(queued_tickets.len()); + for ticket_id in &queued_tickets { + let dir = self.find_ticket_dir(&TicketIdOrSlug::Id(ticket_id.clone()))?; + let ticket = self.ticket_from_dir(&dir)?; + let target = resolve_ready_target(self.target_authority.as_ref(), "local", &ticket)?; + planned.push((dir, target)); + } + + for (dir, target) in planned { + let at = now_utc(); + let mut change = TicketStateChange::new( + TicketWorkflowState::Ready.as_str(), + TicketWorkflowState::Queued.as_str(), + "queued", + self.queued_ready_body(queued_by), + ); + change.author = Some(queued_by.to_string()); + self.apply_workflow_state_change( + &dir, + TicketWorkflowState::Ready, + TicketWorkflowState::Queued, + change, + &[ + ("queued_by", queued_by), + ("queued_at", at.as_str()), + ("repository_id", target.repository_id.as_str()), + ("ref_selector", target.ref_selector.as_str()), + ("queue_root_ticket", requested_ticket.as_str()), + ], + )?; + } + + Ok(TicketQueueOutcome { + requested_ticket, + queued_tickets, + }) } fn close(&self, id: TicketIdOrSlug, resolution: MarkdownText) -> Result<()> { @@ -5395,6 +5507,142 @@ fn ticket_state_resolved(state: TicketWorkflowState) -> bool { ) } +fn dependency_queue_plan( + requested_ticket: &str, + states: &HashMap, + relations: &[TicketRelation], +) -> Result> { + let requested_state = states + .get(requested_ticket) + .copied() + .ok_or_else(|| TicketError::NotFound(requested_ticket.to_owned()))?; + if requested_state != TicketWorkflowState::Ready { + return Err(TicketError::StaleWorkflowState { + expected: TicketWorkflowState::Ready.as_str().to_owned(), + actual: requested_state.as_str().to_owned(), + }); + } + + let mut prerequisites = BTreeMap::>::new(); + for relation in relations { + match relation.kind { + TicketRelationKind::DependsOn => { + prerequisites + .entry(relation.ticket_id.clone()) + .or_default() + .insert(relation.target.clone()); + } + TicketRelationKind::Blocks => { + prerequisites + .entry(relation.target.clone()) + .or_default() + .insert(relation.ticket_id.clone()); + } + TicketRelationKind::Related + | TicketRelationKind::Supersedes + | TicketRelationKind::DuplicateOf => {} + } + } + + fn visit( + ticket: &str, + prerequisites: &BTreeMap>, + marks: &mut BTreeMap, + stack: &mut Vec, + ordered: &mut Vec, + ) -> Result<()> { + match marks.get(ticket).copied() { + Some(2) => return Ok(()), + Some(1) => { + let start = stack.iter().position(|item| item == ticket).unwrap_or(0); + let mut cycle = stack[start..].to_vec(); + cycle.push(ticket.to_owned()); + return Err(TicketError::Conflict(format!( + "ticket dependency cycle detected: {}", + cycle.join(" -> ") + ))); + } + _ => {} + } + + marks.insert(ticket.to_owned(), 1); + stack.push(ticket.to_owned()); + if let Some(dependencies) = prerequisites.get(ticket) { + for dependency in dependencies { + visit(dependency, prerequisites, marks, stack, ordered)?; + } + } + stack.pop(); + marks.insert(ticket.to_owned(), 2); + ordered.push(ticket.to_owned()); + Ok(()) + } + + let mut ordered = Vec::new(); + visit( + requested_ticket, + &prerequisites, + &mut BTreeMap::new(), + &mut Vec::new(), + &mut ordered, + )?; + + fn collect_ready( + ticket: &str, + requested_ticket: &str, + states: &HashMap, + prerequisites: &BTreeMap>, + collected: &mut BTreeSet, + ready: &mut Vec, + ) -> Result<()> { + if !collected.insert(ticket.to_owned()) { + return Ok(()); + } + let state = states + .get(ticket) + .copied() + .ok_or_else(|| TicketError::NotFound(ticket.to_owned()))?; + if ticket != requested_ticket && state == TicketWorkflowState::Planning { + return Err(TicketError::BlockingRelations(format!( + "dependency {ticket} is still planning" + ))); + } + if matches!( + state, + TicketWorkflowState::Done | TicketWorkflowState::Closed + ) { + return Ok(()); + } + if let Some(dependencies) = prerequisites.get(ticket) { + for dependency in dependencies { + collect_ready( + dependency, + requested_ticket, + states, + prerequisites, + collected, + ready, + )?; + } + } + if state == TicketWorkflowState::Ready { + ready.push(ticket.to_owned()); + } + Ok(()) + } + + let mut ready = Vec::new(); + collect_ready( + requested_ticket, + requested_ticket, + states, + &prerequisites, + &mut BTreeSet::new(), + &mut ready, + )?; + Ok(ready) +} + fn relation_view_from_records( meta: &TicketMeta, records: &[TicketRelation], @@ -6892,6 +7140,86 @@ mod tests { } } + fn dependency_relation(ticket_id: &str, target: &str) -> TicketRelation { + TicketRelation { + ticket_id: ticket_id.to_owned(), + kind: TicketRelationKind::DependsOn, + target: target.to_owned(), + note: None, + author: "test".to_owned(), + at: String::new(), + } + } + + #[test] + fn dependency_queue_plan_orders_transitive_ready_dependencies() { + let states = HashMap::from([ + ("root".to_owned(), TicketWorkflowState::Ready), + ("middle".to_owned(), TicketWorkflowState::Ready), + ("leaf".to_owned(), TicketWorkflowState::Ready), + ]); + let relations = [ + dependency_relation("root", "middle"), + dependency_relation("middle", "leaf"), + ]; + + assert_eq!( + dependency_queue_plan("root", &states, &relations).unwrap(), + vec!["leaf", "middle", "root"] + ); + } + + #[test] + fn dependency_queue_plan_stops_at_resolved_dependency() { + let states = HashMap::from([ + ("root".to_owned(), TicketWorkflowState::Ready), + ("done".to_owned(), TicketWorkflowState::Done), + ("planning".to_owned(), TicketWorkflowState::Planning), + ]); + let relations = [ + dependency_relation("root", "done"), + dependency_relation("done", "planning"), + ]; + + assert_eq!( + dependency_queue_plan("root", &states, &relations).unwrap(), + vec!["root"] + ); + } + + #[test] + fn dependency_queue_plan_rejects_transitive_planning_dependency() { + let states = HashMap::from([ + ("root".to_owned(), TicketWorkflowState::Ready), + ("middle".to_owned(), TicketWorkflowState::Queued), + ("leaf".to_owned(), TicketWorkflowState::Planning), + ]); + let relations = [ + dependency_relation("root", "middle"), + dependency_relation("middle", "leaf"), + ]; + + let error = dependency_queue_plan("root", &states, &relations).unwrap_err(); + assert!(matches!(error, TicketError::BlockingRelations(_))); + assert!(error.to_string().contains("leaf")); + } + + #[test] + fn dependency_queue_plan_reports_cycle_path() { + let states = HashMap::from([ + ("root".to_owned(), TicketWorkflowState::Ready), + ("middle".to_owned(), TicketWorkflowState::Ready), + ]); + let relations = [ + dependency_relation("root", "middle"), + dependency_relation("middle", "root"), + ]; + + let error = dependency_queue_plan("root", &states, &relations).unwrap_err(); + assert!(matches!(error, TicketError::Conflict(_))); + assert!(error.to_string().contains("root -> middle -> root")); + } + #[test] fn workspace_projection_queues_ready_ticket_for_orchestrator() { let summary = summary_with_state(TicketWorkflowState::Ready); @@ -6911,7 +7239,7 @@ mod tests { } #[test] - fn workspace_projection_queues_ready_ticket_with_unstarted_dependency_context() { + fn workspace_projection_blocks_ready_ticket_with_planning_dependency() { let summary = summary_with_state(TicketWorkflowState::Ready); let blockers = [blocker_with_state(TicketWorkflowState::Planning)]; let projection = project_ticket_workspace_item(&summary, &blockers, None); @@ -6919,18 +7247,11 @@ mod tests { assert_eq!(projection.kind, TicketWorkspaceRowKind::Ticket); assert_eq!( projection.next_action, - Some(TicketWorkspaceNextAction::QueueForOrchestrator) + Some(TicketWorkspaceNextAction::WaitForOrchestrator) ); - assert!(projection.queue_guard.can_queue_for_orchestrator); + assert!(!projection.queue_guard.can_queue_for_orchestrator); assert!(projection.blocked_reason.is_some()); - assert!(projection.disabled_reason.is_none()); - assert!( - projection - .key_hint - .as_deref() - .unwrap_or_default() - .contains("scheduling context") - ); + assert!(projection.disabled_reason.is_some()); } #[test] @@ -6950,7 +7271,7 @@ mod tests { .key_hint .as_deref() .unwrap_or_default() - .contains("scheduling context") + .contains("orchestration context") ); } @@ -7211,13 +7532,104 @@ state: planning } #[test] - fn sqlite_mark_ready_and_queue_preserve_dependency_context_atomically() { + fn sqlite_queue_cycle_diagnostic_leaves_all_tickets_ready() { + let temp = TempDir::new().unwrap(); + let backend = SqliteTicketBackend::open(temp.path().join("tickets.db"), "workspace-test") + .unwrap() + .with_target_authority(Arc::new(TestTargetAuthority)); + let mut first_input = NewTicket::new("First ready Ticket"); + first_input.workflow_state = Some(TicketWorkflowState::Ready); + first_input.repository_id = Some("main".to_string()); + let first = backend.create(first_input).unwrap(); + let mut second_input = NewTicket::new("Second ready Ticket"); + second_input.workflow_state = Some(TicketWorkflowState::Ready); + second_input.repository_id = Some("main".to_string()); + let second = backend.create(second_input).unwrap(); + for (ticket, target) in [ + (first.id.clone(), second.id.clone()), + (second.id.clone(), first.id.clone()), + ] { + backend + .add_ticket_relation( + TicketIdOrSlug::Id(ticket), + NewTicketRelation { + kind: TicketRelationKind::DependsOn, + target, + note: None, + author: None, + }, + ) + .unwrap(); + } + + let error = backend + .queue_ready(TicketIdOrSlug::Id(first.id.clone()), "orchestrator") + .unwrap_err(); + assert!(matches!(error, TicketError::Conflict(_))); + assert!(error.to_string().contains(" -> ")); + for ticket_id in [first.id, second.id] { + assert_eq!( + backend + .show(TicketIdOrSlug::Id(ticket_id)) + .unwrap() + .meta + .workflow_state, + TicketWorkflowState::Ready + ); + } + } + + #[test] + fn sqlite_queue_target_failure_rolls_back_entire_ready_closure() { + let temp = TempDir::new().unwrap(); + let backend = SqliteTicketBackend::open(temp.path().join("tickets.db"), "workspace-test") + .unwrap() + .with_target_authority(Arc::new(TestTargetAuthority)); + let mut dependency_input = NewTicket::new("Invalid ready dependency"); + dependency_input.workflow_state = Some(TicketWorkflowState::Ready); + dependency_input.repository_id = Some("unknown".to_string()); + let dependency = backend.create(dependency_input).unwrap(); + let mut root_input = NewTicket::new("Queue root"); + root_input.workflow_state = Some(TicketWorkflowState::Ready); + root_input.repository_id = Some("main".to_string()); + let root = backend.create(root_input).unwrap(); + backend + .add_ticket_relation( + TicketIdOrSlug::Id(root.id.clone()), + NewTicketRelation { + kind: TicketRelationKind::DependsOn, + target: dependency.id.clone(), + note: None, + author: None, + }, + ) + .unwrap(); + + let error = backend + .queue_ready(TicketIdOrSlug::Id(root.id.clone()), "orchestrator") + .unwrap_err(); + assert!(matches!(error, TicketError::UnknownTargetRepository(_))); + for ticket_id in [dependency.id, root.id] { + assert_eq!( + backend + .show(TicketIdOrSlug::Id(ticket_id)) + .unwrap() + .meta + .workflow_state, + TicketWorkflowState::Ready + ); + } + } + + #[test] + fn sqlite_queue_atomically_queues_ready_dependency_closure() { let tmp = TempDir::new().unwrap(); let backend = SqliteTicketBackend::open(tmp.path().join("workspace.db"), "workspace-test") .unwrap() .with_target_authority(Arc::new(TestTargetAuthority)); let mut dependency = NewTicket::new("Dependency"); dependency.repository_id = Some("main".to_owned()); + dependency.workflow_state = Some(TicketWorkflowState::Ready); let dependency = backend.create(dependency).unwrap(); let mut implementation = NewTicket::new("Implementation"); implementation.repository_id = Some("main".to_owned()); @@ -7257,14 +7669,28 @@ state: planning .count(), 1 ); - backend + let outcome = backend .queue_ready( TicketIdOrSlug::Id(implementation.id.clone()), "orchestrator", ) .unwrap(); - let queued = backend.show(TicketIdOrSlug::Id(implementation.id)).unwrap(); + assert_eq!(outcome.requested_ticket, implementation.id); + assert_eq!( + outcome.queued_tickets, + vec![dependency.id.clone(), implementation.id.clone()] + ); + let queued = backend + .show(TicketIdOrSlug::Id(implementation.id.clone())) + .unwrap(); + let queued_dependency = backend + .show(TicketIdOrSlug::Id(dependency.id.clone())) + .unwrap(); assert_eq!(queued.meta.workflow_state, TicketWorkflowState::Queued); + assert_eq!( + queued_dependency.meta.workflow_state, + TicketWorkflowState::Queued + ); assert_eq!(queued.meta.queued_by.as_deref(), Some("orchestrator")); assert_eq!(queued.relations.blockers.len(), 1); assert_eq!(queued.relations.blockers[0].blocking_ticket, dependency.id); @@ -8321,7 +8747,7 @@ state: planning } #[test] - fn queue_accepts_unresolved_dependency_and_incoming_blocker_as_context() { + fn queue_rejects_planning_dependency_and_incoming_blocker_without_mutation() { let tmp = TempDir::new().unwrap(); let backend = backend(&tmp); let mut blocked_input = NewTicket::new("Blocked Ready"); @@ -8339,15 +8765,19 @@ state: planning }, ) .unwrap(); - backend + let error = backend .queue_ready(TicketIdOrSlug::Id(blocked.id.clone()), "test") - .unwrap(); - let queued = backend + .unwrap_err(); + assert!(matches!(error, TicketError::BlockingRelations(_))); + let unchanged = backend .show(TicketIdOrSlug::Id(blocked.id.clone())) .unwrap(); - assert_eq!(queued.meta.workflow_state, TicketWorkflowState::Queued); - assert_eq!(queued.relations.blockers.len(), 1); - assert_eq!(queued.relations.blockers[0].blocking_ticket, dependency.id); + assert_eq!(unchanged.meta.workflow_state, TicketWorkflowState::Ready); + assert_eq!(unchanged.relations.blockers.len(), 1); + assert_eq!( + unchanged.relations.blockers[0].blocking_ticket, + dependency.id + ); let mut incoming_input = NewTicket::new("Incoming Blocked Ready"); incoming_input.workflow_state = Some(TicketWorkflowState::Ready); @@ -8364,19 +8794,20 @@ state: planning }, ) .unwrap(); - backend + let error = backend .queue_ready(TicketIdOrSlug::Id(incoming.id.clone()), "test") - .unwrap(); - let queued_incoming = backend + .unwrap_err(); + assert!(matches!(error, TicketError::BlockingRelations(_))); + let unchanged_incoming = backend .show(TicketIdOrSlug::Id(incoming.id.clone())) .unwrap(); assert_eq!( - queued_incoming.meta.workflow_state, - TicketWorkflowState::Queued + unchanged_incoming.meta.workflow_state, + TicketWorkflowState::Ready ); - assert_eq!(queued_incoming.relations.blockers.len(), 1); + assert_eq!(unchanged_incoming.relations.blockers.len(), 1); assert_eq!( - queued_incoming.relations.blockers[0].blocking_ticket, + unchanged_incoming.relations.blockers[0].blocking_ticket, blocker.id ); } diff --git a/crates/ticket/src/tool.rs b/crates/ticket/src/tool.rs index d8763657..f6fb2842 100644 --- a/crates/ticket/src/tool.rs +++ b/crates/ticket/src/tool.rs @@ -142,8 +142,8 @@ const INTAKE_READY_DESCRIPTION: &str = "Record a bounded intake summary and mark The backend applies the same target validation and lock as TicketMarkReady and commits the summary, \ state_changed event, effective target, and planning -> ready transition atomically."; const QUEUE_DESCRIPTION: &str = "Queue a ready Ticket for Orchestrator routing through the typed \ -Ticket backend. The backend performs the gated ready -> queued transition, records queued_by/queued_at, \ -and preserves unresolved blocking relations as Orchestrator scheduling context rather than Queue admission gates."; +Ticket backend. The backend rejects transitive planning dependencies and cycles, atomically queues the \ +requested Ticket plus every transitive ready dependency, and leaves queued or in-progress dependencies unchanged."; const WORKFLOW_STATE_DESCRIPTION: &str = "Transition Ticket `state` through the typed \ Ticket backend with a bounded `state_changed` event. Treat `queued -> inprogress` \ as the implementation acceptance step: implementation side effects should happen only after that \ @@ -316,7 +316,11 @@ impl TicketBackend for TicketToolBackend { self.backend.mark_ready(id, request) } - fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> TicketResult<()> { + fn queue_ready( + &self, + id: TicketIdOrSlug, + queued_by: &str, + ) -> TicketResult { self.backend.queue_ready(id, queued_by) } @@ -1219,12 +1223,22 @@ impl Tool for TicketQueueTool { ) -> Result { let params: TicketQueueParams = parse_input("TicketQueue", input_json)?; let queued_by = default_author(); - self.backend + let outcome = self + .backend .queue_ready(TicketIdOrSlug::Query(params.ticket.clone()), &queued_by) .map_err(|error| backend_error("TicketQueue", error))?; Ok(json_output( - format!("Queued ticket {} for Orchestrator", params.ticket), - json!({ "ticket": params.ticket, "state": "queued", "queued_by": queued_by, "ok": true }), + format!( + "Queued {} ticket(s) for Orchestrator", + outcome.queued_tickets.len() + ), + json!({ + "ticket": outcome.requested_ticket, + "queued_tickets": outcome.queued_tickets, + "state": "queued", + "queued_by": queued_by, + "ok": true + }), )) } } diff --git a/crates/tui/src/workspace_panel.rs b/crates/tui/src/workspace_panel.rs index 1400fff4..bcbfea10 100644 --- a/crates/tui/src/workspace_panel.rs +++ b/crates/tui/src/workspace_panel.rs @@ -2203,7 +2203,7 @@ mod tests { } #[test] - fn workspace_panel_queues_ready_ticket_with_unresolved_relation_context() { + fn workspace_panel_blocks_ready_ticket_with_planning_relation() { let temp = TempDir::new().unwrap(); write_ticket_config(temp.path()); let backend = LocalTicketBackend::new(temp.path().join(".yoi/tickets")); @@ -2233,9 +2233,9 @@ mod tests { .unwrap(); assert_eq!(row.kind, PanelRowKind::Ticket); - assert_eq!(row.next_action, Some(NextUserAction::Queue)); - assert_eq!(row.priority, ActionPriority::ReadyForQueue); - assert!(row.disabled_reason.is_none()); + assert_eq!(row.next_action, Some(NextUserAction::Wait)); + assert_eq!(row.priority, ActionPriority::Background); + assert!(row.disabled_reason.is_some()); assert!( row.ticket .as_ref() @@ -2294,7 +2294,7 @@ mod tests { row.key_hint .as_deref() .unwrap() - .contains("dependency relations remain scheduling context") + .contains("dependency relations remain orchestration context") ); assert!(row.key_hint.as_deref().unwrap().contains(&dependency.id)); } diff --git a/crates/worker/src/feature/builtin/ticket.rs b/crates/worker/src/feature/builtin/ticket.rs index 6a70263b..9001d6bb 100644 --- a/crates/worker/src/feature/builtin/ticket.rs +++ b/crates/worker/src/feature/builtin/ticket.rs @@ -873,12 +873,13 @@ impl WorkspaceHttpTicketBackend { })?), ) .map(TicketBackendOperationResult::Ticket), - TicketBackendOperation::QueueReady { id, .. } => Self::request_unit( + TicketBackendOperation::QueueReady { id, .. } => Self::request( client, WorkspaceRequestMethod::Post, format!("{base}/{}/workflow/queue", Self::ticket_path(&id)), None, - ), + ) + .map(TicketBackendOperationResult::QueueOutcome), TicketBackendOperation::Close { id, resolution } => Self::request_unit( client, WorkspaceRequestMethod::Post, @@ -1099,16 +1100,18 @@ impl TicketBackend for WorkspaceHttpTicketBackend { ) } - fn queue_ready(&self, id: TicketIdOrSlug, queued_by: &str) -> TicketResult<()> { - match self.invoke(TicketBackendOperation::QueueReady { - id, - queued_by: queued_by.to_string(), - })? { - TicketBackendOperationResult::Unit => Ok(()), - other => Err(TicketError::Conflict(format!( - "unexpected ticket backend response: {other:?}" - ))), - } + fn queue_ready( + &self, + id: TicketIdOrSlug, + queued_by: &str, + ) -> TicketResult { + expect_ticket_result!( + self.invoke(TicketBackendOperation::QueueReady { + id, + queued_by: queued_by.to_string(), + }), + TicketBackendOperationResult::QueueOutcome + ) } fn close(&self, id: TicketIdOrSlug, resolution: MarkdownText) -> TicketResult<()> { diff --git a/crates/workspace-server/src/authority.rs b/crates/workspace-server/src/authority.rs index 3ea7decc..66b91ddb 100644 --- a/crates/workspace-server/src/authority.rs +++ b/crates/workspace-server/src/authority.rs @@ -3047,10 +3047,7 @@ VALUES ('workspace-test', 'ticket', 4); assert_eq!(tickets.items[0].record_source, "sqlite_yoi_ticket"); assert_eq!(tickets.items[0].id, "00000000001J2"); assert_eq!(tickets.items[0].state, "ready"); - assert_eq!( - tickets.items[0].workspace_action_priority, - "ready_for_queue" - ); + assert_eq!(tickets.items[0].workspace_action_priority, "background"); let ticket_by_key = authority.ticket(&tickets.items[0].resource_key).unwrap(); assert_eq!(ticket_by_key.id, tickets.items[0].id); diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 1dfca8a4..5fdbb4d9 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -1,4 +1,4 @@ -use std::collections::{BTreeMap, HashMap, HashSet}; +use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; use std::path::{Component, Path, PathBuf}; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex, Weak}; @@ -3951,6 +3951,39 @@ fn reject_unguarded_ticket_completion(operation: &TicketBackendOperation) -> Res Ok(()) } +fn queue_assignment_candidates( + backend: &dyn TicketBackend, + requested_ticket_id: &str, +) -> ticket::Result> { + fn visit( + backend: &dyn TicketBackend, + ticket_id: &str, + visited: &mut BTreeSet, + ready: &mut BTreeSet, + ) -> ticket::Result<()> { + if !visited.insert(ticket_id.to_owned()) { + return Ok(()); + } + let ticket = backend.show(TicketIdOrSlug::Id(ticket_id.to_owned()))?; + if ticket.meta.workflow_state == TicketWorkflowState::Ready { + ready.insert(ticket.meta.id.clone()); + } + for blocker in ticket.relations.blockers { + visit(backend, &blocker.blocking_ticket, visited, ready)?; + } + Ok(()) + } + + let mut ready = BTreeSet::new(); + visit( + backend, + requested_ticket_id, + &mut BTreeSet::new(), + &mut ready, + )?; + Ok(ready.into_iter().collect()) +} + async fn execute_ticket_rest_operation( api: &WorkspaceApi, workspace_id: &str, @@ -3988,27 +4021,34 @@ async fn execute_ticket_rest_operation( .to_string(), ) })?; - let assignment = active_orchestrator_assignment(api, workspace_id, &ticket.meta.id)? - .ok_or_else(|| { - Error::TicketAssignmentConflict( - "Queue requires role=orchestrator assignment to workspace-orchestrator" - .to_string(), - ) - })?; + let candidates = + queue_assignment_candidates(&backend, &ticket.meta.id).map_err(Error::from)?; + let mut assignment_ids = BTreeMap::new(); + for ticket_id in candidates { + let assignment = active_orchestrator_assignment(api, workspace_id, &ticket_id)? + .ok_or_else(|| { + Error::TicketAssignmentConflict(format!( + "Queue requires role=orchestrator assignment for Ticket {ticket_id}" + )) + })?; + assignment_ids.insert(ticket_id, assignment.assignment_id); + } + let assignment_json = serde_json::to_string(&assignment_ids).map_err(|error| { + Error::Config(format!("failed to encode Queue assignment fence: {error}")) + })?; let operation_id = new_id("tqueue"); let fingerprint = Sha256::digest(format!( - "ticket-queue:v1\0{workspace_id}\0{}\0{}\0{}", + "ticket-queue:v2\0{workspace_id}\0{}\0{}\0{assignment_json}", ticket.meta.id, ticket.meta.workflow_state.as_str(), - assignment.assignment_id )) .iter() .map(|byte| format!("{byte:02x}")) .collect::(); event_attributes.extend([ ( - "orchestrator_assignment_id".to_string(), - assignment.assignment_id, + "queue_orchestrator_assignments".to_string(), + assignment_json, ), ( "routing_principal".to_string(), @@ -4029,18 +4069,30 @@ async fn execute_ticket_rest_operation( } let result = execute_ticket_backend_operation(&backend, operation).map_err(Error::from)?; - if is_mutation - && let Some(target) = target - && let Ok(ticket) = backend.show(target) - { - notify_ticket_recipients( - api, - workspace_id, - &ticket.meta.id, - &previous_state, - ticket.meta.workflow_state.as_str(), - source, - ); + if is_mutation { + if let TicketBackendOperationResult::QueueOutcome(outcome) = &result { + for ticket_id in &outcome.queued_tickets { + notify_ticket_recipients( + api, + workspace_id, + ticket_id, + TicketWorkflowState::Ready.as_str(), + TicketWorkflowState::Queued.as_str(), + source.clone(), + ); + } + } else if let Some(target) = target + && let Ok(ticket) = backend.show(target) + { + notify_ticket_recipients( + api, + workspace_id, + &ticket.meta.id, + &previous_state, + ticket.meta.workflow_state.as_str(), + source, + ); + } } Ok(result) } @@ -4364,7 +4416,7 @@ async fn scoped_queue_ticket_record( State(api): State, AxumPath((workspace_id, id)): AxumPath<(String, String)>, headers: HeaderMap, -) -> ApiResult { +) -> ApiResult> { let result = execute_ticket_rest_operation( &api, &workspace_id, @@ -4375,7 +4427,10 @@ async fn scoped_queue_ticket_record( }, ) .await?; - ticket_rest_unit(result) + ticket_rest_result(result, |result| match result { + TicketBackendOperationResult::QueueOutcome(outcome) => Some(outcome), + _ => None, + }) } #[derive(Debug, serde::Deserialize)] @@ -16404,7 +16459,7 @@ mod tests { ); assign_test_orchestrator(&api, &ticket.id); - scoped_queue_ticket_record(State(api.clone()), AxumPath(path), HeaderMap::new()) + let _ = scoped_queue_ticket_record(State(api.clone()), AxumPath(path), HeaderMap::new()) .await .unwrap(); let queued = backend.show(ticket.id.into()).unwrap(); @@ -16422,6 +16477,88 @@ mod tests { assert!(event.attributes.contains_key("routing_request_fingerprint")); } + #[tokio::test] + async fn queue_requires_assignments_for_every_ready_dependency() { + let dir = tempfile::tempdir().unwrap(); + init_clean_git_workspace(dir.path()); + let api = test_api(dir.path()).await; + let backend = browser_ticket_backend(&api).unwrap(); + let mut dependency_input = ticket::NewTicket::new("Ready dependency"); + dependency_input.workflow_state = Some(TicketWorkflowState::Ready); + dependency_input.repository_id = Some(TEST_REPOSITORY_ID.to_string()); + dependency_input.ref_selector = Some("develop".to_string()); + let dependency = backend.create(dependency_input).unwrap(); + let mut root_input = ticket::NewTicket::new("Queue root"); + root_input.workflow_state = Some(TicketWorkflowState::Ready); + root_input.repository_id = Some(TEST_REPOSITORY_ID.to_string()); + root_input.ref_selector = Some("develop".to_string()); + let root = backend.create(root_input).unwrap(); + backend + .add_ticket_relation( + root.id.clone().into(), + ticket::NewTicketRelation { + kind: ticket::TicketRelationKind::DependsOn, + target: dependency.id.clone(), + note: Some("must queue first".to_string()), + author: Some("test".to_string()), + }, + ) + .unwrap(); + + assign_test_orchestrator(&api, &root.id); + let error = scoped_queue_ticket_record( + State(api.clone()), + AxumPath((TEST_WORKSPACE_ID.to_string(), root.id.clone())), + HeaderMap::new(), + ) + .await + .unwrap_err() + .into_response(); + assert_eq!(error.status(), StatusCode::CONFLICT); + assert_eq!( + backend + .show(root.id.clone().into()) + .unwrap() + .meta + .workflow_state, + TicketWorkflowState::Ready + ); + assert_eq!( + backend + .show(dependency.id.clone().into()) + .unwrap() + .meta + .workflow_state, + TicketWorkflowState::Ready + ); + + assign_test_orchestrator(&api, &dependency.id); + let Json(outcome) = scoped_queue_ticket_record( + State(api.clone()), + AxumPath((TEST_WORKSPACE_ID.to_string(), root.id.clone())), + HeaderMap::new(), + ) + .await + .unwrap(); + assert_eq!(outcome.requested_ticket, root.id); + assert_eq!( + outcome.queued_tickets, + vec![dependency.id.clone(), root.id.clone()] + ); + for ticket_id in [dependency.id, root.id] { + let queued = backend.show(ticket_id.clone().into()).unwrap(); + assert_eq!(queued.meta.workflow_state, TicketWorkflowState::Queued); + assert_eq!( + queued + .events + .last() + .and_then(|event| event.attributes.get("orchestrator_assignment_id")) + .cloned(), + Some(format!("orchestrator-{ticket_id}")) + ); + } + } + #[tokio::test] async fn queued_ticket_mutation_succeeds_without_orchestrator() { let dir = tempfile::tempdir().unwrap(); @@ -17200,9 +17337,13 @@ mod tests { workspace_id: TEST_WORKSPACE_ID.to_string(), id: ticket_id.clone(), }; + let mut related_input = ticket::NewTicket::new("Related Browser Ticket"); + related_input.workflow_state = Some(TicketWorkflowState::Ready); + related_input.repository_id = Some(TEST_REPOSITORY_ID.to_string()); + related_input.ref_selector = Some("develop".to_string()); let related_ticket_id = browser_ticket_backend(&api) .unwrap() - .create(ticket::NewTicket::new("Related Browser Ticket")) + .create(related_input) .unwrap() .id; browser_ticket_backend(&api) @@ -17274,6 +17415,7 @@ mod tests { .unwrap(); assert_eq!(ready.state, "ready"); assign_test_orchestrator(&api, &ticket_id); + assign_test_orchestrator(&api, &related_ticket_id); let Json(ready_detail) = scoped_get_ticket(State(api.clone()), AxumPath(path())) .await .unwrap(); diff --git a/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte index 97a3e263..63095682 100644 --- a/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/tickets/[ticketId]/+page.svelte @@ -466,9 +466,9 @@ {busy === "queue" ? "Queueing…" : "Queue ticket"} {#if !ticket.action_eligibility.can_queue} -

Queue requires a valid target, an active Orchestrator assignment, and no active Coder assignment.

+

Queue requires a valid target, an active Orchestrator assignment, no active Coder assignment, and no dependency still in planning.

{:else if ticket.relations.blockers.length > 0} -

Queue records orchestration demand. Dependency relations remain visible so the Orchestrator can decide whether to wait or start work in parallel.

+

Ready dependencies are queued atomically. Queued or in-progress dependencies remain unchanged for the Orchestrator to schedule.

{/if} {/if}