返回 CodeWhale
events.rs
根目录 / crates / tui / src / work_graph / events.rs
1 //! Changes, observations, and receipts — the reducer's entire input/output
2 //! vocabulary.
3 //!
4 //! The reducer never reads clocks or RNG: [`ChangeCtx`] carries the timestamp,
5 //! the session identity used for deterministic ID derivation, and an optional
6 //! idempotency key. Same snapshot + same change + same ctx ⇒ same result,
7 //! always.
8 //!
9 //! Spec-silent shapes chosen minimally here (documented on each type):
10 //! [`WorkNodePatch`], [`WorkGraphProposal`], [`ApprovalRef`], and the
11 //! placeholder observation types that later slices' liveness adapters will
12 //! feed ([`OperationObservation`], [`OwnerState`], [`CancelOutcome`]).
13
14 use serde::{Deserialize, Serialize};
15
16 use super::ids::{ChangeId, ProposalId, WorkEdgeId, WorkNodeId};
17 use super::model::{
18 AcceptanceRequirement, CompatProjectionState, EvidenceRef, IdempotencyKey, NodeState,
19 OperationBinding, Provenance, Ts, WorkActivityEvent, WorkEdge, WorkNode,
20 };
21
22 /// A single mutation of the work graph. The reducer is the only write path;
23 /// UI, tools, and runtime adapters all speak this vocabulary.
24 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
25 #[serde(rename_all = "snake_case")]
26 // Variants intentionally carry full payloads (a node, a proposal) rather than
27 // boxed indirection: changes are transient values, not stored long-term.
28 #[allow(clippy::large_enum_variant)]
29 pub enum WorkGraphChange {
30 AddNode {
31 node: WorkNode,
32 },
33 UpdateNode {
34 id: WorkNodeId,
35 patch: WorkNodePatch,
36 },
37 AddEdge {
38 edge: WorkEdge,
39 },
40 RemoveEdge {
41 id: WorkEdgeId,
42 },
43 BindOperation {
44 node: WorkNodeId,
45 binding: OperationBinding,
46 },
47 ReconcileOperation {
48 node: WorkNodeId,
49 obs: OperationObservation,
50 },
51 AttachEvidence {
52 node: WorkNodeId,
53 evidence: EvidenceRef,
54 },
55 ProposePlanDiff {
56 proposal: WorkGraphProposal,
57 },
58 /// Explicitly retire a pending proposal before a replacement is
59 /// proposed. This keeps repeated "Revise plan" turns reviewable without
60 /// silently rewriting or accumulating stale proposals.
61 WithdrawPlanDiff {
62 proposal_id: ProposalId,
63 },
64 AcceptPlanDiff {
65 proposal_id: ProposalId,
66 approval: ApprovalRef,
67 },
68 Supersede {
69 old: WorkNodeId,
70 replacement: WorkNodeId,
71 },
72 /// Atomically replace the inputs for the legacy Plan/To-do projections.
73 ReplaceCompatProjection {
74 compat: CompatProjectionState,
75 },
76 /// Record the canonical digest of a completed legacy import.
77 SetImportDigest {
78 digest: String,
79 },
80 /// Append one bounded, configuration-only activity receipt.
81 RecordActivity {
82 event: WorkActivityEvent,
83 },
84 /// Evict the oldest ended, non-durable Operation nodes beyond `keep`
85 /// (#6842). See `reducer::prune_ended_operations` for eligibility.
86 PruneEndedOperations {
87 keep: usize,
88 },
89 }
90
91 impl WorkGraphChange {
92 /// Stable discriminant name recorded on receipts. Names only — receipts
93 /// never carry payload text.
94 #[must_use]
95 pub fn kind_name(&self) -> &'static str {
96 match self {
97 WorkGraphChange::AddNode { .. } => "add_node",
98 WorkGraphChange::UpdateNode { .. } => "update_node",
99 WorkGraphChange::AddEdge { .. } => "add_edge",
100 WorkGraphChange::RemoveEdge { .. } => "remove_edge",
101 WorkGraphChange::BindOperation { .. } => "bind_operation",
102 WorkGraphChange::ReconcileOperation { .. } => "reconcile_operation",
103 WorkGraphChange::AttachEvidence { .. } => "attach_evidence",
104 WorkGraphChange::ProposePlanDiff { .. } => "propose_plan_diff",
105 WorkGraphChange::WithdrawPlanDiff { .. } => "withdraw_plan_diff",
106 WorkGraphChange::AcceptPlanDiff { .. } => "accept_plan_diff",
107 WorkGraphChange::Supersede { .. } => "supersede",
108 WorkGraphChange::ReplaceCompatProjection { .. } => "replace_compat_projection",
109 WorkGraphChange::SetImportDigest { .. } => "set_import_digest",
110 WorkGraphChange::RecordActivity { .. } => "record_activity",
111 WorkGraphChange::PruneEndedOperations { .. } => "prune_ended_operations",
112 }
113 }
114 }
115
116 /// Partial update of a node. Spec-silent shape: `Option` per patchable field,
117 /// `None` meaning "leave unchanged". Identity, kind, binding, and evidence
118 /// are deliberately NOT patchable here — they move only through their
119 /// dedicated changes (`BindOperation`, `AttachEvidence`, `Supersede`).
120 #[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
121 pub struct WorkNodePatch {
122 pub title: Option<String>,
123 pub state: Option<NodeState>,
124 pub acceptance: Option<Vec<AcceptanceRequirement>>,
125 pub provenance: Option<Provenance>,
126 }
127
128 /// A reviewable plan diff. Spec-silent shape: explicit added/updated/removed
129 /// sets rather than nested changes, so the whole delta is inspectable before
130 /// acceptance and applies atomically (validated as one unit — no silent
131 /// mutation of objectives, dependencies, acceptance, or scope).
132 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
133 pub struct WorkGraphProposal {
134 pub id: ProposalId,
135 #[serde(default)]
136 pub added_nodes: Vec<WorkNode>,
137 #[serde(default)]
138 pub added_edges: Vec<WorkEdge>,
139 #[serde(default)]
140 pub updated_nodes: Vec<ProposedNodeUpdate>,
141 #[serde(default)]
142 pub removed_nodes: Vec<WorkNodeId>,
143 #[serde(default)]
144 pub removed_edges: Vec<WorkEdgeId>,
145 /// Graph-owned inputs for the legacy Plan/To-do projections. This is part
146 /// of the reviewed scope delta and is applied atomically with the graph
147 /// changes. Older saved proposals deserialize with no replacement.
148 #[serde(default, skip_serializing_if = "Option::is_none")]
149 pub replacement_compat: Option<CompatProjectionState>,
150 }
151
152 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
153 pub struct ProposedNodeUpdate {
154 pub id: WorkNodeId,
155 pub patch: WorkNodePatch,
156 }
157
158 /// Reference to the approval that accepted a plan diff. Spec-silent shape:
159 /// a reference-only string (approval receipt / user action handle), recorded
160 /// on the Approval node the acceptance creates.
161 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
162 pub struct ApprovalRef {
163 pub reference: String,
164 }
165
166 /// Lifecycle state as reported by an operation's owner. Placeholder for the
167 /// liveness slice; present now so reducer signatures are final.
168 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
169 #[serde(rename_all = "snake_case")]
170 pub enum OwnerState {
171 Initializing,
172 Running,
173 Waiting,
174 Completed,
175 /// Terminal: finished and produced output, but with dropped or failed
176 /// slots the operator must be able to distinguish from an ordinary
177 /// success (#5582). Collapsing this into `Completed` let dashboards and
178 /// automation treat a partial workflow as fully accepted.
179 Degraded,
180 Failed,
181 Cancelled,
182 }
183
184 /// Typed cancellation outcomes, mirroring real owner semantics (immediate
185 /// abort vs teardown-wait vs already-finished vs unknown-after-restart).
186 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
187 #[serde(rename_all = "snake_case")]
188 pub enum CancelOutcome {
189 Requested,
190 Acknowledged,
191 Forced,
192 AlreadyFinished,
193 NotFound,
194 StaleUnknown,
195 }
196
197 /// An observation about a bound operation, produced by owner adapters (later
198 /// slice) and consumed by the reducer. The reducer applies these purely; it
199 /// never queries owners itself.
200 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
201 #[serde(rename_all = "snake_case")]
202 pub enum OperationObservation {
203 /// Owner is authoritative. Idempotency key = `(binding, seq)`.
204 OwnerReported {
205 state: OwnerState,
206 seq: u64,
207 at: Ts,
208 /// Bounded logical output receipt. The reference never contains raw
209 /// logs or reasoning; `raw_bytes` preserves the pre-truncation size.
210 #[serde(default, skip_serializing_if = "Option::is_none")]
211 output: Option<EvidenceRef>,
212 },
213 /// No live handle for the binding (e.g. after restart). Never maps to
214 /// Active or Completed — only to Stale (fail toward honesty).
215 OwnerMissing {
216 checked_at: Ts,
217 },
218 CancelUpdate {
219 outcome: CancelOutcome,
220 at: Ts,
221 },
222 }
223
224 /// Compact record of the most recent observation, stored on the binding.
225 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
226 pub struct ObservationSummary {
227 pub owner_state: OwnerState,
228 pub seq: u64,
229 pub observed_at: Ts,
230 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub output: Option<EvidenceRef>,
232 }
233
234 /// Everything ambient the reducer needs, supplied by the caller so the
235 /// reducer itself stays pure: no clock reads, no RNG, no globals.
236 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
237 pub struct ChangeCtx {
238 /// Session identity used for deterministic ID derivation.
239 pub session_id: String,
240 /// Timestamp to record for this change (milliseconds since Unix epoch).
241 pub now: Ts,
242 /// Present for owner-observation changes; duplicates inside the
243 /// snapshot's dedup window become no-op receipts.
244 pub idempotency_key: Option<IdempotencyKey>,
245 }
246
247 /// Receipt for an applied (or deduplicated) change. Bounded history of these
248 /// lives on the snapshot. Receipts carry discriminant names and identifiers
249 /// only — no payload text, no secrets.
250 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
251 pub struct ChangeReceipt {
252 pub change_id: ChangeId,
253 pub revision: u64,
254 pub summary: String,
255 pub applied_at: Ts,
256 pub idempotency_key: Option<IdempotencyKey>,
257 pub no_op: bool,
258 }
259
260 impl ChangeReceipt {
261 #[must_use]
262 pub fn of(change: &WorkGraphChange, revision: u64, ctx: &ChangeCtx) -> Self {
263 let kind = change.kind_name();
264 ChangeReceipt {
265 change_id: ChangeId::derive(&ctx.session_id, &format!("change:{revision}:{kind}")),
266 revision,
267 summary: kind.to_string(),
268 applied_at: ctx.now,
269 idempotency_key: ctx.idempotency_key.clone(),
270 no_op: false,
271 }
272 }
273 }
274
274 lines RUST