返回 CodeWhale
runtime.rs
根目录 / crates / tui / src / work_graph / runtime.rs
1 //! Active-session Work Graph authority and legacy tool adapters.
2
3 use std::path::Path;
4 use std::sync::{Arc, Mutex, MutexGuard};
5
6 use crate::fleet::ledger::FleetLedger;
7 use crate::tools::plan::{PlanSnapshot, PlanState, SharedPlanState, StepStatus};
8 use crate::tools::todo::{SharedTodoList, TodoList, TodoListSnapshot, TodoStatus};
9 use codewhale_lane::LaneRegistry;
10
11 use super::{
12 BindingId, ChangeCtx, CompatPlanMetadata, CompatProjectionState, CompatTodoBinding, EdgeKind,
13 IdempotencyKey, NodeKind, NodeState, OperationBinding, OperationIntent, OperationObservation,
14 OperationOwnerSnapshot, Provenance, ReasoningEffortTier, WorkActivityEvent, WorkEdge,
15 WorkEdgeId, WorkGraph, WorkGraphChange, WorkGraphSnapshot, WorkNode, WorkNodeId, WorkNodePatch,
16 external_identity_is_well_formed, fleet_task_owner_snapshot, import_legacy,
17 lane_owner_snapshot, project_plan, project_todos, validate,
18 };
19
20 pub(crate) const ACTIVE_OPERATION_SUMMARY_START: &str =
21 "<!-- codewhale:active-work-operations:start -->";
22 pub(crate) const ACTIVE_OPERATION_SUMMARY_END: &str =
23 "<!-- codewhale:active-work-operations:end -->";
24
25 #[derive(Debug, Clone, PartialEq)]
26 pub struct WorkRuntimeSnapshot {
27 pub graph: WorkGraphSnapshot,
28 pub todos: TodoListSnapshot,
29 pub plan: PlanSnapshot,
30 }
31
32 #[derive(Debug, Default)]
33 struct ActiveGraph {
34 session_id: Option<String>,
35 snapshot: Option<WorkGraphSnapshot>,
36 pending_publish: bool,
37 }
38
39 /// One active session graph plus the read-only legacy views it publishes.
40 pub struct WorkRuntime {
41 todos: SharedTodoList,
42 plan: SharedPlanState,
43 graph: Mutex<ActiveGraph>,
44 }
45
46 pub type SharedWorkRuntime = Arc<WorkRuntime>;
47
48 impl std::fmt::Debug for WorkRuntime {
49 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
50 let graph = lock_unpoisoned(&self.graph);
51 f.debug_struct("WorkRuntime")
52 .field("session_id", &graph.session_id)
53 .field("has_graph", &graph.snapshot.is_some())
54 .field("pending_publish", &graph.pending_publish)
55 .finish()
56 }
57 }
58
59 #[must_use]
60 pub fn new_shared_work_runtime(todos: SharedTodoList, plan: SharedPlanState) -> SharedWorkRuntime {
61 Arc::new(WorkRuntime {
62 todos,
63 plan,
64 graph: Mutex::new(ActiveGraph::default()),
65 })
66 }
67
68 impl WorkRuntime {
69 #[must_use]
70 pub fn matches_todos(&self, todos: &SharedTodoList) -> bool {
71 Arc::ptr_eq(&self.todos, todos)
72 }
73
74 #[must_use]
75 pub fn matches_plan(&self, plan: &SharedPlanState) -> bool {
76 Arc::ptr_eq(&self.plan, plan)
77 }
78
79 #[must_use]
80 pub fn has_operation_binding(&self, session_id: Option<&str>, external: &str) -> bool {
81 let active = lock_unpoisoned(&self.graph);
82 if let (Some(expected), Some(actual)) = (session_id, active.session_id.as_deref())
83 && expected != actual
84 {
85 return false;
86 }
87 active.snapshot.as_ref().is_some_and(|graph| {
88 graph.nodes.iter().any(|node| {
89 node.binding
90 .as_ref()
91 .is_some_and(|binding| binding.external == external)
92 })
93 })
94 }
95
96 /// Durable owner bindings that still require restore-time confirmation.
97 /// Terminal operations are excluded because their saved owner receipt is
98 /// already sufficient and must not be reopened by a missing live handle.
99 #[must_use]
100 pub fn reconcilable_durable_bindings(&self, session_id: Option<&str>) -> Vec<String> {
101 let active = lock_unpoisoned(&self.graph);
102 if let (Some(expected), Some(actual)) = (session_id, active.session_id.as_deref())
103 && expected != actual
104 {
105 return Vec::new();
106 }
107 active
108 .snapshot
109 .as_ref()
110 .into_iter()
111 .flat_map(|graph| graph.nodes.iter())
112 .filter(|node| node.state.is_live() || node.state == NodeState::Stale)
113 .filter_map(|node| node.binding.as_ref())
114 .filter(|binding| binding.durable)
115 .map(|binding| binding.external.clone())
116 .collect()
117 }
118
119 /// Register an Operation before its owner starts work. The operation is
120 /// first added inert, connected to an Objective/PlanStep, and only then
121 /// advanced to `Initializing`, so every reducer intermediate satisfies
122 /// the no-orphan invariant.
123 pub fn register_operation(
124 &self,
125 session_id: &str,
126 intent: OperationIntent,
127 ) -> Result<WorkNodeId, String> {
128 if !external_identity_is_well_formed(&intent.external) {
129 return Err(format!(
130 "invalid lifecycle binding external {:?}",
131 intent.external
132 ));
133 }
134 let todos = retry_lock(&self.todos, 100)
135 .ok_or_else(|| "To-do state is busy; operation was not registered".to_string())?;
136 let plan = retry_lock(&self.plan, 100)
137 .ok_or_else(|| "Plan state is busy; operation was not registered".to_string())?;
138 let mut active = lock_unpoisoned(&self.graph);
139 let base = graph_for_update(&mut active, session_id, &plan.snapshot(), &todos.snapshot())?;
140 if let Some(existing) = base.nodes.iter().find(|node| {
141 node.binding
142 .as_ref()
143 .is_some_and(|binding| binding.external == intent.external)
144 }) {
145 let binding = existing.binding.as_ref().expect("binding matched above");
146 if binding.durable != intent.durable {
147 return Err(format!(
148 "lifecycle binding {} changed durability",
149 intent.external
150 ));
151 }
152 return Ok(existing.id.clone());
153 }
154
155 let mut graph = WorkGraph::from_snapshot(base);
156 let parent = operation_parent(&mut graph, session_id, &intent.source)?;
157 let node_id = WorkNodeId::derive(session_id, &format!("operation:{}", intent.external));
158 let now = now_ms();
159 apply_change(
160 &mut graph,
161 session_id,
162 &intent.source,
163 WorkGraphChange::AddNode {
164 node: WorkNode {
165 id: node_id.clone(),
166 kind: NodeKind::Operation,
167 title: bounded_operation_title(&intent.title),
168 state: NodeState::Ready,
169 acceptance: intent.acceptance,
170 binding: Some(OperationBinding {
171 external: intent.external,
172 durable: intent.durable,
173 last_observation: None,
174 }),
175 evidence: None,
176 provenance: Provenance::ToolUpdate {
177 tool: intent.source.clone(),
178 call_id: intent.call_id,
179 },
180 created_at: now,
181 updated_at: now,
182 },
183 },
184 )?;
185 ensure_contains(&mut graph, session_id, &intent.source, &parent, &node_id)?;
186 let title = graph
187 .snapshot()
188 .node(&node_id)
189 .map(|node| node.title.clone())
190 .ok_or_else(|| format!("operation {node_id} disappeared during registration"))?;
191 patch_existing_node(
192 &mut graph,
193 session_id,
194 &intent.source,
195 &node_id,
196 title,
197 NodeState::Initializing,
198 )?;
199 // Registration is where finished shell calls accumulate (#6842):
200 // keep only the newest ended, non-durable operations.
201 if super::reducer::evictable_operations(graph.snapshot()).len()
202 > super::model::ENDED_OPERATION_CAP
203 {
204 apply_change(
205 &mut graph,
206 session_id,
207 &intent.source,
208 WorkGraphChange::PruneEndedOperations {
209 keep: super::model::ENDED_OPERATION_CAP,
210 },
211 )?;
212 }
213 let next = graph.into_snapshot();
214 validate_combined(&next, &project_plan(&next), &project_todos(&next))?;
215 active.snapshot = Some(next);
216 active.pending_publish = true;
217 Ok(node_id)
218 }
219
220 /// Retain a crossed approval boundary as provenance attached to an
221 /// operation. The completed Approval node is historical evidence only: it
222 /// does not grant capabilities or change the owner's runtime authority.
223 pub fn record_operation_approval(
224 &self,
225 session_id: &str,
226 external: &str,
227 reference: &str,
228 source: &str,
229 call_id: &str,
230 ) -> Result<(), String> {
231 let todos = retry_lock(&self.todos, 100)
232 .ok_or_else(|| "To-do state is busy; approval was not recorded".to_string())?;
233 let plan = retry_lock(&self.plan, 100)
234 .ok_or_else(|| "Plan state is busy; approval was not recorded".to_string())?;
235 let mut active = lock_unpoisoned(&self.graph);
236 let base = graph_for_update(&mut active, session_id, &plan.snapshot(), &todos.snapshot())?;
237 let operation = base
238 .nodes
239 .iter()
240 .find(|node| {
241 node.kind == NodeKind::Operation
242 && node
243 .binding
244 .as_ref()
245 .is_some_and(|binding| binding.external == external)
246 })
247 .map(|node| node.id.clone())
248 .ok_or_else(|| format!("operation binding {external} is not registered"))?;
249 let approval = WorkNodeId::derive(
250 session_id,
251 &format!("approval:operation:{external}:{reference}"),
252 );
253 let mut graph = WorkGraph::from_snapshot(base);
254 if graph.snapshot().node(&approval).is_none() {
255 let now = now_ms();
256 apply_change(
257 &mut graph,
258 session_id,
259 source,
260 WorkGraphChange::AddNode {
261 node: WorkNode {
262 id: approval.clone(),
263 kind: NodeKind::Approval,
264 title: bounded_operation_title(&format!(
265 "verification approved: {reference}"
266 )),
267 state: NodeState::Completed,
268 acceptance: Vec::new(),
269 binding: None,
270 evidence: None,
271 provenance: Provenance::ToolUpdate {
272 tool: source.to_string(),
273 call_id: call_id.to_string(),
274 },
275 created_at: now,
276 updated_at: now,
277 },
278 },
279 )?;
280 }
281 let edge = WorkEdgeId::derive(
282 session_id,
283 &format!(
284 "requires-approval:{}:{}",
285 operation.as_str(),
286 approval.as_str()
287 ),
288 );
289 if graph.snapshot().edge(&edge).is_none() {
290 apply_change(
291 &mut graph,
292 session_id,
293 source,
294 WorkGraphChange::AddEdge {
295 edge: WorkEdge {
296 id: edge,
297 kind: EdgeKind::RequiresApproval,
298 from: operation,
299 to: approval,
300 },
301 },
302 )?;
303 }
304 let next = graph.into_snapshot();
305 validate_combined(&next, &project_plan(&next), &project_todos(&next))?;
306 active.snapshot = Some(next);
307 active.pending_publish = true;
308 Ok(())
309 }
310
311 /// Apply one owner observation through the reducer. Unknown bindings are
312 /// rejected rather than materialized after the fact: spawn intent must be
313 /// registered before work begins.
314 pub fn reconcile_operation(
315 &self,
316 session_id: &str,
317 snapshot: OperationOwnerSnapshot,
318 ) -> Result<bool, String> {
319 let external = snapshot.external.clone();
320 self.reconcile_observation(session_id, &external, snapshot.into_observation())
321 }
322
323 pub fn reconcile_observation(
324 &self,
325 session_id: &str,
326 external: &str,
327 observation: OperationObservation,
328 ) -> Result<bool, String> {
329 let mut active = lock_unpoisoned(&self.graph);
330 if active
331 .session_id
332 .as_deref()
333 .is_some_and(|id| id != session_id)
334 {
335 return Err(format!(
336 "lifecycle observation for session {session_id} does not match active session"
337 ));
338 }
339 let base = active
340 .snapshot
341 .clone()
342 .ok_or_else(|| "lifecycle observation arrived before graph registration".to_string())?;
343 let Some(next) = apply_observation_to_snapshot(&base, session_id, external, observation)?
344 else {
345 return Ok(false);
346 };
347 active.snapshot = Some(next);
348 active.pending_publish = true;
349 Ok(true)
350 }
351
352 /// Compact graph-owned continuity injected after context compaction.
353 /// It contains identities and state only; no raw output or reasoning.
354 #[must_use]
355 pub fn active_operation_summary(&self, session_id: Option<&str>) -> Option<String> {
356 let active = lock_unpoisoned(&self.graph);
357 if let (Some(expected), Some(actual)) = (session_id, active.session_id.as_deref())
358 && expected != actual
359 {
360 return None;
361 }
362 let graph = active.snapshot.as_ref()?;
363 let operations = graph
364 .nodes
365 .iter()
366 .filter(|node| {
367 node.kind == NodeKind::Operation
368 && (node.state.is_live() || node.state == NodeState::Stale)
369 })
370 .take(24)
371 .collect::<Vec<_>>();
372 if operations.is_empty() {
373 return None;
374 }
375 let mut out = format!(
376 "{ACTIVE_OPERATION_SUMMARY_START}\n## Active Work Graph Operations\n\nOwner records remain authoritative; reconcile before acting on a restored operation.\n"
377 );
378 for node in operations {
379 let external = node
380 .binding
381 .as_ref()
382 .map_or("unbound", |binding| binding.external.as_str());
383 out.push_str(&format!(
384 "- `{}` - {} - {}\n",
385 external,
386 operation_state_label(node.state),
387 prompt_safe_title(&node.title)
388 ));
389 }
390 out.push_str(ACTIVE_OPERATION_SUMMARY_END);
391 Some(out)
392 }
393
394 /// Record one reasoning-effort configuration change as bounded graph
395 /// activity. The event contains typed tiers and a non-secret route
396 /// identity only; there is no field capable of carrying reasoning text.
397 /// When live operations exist, the most recently updated one receives the
398 /// historical link. Later terminalization does not invalidate that link.
399 #[allow(clippy::too_many_arguments)]
400 pub fn record_reasoning_effort_change(
401 &self,
402 session_id: Option<&str>,
403 requested: ReasoningEffortTier,
404 effective: ReasoningEffortTier,
405 provider_kind: crate::config::ProviderKind,
406 provider: &str,
407 endpoint_identity: Option<&str>,
408 model: Option<&str>,
409 ) -> Result<Option<WorkNodeId>, String> {
410 let todos = retry_lock(&self.todos, 100)
411 .ok_or_else(|| "To-do state is busy; effort change was not recorded".to_string())?;
412 let plan = retry_lock(&self.plan, 100)
413 .ok_or_else(|| "Plan state is busy; effort change was not recorded".to_string())?;
414 let mut active = lock_unpoisoned(&self.graph);
415 let session_id = resolved_session_id(&active, session_id);
416 let base = graph_for_update(
417 &mut active,
418 &session_id,
419 &plan.snapshot(),
420 &todos.snapshot(),
421 )?;
422 let operation = base
423 .nodes
424 .iter()
425 .filter(|node| node.kind == NodeKind::Operation && node.state.is_live())
426 .max_by(|left, right| {
427 left.updated_at
428 .cmp(&right.updated_at)
429 .then_with(|| left.id.as_str().cmp(right.id.as_str()))
430 })
431 .map(|node| node.id.clone());
432 let now = now_ms();
433 let mut graph = WorkGraph::from_snapshot(base);
434 graph
435 .apply(
436 WorkGraphChange::RecordActivity {
437 event: WorkActivityEvent::ReasoningEffortChanged {
438 requested,
439 effective,
440 provider_kind: Some(provider_kind),
441 provider: provider.to_string(),
442 endpoint_identity: endpoint_identity.map(str::to_string),
443 model: model.map(str::to_string),
444 ts: now,
445 operation: operation.clone(),
446 },
447 },
448 ChangeCtx {
449 session_id,
450 now,
451 idempotency_key: None,
452 },
453 )
454 .map_err(|err| format!("reasoning effort: {err}"))?;
455 let next = graph.into_snapshot();
456 validate_combined(&next, &project_plan(&next), &project_todos(&next))?;
457 active.snapshot = Some(next);
458 active.pending_publish = true;
459 Ok(operation)
460 }
461
462 /// Apply an `update_plan` payload through the graph and publish both
463 /// legacy projections only after the candidate graph validates.
464 pub async fn apply_plan_update(
465 &self,
466 session_id: &str,
467 tool: &str,
468 plan: &PlanSnapshot,
469 ) -> Result<PlanSnapshot, String> {
470 let todos_guard = self.todos.lock().await;
471 let plan_guard = self.plan.lock().await;
472 let mut active = lock_unpoisoned(&self.graph);
473 let base = graph_for_update(
474 &mut active,
475 session_id,
476 &plan_guard.snapshot(),
477 &todos_guard.snapshot(),
478 )?;
479 let next = update_plan_graph(base, session_id, tool, plan)?;
480 let derived_plan = project_plan(&next);
481 let derived_todos = project_todos(&next);
482 let next_plan = PlanState::from_snapshot(&derived_plan);
483 let next_todos = TodoList::from_snapshot(&derived_todos)?;
484 validate_combined(&next, &next_plan.snapshot(), &next_todos.snapshot())?;
485 active.snapshot = Some(next);
486 active.pending_publish = true;
487 Ok(derived_plan)
488 }
489
490 /// Apply a legacy To-do/checklist payload through the graph and publish
491 /// both projections from the committed candidate.
492 pub async fn apply_todo_update(
493 &self,
494 session_id: &str,
495 tool: &str,
496 todos: &TodoListSnapshot,
497 ) -> Result<TodoListSnapshot, String> {
498 let todos_guard = self.todos.lock().await;
499 let plan_guard = self.plan.lock().await;
500 let mut active = lock_unpoisoned(&self.graph);
501 let base = graph_for_update(
502 &mut active,
503 session_id,
504 &plan_guard.snapshot(),
505 &todos_guard.snapshot(),
506 )?;
507 let next = update_todo_graph(base, session_id, tool, todos)?;
508 let derived_plan = project_plan(&next);
509 let derived_todos = project_todos(&next);
510 let next_plan = PlanState::from_snapshot(&derived_plan);
511 let next_todos = TodoList::from_snapshot(&derived_todos)?;
512 validate_combined(&next, &next_plan.snapshot(), &next_todos.snapshot())?;
513 active.snapshot = Some(next);
514 active.pending_publish = true;
515 Ok(derived_todos)
516 }
517
518 /// Publish the latest validated legacy views after the caller has queued
519 /// their graph-backed session/checkpoint write.
520 pub async fn publish_pending(&self) -> Result<bool, String> {
521 let mut todos = self.todos.lock().await;
522 let mut plan = self.plan.lock().await;
523 let mut active = lock_unpoisoned(&self.graph);
524 if !active.pending_publish {
525 return Ok(false);
526 }
527 let graph = active
528 .snapshot
529 .as_ref()
530 .ok_or_else(|| "pending Work projection has no graph".to_string())?;
531 let derived_plan = project_plan(graph);
532 let derived_todos = project_todos(graph);
533 let next_plan = PlanState::from_snapshot(&derived_plan);
534 let next_todos = TodoList::from_snapshot(&derived_todos)?;
535 validate_combined(graph, &next_plan.snapshot(), &next_todos.snapshot())?;
536 *plan = next_plan;
537 *todos = next_todos;
538 active.pending_publish = false;
539 Ok(true)
540 }
541
542 /// Synchronous counterpart for explicit save/rename/fork commands that
543 /// have already completed their atomic disk write.
544 pub fn publish_pending_sync(&self) -> Result<bool, String> {
545 let mut todos = retry_lock(&self.todos, 100).ok_or_else(|| {
546 "To-do state is busy; saved Work views were not published".to_string()
547 })?;
548 let mut plan = retry_lock(&self.plan, 100)
549 .ok_or_else(|| "Plan state is busy; saved Work views were not published".to_string())?;
550 let mut active = lock_unpoisoned(&self.graph);
551 if !active.pending_publish {
552 return Ok(false);
553 }
554 let graph = active
555 .snapshot
556 .as_ref()
557 .ok_or_else(|| "pending Work projection has no graph".to_string())?;
558 let derived_plan = project_plan(graph);
559 let derived_todos = project_todos(graph);
560 let next_plan = PlanState::from_snapshot(&derived_plan);
561 let next_todos = TodoList::from_snapshot(&derived_todos)?;
562 validate_combined(graph, &next_plan.snapshot(), &next_todos.snapshot())?;
563 *plan = next_plan;
564 *todos = next_todos;
565 active.pending_publish = false;
566 Ok(true)
567 }
568
569 #[must_use]
570 pub fn has_pending_publish(&self) -> bool {
571 lock_unpoisoned(&self.graph).pending_publish
572 }
573
574 /// Latest graph-derived To-do view, including an unpublished transaction.
575 pub async fn current_todos(&self) -> Result<TodoListSnapshot, String> {
576 let projected = {
577 let active = lock_unpoisoned(&self.graph);
578 active.snapshot.as_ref().map(project_todos)
579 };
580 if let Some(projected) = projected {
581 return Ok(projected);
582 }
583 Ok(self.todos.lock().await.snapshot())
584 }
585
586 /// Capture a persistence-ready graph plus fully populated old views.
587 /// Legacy-only in-memory state is imported once and normalized in place.
588 pub fn capture(&self, session_id: Option<&str>) -> Result<Option<WorkRuntimeSnapshot>, String> {
589 self.capture_with_retries(session_id, 100)
590 }
591
592 /// Non-blocking capture for the render/event loop.
593 pub fn try_capture(
594 &self,
595 session_id: Option<&str>,
596 ) -> Result<Option<WorkRuntimeSnapshot>, String> {
597 self.capture_with_retries(session_id, 1)
598 }
599
600 fn capture_with_retries(
601 &self,
602 session_id: Option<&str>,
603 retries: u32,
604 ) -> Result<Option<WorkRuntimeSnapshot>, String> {
605 let todos = retry_lock(&self.todos, retries)
606 .ok_or_else(|| "To-do state is busy; try saving again".to_string())?;
607 let plan = retry_lock(&self.plan, retries)
608 .ok_or_else(|| "Plan state is busy; try saving again".to_string())?;
609 let mut active = lock_unpoisoned(&self.graph);
610 let todos_snapshot = todos.snapshot();
611 let plan_snapshot = plan.snapshot();
612 if todos_snapshot.is_empty()
613 && plan_snapshot.is_empty()
614 && active
615 .snapshot
616 .as_ref()
617 .is_none_or(WorkGraphSnapshot::is_empty)
618 {
619 return Ok(None);
620 }
621 let had_graph = active.snapshot.is_some();
622 let had_pending_publish = active.pending_publish;
623 let session_id = resolved_session_id(&active, session_id);
624 let graph = graph_for_update(&mut active, &session_id, &plan_snapshot, &todos_snapshot)?;
625 let derived_plan = project_plan(&graph);
626 let derived_todos = project_todos(&graph);
627 validate_combined(&graph, &derived_plan, &derived_todos)?;
628 if had_graph
629 && !had_pending_publish
630 && (derived_plan != plan_snapshot || derived_todos != todos_snapshot)
631 {
632 return Err("live Work Graph and legacy views disagree".to_string());
633 }
634 active.snapshot = Some(graph.clone());
635 if !had_graph {
636 active.pending_publish = true;
637 }
638 Ok(Some(WorkRuntimeSnapshot {
639 graph,
640 todos: derived_todos,
641 plan: derived_plan,
642 }))
643 }
644
645 /// Validate and atomically activate persisted state. Sessions without a
646 /// graph are deterministically imported from their complete old views.
647 pub fn restore(
648 &self,
649 session_id: &str,
650 graph: Option<&WorkGraphSnapshot>,
651 todos: &TodoListSnapshot,
652 plan: &PlanSnapshot,
653 ) -> Result<Option<WorkRuntimeSnapshot>, String> {
654 self.restore_internal(session_id, graph, todos, plan, None)
655 }
656
657 /// Restore a saved graph and reconcile workspace-scoped durable owners as
658 /// one candidate transaction. No live state changes until the restored
659 /// graph, owner observations, and both legacy projections all validate.
660 pub fn restore_with_workspace_owner_bindings(
661 &self,
662 session_id: &str,
663 workspace: &Path,
664 graph: Option<&WorkGraphSnapshot>,
665 todos: &TodoListSnapshot,
666 plan: &PlanSnapshot,
667 ) -> Result<Option<WorkRuntimeSnapshot>, String> {
668 self.restore_internal(session_id, graph, todos, plan, Some(workspace))
669 }
670
671 fn restore_internal(
672 &self,
673 session_id: &str,
674 graph: Option<&WorkGraphSnapshot>,
675 todos: &TodoListSnapshot,
676 plan: &PlanSnapshot,
677 workspace: Option<&Path>,
678 ) -> Result<Option<WorkRuntimeSnapshot>, String> {
679 let had_graph = graph.is_some();
680 let graph = match graph {
681 Some(graph) => {
682 validate(graph).map_err(|err| err.to_string())?;
683 graph.clone()
684 }
685 None if todos.is_empty() && plan.is_empty() => WorkGraphSnapshot::new(),
686 None => import_legacy(session_id, plan, todos)?,
687 };
688 let (mut graph, reconciled_ephemeral) =
689 mark_restored_ephemeral_operations_stale(graph, session_id)?;
690 let reconciled_workspace = if let Some(workspace) = workspace {
691 let (reconciled, changed) = reconcile_workspace_snapshot(graph, session_id, workspace)?;
692 graph = reconciled;
693 changed > 0
694 } else {
695 false
696 };
697 let derived_plan = project_plan(&graph);
698 let derived_todos = project_todos(&graph);
699 if graph.is_empty() {
700 if !todos.is_empty() || !plan.is_empty() {
701 return Err("empty Work Graph cannot carry non-empty legacy views".to_string());
702 }
703 } else if graph.import_digest.is_some() && graph.compat.is_empty() {
704 return Err("imported Work Graph is missing compatibility projections".to_string());
705 }
706 validate_combined(&graph, &derived_plan, &derived_todos)?;
707 if had_graph && (&derived_plan != plan || &derived_todos != todos) {
708 return Err("persisted Work Graph and legacy views disagree".to_string());
709 }
710 let next_plan = PlanState::from_snapshot(&derived_plan);
711 let next_todos = TodoList::from_snapshot(&derived_todos)?;
712 let mut todos_guard = retry_lock(&self.todos, 100)
713 .ok_or_else(|| "To-do state is busy; session was not restored".to_string())?;
714 let mut plan_guard = retry_lock(&self.plan, 100)
715 .ok_or_else(|| "Plan state is busy; session was not restored".to_string())?;
716 let mut active = lock_unpoisoned(&self.graph);
717 *todos_guard = next_todos;
718 *plan_guard = next_plan;
719 active.session_id = Some(session_id.to_string());
720 active.snapshot = Some(graph.clone());
721 // A legacy load has already restored its complete old views, but its
722 // newly imported graph still needs one acknowledged graph-bearing
723 // write (and pre-import archive) before the migration is settled.
724 active.pending_publish =
725 reconciled_ephemeral || reconciled_workspace || (!had_graph && !graph.is_empty());
726 if graph.is_empty() {
727 Ok(None)
728 } else {
729 Ok(Some(WorkRuntimeSnapshot {
730 graph,
731 todos: derived_todos,
732 plan: derived_plan,
733 }))
734 }
735 }
736
737 pub fn clear(&self, session_id: Option<&str>) -> bool {
738 let Some(mut todos) = retry_lock(&self.todos, 100) else {
739 return false;
740 };
741 let Some(mut plan) = retry_lock(&self.plan, 100) else {
742 return false;
743 };
744 let mut active = lock_unpoisoned(&self.graph);
745 todos.clear();
746 *plan = PlanState::default();
747 active.session_id = Some(resolved_session_id(&active, session_id));
748 active.snapshot = Some(WorkGraphSnapshot::new());
749 active.pending_publish = false;
750 true
751 }
752 }
753
754 fn apply_observation_to_snapshot(
755 base: &WorkGraphSnapshot,
756 session_id: &str,
757 external: &str,
758 observation: OperationObservation,
759 ) -> Result<Option<WorkGraphSnapshot>, String> {
760 let node = base
761 .nodes
762 .iter()
763 .find(|node| {
764 node.binding
765 .as_ref()
766 .is_some_and(|binding| binding.external == external)
767 })
768 .ok_or_else(|| format!("operation binding {external} is not registered"))?;
769 let stale_owner_recovery = if let OperationObservation::OwnerReported {
770 state, seq, output, ..
771 } = &observation
772 && let Some(previous) = node
773 .binding
774 .as_ref()
775 .and_then(|binding| binding.last_observation.as_ref())
776 {
777 if *seq < previous.seq {
778 return Err(format!(
779 "operation owner {external} sequence regressed from {} to {seq}",
780 previous.seq
781 ));
782 }
783 if *seq == previous.seq {
784 if node.state != NodeState::Stale {
785 return Ok(None);
786 }
787 if *state != previous.owner_state || *output != previous.output {
788 return Err(format!(
789 "operation owner {external} changed observation at sequence {seq}"
790 ));
791 }
792 true
793 } else {
794 false
795 }
796 } else {
797 false
798 };
799 if matches!(&observation, OperationObservation::OwnerMissing { .. })
800 && node.state == NodeState::Stale
801 {
802 return Ok(None);
803 }
804 let idempotency_key = match &observation {
805 OperationObservation::OwnerReported { .. } if stale_owner_recovery => None,
806 OperationObservation::OwnerReported { seq, .. } => Some(IdempotencyKey {
807 binding: BindingId::derive(session_id, &format!("binding:{external}")),
808 seq: *seq,
809 }),
810 OperationObservation::OwnerMissing { .. } | OperationObservation::CancelUpdate { .. } => {
811 None
812 }
813 };
814 let (next, receipt) = super::reducer::apply(
815 base,
816 WorkGraphChange::ReconcileOperation {
817 node: node.id.clone(),
818 obs: observation,
819 },
820 ChangeCtx {
821 session_id: session_id.to_string(),
822 now: now_ms(),
823 idempotency_key,
824 },
825 )
826 .map_err(|err| format!("runtime reconcile: {err}"))?;
827 Ok((!receipt.no_op).then_some(next))
828 }
829
830 fn reconcile_workspace_snapshot(
831 mut graph: WorkGraphSnapshot,
832 session_id: &str,
833 workspace: &Path,
834 ) -> Result<(WorkGraphSnapshot, usize), String> {
835 let candidates = graph
836 .nodes
837 .iter()
838 .filter(|node| node.state.is_live() || node.state == NodeState::Stale)
839 .filter_map(|node| node.binding.as_ref())
840 .filter(|binding| binding.durable)
841 .map(|binding| binding.external.clone())
842 .filter(|external| external.starts_with("fleet:") || external.starts_with("lane:"))
843 .collect::<Vec<_>>();
844 if candidates.is_empty() {
845 return Ok((graph, 0));
846 }
847
848 let fleet = candidates
849 .iter()
850 .any(|external| external.starts_with("fleet:"))
851 .then(|| FleetLedger::open(workspace).and_then(|ledger| ledger.rebuild_state()));
852 if let Some(Err(err)) = fleet.as_ref() {
853 tracing::warn!(
854 workspace = %workspace.display(),
855 error = %err,
856 "Fleet owner store could not be replayed; restored bindings will be stale"
857 );
858 }
859 let lanes = candidates
860 .iter()
861 .any(|external| external.starts_with("lane:"))
862 .then(LaneRegistry::open_default);
863 if let Some(Err(err)) = lanes.as_ref() {
864 tracing::warn!(
865 error = %err,
866 "Lane owner registry could not be opened; restored bindings will be stale"
867 );
868 }
869 let observed_at = now_ms();
870 let mut changed = 0usize;
871 for external in candidates {
872 let observation = if let Some(rest) = external.strip_prefix("fleet:") {
873 rest.split_once('/')
874 .and_then(|(run_id, task_id)| {
875 fleet
876 .as_ref()
877 .and_then(|state| state.as_ref().ok())
878 .and_then(|state| state.tasks.get(&format!("{run_id}:{task_id}")))
879 })
880 .map(|record| fleet_task_owner_snapshot(record, observed_at).into_observation())
881 .unwrap_or(OperationObservation::OwnerMissing {
882 checked_at: observed_at,
883 })
884 } else if let Some(lane_id) = external.strip_prefix("lane:") {
885 lanes
886 .as_ref()
887 .and_then(|registry| registry.as_ref().ok())
888 .and_then(|registry| registry.load(lane_id).ok())
889 .map(|record| lane_owner_snapshot(&record, observed_at).into_observation())
890 .unwrap_or(OperationObservation::OwnerMissing {
891 checked_at: observed_at,
892 })
893 } else {
894 continue;
895 };
896 if let Some(next) =
897 apply_observation_to_snapshot(&graph, session_id, &external, observation)?
898 {
899 graph = next;
900 changed = changed.saturating_add(1);
901 }
902 }
903 Ok((graph, changed))
904 }
905
906 fn operation_parent(
907 graph: &mut WorkGraph,
908 session_id: &str,
909 source: &str,
910 ) -> Result<WorkNodeId, String> {
911 let snapshot = graph.snapshot();
912 // Only a step the current plan or To-do projection still lists may adopt
913 // new work. Replacing a plan keeps the dropped steps as history (their
914 // ids are reused if the plan grows back, so they are not retired), but a
915 // step nobody can see must never become an operation's parent.
916 let projected_step = |node: &&WorkNode, state: NodeState| {
917 node.kind == NodeKind::PlanStep
918 && node.state == state
919 && (snapshot.compat.plan_order.contains(&node.id)
920 || snapshot
921 .compat
922 .todos
923 .iter()
924 .any(|binding| binding.node == node.id))
925 };
926 if let Some(parent) = snapshot
927 .nodes
928 .iter()
929 .find(|node| projected_step(node, NodeState::Active))
930 .or_else(|| {
931 snapshot
932 .nodes
933 .iter()
934 .find(|node| projected_step(node, NodeState::Ready))
935 })
936 .or_else(|| {
937 snapshot
938 .nodes
939 .iter()
940 .find(|node| node.kind == NodeKind::Objective)
941 })
942 {
943 return Ok(parent.id.clone());
944 }
945
946 let now = now_ms();
947 let id = WorkNodeId::derive(session_id, "objective:runtime-operations");
948 apply_change(
949 graph,
950 session_id,
951 source,
952 WorkGraphChange::AddNode {
953 node: WorkNode {
954 id: id.clone(),
955 kind: NodeKind::Objective,
956 title: "Runtime operations".to_string(),
957 state: NodeState::Ready,
958 acceptance: Vec::new(),
959 binding: None,
960 evidence: None,
961 provenance: Provenance::RuntimeReconcile {
962 source: source.to_string(),
963 observed_at: now,
964 },
965 created_at: now,
966 updated_at: now,
967 },
968 },
969 )?;
970 Ok(id)
971 }
972
973 fn mark_restored_ephemeral_operations_stale(
974 graph: WorkGraphSnapshot,
975 session_id: &str,
976 ) -> Result<(WorkGraphSnapshot, bool), String> {
977 let candidates = graph
978 .nodes
979 .iter()
980 .filter_map(|node| {
981 let binding = node.binding.as_ref()?;
982 (!binding.durable && node.state.is_live()).then(|| node.id.clone())
983 })
984 .collect::<Vec<_>>();
985 if candidates.is_empty() {
986 return Ok((graph, false));
987 }
988 let mut graph = WorkGraph::from_snapshot(graph);
989 for node in candidates {
990 graph
991 .apply(
992 WorkGraphChange::ReconcileOperation {
993 node,
994 obs: OperationObservation::OwnerMissing {
995 checked_at: now_ms(),
996 },
997 },
998 ChangeCtx {
999 session_id: session_id.to_string(),
1000 now: now_ms(),
1001 idempotency_key: None,
1002 },
1003 )
1004 .map_err(|err| format!("restart reconcile: {err}"))?;
1005 }
1006 Ok((graph.into_snapshot(), true))
1007 }
1008
1009 fn bounded_operation_title(title: &str) -> String {
1010 let normalized = title.split_whitespace().collect::<Vec<_>>().join(" ");
1011 let mut chars = normalized.chars();
1012 let bounded = chars.by_ref().take(180).collect::<String>();
1013 if chars.next().is_some() {
1014 format!("{bounded}...")
1015 } else if bounded.is_empty() {
1016 "Runtime operation".to_string()
1017 } else {
1018 bounded
1019 }
1020 }
1021
1022 fn prompt_safe_title(title: &str) -> String {
1023 bounded_operation_title(&title.replace('`', "'"))
1024 }
1025
1026 const fn operation_state_label(state: NodeState) -> &'static str {
1027 match state {
1028 NodeState::Ready => "ready",
1029 NodeState::Initializing => "initializing",
1030 NodeState::Active => "running",
1031 NodeState::Waiting => "waiting",
1032 NodeState::Blocked => "blocked",
1033 NodeState::Completed => "completed",
1034 NodeState::Verified => "verified",
1035 NodeState::Stale => "stale",
1036 NodeState::Superseded => "superseded",
1037 NodeState::Cancelled => "cancelled",
1038 NodeState::Failed => "failed",
1039 }
1040 }
1041
1042 fn graph_for_update(
1043 active: &mut ActiveGraph,
1044 session_id: &str,
1045 plan: &PlanSnapshot,
1046 todos: &TodoListSnapshot,
1047 ) -> Result<WorkGraphSnapshot, String> {
1048 match active.session_id.as_deref() {
1049 // App session transitions are already blocked while runtime work is
1050 // active. Rebind the authority namespace without re-keying graph IDs
1051 // so save-as/fork/new-session flows keep one coherent snapshot.
1052 Some(active_id) if active_id != session_id => {
1053 active.session_id = Some(session_id.to_string());
1054 }
1055 None => active.session_id = Some(session_id.to_string()),
1056 Some(_) => {}
1057 }
1058 if let Some(snapshot) = active.snapshot.as_ref() {
1059 validate(snapshot).map_err(|err| err.to_string())?;
1060 return Ok(snapshot.clone());
1061 }
1062 let graph = if plan.is_empty() && todos.is_empty() {
1063 WorkGraphSnapshot::new()
1064 } else {
1065 import_legacy(session_id, plan, todos)?
1066 };
1067 active.snapshot = Some(graph.clone());
1068 Ok(graph)
1069 }
1070
1071 fn update_plan_graph(
1072 base: WorkGraphSnapshot,
1073 session_id: &str,
1074 tool: &str,
1075 plan: &PlanSnapshot,
1076 ) -> Result<WorkGraphSnapshot, String> {
1077 let mut graph = WorkGraph::from_snapshot(base);
1078 let objective = ensure_objective(&mut graph, session_id, tool, plan)?;
1079 let desired_active_alias = plan.items.iter().enumerate().find_map(|(index, item)| {
1080 (item.status == StepStatus::InProgress)
1081 .then(|| graph.snapshot().compat.plan_order.get(index).cloned())
1082 .flatten()
1083 .filter(|node| {
1084 graph
1085 .snapshot()
1086 .compat
1087 .todos
1088 .iter()
1089 .any(|binding| &binding.node == node)
1090 })
1091 });
1092 if desired_active_alias.is_some() {
1093 deactivate_projected_todos(&mut graph, session_id, tool)?;
1094 }
1095
1096 let mut order = Vec::with_capacity(plan.items.len());
1097 for (index, item) in plan.items.iter().enumerate() {
1098 let id = graph
1099 .snapshot()
1100 .compat
1101 .plan_order
1102 .get(index)
1103 .cloned()
1104 .unwrap_or_else(|| WorkNodeId::derive(session_id, &format!("plan:{index}")));
1105 let provenance = tool_provenance(graph.snapshot(), tool);
1106 upsert_node(
1107 &mut graph,
1108 session_id,
1109 tool,
1110 WorkNode {
1111 id: id.clone(),
1112 kind: NodeKind::PlanStep,
1113 title: item.step.trim().to_string(),
1114 state: plan_node_state(&item.status),
1115 acceptance: Vec::new(),
1116 binding: None,
1117 evidence: None,
1118 provenance,
1119 created_at: now_ms(),
1120 updated_at: now_ms(),
1121 },
1122 )?;
1123 ensure_contains(&mut graph, session_id, tool, &objective, &id)?;
1124 order.push(id);
1125 }
1126 let mut compat = graph.snapshot().compat.clone();
1127 compat.plan = CompatPlanMetadata::from_plan_snapshot(plan);
1128 compat.plan_order = order;
1129 compat.todos.retain(|binding| {
1130 binding.plan_index.is_none_or(|index| {
1131 usize::try_from(index)
1132 .ok()
1133 .is_some_and(|i| i < plan.items.len())
1134 })
1135 });
1136 for binding in &mut compat.todos {
1137 if let Some(index) = binding.plan_index
1138 && let Some(node) = compat
1139 .plan_order
1140 .get(usize::try_from(index).unwrap_or(usize::MAX))
1141 {
1142 binding.node.clone_from(node);
1143 }
1144 }
1145 apply_change(
1146 &mut graph,
1147 session_id,
1148 tool,
1149 WorkGraphChange::ReplaceCompatProjection { compat },
1150 )?;
1151 Ok(graph.into_snapshot())
1152 }
1153
1154 fn update_todo_graph(
1155 base: WorkGraphSnapshot,
1156 session_id: &str,
1157 tool: &str,
1158 todos: &TodoListSnapshot,
1159 ) -> Result<WorkGraphSnapshot, String> {
1160 let mut graph = WorkGraph::from_snapshot(base);
1161 deactivate_projected_todos(&mut graph, session_id, tool)?;
1162 let current_plan = project_plan(graph.snapshot());
1163 let objective = ensure_objective(&mut graph, session_id, tool, &current_plan)?;
1164 let plan_order = graph.snapshot().compat.plan_order.clone();
1165 let mut bindings = Vec::with_capacity(todos.items.len());
1166 for item in &todos.items {
1167 let title = item.content.trim().to_string();
1168 let alias = graph
1169 .snapshot()
1170 .compat
1171 .todos
1172 .iter()
1173 .find(|binding| binding.legacy_id == item.id)
1174 .and_then(|binding| {
1175 binding
1176 .plan_index
1177 .map(|index| (index, binding.node.clone()))
1178 })
1179 .filter(|(index, node)| {
1180 plan_order.get(usize::try_from(*index).unwrap_or(usize::MAX)) == Some(node)
1181 });
1182 let (node, plan_index) = if let Some((index, node)) = alias {
1183 patch_existing_node(
1184 &mut graph,
1185 session_id,
1186 tool,
1187 &node,
1188 title,
1189 todo_node_state(item.status),
1190 )?;
1191 (node, Some(index))
1192 } else {
1193 let node = graph
1194 .snapshot()
1195 .compat
1196 .todos
1197 .iter()
1198 .find(|binding| binding.legacy_id == item.id && binding.plan_index.is_none())
1199 .map(|binding| binding.node.clone())
1200 .unwrap_or_else(|| WorkNodeId::derive(session_id, &format!("todo:{}", item.id)));
1201 let desired = todo_node_state(item.status);
1202 let provenance = tool_provenance(graph.snapshot(), tool);
1203 upsert_node(
1204 &mut graph,
1205 session_id,
1206 tool,
1207 WorkNode {
1208 id: node.clone(),
1209 kind: NodeKind::PlanStep,
1210 title,
1211 state: if desired == NodeState::Active {
1212 NodeState::Ready
1213 } else {
1214 desired
1215 },
1216 acceptance: Vec::new(),
1217 binding: None,
1218 evidence: None,
1219 provenance,
1220 created_at: now_ms(),
1221 updated_at: now_ms(),
1222 },
1223 )?;
1224 ensure_contains(&mut graph, session_id, tool, &objective, &node)?;
1225 if desired == NodeState::Active {
1226 let clean_title = graph
1227 .snapshot()
1228 .node(&node)
1229 .map(|node| node.title.clone())
1230 .ok_or_else(|| format!("node {node} not found after insert"))?;
1231 patch_existing_node(&mut graph, session_id, tool, &node, clean_title, desired)?;
1232 }
1233 (node, None)
1234 };
1235 bindings.push(CompatTodoBinding {
1236 legacy_id: item.id,
1237 node,
1238 plan_index,
1239 });
1240 }
1241 let mut compat = graph.snapshot().compat.clone();
1242 compat.todos = bindings;
1243 apply_change(
1244 &mut graph,
1245 session_id,
1246 tool,
1247 WorkGraphChange::ReplaceCompatProjection { compat },
1248 )?;
1249 Ok(graph.into_snapshot())
1250 }
1251
1252 fn ensure_objective(
1253 graph: &mut WorkGraph,
1254 session_id: &str,
1255 tool: &str,
1256 plan: &PlanSnapshot,
1257 ) -> Result<WorkNodeId, String> {
1258 let id = graph
1259 .snapshot()
1260 .nodes
1261 .iter()
1262 .find(|node| node.kind == NodeKind::Objective)
1263 .map(|node| node.id.clone())
1264 .unwrap_or_else(|| WorkNodeId::derive(session_id, "objective"));
1265 let title = plan
1266 .objective
1267 .as_deref()
1268 .or(plan.title.as_deref())
1269 .unwrap_or("Session work")
1270 .to_string();
1271 upsert_node(
1272 graph,
1273 session_id,
1274 tool,
1275 WorkNode {
1276 id: id.clone(),
1277 kind: NodeKind::Objective,
1278 title,
1279 state: NodeState::Ready,
1280 acceptance: Vec::new(),
1281 binding: None,
1282 evidence: None,
1283 provenance: tool_provenance(graph.snapshot(), tool),
1284 created_at: now_ms(),
1285 updated_at: now_ms(),
1286 },
1287 )?;
1288 Ok(id)
1289 }
1290
1291 fn upsert_node(
1292 graph: &mut WorkGraph,
1293 session_id: &str,
1294 tool: &str,
1295 node: WorkNode,
1296 ) -> Result<(), String> {
1297 if let Some(existing) = graph.snapshot().node(&node.id) {
1298 if existing.kind != node.kind {
1299 return Err(format!("node {} changed kind", node.id));
1300 }
1301 patch_existing_node(graph, session_id, tool, &node.id, node.title, node.state)
1302 } else {
1303 apply_change(graph, session_id, tool, WorkGraphChange::AddNode { node })
1304 }
1305 }
1306
1307 fn patch_existing_node(
1308 graph: &mut WorkGraph,
1309 session_id: &str,
1310 tool: &str,
1311 id: &WorkNodeId,
1312 title: String,
1313 state: NodeState,
1314 ) -> Result<(), String> {
1315 let current = graph
1316 .snapshot()
1317 .node(id)
1318 .ok_or_else(|| format!("node {id} not found"))?;
1319 if current.title == title && current.state == state {
1320 // Semantic no-ops must stay no-ops. In particular, terminal nodes are
1321 // immutable, so refreshing another To-do item must not try to rewrite
1322 // a settled sibling merely to stamp newer tool provenance.
1323 return Ok(());
1324 }
1325 let provenance = tool_provenance(graph.snapshot(), tool);
1326 apply_change(
1327 graph,
1328 session_id,
1329 tool,
1330 WorkGraphChange::UpdateNode {
1331 id: id.clone(),
1332 patch: WorkNodePatch {
1333 title: Some(title),
1334 state: Some(state),
1335 provenance: Some(provenance),
1336 ..WorkNodePatch::default()
1337 },
1338 },
1339 )
1340 }
1341
1342 fn ensure_contains(
1343 graph: &mut WorkGraph,
1344 session_id: &str,
1345 tool: &str,
1346 parent: &WorkNodeId,
1347 child: &WorkNodeId,
1348 ) -> Result<(), String> {
1349 let id = WorkEdgeId::derive(
1350 session_id,
1351 &format!("contains:{}:{}", parent.as_str(), child.as_str()),
1352 );
1353 if graph.snapshot().edge(&id).is_some() {
1354 return Ok(());
1355 }
1356 apply_change(
1357 graph,
1358 session_id,
1359 tool,
1360 WorkGraphChange::AddEdge {
1361 edge: WorkEdge {
1362 id,
1363 kind: EdgeKind::Contains,
1364 from: parent.clone(),
1365 to: child.clone(),
1366 },
1367 },
1368 )
1369 }
1370
1371 fn deactivate_projected_todos(
1372 graph: &mut WorkGraph,
1373 session_id: &str,
1374 tool: &str,
1375 ) -> Result<(), String> {
1376 let active = graph
1377 .snapshot()
1378 .compat
1379 .todos
1380 .iter()
1381 .filter_map(|binding| {
1382 graph
1383 .snapshot()
1384 .node(&binding.node)
1385 .filter(|node| node.state == NodeState::Active)
1386 .map(|node| (node.id.clone(), node.title.clone()))
1387 })
1388 .collect::<Vec<_>>();
1389 for (id, title) in active {
1390 patch_existing_node(graph, session_id, tool, &id, title, NodeState::Ready)?;
1391 }
1392 Ok(())
1393 }
1394
1395 fn apply_change(
1396 graph: &mut WorkGraph,
1397 session_id: &str,
1398 tool: &str,
1399 change: WorkGraphChange,
1400 ) -> Result<(), String> {
1401 graph
1402 .apply(
1403 change,
1404 ChangeCtx {
1405 session_id: session_id.to_string(),
1406 now: now_ms(),
1407 idempotency_key: None,
1408 },
1409 )
1410 .map(|_| ())
1411 .map_err(|err| format!("{tool}: {err}"))
1412 }
1413
1414 fn validate_combined(
1415 graph: &WorkGraphSnapshot,
1416 plan: &PlanSnapshot,
1417 todos: &TodoListSnapshot,
1418 ) -> Result<(), String> {
1419 validate(graph).map_err(|err| err.to_string())?;
1420 if &project_plan(graph) != plan {
1421 return Err("Work Graph Plan projection is inconsistent".to_string());
1422 }
1423 if &project_todos(graph) != todos {
1424 return Err("Work Graph To-do projection is inconsistent".to_string());
1425 }
1426 TodoList::from_snapshot(todos)?;
1427 Ok(())
1428 }
1429
1430 fn tool_provenance(snapshot: &WorkGraphSnapshot, tool: &str) -> Provenance {
1431 Provenance::ToolUpdate {
1432 tool: tool.to_string(),
1433 call_id: format!("{tool}:{}", snapshot.revision.saturating_add(1)),
1434 }
1435 }
1436
1437 fn plan_node_state(status: &StepStatus) -> NodeState {
1438 match status {
1439 StepStatus::Pending => NodeState::Ready,
1440 StepStatus::InProgress => NodeState::Active,
1441 StepStatus::Completed => NodeState::Completed,
1442 }
1443 }
1444
1445 fn todo_node_state(status: TodoStatus) -> NodeState {
1446 match status {
1447 TodoStatus::Pending => NodeState::Ready,
1448 TodoStatus::InProgress => NodeState::Active,
1449 TodoStatus::Completed => NodeState::Completed,
1450 TodoStatus::Cancelled => NodeState::Cancelled,
1451 }
1452 }
1453
1454 fn resolved_session_id(active: &ActiveGraph, requested: Option<&str>) -> String {
1455 requested
1456 .map(str::to_string)
1457 .or_else(|| active.session_id.clone())
1458 .unwrap_or_else(|| "unsaved-work".to_string())
1459 }
1460
1461 fn now_ms() -> i64 {
1462 chrono::Utc::now().timestamp_millis()
1463 }
1464
1465 fn lock_unpoisoned<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
1466 mutex
1467 .lock()
1468 .unwrap_or_else(std::sync::PoisonError::into_inner)
1469 }
1470
1471 // These callers are synchronous methods reached from tool code on the Tokio
1472 // runtime, so this contention retry must not park the worker with
1473 // `thread::sleep`. `yield_now` hands the thread to the holder long enough for
1474 // its microsecond-scale critical section; if the lock is still contended the
1475 // caller surfaces the existing "state is busy" error instead of stalling.
1476 fn retry_lock<T>(
1477 mutex: &tokio::sync::Mutex<T>,
1478 retries: u32,
1479 ) -> Option<tokio::sync::MutexGuard<'_, T>> {
1480 for _ in 0..retries {
1481 if let Ok(guard) = mutex.try_lock() {
1482 return Some(guard);
1483 }
1484 std::thread::yield_now();
1485 }
1486 None
1487 }
1488
1489 #[cfg(test)]
1490 mod tests {
1491 use super::*;
1492 use crate::tools::plan::PlanItemArg;
1493 use crate::work_graph::{EvidenceKind, EvidenceRef, OwnerState};
1494
1495 #[test]
1496 fn operation_lifecycle_is_registered_idempotent_and_receipt_only() {
1497 let runtime = new_shared_work_runtime(
1498 crate::tools::todo::new_shared_todo_list(),
1499 crate::tools::plan::new_shared_plan_state(),
1500 );
1501 let intent = OperationIntent::new(
1502 "shell:shell_test",
1503 "silent `owner` command",
1504 false,
1505 "exec_shell",
1506 "shell_test",
1507 );
1508 let node_id = runtime
1509 .register_operation("session", intent.clone())
1510 .expect("register before spawn");
1511 assert_eq!(
1512 runtime.register_operation("session", intent),
1513 Ok(node_id.clone()),
1514 "repeat spawn intent must not duplicate the operation"
1515 );
1516 let initialized = runtime
1517 .capture(Some("session"))
1518 .expect("capture")
1519 .expect("graph");
1520 let node = initialized.graph.node(&node_id).expect("operation node");
1521 assert_eq!(node.state, NodeState::Initializing);
1522 assert!(
1523 initialized
1524 .graph
1525 .edges
1526 .iter()
1527 .any(|edge| { edge.kind == EdgeKind::Contains && edge.to == node_id })
1528 );
1529
1530 let output = EvidenceRef::new(
1531 EvidenceKind::Receipt {
1532 owner: "shell".to_string(),
1533 },
1534 "shell:shell_test:output",
1535 Some(4_096),
1536 true,
1537 )
1538 .expect("safe logical receipt");
1539 assert_eq!(
1540 runtime.reconcile_operation(
1541 "session",
1542 OperationOwnerSnapshot::new("shell:shell_test", OwnerState::Running, 7, 10,)
1543 .with_output(output),
1544 ),
1545 Ok(true)
1546 );
1547 assert_eq!(
1548 runtime.reconcile_operation(
1549 "session",
1550 OperationOwnerSnapshot::new("shell:shell_test", OwnerState::Completed, 7, 11,),
1551 ),
1552 Ok(false),
1553 "the same binding sequence is an idempotent no-op"
1554 );
1555 assert!(
1556 runtime
1557 .reconcile_operation(
1558 "session",
1559 OperationOwnerSnapshot::new("shell:unknown", OwnerState::Running, 1, 12,),
1560 )
1561 .expect_err("unknown owner must not materialize after spawn")
1562 .contains("not registered")
1563 );
1564 let running = runtime
1565 .capture(Some("session"))
1566 .expect("capture running")
1567 .expect("graph");
1568 let binding = running
1569 .graph
1570 .node(&node_id)
1571 .and_then(|node| node.binding.as_ref())
1572 .expect("binding");
1573 assert_eq!(
1574 running.graph.node(&node_id).map(|node| node.state),
1575 Some(NodeState::Active)
1576 );
1577 assert_eq!(
1578 binding
1579 .last_observation
1580 .as_ref()
1581 .and_then(|obs| obs.output.as_ref())
1582 .and_then(EvidenceRef::raw_bytes),
1583 Some(4_096)
1584 );
1585 let summary = runtime
1586 .active_operation_summary(Some("session"))
1587 .expect("compaction re-anchor");
1588 assert!(summary.contains("shell:shell_test"), "{summary}");
1589 assert!(summary.contains("silent 'owner' command"), "{summary}");
1590 assert!(!summary.contains("4,096"), "{summary}");
1591
1592 assert_eq!(
1593 runtime.record_reasoning_effort_change(
1594 Some("session"),
1595 ReasoningEffortTier::Low,
1596 ReasoningEffortTier::High,
1597 crate::config::ProviderKind::Moonshot,
1598 "moonshot",
1599 Some(crate::config::DEFAULT_MOONSHOT_BASE_URL),
1600 Some("kimi-k2.5"),
1601 ),
1602 Ok(Some(node_id.clone()))
1603 );
1604 let activity = runtime
1605 .capture(Some("session"))
1606 .expect("capture effort activity")
1607 .expect("graph")
1608 .graph
1609 .activities
1610 .last()
1611 .cloned()
1612 .expect("effort activity");
1613 let ts = match &activity {
1614 WorkActivityEvent::ReasoningEffortChanged { ts, .. } => *ts,
1615 };
1616 assert_eq!(
1617 activity,
1618 WorkActivityEvent::ReasoningEffortChanged {
1619 requested: ReasoningEffortTier::Low,
1620 effective: ReasoningEffortTier::High,
1621 provider_kind: Some(crate::config::ProviderKind::Moonshot),
1622 provider: "moonshot".to_string(),
1623 endpoint_identity: Some(crate::config::DEFAULT_MOONSHOT_BASE_URL.to_string()),
1624 model: Some("kimi-k2.5".to_string()),
1625 ts,
1626 operation: Some(node_id.clone()),
1627 }
1628 );
1629
1630 runtime
1631 .record_operation_approval(
1632 "session",
1633 "shell:shell_test",
1634 "operate-verification:shell_test",
1635 "exec_shell",
1636 "approval_test",
1637 )
1638 .expect("approval provenance");
1639 let approved = runtime
1640 .capture(Some("session"))
1641 .expect("capture approval")
1642 .expect("graph");
1643 assert!(
1644 approved
1645 .graph
1646 .nodes
1647 .iter()
1648 .any(|node| node.kind == NodeKind::Approval)
1649 );
1650 assert!(
1651 approved
1652 .graph
1653 .edges
1654 .iter()
1655 .any(|edge| edge.kind == EdgeKind::RequiresApproval)
1656 );
1657
1658 runtime
1659 .reconcile_operation(
1660 "session",
1661 OperationOwnerSnapshot::new("shell:shell_test", OwnerState::Completed, 8, 13),
1662 )
1663 .expect("terminal owner report");
1664 runtime
1665 .reconcile_observation(
1666 "session",
1667 "shell:shell_test",
1668 OperationObservation::CancelUpdate {
1669 outcome: super::super::CancelOutcome::AlreadyFinished,
1670 at: 14,
1671 },
1672 )
1673 .expect("already-finished cancellation receipt");
1674 assert_eq!(
1675 runtime
1676 .capture(Some("session"))
1677 .expect("capture terminal")
1678 .expect("graph")
1679 .graph
1680 .node(&node_id)
1681 .map(|node| node.state),
1682 Some(NodeState::Completed),
1683 "already-finished cancellation must not rewrite owner state"
1684 );
1685 }
1686
1687 #[test]
1688 fn restore_stales_only_live_ephemeral_operations() {
1689 let runtime = new_shared_work_runtime(
1690 crate::tools::todo::new_shared_todo_list(),
1691 crate::tools::plan::new_shared_plan_state(),
1692 );
1693 for id in ["shell:live", "shell:done"] {
1694 runtime
1695 .register_operation(
1696 "session",
1697 OperationIntent::new(id, id, false, "exec_shell", id),
1698 )
1699 .expect("register shell");
1700 }
1701 runtime
1702 .reconcile_operation(
1703 "session",
1704 OperationOwnerSnapshot::new("shell:live", OwnerState::Running, 1, 1),
1705 )
1706 .expect("live shell");
1707 runtime
1708 .reconcile_operation(
1709 "session",
1710 OperationOwnerSnapshot::new("shell:done", OwnerState::Completed, 1, 1),
1711 )
1712 .expect("completed shell");
1713 let saved = runtime
1714 .capture(Some("session"))
1715 .expect("capture")
1716 .expect("saved graph");
1717
1718 let restored = new_shared_work_runtime(
1719 crate::tools::todo::new_shared_todo_list(),
1720 crate::tools::plan::new_shared_plan_state(),
1721 );
1722 restored
1723 .restore("session", Some(&saved.graph), &saved.todos, &saved.plan)
1724 .expect("restore graph");
1725 let graph = restored
1726 .capture(Some("session"))
1727 .expect("capture restored")
1728 .expect("restored graph")
1729 .graph;
1730 let state_for = |external: &str| {
1731 graph
1732 .nodes
1733 .iter()
1734 .find(|node| {
1735 node.binding
1736 .as_ref()
1737 .is_some_and(|binding| binding.external == external)
1738 })
1739 .map(|node| node.state)
1740 };
1741 assert_eq!(state_for("shell:live"), Some(NodeState::Stale));
1742 assert_eq!(state_for("shell:done"), Some(NodeState::Completed));
1743 assert!(restored.has_pending_publish());
1744 }
1745
1746 #[test]
1747 fn durable_owner_same_sequence_recovers_stale_once_and_rejects_regression() {
1748 let runtime = new_shared_work_runtime(
1749 crate::tools::todo::new_shared_todo_list(),
1750 crate::tools::plan::new_shared_plan_state(),
1751 );
1752 runtime
1753 .register_operation(
1754 "session",
1755 OperationIntent::new(
1756 "task:task_restore",
1757 "restored task",
1758 true,
1759 "task_create",
1760 "task_restore",
1761 ),
1762 )
1763 .expect("register durable owner");
1764 let running = OperationOwnerSnapshot::new("task:task_restore", OwnerState::Running, 7, 10);
1765 assert_eq!(
1766 runtime.reconcile_operation("session", running.clone()),
1767 Ok(true)
1768 );
1769 assert_eq!(
1770 runtime.reconcile_observation(
1771 "session",
1772 "task:task_restore",
1773 OperationObservation::OwnerMissing { checked_at: 11 },
1774 ),
1775 Ok(true)
1776 );
1777 assert_eq!(
1778 runtime
1779 .capture(Some("session"))
1780 .expect("capture stale")
1781 .expect("graph")
1782 .graph
1783 .nodes
1784 .iter()
1785 .find(|node| {
1786 node.binding
1787 .as_ref()
1788 .is_some_and(|binding| binding.external == "task:task_restore")
1789 })
1790 .map(|node| node.state),
1791 Some(NodeState::Stale)
1792 );
1793
1794 let replay = OperationOwnerSnapshot::new("task:task_restore", OwnerState::Running, 7, 12);
1795 assert_eq!(
1796 runtime.reconcile_operation("session", replay.clone()),
1797 Ok(true)
1798 );
1799 assert_eq!(
1800 runtime.reconcile_operation("session", replay),
1801 Ok(false),
1802 "same-sequence recovery must happen only while the node is stale"
1803 );
1804 assert_eq!(
1805 runtime
1806 .capture(Some("session"))
1807 .expect("capture recovered")
1808 .expect("graph")
1809 .graph
1810 .nodes
1811 .iter()
1812 .find(|node| {
1813 node.binding
1814 .as_ref()
1815 .is_some_and(|binding| binding.external == "task:task_restore")
1816 })
1817 .map(|node| node.state),
1818 Some(NodeState::Active)
1819 );
1820 assert_eq!(
1821 runtime.reconcile_operation(
1822 "session",
1823 OperationOwnerSnapshot::new("task:task_restore", OwnerState::Waiting, 7, 13,),
1824 ),
1825 Ok(false),
1826 "ordinary same-key duplicates retain the reducer no-op contract"
1827 );
1828 runtime
1829 .reconcile_observation(
1830 "session",
1831 "task:task_restore",
1832 OperationObservation::OwnerMissing { checked_at: 14 },
1833 )
1834 .expect("mark owner missing again");
1835 assert!(
1836 runtime
1837 .reconcile_operation(
1838 "session",
1839 OperationOwnerSnapshot::new("task:task_restore", OwnerState::Waiting, 7, 15,),
1840 )
1841 .expect_err("inconsistent replay cannot revive a stale node")
1842 .contains("changed observation")
1843 );
1844 assert!(
1845 runtime
1846 .reconcile_operation(
1847 "session",
1848 OperationOwnerSnapshot::new("task:task_restore", OwnerState::Running, 6, 16,),
1849 )
1850 .expect_err("owner sequence cannot regress")
1851 .contains("sequence regressed")
1852 );
1853 }
1854
1855 #[tokio::test]
1856 async fn replaced_plan_steps_never_parent_new_operations() {
1857 let runtime = new_shared_work_runtime(
1858 crate::tools::todo::new_shared_todo_list(),
1859 crate::tools::plan::new_shared_plan_state(),
1860 );
1861 let step = |step: &str, status: StepStatus| PlanItemArg {
1862 step: step.to_string(),
1863 status,
1864 };
1865 let three = PlanSnapshot {
1866 items: vec![
1867 step("first", StepStatus::Pending),
1868 step("second", StepStatus::Pending),
1869 step("third", StepStatus::InProgress),
1870 ],
1871 ..PlanSnapshot::default()
1872 };
1873 runtime
1874 .apply_plan_update("session", "update_plan", &three)
1875 .await
1876 .expect("three-step plan");
1877 let one = PlanSnapshot {
1878 items: vec![step("only", StepStatus::Pending)],
1879 ..PlanSnapshot::default()
1880 };
1881 runtime
1882 .apply_plan_update("session", "update_plan", &one)
1883 .await
1884 .expect("replacement plan");
1885
1886 let operation = runtime
1887 .register_operation(
1888 "session",
1889 OperationIntent::new(
1890 "shell:after",
1891 "after replacement",
1892 false,
1893 "exec_shell",
1894 "c1",
1895 ),
1896 )
1897 .expect("register operation");
1898 let graph = runtime
1899 .capture(Some("session"))
1900 .expect("capture")
1901 .expect("graph")
1902 .graph;
1903 assert_eq!(graph.compat.plan_order.len(), 1);
1904 let parent = graph
1905 .edges
1906 .iter()
1907 .find(|edge| edge.kind == EdgeKind::Contains && edge.to == operation)
1908 .map(|edge| edge.from.clone());
1909 assert_eq!(
1910 parent,
1911 Some(graph.compat.plan_order[0].clone()),
1912 "an orphaned step from the replaced plan must not adopt new work"
1913 );
1914 }
1915
1916 #[test]
1917 fn legacy_restore_stays_pending_until_first_graph_bearing_write() {
1918 let todos = crate::tools::todo::new_shared_todo_list();
1919 let plan = crate::tools::plan::new_shared_plan_state();
1920 let runtime = new_shared_work_runtime(todos.clone(), plan.clone());
1921 let legacy_plan = PlanSnapshot {
1922 items: vec![PlanItemArg {
1923 step: "Migrate once".to_string(),
1924 status: StepStatus::InProgress,
1925 }],
1926 ..PlanSnapshot::default()
1927 };
1928
1929 runtime
1930 .restore(
1931 "legacy-session",
1932 None,
1933 &TodoListSnapshot::default(),
1934 &legacy_plan,
1935 )
1936 .expect("restore legacy state");
1937 assert!(runtime.has_pending_publish());
1938 let captured = runtime
1939 .capture(Some("legacy-session"))
1940 .expect("capture imported graph")
1941 .expect("state");
1942 assert!(captured.graph.import_digest.is_some());
1943 assert_eq!(plan.blocking_lock().snapshot(), legacy_plan);
1944 assert_eq!(runtime.publish_pending_sync(), Ok(true));
1945 assert!(!runtime.has_pending_publish());
1946 assert!(todos.blocking_lock().snapshot().is_empty());
1947 }
1948
1949 /// #6842: one Operation node per finished shell call grew without bound.
1950 #[test]
1951 fn ended_shell_operations_are_capped_oldest_first() {
1952 let runtime = new_shared_work_runtime(
1953 crate::tools::todo::new_shared_todo_list(),
1954 crate::tools::plan::new_shared_plan_state(),
1955 );
1956 let total = super::super::model::ENDED_OPERATION_CAP + 44;
1957 let mut ids = Vec::new();
1958 for i in 0..total {
1959 let external = format!("shell:cap_{i}");
1960 let intent =
1961 OperationIntent::new(&external, "ls", false, "exec_shell", format!("c{i}"));
1962 ids.push(
1963 runtime
1964 .register_operation("session", intent)
1965 .expect("register"),
1966 );
1967 for (seq, state) in [(1, OwnerState::Running), (2, OwnerState::Completed)] {
1968 runtime
1969 .reconcile_operation(
1970 "session",
1971 OperationOwnerSnapshot::new(
1972 &external,
1973 state,
1974 seq,
1975 i as i64 * 10 + seq as i64,
1976 ),
1977 )
1978 .expect("reconcile");
1979 }
1980 }
1981 let graph = runtime
1982 .capture(Some("session"))
1983 .expect("capture")
1984 .expect("graph")
1985 .graph;
1986 crate::work_graph::validate(&graph).expect("pruned graph validates");
1987 let operations = graph
1988 .nodes
1989 .iter()
1990 .filter(|node| node.kind == NodeKind::Operation)
1991 .count();
1992 // Pruned to the cap at each registration; the last call then ended.
1993 assert!(
1994 operations <= super::super::model::ENDED_OPERATION_CAP + 1,
1995 "{operations}"
1996 );
1997 assert!(graph.node(&ids[0]).is_none(), "oldest ended call evicted");
1998 assert!(graph.node(&ids[total - 1]).is_some(), "newest call kept");
1999 assert!(
2000 graph
2001 .edges
2002 .iter()
2003 .all(|edge| graph.node(&edge.from).is_some() && graph.node(&edge.to).is_some())
2004 );
2005 }
2006 }
2007
2007 lines RUST