1use 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)]
21pub struct SourceRef {
23 pub kind: String,
25 pub id: String,
27 pub metadata: Option<BTreeMap<String, JsonValue>>,
29}
30
31impl SourceRef {
32 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)]
43pub struct ScheduledReadinessRequested {
45 pub schedule_id: String,
47 pub subject_refs: Vec<SourceRef>,
49 pub ready_at_ms: u64,
51 pub deadline_ms: Option<u64>,
53 pub reason: Option<String>,
55 pub policy_refs: Vec<SourceRef>,
57 pub source_refs: Vec<SourceRef>,
59 pub metadata: Option<BTreeMap<String, JsonValue>>,
61}
62
63#[derive(Debug, Clone, PartialEq)]
64pub struct ScheduledReadinessClock {
66 pub clock_id: String,
68 pub now_ms: u64,
70 pub source_refs: Vec<SourceRef>,
72 pub metadata: Option<BTreeMap<String, JsonValue>>,
74}
75
76#[derive(Debug, Clone, PartialEq)]
77pub struct ScheduledReadinessPending {
79 pub schedule_id: String,
81 pub subject_refs: Vec<SourceRef>,
83 pub ready_at_ms: u64,
85 pub deadline_ms: Option<u64>,
87 pub now_ms: Option<u64>,
89 pub source_refs: Vec<SourceRef>,
91 pub metadata: Option<BTreeMap<String, JsonValue>>,
93}
94
95#[derive(Debug, Clone, PartialEq)]
96pub struct ScheduledReadinessReady {
98 pub schedule_id: String,
100 pub subject_refs: Vec<SourceRef>,
102 pub ready_at_ms: u64,
104 pub deadline_ms: Option<u64>,
106 pub now_ms: u64,
108 pub source_refs: Vec<SourceRef>,
110 pub metadata: Option<BTreeMap<String, JsonValue>>,
112}
113
114#[derive(Debug, Clone, PartialEq)]
115pub struct ScheduledReadinessOverdue {
117 pub schedule_id: String,
119 pub subject_refs: Vec<SourceRef>,
121 pub ready_at_ms: u64,
123 pub deadline_ms: u64,
125 pub now_ms: u64,
127 pub source_refs: Vec<SourceRef>,
129 pub metadata: Option<BTreeMap<String, JsonValue>>,
131}
132
133#[derive(Debug, Clone, Copy, PartialEq, Eq)]
134pub enum ScheduledReadinessStatusState {
136 Pending,
138 Ready,
140 Overdue,
142 Issue,
144}
145
146#[derive(Debug, Clone, PartialEq)]
147pub struct ScheduledReadinessStatus {
149 pub status_id: String,
151 pub schedule_id: String,
153 pub state: ScheduledReadinessStatusState,
155 pub subject_refs: Vec<SourceRef>,
157 pub ready_at_ms: Option<u64>,
159 pub deadline_ms: Option<u64>,
161 pub now_ms: Option<u64>,
163 pub source_refs: Vec<SourceRef>,
165 pub issue_codes: Vec<String>,
167 pub metadata: Option<BTreeMap<String, JsonValue>>,
169}
170
171#[derive(Debug, Clone, PartialEq)]
172pub struct ScheduledReadinessAuditRecord {
174 pub id: String,
176 pub kind: String,
178 pub subject_id: Option<String>,
180 pub source_refs: Vec<SourceRef>,
182 pub metadata: Option<BTreeMap<String, JsonValue>>,
184}
185
186#[derive(Debug, Clone, PartialEq, Default)]
187pub struct ScheduledReadinessViews {
189 pub schedules_by_id: BTreeMap<String, ScheduledReadinessRequested>,
191 pub pending_by_id: BTreeMap<String, ScheduledReadinessPending>,
193 pub ready_by_id: BTreeMap<String, ScheduledReadinessReady>,
195 pub overdue_by_id: BTreeMap<String, ScheduledReadinessOverdue>,
197 pub status_by_id: BTreeMap<String, ScheduledReadinessStatus>,
199 pub now_ms: Option<u64>,
201}
202
203#[derive(Clone)]
204pub struct ScheduledReadinessBundle {
206 pub pending: Node<ScheduledReadinessPending>,
208 pub ready: Node<ScheduledReadinessReady>,
210 pub overdue: Node<ScheduledReadinessOverdue>,
212 pub status: Node<ScheduledReadinessStatus>,
214 pub issues: Node<DataIssue>,
216 pub audit: Node<ScheduledReadinessAuditRecord>,
218 pub views: Node<ScheduledReadinessViews>,
220}
221
222#[derive(Clone)]
223pub struct ScheduledReadinessOptions {
225 pub name: Option<String>,
227 pub schedules: Vec<Node<ScheduledReadinessRequested>>,
229 pub clocks: Vec<Node<ScheduledReadinessClock>>,
231}
232
233impl ScheduledReadinessOptions {
234 pub fn new(schedules: Vec<Node<ScheduledReadinessRequested>>) -> Self {
236 Self {
237 name: None,
238 schedules,
239 clocks: Vec::new(),
240 }
241 }
242
243 pub fn named(mut self, name: impl Into<String>) -> Self {
245 self.name = Some(name.into());
246 self
247 }
248
249 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
281pub 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
405pub 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
804pub 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}