Skip to main content

graphrefly/process/
work_queue.rs

1//! Optional ProcessBundle-over-workQueue recipe (D349/D353).
2//!
3//! Queue records are mapped to graph-visible process evidence. Queue completion
4//! remains disposition evidence, not process completion or domain truth.
5
6use std::collections::{BTreeMap, BTreeSet};
7
8use crate::ctx::Ctx;
9use crate::graph::{Graph, GraphNodeOpts};
10use crate::identity::canonical_tuple_key;
11use crate::messaging::DataIssue;
12use crate::node::Node;
13use crate::process::ProcessEffectRequest;
14use crate::work_queue::{WorkQueueCommand, WorkQueueRecord};
15
16#[derive(Debug, Clone, PartialEq)]
17/// `ProcessQueuedEffectPayload` data container.
18pub struct ProcessQueuedEffectPayload<TEffect> {
19    /// `kind` field for kind.
20    pub kind: String,
21    /// `effect` field for effect.
22    pub effect: ProcessEffectRequest<TEffect>,
23    /// `idempotency_key` field for idempotency key.
24    pub idempotency_key: Option<String>,
25    /// `source_refs` field for source refs.
26    pub source_refs: Vec<String>,
27    /// `policy_refs` field for policy refs.
28    pub policy_refs: Vec<String>,
29    /// `metadata` field for metadata.
30    pub metadata: Option<String>,
31}
32
33impl<TEffect> ProcessQueuedEffectPayload<TEffect> {
34    /// Creates or computes `new`.
35    pub fn new(effect: ProcessEffectRequest<TEffect>) -> Self {
36        Self {
37            kind: "process-queued-effect".to_owned(),
38            effect,
39            idempotency_key: None,
40            source_refs: Vec::new(),
41            policy_refs: Vec::new(),
42            metadata: None,
43        }
44    }
45}
46
47#[derive(Debug, Clone, PartialEq, Eq)]
48/// `ProcessQueueEvidence` data container.
49pub struct ProcessQueueEvidence {
50    /// `kind` field for kind.
51    pub kind: String,
52    /// `evidence_id` field for evidence id.
53    pub evidence_id: String,
54    /// `effect_id` field for effect id.
55    pub effect_id: String,
56    /// `effect_type` field for effect type.
57    pub effect_type: String,
58    /// `work_id` field for work id.
59    pub work_id: String,
60    /// `queue_record_kind` field for queue record kind.
61    pub queue_record_kind: String,
62    /// `result` field for result.
63    pub result: Option<String>,
64    /// `error` field for error.
65    pub error: Option<String>,
66    /// `recorded_at_ms` field for recorded at ms.
67    pub recorded_at_ms: u64,
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
71/// `ProcessQueueStatus` data container.
72pub struct ProcessQueueStatus {
73    /// `kind` field for kind.
74    pub kind: String,
75    /// `state` field for state.
76    pub state: String,
77    /// `effect_id` field for effect id.
78    pub effect_id: Option<String>,
79    /// `effect_type` field for effect type.
80    pub effect_type: Option<String>,
81    /// `work_id` field for work id.
82    pub work_id: Option<String>,
83    /// `queue_record_kind` field for queue record kind.
84    pub queue_record_kind: Option<String>,
85    /// `evidence_id` field for evidence id.
86    pub evidence_id: Option<String>,
87    /// `issue_codes` field for issue codes.
88    pub issue_codes: Vec<String>,
89}
90
91#[derive(Debug, Clone, PartialEq, Eq)]
92/// `ProcessQueueAuditRecord` data container.
93pub struct ProcessQueueAuditRecord {
94    /// `kind` field for kind.
95    pub kind: String,
96    /// `seq` field for seq.
97    pub seq: u64,
98    /// `outcome` field for outcome.
99    pub outcome: String,
100    /// `effect_id` field for effect id.
101    pub effect_id: Option<String>,
102    /// `effect_type` field for effect type.
103    pub effect_type: Option<String>,
104    /// `work_id` field for work id.
105    pub work_id: Option<String>,
106    /// `queue_record_kind` field for queue record kind.
107    pub queue_record_kind: Option<String>,
108    /// `evidence_id` field for evidence id.
109    pub evidence_id: Option<String>,
110}
111
112#[derive(Clone)]
113/// `ProcessWorkQueueRecipeOptions` data container.
114pub struct ProcessWorkQueueRecipeOptions<TEffect> {
115    /// `name` field for name.
116    pub name: String,
117    /// `effect_requests` field for effect requests.
118    pub effect_requests: Option<Node<ProcessEffectRequest<TEffect>>>,
119    /// `records` field for records.
120    pub records: Node<WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>>,
121}
122
123impl<TEffect> ProcessWorkQueueRecipeOptions<TEffect> {
124    /// Creates or computes `new`.
125    pub fn new(records: Node<WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>>) -> Self {
126        Self {
127            name: "processWorkQueue".to_owned(),
128            effect_requests: None,
129            records,
130        }
131    }
132
133    /// Updates or reads `named`.
134    pub fn named(mut self, name: impl Into<String>) -> Self {
135        self.name = name.into();
136        self
137    }
138
139    /// Updates or reads `with_effect_requests`.
140    pub fn with_effect_requests(
141        mut self,
142        effect_requests: Node<ProcessEffectRequest<TEffect>>,
143    ) -> Self {
144        self.effect_requests = Some(effect_requests);
145        self
146    }
147}
148
149#[derive(Clone)]
150/// `ProcessWorkQueueRecipeBundle` data container.
151pub struct ProcessWorkQueueRecipeBundle<TEffect> {
152    /// `submit_commands` field for submit commands.
153    pub submit_commands: Option<Node<WorkQueueCommand<ProcessQueuedEffectPayload<TEffect>>>>,
154    /// `evidence` field for evidence.
155    pub evidence: Node<ProcessQueueEvidence>,
156    /// `status` field for status.
157    pub status: Node<ProcessQueueStatus>,
158    /// `issues` field for issues.
159    pub issues: Node<DataIssue>,
160    /// `audit` field for audit.
161    pub audit: Node<ProcessQueueAuditRecord>,
162}
163
164#[derive(Clone)]
165enum ProcessQueueFact {
166    Evidence(ProcessQueueEvidence),
167    Status(ProcessQueueStatus),
168    Issue(DataIssue),
169    Audit(ProcessQueueAuditRecord),
170}
171
172#[derive(Clone)]
173struct ProcessQueueState<TEffect> {
174    payloads: BTreeMap<String, ProcessQueuedEffectPayload<TEffect>>,
175    terminal_records: BTreeSet<String>,
176    audit_seq: u64,
177}
178
179impl<TEffect> Default for ProcessQueueState<TEffect> {
180    fn default() -> Self {
181        Self {
182            payloads: BTreeMap::new(),
183            terminal_records: BTreeSet::new(),
184            audit_seq: 0,
185        }
186    }
187}
188
189/// Creates or computes `process_work_queue_recipe`.
190pub fn process_work_queue_recipe<TEffect: Clone + 'static>(
191    graph: &Graph,
192    opts: ProcessWorkQueueRecipeOptions<TEffect>,
193) -> ProcessWorkQueueRecipeBundle<TEffect> {
194    let name = opts.name.clone();
195    let submit_commands = opts.effect_requests.as_ref().map(|effect_requests| {
196        process_effect_submit_commands(
197            graph,
198            effect_requests.clone(),
199            format!("{name}/submitCommands"),
200        )
201    });
202    let runtime = graph.node_opts::<ProcessQueueFact, _>(
203        vec![opts.records.erased()],
204        move |ctx| {
205            let mut state = ctx
206                .state_get::<ProcessQueueState<TEffect>>()
207                .map(|state| (*state).clone())
208                .unwrap_or_default();
209            for record in ctx.batch::<WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>>(0) {
210                reduce_record(ctx, &mut state, &record);
211            }
212            ctx.state_set(state);
213            ctx.state_persist(true);
214        },
215        GraphNodeOpts::named(format!("{name}/runtime")),
216    );
217    ProcessWorkQueueRecipeBundle {
218        submit_commands,
219        evidence: project(
220            graph,
221            &runtime,
222            format!("{name}/evidence"),
223            |fact| match fact {
224                ProcessQueueFact::Evidence(evidence) => Some(evidence.clone()),
225                _ => None,
226            },
227        ),
228        status: project(
229            graph,
230            &runtime,
231            format!("{name}/status"),
232            |fact| match fact {
233                ProcessQueueFact::Status(status) => Some(status.clone()),
234                _ => None,
235            },
236        ),
237        issues: project(
238            graph,
239            &runtime,
240            format!("{name}/issues"),
241            |fact| match fact {
242                ProcessQueueFact::Issue(issue) => Some(issue.clone()),
243                _ => None,
244            },
245        ),
246        audit: project(
247            graph,
248            &runtime,
249            format!("{name}/audit"),
250            |fact| match fact {
251                ProcessQueueFact::Audit(audit) => Some(audit.clone()),
252                _ => None,
253            },
254        ),
255    }
256}
257
258/// Creates or computes `process_effect_submit_command`.
259pub fn process_effect_submit_command<TEffect: Clone>(
260    effect: ProcessEffectRequest<TEffect>,
261) -> WorkQueueCommand<ProcessQueuedEffectPayload<TEffect>> {
262    let command_id = format!("{}:process-work-queue-submit", effect.id);
263    let mut payload = ProcessQueuedEffectPayload::new(effect.clone());
264    payload.idempotency_key = Some(effect.id.clone());
265    WorkQueueCommand::Submit {
266        payload,
267        command_id,
268        queue_id: None,
269        idempotency_key: Some(effect.id),
270    }
271}
272
273/// Creates or computes `process_effect_submit_commands`.
274pub fn process_effect_submit_commands<TEffect: Clone + 'static>(
275    graph: &Graph,
276    effect_requests: Node<ProcessEffectRequest<TEffect>>,
277    name: impl Into<String>,
278) -> Node<WorkQueueCommand<ProcessQueuedEffectPayload<TEffect>>> {
279    graph.node_opts::<WorkQueueCommand<ProcessQueuedEffectPayload<TEffect>>, _>(
280        vec![effect_requests.erased()],
281        move |ctx| {
282            for effect in ctx.batch::<ProcessEffectRequest<TEffect>>(0) {
283                ctx.emit(process_effect_submit_command((*effect).clone()));
284            }
285        },
286        GraphNodeOpts::named(name.into()),
287    )
288}
289
290fn reduce_record<TEffect: Clone + 'static>(
291    ctx: &Ctx,
292    state: &mut ProcessQueueState<TEffect>,
293    record: &WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>,
294) {
295    match record {
296        WorkQueueRecord::WorkAdmitted {
297            work_id, payload, ..
298        } => {
299            if payload.kind == "process-queued-effect" {
300                state.payloads.insert(work_id.clone(), payload.clone());
301            } else {
302                emit_issue(ctx, state, record, "process-queue-malformed-payload");
303            }
304        }
305        _ if is_terminal_record(record) => {
306            let key = canonical_tuple_key(&[record_kind(record), &record.record_seq().to_string()]);
307            if !state.terminal_records.insert(key) {
308                return;
309            }
310            let Some(payload) = state.payloads.get(record.work_id()).cloned() else {
311                emit_issue(ctx, state, record, "process-queue-record-without-payload");
312                return;
313            };
314            let evidence = evidence_from_record(record, &payload);
315            ctx.emit(ProcessQueueFact::Evidence(evidence.clone()));
316            ctx.emit(ProcessQueueFact::Status(status_from_evidence(
317                &evidence, &payload,
318            )));
319            state.audit_seq += 1;
320            ctx.emit(ProcessQueueFact::Audit(audit_from_evidence(
321                state.audit_seq,
322                &evidence,
323                &payload,
324            )));
325        }
326        _ => {}
327    }
328}
329
330fn is_terminal_record<TEffect>(record: &WorkQueueRecord<TEffect>) -> bool {
331    matches!(
332        record,
333        WorkQueueRecord::WorkCompleted { .. }
334            | WorkQueueRecord::AttemptCompleted { .. }
335            | WorkQueueRecord::AttemptFailed { .. }
336            | WorkQueueRecord::WorkDeadLettered { .. }
337            | WorkQueueRecord::WorkCanceled { .. }
338    )
339}
340
341fn evidence_from_record<TEffect>(
342    record: &WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>,
343    payload: &ProcessQueuedEffectPayload<TEffect>,
344) -> ProcessQueueEvidence {
345    let (result, error, recorded_at_ms) = match record {
346        WorkQueueRecord::WorkCompleted {
347            result,
348            recorded_at_ms,
349            ..
350        }
351        | WorkQueueRecord::AttemptCompleted {
352            result,
353            recorded_at_ms,
354            ..
355        } => (result.clone(), None, *recorded_at_ms),
356        WorkQueueRecord::AttemptFailed {
357            error,
358            recorded_at_ms,
359            ..
360        } => (None, error.clone(), *recorded_at_ms),
361        WorkQueueRecord::WorkDeadLettered { recorded_at_ms, .. } => {
362            (None, Some("dead-lettered".to_owned()), *recorded_at_ms)
363        }
364        WorkQueueRecord::WorkCanceled { canceled_at_ms, .. } => {
365            (None, Some("canceled".to_owned()), *canceled_at_ms)
366        }
367        _ => (None, None, 0),
368    };
369    ProcessQueueEvidence {
370        kind: "process-queue-evidence".to_owned(),
371        evidence_id: format!("work-queue:{}", record.record_seq()),
372        effect_id: payload.effect.id.clone(),
373        effect_type: payload.effect.effect_type.clone(),
374        work_id: record.work_id().to_owned(),
375        queue_record_kind: record_kind(record).to_owned(),
376        result,
377        error,
378        recorded_at_ms,
379    }
380}
381
382fn status_from_evidence<TEffect>(
383    evidence: &ProcessQueueEvidence,
384    payload: &ProcessQueuedEffectPayload<TEffect>,
385) -> ProcessQueueStatus {
386    ProcessQueueStatus {
387        kind: "process-queue-status".to_owned(),
388        state: "evidence-recorded".to_owned(),
389        effect_id: Some(payload.effect.id.clone()),
390        effect_type: Some(payload.effect.effect_type.clone()),
391        work_id: Some(evidence.work_id.clone()),
392        queue_record_kind: Some(evidence.queue_record_kind.clone()),
393        evidence_id: Some(evidence.evidence_id.clone()),
394        issue_codes: Vec::new(),
395    }
396}
397
398fn audit_from_evidence<TEffect>(
399    seq: u64,
400    evidence: &ProcessQueueEvidence,
401    payload: &ProcessQueuedEffectPayload<TEffect>,
402) -> ProcessQueueAuditRecord {
403    ProcessQueueAuditRecord {
404        kind: "process-queue-audit".to_owned(),
405        seq,
406        outcome: "mapped".to_owned(),
407        effect_id: Some(payload.effect.id.clone()),
408        effect_type: Some(payload.effect.effect_type.clone()),
409        work_id: Some(evidence.work_id.clone()),
410        queue_record_kind: Some(evidence.queue_record_kind.clone()),
411        evidence_id: Some(evidence.evidence_id.clone()),
412    }
413}
414
415fn emit_issue<TEffect: Clone + 'static>(
416    ctx: &Ctx,
417    state: &mut ProcessQueueState<TEffect>,
418    record: &WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>,
419    code: &str,
420) {
421    let issue = DataIssue {
422        kind: "issue".to_owned(),
423        code: code.to_owned(),
424        message: format!(
425            "Process workQueue recipe could not map workQueue record '{}'",
426            record_kind(record)
427        ),
428        severity: "error".to_owned(),
429        source: "process.workQueue".to_owned(),
430        topic: None,
431        details: Some(format!(
432            "work_id={};record_seq={}",
433            record.work_id(),
434            record.record_seq()
435        )),
436    };
437    ctx.emit(ProcessQueueFact::Issue(issue.clone()));
438    ctx.emit(ProcessQueueFact::Status(ProcessQueueStatus {
439        kind: "process-queue-status".to_owned(),
440        state: "mapping-issue".to_owned(),
441        effect_id: None,
442        effect_type: None,
443        work_id: Some(record.work_id().to_owned()),
444        queue_record_kind: Some(record_kind(record).to_owned()),
445        evidence_id: None,
446        issue_codes: vec![issue.code.clone()],
447    }));
448    state.audit_seq += 1;
449    ctx.emit(ProcessQueueFact::Audit(ProcessQueueAuditRecord {
450        kind: "process-queue-audit".to_owned(),
451        seq: state.audit_seq,
452        outcome: "issue".to_owned(),
453        effect_id: None,
454        effect_type: None,
455        work_id: Some(record.work_id().to_owned()),
456        queue_record_kind: Some(record_kind(record).to_owned()),
457        evidence_id: None,
458    }));
459}
460
461fn project<TIn: Clone + 'static, TOut: 'static>(
462    graph: &Graph,
463    source: &Node<TIn>,
464    name: String,
465    pick: impl Fn(&TIn) -> Option<TOut> + 'static,
466) -> Node<TOut> {
467    graph.node_opts::<TOut, _>(
468        vec![source.erased()],
469        move |ctx| {
470            for fact in ctx.batch::<TIn>(0) {
471                if let Some(value) = pick(&fact) {
472                    ctx.emit(value);
473                }
474            }
475        },
476        GraphNodeOpts::named(name),
477    )
478}
479
480fn record_kind<TEffect>(record: &WorkQueueRecord<TEffect>) -> &'static str {
481    match record {
482        WorkQueueRecord::WorkAdmitted { .. } => "work-admitted",
483        WorkQueueRecord::AdmissionDeduped { .. } => "admission-deduped",
484        WorkQueueRecord::WorkScheduled { .. } => "work-scheduled",
485        WorkQueueRecord::WorkClaimed { .. } => "work-claimed",
486        WorkQueueRecord::LeaseRenewed { .. } => "lease-renewed",
487        WorkQueueRecord::WorkReleased { .. } => "work-released",
488        WorkQueueRecord::LeaseExpired { .. } => "lease-expired",
489        WorkQueueRecord::AttemptCompleted { .. } => "attempt-completed",
490        WorkQueueRecord::WorkCompleted { .. } => "work-completed",
491        WorkQueueRecord::AttemptFailed { .. } => "attempt-failed",
492        WorkQueueRecord::RetryScheduled { .. } => "retry-scheduled",
493        WorkQueueRecord::WorkDeadLettered { .. } => "work-dead-lettered",
494        WorkQueueRecord::WorkCanceled { .. } => "work-canceled",
495    }
496}