1use std::collections::{BTreeMap, HashMap};
8
9use serde_json::Value;
10
11use crate::node::StepNode;
12use crate::{CapabilityManifest, CapabilityPin, GuardKind};
13
14#[derive(Debug, thiserror::Error)]
16pub enum NodeError {
17 #[error("unknown node type '{0}' (not registered)")]
19 UnknownType(String),
20 #[error("node '{node_type}' has invalid config: {reason}")]
22 InvalidConfig {
23 node_type: String,
25 reason: String,
27 },
28 #[error("invalid capability manifest: {0}")]
30 InvalidCapability(String),
31 #[error("capability '{0}' is already registered")]
33 DuplicateCapability(String),
34}
35
36pub type StepFactory = fn(config: &Value) -> Result<Box<dyn StepNode>, NodeError>;
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
41#[serde(rename_all = "lowercase")]
42pub enum FieldType {
43 String,
45 Number,
47 Bool,
49 Array,
51 Object,
53 Any,
55}
56
57#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
59pub struct FieldSpec {
60 pub key: String,
62 pub ty: FieldType,
64 pub required: bool,
66}
67
68impl NodeSchema {
69 pub fn to_json_schema(&self) -> Value {
71 let mut properties = serde_json::Map::new();
72 let mut required = Vec::new();
73 for field in &self.fields {
74 let schema = match field.ty {
75 FieldType::String => serde_json::json!({"type":"string"}),
76 FieldType::Number => serde_json::json!({"type":"number"}),
77 FieldType::Bool => serde_json::json!({"type":"boolean"}),
78 FieldType::Array => serde_json::json!({"type":"array"}),
79 FieldType::Object => serde_json::json!({"type":"object"}),
80 FieldType::Any => serde_json::json!({}),
81 };
82 properties.insert(field.key.clone(), schema);
83 if field.required {
84 required.push(field.key.clone());
85 }
86 }
87 serde_json::json!({"type":"object","properties":properties,"required":required,"additionalProperties":false})
88 }
89}
90
91impl FieldSpec {
92 pub fn required(key: impl Into<String>, ty: FieldType) -> Self {
94 Self {
95 key: key.into(),
96 ty,
97 required: true,
98 }
99 }
100
101 pub fn optional(key: impl Into<String>, ty: FieldType) -> Self {
103 Self {
104 key: key.into(),
105 ty,
106 required: false,
107 }
108 }
109}
110
111#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
113pub struct NodeSchema {
114 pub fields: Vec<FieldSpec>,
116}
117
118#[derive(Debug, Clone)]
121pub struct NodeRegistry {
122 steps: HashMap<String, StepFactory>,
123 ingress: std::collections::HashSet<String>,
124 fan_out: std::collections::HashSet<String>,
125 guard_kinds: HashMap<String, GuardKind>,
126 capabilities: BTreeMap<CapabilityPin, CapabilityManifest>,
127 capability_steps: BTreeMap<CapabilityPin, StepFactory>,
128 capability_schemas: BTreeMap<CapabilityPin, NodeSchema>,
129 capability_guards: BTreeMap<CapabilityPin, GuardKind>,
130 capability_fan_out: std::collections::BTreeSet<CapabilityPin>,
131 schemas: HashMap<String, NodeSchema>,
132}
133
134impl NodeRegistry {
135 pub fn empty() -> Self {
137 Self {
138 steps: HashMap::new(),
139 ingress: Default::default(),
140 fan_out: Default::default(),
141 guard_kinds: Default::default(),
142 capabilities: Default::default(),
143 capability_steps: Default::default(),
144 capability_schemas: Default::default(),
145 capability_guards: Default::default(),
146 capability_fan_out: Default::default(),
147 schemas: HashMap::new(),
148 }
149 }
150
151 pub fn with_builtins() -> Self {
153 let mut r = Self::empty();
154 crate::builtins::register_builtins(&mut r);
155 r
156 }
157
158 pub fn register_step(&mut self, node_type: impl Into<String>, factory: StepFactory) {
160 self.steps.insert(node_type.into(), factory);
161 }
162
163 pub fn register_ingress(&mut self, node_type: impl Into<String>) {
165 self.ingress.insert(node_type.into());
166 }
167
168 pub fn register_fan_out(&mut self, node_type: impl Into<String>) {
171 self.fan_out.insert(node_type.into());
172 }
173
174 pub fn is_fan_out_capable(&self, node_type: &str) -> bool {
176 self.fan_out.contains(node_type)
177 }
178
179 pub fn register_side_effect_guard(
181 &mut self,
182 node_type: impl Into<String>,
183 ) -> Result<(), NodeError> {
184 self.register_guard(node_type, GuardKind::Authorization)
185 }
186
187 pub fn register_guard(
189 &mut self,
190 node_type: impl Into<String>,
191 kind: GuardKind,
192 ) -> Result<(), NodeError> {
193 let node_type = node_type.into();
194 let mut manifests = self
195 .capabilities
196 .values()
197 .filter(|manifest| manifest.id == node_type);
198 if let Some(manifest) = manifests.next() {
199 if manifests.next().is_some() {
200 return Err(NodeError::InvalidCapability(format!(
201 "guard capability '{node_type}' has multiple versions; register an exact implementation"
202 )));
203 }
204 if manifest.kind != crate::CapabilityKind::Guard {
205 return Err(NodeError::InvalidCapability(format!(
206 "capability '{node_type}' manifest is not a guard"
207 )));
208 }
209 if manifest.guard_kind != Some(kind) {
210 return Err(NodeError::InvalidCapability(format!(
211 "guard capability '{node_type}' implementation guard {:?} does not match manifest guard {:?}",
212 kind, manifest.guard_kind
213 )));
214 }
215 }
216 self.guard_kinds.insert(node_type, kind);
217 Ok(())
218 }
219
220 pub fn guard_kind(&self, node_type: &str) -> Option<GuardKind> {
222 self.capability(node_type)
223 .and_then(|manifest| manifest.guard_kind)
224 .or_else(|| self.guard_kinds.get(node_type).copied())
225 }
226
227 pub fn is_side_effect_guard(&self, node_type: &str) -> bool {
229 self.guard_kind(node_type).is_some()
230 }
231
232 pub fn is_ingress(&self, node_type: &str) -> bool {
234 self.ingress.contains(node_type) || node_type.starts_with("ingress.")
235 }
236
237 pub fn is_step(&self, node_type: &str) -> bool {
239 self.steps.contains_key(node_type)
240 || self
241 .capabilities
242 .values()
243 .any(|manifest| manifest.id == node_type)
244 }
245
246 pub fn build_step(
248 &self,
249 node_type: &str,
250 config: &Value,
251 ) -> Result<Box<dyn StepNode>, NodeError> {
252 let factory = self
253 .steps
254 .get(node_type)
255 .ok_or_else(|| NodeError::UnknownType(node_type.to_string()))?;
256 factory(config)
257 }
258
259 pub fn known_step_types(&self) -> impl Iterator<Item = &str> {
261 self.steps.keys().map(|s| s.as_str())
262 }
263
264 pub fn known_node_types(&self) -> impl Iterator<Item = &str> {
266 self.steps
267 .keys()
268 .chain(self.ingress.iter())
269 .map(String::as_str)
270 }
271
272 pub fn register_schema(&mut self, node_type: impl Into<String>, schema: NodeSchema) {
274 self.schemas.insert(node_type.into(), schema);
275 }
276
277 pub fn schema(&self, node_type: &str) -> Option<&NodeSchema> {
279 self.schemas.get(node_type)
280 }
281
282 pub fn authoring_schemas(&self) -> Vec<(String, Value)> {
284 let mut entries = self
285 .schemas
286 .iter()
287 .map(|(id, schema)| (id.clone(), schema.to_json_schema()))
288 .collect::<Vec<_>>();
289 entries.sort_by(|left, right| left.0.cmp(&right.0));
290 entries
291 }
292
293 pub fn register_capability(&mut self, manifest: CapabilityManifest) -> Result<(), NodeError> {
295 manifest
296 .validate()
297 .map_err(|error| NodeError::InvalidCapability(error.to_string()))?;
298 let pin = CapabilityPin {
299 id: manifest.id.clone(),
300 contract_version: manifest.contract_version.clone(),
301 content_digest: manifest.content_digest.clone(),
302 };
303 if self.capabilities.contains_key(&pin) {
304 return Err(NodeError::DuplicateCapability(manifest.id));
305 }
306 if let Some(guard) = self.guard_kinds.get(&manifest.id).copied() {
307 if self
308 .capabilities
309 .values()
310 .any(|registered| registered.id == manifest.id)
311 {
312 return Err(NodeError::InvalidCapability(format!(
313 "guard capability '{}' has multiple versions; register an exact implementation",
314 manifest.id
315 )));
316 }
317 if manifest.kind != crate::CapabilityKind::Guard {
318 return Err(NodeError::InvalidCapability(format!(
319 "capability '{}' manifest is not a guard",
320 manifest.id
321 )));
322 }
323 if manifest.guard_kind != Some(guard) {
324 return Err(NodeError::InvalidCapability(format!(
325 "guard capability '{}' manifest guard {:?} does not match implementation guard {:?}",
326 manifest.id, manifest.guard_kind, guard
327 )));
328 }
329 }
330 self.capabilities.insert(pin, manifest);
331 Ok(())
332 }
333
334 pub fn register_capability_implementation(
336 &mut self,
337 pin: CapabilityPin,
338 factory: StepFactory,
339 schema: Option<NodeSchema>,
340 guard: Option<GuardKind>,
341 fan_out: bool,
342 ) -> Result<(), NodeError> {
343 let manifest = self.capabilities.get(&pin).ok_or_else(|| {
344 NodeError::InvalidCapability(format!(
345 "capability '{}' implementation has no registered manifest at {} ({})",
346 pin.id, pin.contract_version, pin.content_digest
347 ))
348 })?;
349 if guard != manifest.guard_kind {
350 return Err(NodeError::InvalidCapability(format!(
351 "capability '{}' implementation guard {:?} does not match manifest guard {:?}",
352 pin.id, guard, manifest.guard_kind
353 )));
354 }
355 if self.capability_steps.contains_key(&pin) {
356 return Err(NodeError::DuplicateCapability(pin.id));
357 }
358 self.capability_steps.insert(pin.clone(), factory);
359 if let Some(schema) = schema {
360 self.capability_schemas.insert(pin.clone(), schema);
361 }
362 if let Some(guard) = guard {
363 self.capability_guards.insert(pin.clone(), guard);
364 }
365 if fan_out {
366 self.capability_fan_out.insert(pin);
367 }
368 Ok(())
369 }
370
371 pub fn capability(&self, id: &str) -> Option<&CapabilityManifest> {
373 let mut matches = self
374 .capabilities
375 .values()
376 .filter(|manifest| manifest.id == id);
377 let manifest = matches.next()?;
378 matches.next().is_none().then_some(manifest)
379 }
380
381 pub fn capability_by_pin(&self, pin: &CapabilityPin) -> Option<&CapabilityManifest> {
383 self.capabilities.get(pin)
384 }
385
386 pub fn for_capability_pins(&self, pins: &[CapabilityPin]) -> Result<Self, NodeError> {
388 let mut selected = self.clone();
389 selected.capabilities.clear();
390 selected
391 .fan_out
392 .retain(|node_type| !pins.iter().any(|pin| pin.id == *node_type));
393 for pin in pins {
394 selected.guard_kinds.remove(&pin.id);
395 let manifest = self.capability_by_pin(pin).ok_or_else(|| {
396 NodeError::InvalidCapability(format!(
397 "capability '{}' is unavailable at {} ({})",
398 pin.id, pin.contract_version, pin.content_digest
399 ))
400 })?;
401 selected.capabilities.insert(pin.clone(), manifest.clone());
402 let versions = self
403 .capabilities
404 .keys()
405 .filter(|candidate| candidate.id == pin.id)
406 .count();
407 if let Some(factory) = self.capability_steps.get(pin) {
408 selected.steps.insert(pin.id.clone(), *factory);
409 } else if versions > 1 && self.steps.contains_key(&pin.id) {
410 return Err(NodeError::InvalidCapability(format!(
411 "capability '{}' has multiple versions but no executable implementation for {} ({})",
412 pin.id, pin.contract_version, pin.content_digest
413 )));
414 }
415 if let Some(schema) = self.capability_schemas.get(pin) {
416 selected.schemas.insert(pin.id.clone(), schema.clone());
417 } else if versions > 1 && self.schemas.contains_key(&pin.id) {
418 return Err(NodeError::InvalidCapability(format!(
419 "capability '{}' has multiple versions but no authoring schema for {} ({})",
420 pin.id, pin.contract_version, pin.content_digest
421 )));
422 }
423 if let Some(manifest_guard) = manifest.guard_kind {
424 let guard = match self.capability_guards.get(pin).copied() {
425 Some(guard) => guard,
426 None => match self.guard_kinds.get(&pin.id).copied() {
427 Some(guard) if versions == 1 && guard == manifest_guard => guard,
428 Some(guard) if versions == 1 => {
429 return Err(NodeError::InvalidCapability(format!(
430 "guard capability '{}' implementation guard {:?} does not match manifest guard {:?}",
431 pin.id, guard, manifest_guard
432 )));
433 }
434 Some(_) => {
435 return Err(NodeError::InvalidCapability(format!(
436 "guard capability '{}' has multiple versions but no exact guard implementation for {} ({})",
437 pin.id, pin.contract_version, pin.content_digest
438 )));
439 }
440 None if !self.steps.contains_key(&pin.id) => manifest_guard,
441 None => {
442 return Err(NodeError::InvalidCapability(format!(
443 "guard capability '{}' has no guard implementation for {} ({})",
444 pin.id, pin.contract_version, pin.content_digest
445 )));
446 }
447 },
448 };
449 selected.guard_kinds.insert(pin.id.clone(), guard);
450 } else if self.steps.contains_key(&pin.id) && self.guard_kinds.contains_key(&pin.id) {
451 return Err(NodeError::InvalidCapability(format!(
452 "capability '{}' implementation declares a guard role but its manifest is not a guard",
453 pin.id
454 )));
455 }
456 if self.capability_fan_out.contains(pin) {
457 selected.fan_out.insert(pin.id.clone());
458 }
459 }
460 Ok(selected)
461 }
462
463 pub fn capability_manifests(&self) -> impl Iterator<Item = &CapabilityManifest> {
465 self.capabilities.values()
466 }
467}
468
469#[cfg(test)]
470mod tests {
471 use super::*;
472 use crate::{CapabilityKind, Effect, IdempotencyMode};
473
474 struct Noop;
475
476 #[async_trait::async_trait]
477 impl StepNode for Noop {
478 async fn process(
479 &self,
480 event: &crate::Event,
481 _ctx: &crate::WorkflowContext,
482 ) -> crate::StepResult {
483 crate::StepResult::Pass(event.clone())
484 }
485 }
486
487 fn noop(_: &Value) -> Result<Box<dyn StepNode>, NodeError> {
488 Ok(Box::new(Noop))
489 }
490
491 #[test]
492 fn authoring_catalog_includes_steps_and_ingress() {
493 let mut registry = NodeRegistry::empty();
494 registry.register_step("transform.test", noop);
495 registry.register_ingress("ingress.test");
496 let types = registry
497 .known_node_types()
498 .collect::<std::collections::BTreeSet<_>>();
499 assert_eq!(
500 types,
501 std::collections::BTreeSet::from(["ingress.test", "transform.test"])
502 );
503 }
504
505 #[test]
506 fn capability_versions_are_indexed_and_selected_by_full_pin() {
507 let manifest = |version: &str, digest: &str| {
508 CapabilityManifest::action(
509 "action.versioned",
510 version,
511 digest,
512 Effect::ExternalWrite,
513 IdempotencyMode::Native,
514 true,
515 )
516 };
517 let v1 = manifest("1", "digest-v1");
518 let v2 = manifest("2", "digest-v2");
519 let pin = CapabilityPin {
520 id: v1.id.clone(),
521 contract_version: v1.contract_version.clone(),
522 content_digest: v1.content_digest.clone(),
523 };
524 let mut registry = NodeRegistry::empty();
525 registry.register_capability(v1.clone()).unwrap();
526 registry.register_capability(v2).unwrap();
527
528 assert!(registry.capability("action.versioned").is_none());
529 assert_eq!(registry.capability_by_pin(&pin), Some(&v1));
530 let selected = registry.for_capability_pins(&[pin]).unwrap();
531 assert_eq!(selected.capability("action.versioned"), Some(&v1));
532 }
533
534 #[test]
535 fn versioned_capability_selects_exact_factory_schema_and_guard() {
536 let manifest = |version: &str, digest: &str| {
537 let mut manifest = CapabilityManifest::action(
538 "guard.versioned",
539 version,
540 digest,
541 Effect::Pure,
542 IdempotencyMode::Native,
543 true,
544 );
545 manifest.kind = CapabilityKind::Guard;
546 manifest.guard_kind = Some(GuardKind::Authorization);
547 manifest
548 };
549 let v1 = manifest("1", "digest-v1");
550 let v2 = manifest("2", "digest-v2");
551 let pin = CapabilityPin {
552 id: v1.id.clone(),
553 contract_version: v1.contract_version.clone(),
554 content_digest: v1.content_digest.clone(),
555 };
556 let schema = NodeSchema {
557 fields: vec![FieldSpec::required("approved", FieldType::Bool)],
558 };
559 let mut registry = NodeRegistry::empty();
560 registry.register_capability(v1).unwrap();
561 registry.register_capability(v2).unwrap();
562 let mismatch = registry
563 .register_capability_implementation(
564 pin.clone(),
565 noop,
566 None,
567 Some(GuardKind::Freshness),
568 false,
569 )
570 .unwrap_err();
571 assert!(matches!(mismatch, NodeError::InvalidCapability(_)));
572 registry
573 .register_capability_implementation(
574 pin.clone(),
575 noop,
576 Some(schema.clone()),
577 Some(GuardKind::Authorization),
578 false,
579 )
580 .unwrap();
581
582 let selected = registry.for_capability_pins(&[pin]).unwrap();
583 assert_eq!(selected.schema("guard.versioned"), Some(&schema));
584 assert_eq!(
585 selected.guard_kind("guard.versioned"),
586 Some(GuardKind::Authorization)
587 );
588 assert!(selected.build_step("guard.versioned", &Value::Null).is_ok());
589 }
590
591 #[test]
592 fn versioned_capability_selects_fan_out_by_full_pin() {
593 let manifest = |version: &str, digest: &str| {
594 CapabilityManifest::action(
595 "transform.versioned",
596 version,
597 digest,
598 Effect::Pure,
599 IdempotencyMode::Native,
600 true,
601 )
602 };
603 let v1 = manifest("1", "digest-v1");
604 let v2 = manifest("2", "digest-v2");
605 let pin = |manifest: &CapabilityManifest| CapabilityPin {
606 id: manifest.id.clone(),
607 contract_version: manifest.contract_version.clone(),
608 content_digest: manifest.content_digest.clone(),
609 };
610 let mut registry = NodeRegistry::empty();
611 registry.register_capability(v1.clone()).unwrap();
612 registry.register_capability(v2.clone()).unwrap();
613 registry
614 .register_capability_implementation(pin(&v1), noop, None, None, true)
615 .unwrap();
616 registry
617 .register_capability_implementation(pin(&v2), noop, None, None, false)
618 .unwrap();
619
620 assert!(registry
621 .for_capability_pins(&[pin(&v1)])
622 .unwrap()
623 .is_fan_out_capable("transform.versioned"));
624 assert!(!registry
625 .for_capability_pins(&[pin(&v2)])
626 .unwrap()
627 .is_fan_out_capable("transform.versioned"));
628 }
629
630 #[test]
631 fn legacy_guard_registration_must_match_the_manifest() {
632 let mut manifest = CapabilityManifest::action(
633 "guard.legacy",
634 "1",
635 "digest",
636 Effect::Pure,
637 IdempotencyMode::None,
638 false,
639 );
640 manifest.kind = CapabilityKind::Guard;
641 manifest.guard_kind = Some(GuardKind::Authorization);
642 let pin = CapabilityPin {
643 id: manifest.id.clone(),
644 contract_version: manifest.contract_version.clone(),
645 content_digest: manifest.content_digest.clone(),
646 };
647 let mut persisted = NodeRegistry::empty();
648 persisted.register_capability(manifest.clone()).unwrap();
649 assert_eq!(
650 persisted
651 .for_capability_pins(std::slice::from_ref(&pin))
652 .unwrap()
653 .guard_kind(&manifest.id),
654 manifest.guard_kind
655 );
656
657 let mut legacy = persisted;
658 legacy.register_step(&manifest.id, noop);
659 assert!(matches!(
660 legacy.register_guard(&manifest.id, GuardKind::Freshness),
661 Err(NodeError::InvalidCapability(reason)) if reason.contains("does not match")
662 ));
663 legacy
664 .register_guard(&manifest.id, GuardKind::Authorization)
665 .unwrap();
666 assert_eq!(
667 legacy
668 .for_capability_pins(&[pin])
669 .unwrap()
670 .guard_kind(&manifest.id),
671 manifest.guard_kind
672 );
673
674 let action = CapabilityManifest::action(
675 "action.legacy",
676 "1",
677 "digest",
678 Effect::Pure,
679 IdempotencyMode::None,
680 false,
681 );
682 let action_pin = CapabilityPin {
683 id: action.id.clone(),
684 contract_version: action.contract_version.clone(),
685 content_digest: action.content_digest.clone(),
686 };
687 let mut non_guard = NodeRegistry::empty();
688 non_guard.register_capability(action.clone()).unwrap();
689 assert!(non_guard
690 .for_capability_pins(std::slice::from_ref(&action_pin))
691 .is_ok());
692 non_guard.register_step(&action.id, noop);
693 assert!(matches!(
694 non_guard.register_guard(&action.id, GuardKind::Authorization),
695 Err(NodeError::InvalidCapability(reason)) if reason.contains("manifest is not a guard")
696 ));
697 assert!(non_guard
698 .for_capability_pins(&[action_pin])
699 .unwrap()
700 .guard_kind(&action.id)
701 .is_none());
702 }
703
704 #[test]
705 fn manifest_registration_must_match_a_legacy_guard() {
706 let manifest = |version: &str, kind| {
707 let mut manifest = CapabilityManifest::action(
708 "guard.reverse",
709 version,
710 format!("digest-{version}"),
711 Effect::Pure,
712 IdempotencyMode::None,
713 false,
714 );
715 manifest.kind = CapabilityKind::Guard;
716 manifest.guard_kind = Some(kind);
717 manifest
718 };
719 let mut registry = NodeRegistry::empty();
720 registry
721 .register_guard("guard.reverse", GuardKind::Authorization)
722 .unwrap();
723 assert!(matches!(
724 registry.register_capability(manifest("1", GuardKind::Freshness)),
725 Err(NodeError::InvalidCapability(reason)) if reason.contains("does not match")
726 ));
727 registry
728 .register_capability(manifest("1", GuardKind::Authorization))
729 .unwrap();
730 assert!(matches!(
731 registry.register_capability(manifest("2", GuardKind::Authorization)),
732 Err(NodeError::InvalidCapability(reason)) if reason.contains("multiple versions")
733 ));
734
735 let mut non_guard = NodeRegistry::empty();
736 non_guard
737 .register_guard("action.reverse", GuardKind::Authorization)
738 .unwrap();
739 assert!(matches!(
740 non_guard.register_capability(CapabilityManifest::action(
741 "action.reverse",
742 "1",
743 "digest",
744 Effect::Pure,
745 IdempotencyMode::None,
746 false,
747 )),
748 Err(NodeError::InvalidCapability(reason)) if reason.contains("manifest is not a guard")
749 ));
750 }
751}