返回 CodeWhale
reducer.rs
根目录 / crates / tui / src / work_graph / reducer.rs
1 //! The only write path for the work graph.
2 //!
3 //! [`apply`] is pure and deterministic: no clock reads, no RNG, no I/O — the
4 //! caller-supplied [`ChangeCtx`] carries timestamps and session identity, and
5 //! every derived ID comes from SHA-256 over `(session_id, discriminator)`.
6 //! The same snapshot + change + ctx always produce byte-identical results.
7 //!
8 //! Contract per change:
9 //! 1. idempotency-key duplicates are acknowledged as no-op receipts;
10 //! 2. the change is applied to a copy;
11 //! 3. the candidate is validated fail-closed — on any violation the input
12 //! snapshot is untouched and the caller gets the full report;
13 //! 4. the revision increments exactly once;
14 //! 5. the receipt is pushed onto the bounded history.
15 //!
16 //! Callers on `Ok`: derive compat snapshots → validate combined → persist →
17 //! publish projections (publish AFTER the persist enqueue, never before).
18 //! On `Err`: surface the report; no state changed.
19
20 use super::events::{
21 CancelOutcome, ChangeCtx, ChangeReceipt, ObservationSummary, OperationObservation, OwnerState,
22 WorkGraphChange, WorkGraphProposal, WorkNodePatch,
23 };
24 use super::ids::{ChangeId, WorkEdgeId, WorkNodeId};
25 use super::model::{
26 EdgeKind, EvidenceRef, NodeKind, NodeState, OperationBinding, Provenance, WorkActivityEvent,
27 WorkEdge, WorkGraphSnapshot, WorkNode,
28 };
29 use super::validate::{ValidationCode, ValidationReport, validate};
30
31 /// Apply one change, producing the next snapshot and a receipt, or a
32 /// validation report with the input snapshot untouched.
33 pub fn apply(
34 g: &WorkGraphSnapshot,
35 change: WorkGraphChange,
36 ctx: ChangeCtx,
37 ) -> Result<(WorkGraphSnapshot, ChangeReceipt), ValidationReport> {
38 if let Some(key) = &ctx.idempotency_key
39 && g.seen_keys.contains(key)
40 {
41 // Duplicate runtime event: acknowledge without effect.
42 let receipt = ChangeReceipt {
43 change_id: ChangeId::derive(
44 &ctx.session_id,
45 &format!("noop:{}:{}", key.binding.as_str(), key.seq),
46 ),
47 revision: g.revision,
48 summary: format!("{} (duplicate)", change.kind_name()),
49 applied_at: ctx.now,
50 idempotency_key: Some(key.clone()),
51 no_op: true,
52 };
53 return Ok((g.clone(), receipt));
54 }
55
56 let mut next = apply_pure(g, &change, &ctx)?;
57 validate(&next)?; // fail closed: `g` unchanged on Err
58 next.revision = g.revision + 1; // exactly once
59 let receipt = ChangeReceipt::of(&change, next.revision, &ctx);
60 next.history.push_bounded(receipt.clone());
61 if let Some(key) = ctx.idempotency_key {
62 next.seen_keys.insert(key);
63 }
64 Ok((next, receipt))
65 }
66
67 fn structural(message: impl Into<String>) -> ValidationReport {
68 ValidationReport::single(ValidationCode::Structural, message)
69 }
70
71 fn apply_pure(
72 g: &WorkGraphSnapshot,
73 change: &WorkGraphChange,
74 ctx: &ChangeCtx,
75 ) -> Result<WorkGraphSnapshot, ValidationReport> {
76 let mut next = g.clone();
77 match change {
78 WorkGraphChange::AddNode { node } => {
79 add_node(&mut next, node.clone())?;
80 }
81 WorkGraphChange::UpdateNode { id, patch } => {
82 patch_node(&mut next, id, patch, ctx.now)?;
83 }
84 WorkGraphChange::AddEdge { edge } => {
85 add_edge(&mut next, edge.clone())?;
86 }
87 WorkGraphChange::RemoveEdge { id } => {
88 let before = next.edges.len();
89 next.edges.retain(|e| &e.id != id);
90 if next.edges.len() == before {
91 return Err(structural(format!("edge {id} not found")));
92 }
93 }
94 WorkGraphChange::BindOperation { node, binding } => {
95 let now = ctx.now;
96 let target = next
97 .node_mut(node)
98 .ok_or_else(|| structural(format!("node {node} not found")))?;
99 target.binding = Some(binding.clone());
100 target.updated_at = now;
101 }
102 WorkGraphChange::ReconcileOperation { node, obs } => {
103 reconcile(&mut next, node, obs, ctx)?;
104 }
105 WorkGraphChange::AttachEvidence { node, evidence } => {
106 attach_evidence(&mut next, node, evidence, ctx)?;
107 }
108 WorkGraphChange::ProposePlanDiff { proposal } => {
109 if next.proposals.iter().any(|p| p.id == proposal.id) {
110 return Err(structural(format!("duplicate proposal {}", proposal.id)));
111 }
112 validate_proposal(&next, proposal, ctx)?;
113 next.proposals.push(proposal.clone());
114 }
115 WorkGraphChange::WithdrawPlanDiff { proposal_id } => {
116 let before = next.proposals.len();
117 next.proposals
118 .retain(|proposal| proposal.id != *proposal_id);
119 if next.proposals.len() == before {
120 return Err(structural(format!("proposal {proposal_id} not found")));
121 }
122 }
123 WorkGraphChange::AcceptPlanDiff {
124 proposal_id,
125 approval,
126 } => {
127 let index = next
128 .proposals
129 .iter()
130 .position(|p| p.id == *proposal_id)
131 .ok_or_else(|| structural(format!("proposal {proposal_id} not found")))?;
132 let proposal = next.proposals.remove(index);
133 apply_proposal(&mut next, &proposal, ctx)?;
134 // Record the acceptance as an Approval node so the review action
135 // itself is part of the graph's history.
136 let approval_node = WorkNode {
137 id: WorkNodeId::derive(
138 &ctx.session_id,
139 &format!("approval:{}", proposal_id.as_str()),
140 ),
141 kind: NodeKind::Approval,
142 title: format!("plan diff approved: {}", approval.reference),
143 state: NodeState::Completed,
144 acceptance: Vec::new(),
145 binding: None,
146 evidence: None,
147 provenance: Provenance::UserEdit {
148 proposal_id: proposal_id.clone(),
149 },
150 created_at: ctx.now,
151 updated_at: ctx.now,
152 };
153 add_node(&mut next, approval_node)?;
154 }
155 WorkGraphChange::Supersede { old, replacement } => {
156 if old == replacement {
157 return Err(structural("node cannot supersede itself"));
158 }
159 if next.node(replacement).is_none() {
160 return Err(structural(format!(
161 "replacement node {replacement} not found"
162 )));
163 }
164 let now = ctx.now;
165 {
166 let old_node = next
167 .node_mut(old)
168 .ok_or_else(|| structural(format!("node {old} not found")))?;
169 // Explicit supersede is the sanctioned way past V9's
170 // terminal-state protection.
171 old_node.state = NodeState::Superseded;
172 old_node.updated_at = now;
173 }
174 let edge = WorkEdge {
175 id: WorkEdgeId::derive(
176 &ctx.session_id,
177 &format!("supersedes:{}:{}", replacement.as_str(), old.as_str()),
178 ),
179 kind: EdgeKind::Supersedes,
180 from: replacement.clone(),
181 to: old.clone(),
182 };
183 add_edge(&mut next, edge)?;
184 }
185 WorkGraphChange::ReplaceCompatProjection { compat } => {
186 next.compat = compat.clone();
187 }
188 WorkGraphChange::SetImportDigest { digest } => {
189 if digest.is_empty() {
190 return Err(structural("legacy import digest cannot be empty"));
191 }
192 next.import_digest = Some(digest.clone());
193 }
194 WorkGraphChange::RecordActivity { event } => {
195 let operation = match event {
196 WorkActivityEvent::ReasoningEffortChanged { operation, .. } => operation,
197 };
198 if let Some(operation) = operation {
199 let node = next.node(operation).ok_or_else(|| {
200 structural(format!("activity references missing operation {operation}"))
201 })?;
202 if node.kind != NodeKind::Operation || !node.state.is_live() {
203 return Err(structural(format!(
204 "activity operation {operation} is not live"
205 )));
206 }
207 }
208 next.activities.push_bounded(event.clone());
209 }
210 WorkGraphChange::PruneEndedOperations { keep } => {
211 prune_ended_operations(&mut next, *keep);
212 }
213 }
214 Ok(next)
215 }
216
217 /// Operation nodes that may be evicted: ended (not live, not `Stale` — a
218 /// stale shell can still be recovered by a late owner report), bound to a
219 /// non-durable owner, touched by no edge except an incoming `Contains`, and
220 /// referenced by no activity or pending proposal.
221 pub(crate) fn evictable_operations(g: &WorkGraphSnapshot) -> Vec<&WorkNode> {
222 use std::collections::HashSet;
223 let mut pinned: HashSet<&WorkNodeId> = HashSet::new();
224 for edge in &g.edges {
225 if edge.kind != EdgeKind::Contains {
226 pinned.insert(&edge.from);
227 pinned.insert(&edge.to);
228 } else {
229 pinned.insert(&edge.from);
230 }
231 }
232 for activity in g.activities.iter() {
233 let WorkActivityEvent::ReasoningEffortChanged { operation, .. } = activity;
234 if let Some(operation) = operation {
235 pinned.insert(operation);
236 }
237 }
238 for proposal in &g.proposals {
239 pinned.extend(proposal.removed_nodes.iter());
240 pinned.extend(proposal.updated_nodes.iter().map(|update| &update.id));
241 for edge in &proposal.added_edges {
242 pinned.insert(&edge.from);
243 pinned.insert(&edge.to);
244 }
245 }
246 g.nodes
247 .iter()
248 .filter(|node| {
249 node.kind == NodeKind::Operation
250 && matches!(
251 node.state,
252 NodeState::Completed
253 | NodeState::Failed
254 | NodeState::Cancelled
255 | NodeState::Superseded
256 | NodeState::Verified
257 )
258 && node
259 .binding
260 .as_ref()
261 .is_some_and(|binding| !binding.durable)
262 && !pinned.contains(&node.id)
263 })
264 .collect()
265 }
266
267 /// Keep the newest `keep` evictable operations (by `updated_at`, then
268 /// insertion order, so the result is deterministic) and remove the rest with
269 /// their `Contains` edges.
270 fn prune_ended_operations(next: &mut WorkGraphSnapshot, keep: usize) {
271 let positions: std::collections::HashMap<&WorkNodeId, usize> = next
272 .nodes
273 .iter()
274 .enumerate()
275 .map(|(position, node)| (&node.id, position))
276 .collect();
277 let mut evictable: Vec<(i64, usize, WorkNodeId)> = evictable_operations(next)
278 .into_iter()
279 .map(|node| (node.updated_at, positions[&node.id], node.id.clone()))
280 .collect();
281 drop(positions);
282 let Some(excess) = evictable.len().checked_sub(keep).filter(|n| *n > 0) else {
283 return;
284 };
285 evictable.sort();
286 let evicted: std::collections::HashSet<WorkNodeId> = evictable
287 .into_iter()
288 .take(excess)
289 .map(|(_, _, id)| id)
290 .collect();
291 next.nodes.retain(|node| !evicted.contains(&node.id));
292 next.edges
293 .retain(|edge| !evicted.contains(&edge.from) && !evicted.contains(&edge.to));
294 }
295
296 fn add_node(next: &mut WorkGraphSnapshot, node: WorkNode) -> Result<(), ValidationReport> {
297 if next.node(&node.id).is_some() {
298 return Err(structural(format!("duplicate node {}", node.id)));
299 }
300 next.nodes.push(node);
301 Ok(())
302 }
303
304 fn add_edge(next: &mut WorkGraphSnapshot, edge: WorkEdge) -> Result<(), ValidationReport> {
305 if next.edge(&edge.id).is_some() {
306 return Err(structural(format!("duplicate edge {}", edge.id)));
307 }
308 for endpoint in [&edge.from, &edge.to] {
309 if next.node(endpoint).is_none() {
310 return Err(structural(format!(
311 "edge {} references missing node {endpoint}",
312 edge.id
313 )));
314 }
315 }
316 next.edges.push(edge);
317 Ok(())
318 }
319
320 /// V9 at the write path: terminal states are never overwritten by patches;
321 /// only the explicit `Supersede` change (or a reconcile-rule change) may move
322 /// a node out of a terminal state.
323 fn patch_node(
324 next: &mut WorkGraphSnapshot,
325 id: &WorkNodeId,
326 patch: &WorkNodePatch,
327 now: i64,
328 ) -> Result<(), ValidationReport> {
329 let node = next
330 .node_mut(id)
331 .ok_or_else(|| structural(format!("node {id} not found")))?;
332 // Only an actual transition OUT of a terminal state is forbidden.
333 // Re-asserting the state a node already holds is a no-op, and rejecting it
334 // broke the only tool that writes here: `work_update` replaces the whole
335 // todo list on every call, so once an item is cancelled every later call
336 // re-sends it as cancelled and the entire update was refused. The model was
337 // then told to "use Supersede", which `work_update` does not expose — a
338 // dead end whose only escape was silently dropping the item from the list.
339 if node.state.is_terminal() && patch.state.is_some_and(|s| s != node.state) {
340 return Err(ValidationReport::single(
341 ValidationCode::V9,
342 format!(
343 "node {id} is terminal ({:?}) and cannot move to {:?}. \
344 Re-sending its current state is fine; changing it needs Supersede.",
345 node.state,
346 patch.state.expect("checked above")
347 ),
348 ));
349 }
350 if let Some(title) = &patch.title {
351 node.title = title.clone();
352 }
353 if let Some(state) = patch.state {
354 node.state = state;
355 }
356 if let Some(acceptance) = &patch.acceptance {
357 node.acceptance = acceptance.clone();
358 }
359 if let Some(provenance) = &patch.provenance {
360 node.provenance = provenance.clone();
361 }
362 node.updated_at = now;
363 Ok(())
364 }
365
366 /// Pure application of an owner observation. The owner is authoritative for
367 /// lifecycle; the graph never invents liveness:
368 /// - a missing owner maps to `Stale` — NEVER to Active or Completed;
369 /// - a terminal owner lifecycle maps to `Completed` — never `Verified`
370 /// (verification only ever comes from evidence, V4);
371 /// - nodes already in a terminal state keep it (V9); only the observation
372 /// summary is updated.
373 fn reconcile(
374 next: &mut WorkGraphSnapshot,
375 id: &WorkNodeId,
376 obs: &OperationObservation,
377 ctx: &ChangeCtx,
378 ) -> Result<(), ValidationReport> {
379 let now = ctx.now;
380 let node = next
381 .node_mut(id)
382 .ok_or_else(|| structural(format!("node {id} not found")))?;
383 let binding = node
384 .binding
385 .as_mut()
386 .ok_or_else(|| structural(format!("node {id} has no operation binding")))?;
387
388 let new_state = match obs {
389 OperationObservation::OwnerReported {
390 state,
391 seq,
392 at,
393 output,
394 } => {
395 binding.last_observation = Some(ObservationSummary {
396 owner_state: *state,
397 seq: *seq,
398 observed_at: *at,
399 output: output.clone(),
400 });
401 Some(match state {
402 OwnerState::Initializing => NodeState::Initializing,
403 OwnerState::Running => NodeState::Active,
404 OwnerState::Waiting => NodeState::Waiting,
405 // Degraded ended like Completed at the graph level — where
406 // Completed already means "ended, NOT done" and acceptance
407 // is verification's job — but the observation summary above
408 // keeps the owner's Degraded truth for every reader (#5582).
409 OwnerState::Completed | OwnerState::Degraded => NodeState::Completed,
410 OwnerState::Failed => NodeState::Failed,
411 OwnerState::Cancelled => NodeState::Cancelled,
412 })
413 }
414 OperationObservation::OwnerMissing { .. } => Some(NodeState::Stale),
415 OperationObservation::CancelUpdate { outcome, .. } => match outcome {
416 // In-flight acknowledgements: record only, no state claim yet.
417 CancelOutcome::Requested
418 | CancelOutcome::Acknowledged
419 | CancelOutcome::AlreadyFinished => None,
420 CancelOutcome::Forced => Some(NodeState::Cancelled),
421 CancelOutcome::NotFound | CancelOutcome::StaleUnknown => Some(NodeState::Stale),
422 },
423 };
424 if let Some(state) = new_state
425 && !node.state.is_terminal()
426 {
427 node.state = state;
428 }
429 node.updated_at = now;
430 Ok(())
431 }
432
433 /// Materialize evidence as an Evidence node plus a `Verifies` edge onto the
434 /// target, both with deterministically derived IDs (so the same evidence
435 /// reference attaches at most once — a repeat is a structural rejection, not
436 /// a duplicate node).
437 fn attach_evidence(
438 next: &mut WorkGraphSnapshot,
439 target: &WorkNodeId,
440 evidence: &EvidenceRef,
441 ctx: &ChangeCtx,
442 ) -> Result<(), ValidationReport> {
443 if next.node(target).is_none() {
444 return Err(structural(format!("node {target} not found")));
445 }
446 let evidence_id = WorkNodeId::derive(
447 &ctx.session_id,
448 &format!("evidence:{}:{}", target.as_str(), evidence.reference()),
449 );
450 let node = WorkNode {
451 id: evidence_id.clone(),
452 kind: NodeKind::Evidence,
453 title: format!("evidence: {}", evidence.reference()),
454 state: NodeState::Completed,
455 acceptance: Vec::new(),
456 binding: None,
457 evidence: Some(evidence.clone()),
458 provenance: Provenance::RuntimeReconcile {
459 source: "attach_evidence".to_string(),
460 observed_at: ctx.now,
461 },
462 created_at: ctx.now,
463 updated_at: ctx.now,
464 };
465 add_node(next, node)?;
466 let edge = WorkEdge {
467 id: WorkEdgeId::derive(
468 &ctx.session_id,
469 &format!("verifies:{}:{}", evidence_id.as_str(), target.as_str()),
470 ),
471 kind: EdgeKind::Verifies,
472 from: evidence_id,
473 to: target.clone(),
474 };
475 add_edge(next, edge)
476 }
477
478 /// Apply an accepted proposal atomically: nodes first (so added edges may
479 /// reference them), then edges, then patches, then removals. Any failure
480 /// rejects the whole acceptance (the caller's snapshot stays untouched).
481 fn apply_proposal(
482 next: &mut WorkGraphSnapshot,
483 proposal: &WorkGraphProposal,
484 ctx: &ChangeCtx,
485 ) -> Result<(), ValidationReport> {
486 for node in &proposal.added_nodes {
487 add_node(next, node.clone())?;
488 }
489 for edge in &proposal.added_edges {
490 add_edge(next, edge.clone())?;
491 }
492 for update in &proposal.updated_nodes {
493 patch_node(next, &update.id, &update.patch, ctx.now)?;
494 }
495 for edge_id in &proposal.removed_edges {
496 let before = next.edges.len();
497 next.edges.retain(|e| &e.id != edge_id);
498 if next.edges.len() == before {
499 return Err(structural(format!("edge {edge_id} not found")));
500 }
501 }
502 for node_id in &proposal.removed_nodes {
503 if next
504 .edges
505 .iter()
506 .any(|edge| edge.from == *node_id || edge.to == *node_id)
507 {
508 return Err(structural(format!(
509 "node {node_id} still has edges after proposed removals"
510 )));
511 }
512 let before = next.nodes.len();
513 next.nodes.retain(|node| node.id != *node_id);
514 if next.nodes.len() == before {
515 return Err(structural(format!("node {node_id} not found")));
516 }
517 }
518 if let Some(compat) = &proposal.replacement_compat {
519 next.compat.clone_from(compat);
520 }
521 Ok(())
522 }
523
524 /// Reject malformed or invariant-breaking plan edits before they become
525 /// user-reviewable. Acceptance reruns the same atomic application and the
526 /// outer reducer validation, so reviewed and accepted semantics cannot drift.
527 fn validate_proposal(
528 current: &WorkGraphSnapshot,
529 proposal: &WorkGraphProposal,
530 ctx: &ChangeCtx,
531 ) -> Result<(), ValidationReport> {
532 preview_plan_diff(current, proposal, ctx).map(|_| ())
533 }
534
535 pub(super) fn preview_plan_diff(
536 current: &WorkGraphSnapshot,
537 proposal: &WorkGraphProposal,
538 ctx: &ChangeCtx,
539 ) -> Result<WorkGraphSnapshot, ValidationReport> {
540 let mut candidate = current.clone();
541 apply_proposal(&mut candidate, proposal, ctx)?;
542 validate(&candidate)?;
543 Ok(candidate)
544 }
545
545 lines RUST