Skip to main content

graphrefly/
scheduled_readiness.rs

1//! Graph-visible scheduled readiness projector (D424/D432/D433).
2//!
3//! The projector is intentionally passive: explicit schedule facts plus explicit
4//! clock facts produce eligibility/deadline visibility. It never runs timers,
5//! claims work, executes providers, or mutates domain lifecycle records.
6
7use std::cell::RefCell;
8use std::collections::{BTreeMap, BTreeSet};
9use std::rc::Rc;
10
11use serde_json::{Map, Value};
12
13use crate::ctx::Ctx;
14use crate::graph::{Graph, GraphNodeOpts};
15use crate::identity::{canonical_tuple_key, compound_tuple_key};
16use crate::json::{stable_json_string, JsonValue};
17use crate::messaging::DataIssue;
18use crate::node::{Node, NodeOpts};
19
20#[derive(Debug, Clone, PartialEq)]
21/// `SourceRef` data container.
22pub struct SourceRef {
23    /// `kind` field for kind.
24    pub kind: String,
25    /// `id` field for id.
26    pub id: String,
27    /// `metadata` field for metadata.
28    pub metadata: Option<BTreeMap<String, JsonValue>>,
29}
30
31impl SourceRef {
32    /// Creates or computes `new`.
33    pub fn new(kind: impl Into<String>, id: impl Into<String>) -> Self {
34        Self {
35            kind: kind.into(),
36            id: id.into(),
37            metadata: None,
38        }
39    }
40}
41
42#[derive(Debug, Clone, PartialEq)]
43/// `ScheduledReadinessRequested` data container.
44pub struct ScheduledReadinessRequested {
45    /// `schedule_id` field for schedule id.
46    pub schedule_id: String,
47    /// `subject_refs` field for subject refs.
48    pub subject_refs: Vec<SourceRef>,
49    /// `ready_at_ms` field for ready at ms.
50    pub ready_at_ms: u64,
51    /// `deadline_ms` field for deadline ms.
52    pub deadline_ms: Option<u64>,
53    /// `reason` field for reason.
54    pub reason: Option<String>,
55    /// `policy_refs` field for policy refs.
56    pub policy_refs: Vec<SourceRef>,
57    /// `source_refs` field for source refs.
58    pub source_refs: Vec<SourceRef>,
59    /// `metadata` field for metadata.
60    pub metadata: Option<BTreeMap<String, JsonValue>>,
61}
62
63#[derive(Debug, Clone, PartialEq)]
64/// `ScheduledReadinessClock` data container.
65pub struct ScheduledReadinessClock {
66    /// `clock_id` field for clock id.
67    pub clock_id: String,
68    /// `now_ms` field for now ms.
69    pub now_ms: u64,
70    /// `source_refs` field for source refs.
71    pub source_refs: Vec<SourceRef>,
72    /// `metadata` field for metadata.
73    pub metadata: Option<BTreeMap<String, JsonValue>>,
74}
75
76#[derive(Debug, Clone, PartialEq)]
77/// `ScheduledReadinessPending` data container.
78pub struct ScheduledReadinessPending {
79    /// `schedule_id` field for schedule id.
80    pub schedule_id: String,
81    /// `subject_refs` field for subject refs.
82    pub subject_refs: Vec<SourceRef>,
83    /// `ready_at_ms` field for ready at ms.
84    pub ready_at_ms: u64,
85    /// `deadline_ms` field for deadline ms.
86    pub deadline_ms: Option<u64>,
87    /// `now_ms` field for now ms.
88    pub now_ms: Option<u64>,
89    /// `source_refs` field for source refs.
90    pub source_refs: Vec<SourceRef>,
91    /// `metadata` field for metadata.
92    pub metadata: Option<BTreeMap<String, JsonValue>>,
93}
94
95#[derive(Debug, Clone, PartialEq)]
96/// `ScheduledReadinessReady` data container.
97pub struct ScheduledReadinessReady {
98    /// `schedule_id` field for schedule id.
99    pub schedule_id: String,
100    /// `subject_refs` field for subject refs.
101    pub subject_refs: Vec<SourceRef>,
102    /// `ready_at_ms` field for ready at ms.
103    pub ready_at_ms: u64,
104    /// `deadline_ms` field for deadline ms.
105    pub deadline_ms: Option<u64>,
106    /// `now_ms` field for now ms.
107    pub now_ms: u64,
108    /// `source_refs` field for source refs.
109    pub source_refs: Vec<SourceRef>,
110    /// `metadata` field for metadata.
111    pub metadata: Option<BTreeMap<String, JsonValue>>,
112}
113
114#[derive(Debug, Clone, PartialEq)]
115/// `ScheduledReadinessOverdue` data container.
116pub struct ScheduledReadinessOverdue {
117    /// `schedule_id` field for schedule id.
118    pub schedule_id: String,
119    /// `subject_refs` field for subject refs.
120    pub subject_refs: Vec<SourceRef>,
121    /// `ready_at_ms` field for ready at ms.
122    pub ready_at_ms: u64,
123    /// `deadline_ms` field for deadline ms.
124    pub deadline_ms: u64,
125    /// `now_ms` field for now ms.
126    pub now_ms: u64,
127    /// `source_refs` field for source refs.
128    pub source_refs: Vec<SourceRef>,
129    /// `metadata` field for metadata.
130    pub metadata: Option<BTreeMap<String, JsonValue>>,
131}
132
133#[derive(Debug, Clone, Copy, PartialEq, Eq)]
134/// `ScheduledReadinessStatusState` variants.
135pub enum ScheduledReadinessStatusState {
136    /// `Pending` variant.
137    Pending,
138    /// `Ready` variant.
139    Ready,
140    /// `Overdue` variant.
141    Overdue,
142    /// `Issue` variant.
143    Issue,
144}
145
146#[derive(Debug, Clone, PartialEq)]
147/// `ScheduledReadinessStatus` data container.
148pub struct ScheduledReadinessStatus {
149    /// `status_id` field for status id.
150    pub status_id: String,
151    /// `schedule_id` field for schedule id.
152    pub schedule_id: String,
153    /// `state` field for state.
154    pub state: ScheduledReadinessStatusState,
155    /// `subject_refs` field for subject refs.
156    pub subject_refs: Vec<SourceRef>,
157    /// `ready_at_ms` field for ready at ms.
158    pub ready_at_ms: Option<u64>,
159    /// `deadline_ms` field for deadline ms.
160    pub deadline_ms: Option<u64>,
161    /// `now_ms` field for now ms.
162    pub now_ms: Option<u64>,
163    /// `source_refs` field for source refs.
164    pub source_refs: Vec<SourceRef>,
165    /// `issue_codes` field for issue codes.
166    pub issue_codes: Vec<String>,
167    /// `metadata` field for metadata.
168    pub metadata: Option<BTreeMap<String, JsonValue>>,
169}
170
171#[derive(Debug, Clone, PartialEq)]
172/// `ScheduledReadinessAuditRecord` data container.
173pub struct ScheduledReadinessAuditRecord {
174    /// `id` field for id.
175    pub id: String,
176    /// `kind` field for kind.
177    pub kind: String,
178    /// `subject_id` field for subject id.
179    pub subject_id: Option<String>,
180    /// `source_refs` field for source refs.
181    pub source_refs: Vec<SourceRef>,
182    /// `metadata` field for metadata.
183    pub metadata: Option<BTreeMap<String, JsonValue>>,
184}
185
186#[derive(Debug, Clone, PartialEq, Default)]
187/// `ScheduledReadinessViews` data container.
188pub struct ScheduledReadinessViews {
189    /// `schedules_by_id` field for schedules by id.
190    pub schedules_by_id: BTreeMap<String, ScheduledReadinessRequested>,
191    /// `pending_by_id` field for pending by id.
192    pub pending_by_id: BTreeMap<String, ScheduledReadinessPending>,
193    /// `ready_by_id` field for ready by id.
194    pub ready_by_id: BTreeMap<String, ScheduledReadinessReady>,
195    /// `overdue_by_id` field for overdue by id.
196    pub overdue_by_id: BTreeMap<String, ScheduledReadinessOverdue>,
197    /// `status_by_id` field for status by id.
198    pub status_by_id: BTreeMap<String, ScheduledReadinessStatus>,
199    /// `now_ms` field for now ms.
200    pub now_ms: Option<u64>,
201}
202
203#[derive(Clone)]
204/// `ScheduledReadinessBundle` data container.
205pub struct ScheduledReadinessBundle {
206    /// `pending` field for pending.
207    pub pending: Node<ScheduledReadinessPending>,
208    /// `ready` field for ready.
209    pub ready: Node<ScheduledReadinessReady>,
210    /// `overdue` field for overdue.
211    pub overdue: Node<ScheduledReadinessOverdue>,
212    /// `status` field for status.
213    pub status: Node<ScheduledReadinessStatus>,
214    /// `issues` field for issues.
215    pub issues: Node<DataIssue>,
216    /// `audit` field for audit.
217    pub audit: Node<ScheduledReadinessAuditRecord>,
218    /// `views` field for views.
219    pub views: Node<ScheduledReadinessViews>,
220}
221
222#[derive(Clone)]
223/// `ScheduledReadinessOptions` data container.
224pub struct ScheduledReadinessOptions {
225    /// `name` field for name.
226    pub name: Option<String>,
227    /// `schedules` field for schedules.
228    pub schedules: Vec<Node<ScheduledReadinessRequested>>,
229    /// `clocks` field for clocks.
230    pub clocks: Vec<Node<ScheduledReadinessClock>>,
231}
232
233impl ScheduledReadinessOptions {
234    /// Creates or computes `new`.
235    pub fn new(schedules: Vec<Node<ScheduledReadinessRequested>>) -> Self {
236        Self {
237            name: None,
238            schedules,
239            clocks: Vec::new(),
240        }
241    }
242
243    /// Updates or reads `named`.
244    pub fn named(mut self, name: impl Into<String>) -> Self {
245        self.name = Some(name.into());
246        self
247    }
248
249    /// Updates or reads `with_clocks`.
250    pub fn with_clocks(mut self, clocks: Vec<Node<ScheduledReadinessClock>>) -> Self {
251        self.clocks = clocks;
252        self
253    }
254}
255
256#[derive(Clone)]
257enum ScheduledReadinessFact {
258    Pending(ScheduledReadinessPending),
259    Ready(ScheduledReadinessReady),
260    Overdue(ScheduledReadinessOverdue),
261    Status(ScheduledReadinessStatus),
262    Issue(DataIssue),
263    Audit(ScheduledReadinessAuditRecord),
264    Views(ScheduledReadinessViews),
265}
266
267#[derive(Default)]
268struct ScheduledReadinessState {
269    schedules: BTreeMap<String, ScheduledReadinessRequested>,
270    pending_by_id: BTreeMap<String, ScheduledReadinessPending>,
271    ready_by_id: BTreeMap<String, ScheduledReadinessReady>,
272    overdue_by_id: BTreeMap<String, ScheduledReadinessOverdue>,
273    status_by_id: BTreeMap<String, ScheduledReadinessStatus>,
274    emitted_keys: BTreeSet<String>,
275    issue_keys: BTreeSet<String>,
276    audit_seq: u64,
277    now_ms: Option<u64>,
278    clock_source_refs: Vec<SourceRef>,
279}
280
281/// Creates or computes `scheduled_readiness_projector`.
282pub fn scheduled_readiness_projector(
283    graph: &Graph,
284    opts: ScheduledReadinessOptions,
285) -> ScheduledReadinessBundle {
286    let name = opts
287        .name
288        .clone()
289        .unwrap_or_else(|| "scheduledReadiness".to_owned());
290    let schedule_count = opts.schedules.len();
291    let mut deps = Vec::with_capacity(opts.schedules.len() + opts.clocks.len());
292    deps.extend(opts.schedules.iter().map(Node::erased));
293    deps.extend(opts.clocks.iter().map(Node::erased));
294    let state = Rc::new(RefCell::new(ScheduledReadinessState::default()));
295    let runtime = graph.node_opts::<ScheduledReadinessFact, _>(
296        deps,
297        {
298            let state = state.clone();
299            move |ctx| {
300                let mut state = state.borrow_mut();
301                for index in 0..schedule_count {
302                    for schedule in ctx.batch::<ScheduledReadinessRequested>(index) {
303                        retain_schedule(ctx, &mut state, (*schedule).clone());
304                    }
305                }
306                for index in schedule_count..(schedule_count + opts.clocks.len()) {
307                    for clock in ctx.batch::<ScheduledReadinessClock>(index) {
308                        retain_clock(ctx, &mut state, (*clock).clone());
309                    }
310                }
311                evaluate_schedules(ctx, &mut state);
312                emit_fact(
313                    ctx,
314                    ScheduledReadinessFact::Views(ScheduledReadinessViews {
315                        schedules_by_id: state.schedules.clone(),
316                        pending_by_id: state.pending_by_id.clone(),
317                        ready_by_id: state.ready_by_id.clone(),
318                        overdue_by_id: state.overdue_by_id.clone(),
319                        status_by_id: state.status_by_id.clone(),
320                        now_ms: state.now_ms,
321                    }),
322                );
323            }
324        },
325        {
326            let mut node_opts = node_opts(format!("{name}/runtime"), "scheduledReadinessProjector");
327            node_opts.node.partial = true;
328            node_opts
329        },
330    );
331    ScheduledReadinessBundle {
332        pending: project_fact(
333            graph,
334            &runtime,
335            format!("{name}/pending"),
336            "scheduledReadinessPending",
337            |fact| match fact {
338                ScheduledReadinessFact::Pending(value) => Some(value.clone()),
339                _ => None,
340            },
341        ),
342        ready: project_fact(
343            graph,
344            &runtime,
345            format!("{name}/ready"),
346            "scheduledReadinessReady",
347            |fact| match fact {
348                ScheduledReadinessFact::Ready(value) => Some(value.clone()),
349                _ => None,
350            },
351        ),
352        overdue: project_fact(
353            graph,
354            &runtime,
355            format!("{name}/overdue"),
356            "scheduledReadinessOverdue",
357            |fact| match fact {
358                ScheduledReadinessFact::Overdue(value) => Some(value.clone()),
359                _ => None,
360            },
361        ),
362        status: project_fact(
363            graph,
364            &runtime,
365            format!("{name}/status"),
366            "scheduledReadinessStatus",
367            |fact| match fact {
368                ScheduledReadinessFact::Status(value) => Some(value.clone()),
369                _ => None,
370            },
371        ),
372        issues: project_fact(
373            graph,
374            &runtime,
375            format!("{name}/issues"),
376            "scheduledReadinessIssues",
377            |fact| match fact {
378                ScheduledReadinessFact::Issue(value) => Some(value.clone()),
379                _ => None,
380            },
381        ),
382        audit: project_fact(
383            graph,
384            &runtime,
385            format!("{name}/audit"),
386            "scheduledReadinessAudit",
387            |fact| match fact {
388                ScheduledReadinessFact::Audit(value) => Some(value.clone()),
389                _ => None,
390            },
391        ),
392        views: project_fact(
393            graph,
394            &runtime,
395            format!("{name}/views"),
396            "scheduledReadinessViews",
397            |fact| match fact {
398                ScheduledReadinessFact::Views(value) => Some(value.clone()),
399                _ => None,
400            },
401        ),
402    }
403}
404
405/// Creates or computes `parse_scheduled_readiness_requested`.
406pub fn parse_scheduled_readiness_requested(
407    value: &JsonValue,
408) -> Result<ScheduledReadinessRequested, Box<DataIssue>> {
409    let Some(object) = value.as_object() else {
410        return Err(Box::new(readiness_issue(
411            "scheduled-readiness-malformed-schedule",
412            "Scheduled readiness requires scheduleId, subjectRefs, and readyAtMs.",
413            "unknown-scheduled-readiness",
414            None,
415            "error",
416        )));
417    };
418    let schedule_id = string_field(object, "scheduleId")
419        .filter(|value| !value.is_empty())
420        .unwrap_or("unknown-scheduled-readiness");
421    if object.contains_key("subjectRef") || object.contains_key("notBeforeMs") {
422        return Err(Box::new(readiness_issue(
423            "scheduled-readiness-malformed-schedule",
424            "Scheduled readiness v1 rejects stale subjectRef/notBeforeMs aliases.",
425            schedule_id,
426            None,
427            "error",
428        )));
429    }
430    let ready_at_ms = u64_field(object, "readyAtMs");
431    let deadline_ms = optional_u64_field(object, "deadlineMs");
432    let subject_refs = source_refs_field(object, "subjectRefs");
433    if object.get("kind").and_then(Value::as_str) != Some("scheduled-readiness-requested")
434        || string_field(object, "scheduleId").is_none_or(str::is_empty)
435        || object
436            .get("subjectRefs")
437            .is_none_or(|value| !value.is_array())
438        || ready_at_ms.is_none()
439        || deadline_ms.is_err()
440    {
441        return Err(Box::new(readiness_issue(
442            "scheduled-readiness-malformed-schedule",
443            "Scheduled readiness requires scheduleId, subjectRefs, and finite readyAtMs.",
444            schedule_id,
445            None,
446            "error",
447        )));
448    }
449    Ok(ScheduledReadinessRequested {
450        schedule_id: schedule_id.to_owned(),
451        subject_refs,
452        ready_at_ms: ready_at_ms.expect("checked readyAtMs"),
453        deadline_ms: deadline_ms.expect("checked deadlineMs"),
454        reason: string_field(object, "reason").map(str::to_owned),
455        policy_refs: source_refs_field(object, "policyRefs"),
456        source_refs: source_refs_field(object, "sourceRefs"),
457        metadata: object_metadata(object.get("metadata")),
458    })
459}
460
461fn retain_schedule(
462    ctx: &Ctx,
463    state: &mut ScheduledReadinessState,
464    schedule: ScheduledReadinessRequested,
465) {
466    let schedule = sanitize_schedule(schedule);
467    let schedule_id = schedule.schedule_id.clone();
468    match state.schedules.get(&schedule_id) {
469        None => {
470            state.schedules.insert(schedule_id, schedule);
471        }
472        Some(existing) if schedule_identity(existing) == schedule_identity(&schedule) => {}
473        Some(existing) => {
474            let existing = existing.clone();
475            let issue = readiness_issue(
476                "scheduled-readiness-schedule-conflict",
477                "Scheduled readiness scheduleId was replayed with conflicting material; first valid schedule retained.",
478                &schedule_id,
479                Some(format!(
480                    "existingReadyAtMs={}; incomingReadyAtMs={}",
481                    existing.ready_at_ms, schedule.ready_at_ms
482                )),
483                "error",
484            );
485            emit_issue(ctx, state, issue.clone());
486            emit_status(
487                ctx,
488                state,
489                ScheduledReadinessStatus {
490                    status_id: format!("{schedule_id}:scheduled-readiness-status:issue"),
491                    schedule_id,
492                    state: ScheduledReadinessStatusState::Issue,
493                    subject_refs: existing.subject_refs.clone(),
494                    ready_at_ms: Some(existing.ready_at_ms),
495                    deadline_ms: existing.deadline_ms,
496                    now_ms: state.now_ms,
497                    source_refs: schedule_source_refs(&existing, &state.clock_source_refs),
498                    issue_codes: vec![issue.code],
499                    metadata: None,
500                },
501            );
502        }
503    }
504}
505
506fn retain_clock(ctx: &Ctx, state: &mut ScheduledReadinessState, clock: ScheduledReadinessClock) {
507    if let Some(previous) = state.now_ms {
508        if clock.now_ms < previous {
509            emit_issue(
510                ctx,
511                state,
512                readiness_issue(
513                    "scheduled-readiness-clock-rollback",
514                    "Scheduled readiness clock facts must be monotonic; rollback was ignored.",
515                    &clock.clock_id,
516                    Some(format!("nowMs={}; previousNowMs={previous}", clock.now_ms)),
517                    "warning",
518                ),
519            );
520            return;
521        }
522    }
523    state.now_ms = Some(clock.now_ms);
524    state.clock_source_refs = canonical_source_refs(clock.source_refs);
525}
526
527fn evaluate_schedules(ctx: &Ctx, state: &mut ScheduledReadinessState) {
528    for schedule in state.schedules.values().cloned().collect::<Vec<_>>() {
529        let source_refs = schedule_source_refs(&schedule, &state.clock_source_refs);
530        let metadata = readiness_metadata(&schedule);
531        if state
532            .now_ms
533            .is_none_or(|now_ms| now_ms < schedule.ready_at_ms)
534        {
535            let pending = ScheduledReadinessPending {
536                schedule_id: schedule.schedule_id.clone(),
537                subject_refs: canonical_source_refs(schedule.subject_refs.clone()),
538                ready_at_ms: schedule.ready_at_ms,
539                deadline_ms: schedule.deadline_ms,
540                now_ms: state.now_ms,
541                source_refs: source_refs.clone(),
542                metadata: metadata.clone(),
543            };
544            emit_pending(ctx, state, pending);
545            emit_status(
546                ctx,
547                state,
548                status_for(
549                    &schedule,
550                    ScheduledReadinessStatusState::Pending,
551                    source_refs,
552                    state.now_ms,
553                ),
554            );
555            continue;
556        }
557        let now_ms = state.now_ms.expect("ready branch has clock");
558        let ready = ScheduledReadinessReady {
559            schedule_id: schedule.schedule_id.clone(),
560            subject_refs: canonical_source_refs(schedule.subject_refs.clone()),
561            ready_at_ms: schedule.ready_at_ms,
562            deadline_ms: schedule.deadline_ms,
563            now_ms,
564            source_refs: source_refs.clone(),
565            metadata: metadata.clone(),
566        };
567        emit_ready(ctx, state, ready);
568        emit_status(
569            ctx,
570            state,
571            status_for(
572                &schedule,
573                ScheduledReadinessStatusState::Ready,
574                source_refs.clone(),
575                Some(now_ms),
576            ),
577        );
578        if schedule
579            .deadline_ms
580            .is_some_and(|deadline| now_ms > deadline)
581        {
582            let overdue = ScheduledReadinessOverdue {
583                schedule_id: schedule.schedule_id.clone(),
584                subject_refs: canonical_source_refs(schedule.subject_refs.clone()),
585                ready_at_ms: schedule.ready_at_ms,
586                deadline_ms: schedule.deadline_ms.expect("checked deadline"),
587                now_ms,
588                source_refs,
589                metadata,
590            };
591            emit_overdue(ctx, state, overdue);
592            emit_status(
593                ctx,
594                state,
595                status_for(
596                    &schedule,
597                    ScheduledReadinessStatusState::Overdue,
598                    schedule_source_refs(&schedule, &state.clock_source_refs),
599                    Some(now_ms),
600                ),
601            );
602        }
603    }
604}
605
606fn emit_pending(
607    ctx: &Ctx,
608    state: &mut ScheduledReadinessState,
609    pending: ScheduledReadinessPending,
610) {
611    state
612        .pending_by_id
613        .insert(pending.schedule_id.clone(), pending.clone());
614    let key = format!("pending:{}", pending.schedule_id);
615    if state.emitted_keys.insert(key) {
616        emit_fact(ctx, ScheduledReadinessFact::Pending(pending));
617    }
618}
619
620fn emit_ready(ctx: &Ctx, state: &mut ScheduledReadinessState, ready: ScheduledReadinessReady) {
621    let key = format!("ready:{}", ready.schedule_id);
622    state
623        .ready_by_id
624        .insert(ready.schedule_id.clone(), ready.clone());
625    state.pending_by_id.remove(&ready.schedule_id);
626    if state.emitted_keys.insert(key) {
627        emit_audit(
628            ctx,
629            state,
630            "scheduled-readiness-ready",
631            Some(ready.schedule_id.clone()),
632            ready.source_refs.clone(),
633            Some(BTreeMap::from([
634                ("nowMs".to_owned(), JsonValue::from(ready.now_ms)),
635                ("readyAtMs".to_owned(), JsonValue::from(ready.ready_at_ms)),
636            ])),
637        );
638        emit_fact(ctx, ScheduledReadinessFact::Ready(ready));
639    }
640}
641
642fn emit_overdue(
643    ctx: &Ctx,
644    state: &mut ScheduledReadinessState,
645    overdue: ScheduledReadinessOverdue,
646) {
647    let key = format!("overdue:{}", overdue.schedule_id);
648    state
649        .overdue_by_id
650        .insert(overdue.schedule_id.clone(), overdue.clone());
651    if state.emitted_keys.insert(key) {
652        emit_audit(
653            ctx,
654            state,
655            "scheduled-readiness-overdue",
656            Some(overdue.schedule_id.clone()),
657            overdue.source_refs.clone(),
658            Some(BTreeMap::from([
659                ("nowMs".to_owned(), JsonValue::from(overdue.now_ms)),
660                ("readyAtMs".to_owned(), JsonValue::from(overdue.ready_at_ms)),
661                (
662                    "deadlineMs".to_owned(),
663                    JsonValue::from(overdue.deadline_ms),
664                ),
665            ])),
666        );
667        emit_fact(ctx, ScheduledReadinessFact::Overdue(overdue));
668    }
669}
670
671fn emit_status(ctx: &Ctx, state: &mut ScheduledReadinessState, status: ScheduledReadinessStatus) {
672    state
673        .status_by_id
674        .insert(status.schedule_id.clone(), status.clone());
675    let key = compound_tuple_key("status", &[&status_identity(&status)]);
676    if state.emitted_keys.insert(key) {
677        emit_fact(ctx, ScheduledReadinessFact::Status(status));
678    }
679}
680
681fn emit_issue(ctx: &Ctx, state: &mut ScheduledReadinessState, issue: DataIssue) {
682    let key = canonical_tuple_key(&[
683        &issue.source,
684        &issue.code,
685        issue.details.as_deref().unwrap_or(""),
686    ]);
687    if state.issue_keys.insert(key) {
688        emit_fact(ctx, ScheduledReadinessFact::Issue(issue));
689    }
690}
691
692fn emit_audit(
693    ctx: &Ctx,
694    state: &mut ScheduledReadinessState,
695    kind: impl Into<String>,
696    subject_id: Option<String>,
697    source_refs: Vec<SourceRef>,
698    metadata: Option<BTreeMap<String, JsonValue>>,
699) {
700    state.audit_seq += 1;
701    emit_fact(
702        ctx,
703        ScheduledReadinessFact::Audit(ScheduledReadinessAuditRecord {
704            id: compound_tuple_key("scheduled-readiness-audit", &[&state.audit_seq.to_string()]),
705            kind: kind.into(),
706            subject_id,
707            source_refs: canonical_source_refs(source_refs),
708            metadata: sanitize_metadata(metadata),
709        }),
710    );
711}
712
713fn status_for(
714    schedule: &ScheduledReadinessRequested,
715    state: ScheduledReadinessStatusState,
716    source_refs: Vec<SourceRef>,
717    now_ms: Option<u64>,
718) -> ScheduledReadinessStatus {
719    let status_name = match state {
720        ScheduledReadinessStatusState::Pending => "pending",
721        ScheduledReadinessStatusState::Ready => "ready",
722        ScheduledReadinessStatusState::Overdue => "overdue",
723        ScheduledReadinessStatusState::Issue => "issue",
724    };
725    ScheduledReadinessStatus {
726        status_id: compound_tuple_key(
727            "scheduled-readiness-status",
728            &[&schedule.schedule_id, status_name],
729        ),
730        schedule_id: schedule.schedule_id.clone(),
731        state,
732        subject_refs: canonical_source_refs(schedule.subject_refs.clone()),
733        ready_at_ms: Some(schedule.ready_at_ms),
734        deadline_ms: schedule.deadline_ms,
735        now_ms,
736        source_refs,
737        issue_codes: Vec::new(),
738        metadata: schedule
739            .reason
740            .as_ref()
741            .map(|reason| BTreeMap::from([("reason".to_owned(), JsonValue::from(reason.clone()))])),
742    }
743}
744
745fn sanitize_schedule(mut schedule: ScheduledReadinessRequested) -> ScheduledReadinessRequested {
746    schedule.subject_refs = canonical_source_refs(schedule.subject_refs);
747    schedule.policy_refs = canonical_source_refs(schedule.policy_refs);
748    schedule.source_refs = canonical_source_refs(schedule.source_refs);
749    schedule.metadata = sanitize_metadata(schedule.metadata);
750    schedule
751}
752
753fn readiness_metadata(
754    schedule: &ScheduledReadinessRequested,
755) -> Option<BTreeMap<String, JsonValue>> {
756    let mut metadata = schedule.metadata.clone().unwrap_or_default();
757    if let Some(reason) = &schedule.reason {
758        metadata.insert("reason".to_owned(), JsonValue::from(reason.clone()));
759    }
760    sanitize_metadata((!metadata.is_empty()).then_some(metadata))
761}
762
763fn schedule_source_refs(
764    schedule: &ScheduledReadinessRequested,
765    clock_source_refs: &[SourceRef],
766) -> Vec<SourceRef> {
767    let mut refs = vec![SourceRef::new(
768        "scheduled-readiness",
769        schedule.schedule_id.clone(),
770    )];
771    refs.extend(schedule.source_refs.clone());
772    refs.extend(schedule.policy_refs.clone());
773    refs.extend(clock_source_refs.iter().cloned());
774    canonical_source_refs(refs)
775}
776
777fn schedule_identity(schedule: &ScheduledReadinessRequested) -> String {
778    let value = serde_json::json!({
779        "scheduleId": schedule.schedule_id,
780        "subjectRefs": source_refs_json(&schedule.subject_refs),
781        "readyAtMs": schedule.ready_at_ms,
782        "deadlineMs": schedule.deadline_ms,
783        "reason": schedule.reason,
784        "policyRefs": source_refs_json(&schedule.policy_refs),
785        "sourceRefs": source_refs_json(&schedule.source_refs),
786        "metadata": schedule.metadata,
787    });
788    stable_json_string(&value).unwrap_or_else(|_| format!("{schedule:?}"))
789}
790
791fn status_identity(status: &ScheduledReadinessStatus) -> String {
792    let value = serde_json::json!({
793        "statusId": status.status_id,
794        "scheduleId": status.schedule_id,
795        "state": format!("{:?}", status.state),
796        "readyAtMs": status.ready_at_ms,
797        "deadlineMs": status.deadline_ms,
798        "nowMs": status.now_ms,
799        "issueCodes": status.issue_codes,
800    });
801    stable_json_string(&value).unwrap_or_else(|_| format!("{status:?}"))
802}
803
804/// Creates or computes `readiness_issue`.
805pub fn readiness_issue(
806    code: impl Into<String>,
807    message: impl Into<String>,
808    subject_id: &str,
809    details: Option<String>,
810    severity: impl Into<String>,
811) -> DataIssue {
812    DataIssue {
813        kind: "data-issue".to_owned(),
814        code: code.into(),
815        message: message.into(),
816        severity: severity.into(),
817        source: format!("scheduled-readiness:{subject_id}"),
818        topic: None,
819        details,
820    }
821}
822
823fn canonical_source_refs(refs: Vec<SourceRef>) -> Vec<SourceRef> {
824    let mut refs = refs
825        .into_iter()
826        .filter_map(|mut source_ref| {
827            if source_ref.kind.is_empty() || source_ref.id.is_empty() {
828                return None;
829            }
830            source_ref.metadata = sanitize_metadata(source_ref.metadata);
831            Some(source_ref)
832        })
833        .collect::<Vec<_>>();
834    refs.sort_by_key(source_ref_sort_key);
835    let mut seen = BTreeSet::new();
836    let mut out = Vec::new();
837    for source_ref in refs {
838        let key = canonical_tuple_key(&[&source_ref.kind, &source_ref.id]);
839        if seen.insert(key) {
840            out.push(source_ref);
841        }
842    }
843    out
844}
845
846fn source_ref_sort_key(source_ref: &SourceRef) -> (String, String, String) {
847    let metadata = source_ref
848        .metadata
849        .as_ref()
850        .map(|metadata| {
851            stable_json_string(&JsonValue::Object(
852                metadata
853                    .iter()
854                    .map(|(key, value)| (key.clone(), value.clone()))
855                    .collect(),
856            ))
857        })
858        .transpose()
859        .ok()
860        .flatten()
861        .unwrap_or_default();
862    (source_ref.kind.clone(), source_ref.id.clone(), metadata)
863}
864
865fn sanitize_metadata(
866    metadata: Option<BTreeMap<String, JsonValue>>,
867) -> Option<BTreeMap<String, JsonValue>> {
868    let mut out = BTreeMap::new();
869    for (key, value) in metadata.unwrap_or_default() {
870        if is_runtime_metadata_key(&key) {
871            continue;
872        }
873        out.insert(key, value);
874    }
875    (!out.is_empty()).then_some(out)
876}
877
878fn is_runtime_metadata_key(key: &str) -> bool {
879    let lower = key.to_ascii_lowercase();
880    matches!(
881        lower.as_str(),
882        "apikey"
883            | "api_key"
884            | "secret"
885            | "client"
886            | "transport"
887            | "subprocess"
888            | "sdk"
889            | "oauth"
890            | "credential"
891            | "credentials"
892            | "accesstoken"
893            | "access_token"
894            | "refreshtoken"
895            | "refresh_token"
896            | "token"
897            | "password"
898            | "authorization"
899            | "cookie"
900            | "stdout"
901            | "stderr"
902            | "stack"
903            | "stacktrace"
904            | "providerraw"
905            | "provider_raw"
906            | "rawresponse"
907            | "raw_response"
908            | "diff"
909            | "patch"
910            | "filecontents"
911            | "file_contents"
912            | "binary"
913            | "media"
914    )
915}
916
917fn source_refs_json(refs: &[SourceRef]) -> JsonValue {
918    JsonValue::Array(
919        refs.iter()
920            .map(|source_ref| {
921                let mut object = Map::new();
922                object.insert("kind".to_owned(), JsonValue::from(source_ref.kind.clone()));
923                object.insert("id".to_owned(), JsonValue::from(source_ref.id.clone()));
924                if let Some(metadata) = &source_ref.metadata {
925                    object.insert(
926                        "metadata".to_owned(),
927                        JsonValue::Object(
928                            metadata
929                                .iter()
930                                .map(|(key, value)| (key.clone(), value.clone()))
931                                .collect(),
932                        ),
933                    );
934                }
935                JsonValue::Object(object)
936            })
937            .collect(),
938    )
939}
940
941fn string_field<'a>(object: &'a Map<String, Value>, key: &str) -> Option<&'a str> {
942    object.get(key).and_then(Value::as_str)
943}
944
945fn u64_field(object: &Map<String, Value>, key: &str) -> Option<u64> {
946    object.get(key).and_then(Value::as_u64)
947}
948
949fn optional_u64_field(object: &Map<String, Value>, key: &str) -> Result<Option<u64>, ()> {
950    match object.get(key) {
951        None => Ok(None),
952        Some(value) => value.as_u64().map(Some).ok_or(()),
953    }
954}
955
956fn source_refs_field(object: &Map<String, Value>, key: &str) -> Vec<SourceRef> {
957    let Some(values) = object.get(key).and_then(Value::as_array) else {
958        return Vec::new();
959    };
960    canonical_source_refs(
961        values
962            .iter()
963            .filter_map(|value| {
964                let object = value.as_object()?;
965                Some(SourceRef {
966                    kind: string_field(object, "kind")?.to_owned(),
967                    id: string_field(object, "id")?.to_owned(),
968                    metadata: object_metadata(object.get("metadata")),
969                })
970            })
971            .collect(),
972    )
973}
974
975fn object_metadata(value: Option<&Value>) -> Option<BTreeMap<String, JsonValue>> {
976    let object = value?.as_object()?;
977    sanitize_metadata(Some(
978        object
979            .iter()
980            .map(|(key, value)| (key.clone(), value.clone()))
981            .collect(),
982    ))
983}
984
985fn project_fact<T: Clone + 'static>(
986    graph: &Graph,
987    runtime: &Node<ScheduledReadinessFact>,
988    name: String,
989    factory: &'static str,
990    select: impl Fn(&ScheduledReadinessFact) -> Option<T> + 'static,
991) -> Node<T> {
992    graph.node_opts::<T, _>(
993        vec![runtime.erased()],
994        move |ctx| {
995            for fact in ctx.batch::<ScheduledReadinessFact>(0) {
996                if let Some(value) = select(&fact) {
997                    ctx.emit(value);
998                }
999            }
1000        },
1001        node_opts(name, factory),
1002    )
1003}
1004
1005fn emit_fact(ctx: &Ctx, fact: ScheduledReadinessFact) {
1006    ctx.emit(fact);
1007}
1008
1009fn node_opts(name: impl Into<String>, factory: impl Into<String>) -> GraphNodeOpts {
1010    let mut opts = GraphNodeOpts::named(name);
1011    opts.node = NodeOpts {
1012        factory: Some(factory.into()),
1013        complete_when_deps_complete: false,
1014        error_when_deps_error: false,
1015        ..opts.node
1016    };
1017    opts
1018}