1use 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)]
17pub struct ProcessQueuedEffectPayload<TEffect> {
19 pub kind: String,
21 pub effect: ProcessEffectRequest<TEffect>,
23 pub idempotency_key: Option<String>,
25 pub source_refs: Vec<String>,
27 pub policy_refs: Vec<String>,
29 pub metadata: Option<String>,
31}
32
33impl<TEffect> ProcessQueuedEffectPayload<TEffect> {
34 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)]
48pub struct ProcessQueueEvidence {
50 pub kind: String,
52 pub evidence_id: String,
54 pub effect_id: String,
56 pub effect_type: String,
58 pub work_id: String,
60 pub queue_record_kind: String,
62 pub result: Option<String>,
64 pub error: Option<String>,
66 pub recorded_at_ms: u64,
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct ProcessQueueStatus {
73 pub kind: String,
75 pub state: String,
77 pub effect_id: Option<String>,
79 pub effect_type: Option<String>,
81 pub work_id: Option<String>,
83 pub queue_record_kind: Option<String>,
85 pub evidence_id: Option<String>,
87 pub issue_codes: Vec<String>,
89}
90
91#[derive(Debug, Clone, PartialEq, Eq)]
92pub struct ProcessQueueAuditRecord {
94 pub kind: String,
96 pub seq: u64,
98 pub outcome: String,
100 pub effect_id: Option<String>,
102 pub effect_type: Option<String>,
104 pub work_id: Option<String>,
106 pub queue_record_kind: Option<String>,
108 pub evidence_id: Option<String>,
110}
111
112#[derive(Clone)]
113pub struct ProcessWorkQueueRecipeOptions<TEffect> {
115 pub name: String,
117 pub effect_requests: Option<Node<ProcessEffectRequest<TEffect>>>,
119 pub records: Node<WorkQueueRecord<ProcessQueuedEffectPayload<TEffect>>>,
121}
122
123impl<TEffect> ProcessWorkQueueRecipeOptions<TEffect> {
124 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 pub fn named(mut self, name: impl Into<String>) -> Self {
135 self.name = name.into();
136 self
137 }
138
139 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)]
150pub struct ProcessWorkQueueRecipeBundle<TEffect> {
152 pub submit_commands: Option<Node<WorkQueueCommand<ProcessQueuedEffectPayload<TEffect>>>>,
154 pub evidence: Node<ProcessQueueEvidence>,
156 pub status: Node<ProcessQueueStatus>,
158 pub issues: Node<DataIssue>,
160 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
189pub 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
258pub 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
273pub 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}