1use std::fmt;
12use std::sync::Arc;
13
14use indexmap::IndexMap;
15use serde::{Deserialize, Serialize};
16
17use crate::case::{CaseRef, Versioned};
18use crate::command::{CommandBatch, CommandPolicy};
19use crate::error::{ErasedCallError, ErasureError, ExecutionError, StoreError};
20use crate::event::{ArtifactRef, Commit, CommittedEvent, OperationalReceipt, ReceiptEvent};
21use crate::flow::{
22 ConfirmationSubject, DomainEnumeration, InteractionRequirement, ObligationId, PhaseOwnership,
23 StartBehaviour, StartPrecondition, ViewOf, WorkflowDefinition, WorkflowExecutor,
24 WorkflowNotice, WorkflowView, WritingStage,
25};
26use crate::ids::{AccountId, CaseId, WorkflowKey, WorkflowVersion};
27use crate::interaction::InteractionSpec;
28use crate::locale::Locale;
29use crate::operation::{GlossaryTerm, OperationSpec};
30use crate::target::ResolvedAct;
31
32#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
34pub struct ErasedObligation {
35 pub id: ObligationId,
37 pub value: serde_json::Value,
39 #[serde(default, skip_serializing_if = "Option::is_none")]
46 pub sentence: Option<crate::locale::LocalizedText>,
47 #[serde(default, skip_serializing_if = "Option::is_none")]
50 pub act: Option<super::ObligationAct>,
51}
52
53#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
56pub struct ErasedWorkflowView {
57 pub case_ref: CaseRef,
59 pub workflow_version: WorkflowVersion,
61 pub phase: serde_json::Value,
63 pub phase_ownership: PhaseOwnership,
65 pub obligations: Vec<ErasedObligation>,
67 #[serde(default, skip_serializing_if = "Option::is_none")]
69 pub blocking_interaction: Option<InteractionRequirement>,
70 #[serde(default)]
72 pub notices: Vec<WorkflowNotice>,
73 #[serde(default, skip_serializing_if = "Option::is_none")]
75 pub outcome: Option<serde_json::Value>,
76 #[serde(default, skip_serializing_if = "Vec::is_empty")]
84 pub state: Vec<crate::flow::StateField>,
85}
86
87impl ErasedWorkflowView {
88 #[must_use]
90 pub fn is_complete(&self) -> bool {
91 self.outcome.is_some()
92 }
93
94 #[must_use]
96 pub fn is_user_owned(&self) -> bool {
97 self.phase_ownership == PhaseOwnership::User
98 }
99
100 #[must_use]
102 pub fn obligation_ids(&self) -> Vec<&ObligationId> {
103 self.obligations.iter().map(|o| &o.id).collect()
104 }
105}
106
107pub trait ErasedWorkflow: Send + Sync {
112 fn key(&self) -> WorkflowKey;
114
115 fn version(&self) -> WorkflowVersion;
117
118 fn project(
120 &self,
121 case_ref: CaseRef,
122 state: Option<&serde_json::Value>,
123 ) -> Result<ErasedWorkflowView, ErasureError>;
124
125 fn operations(
127 &self,
128 case_ref: CaseRef,
129 state: Option<&serde_json::Value>,
130 ) -> Result<Vec<OperationSpec>, ErasureError>;
131 fn summary(&self) -> Option<String>;
133 fn glossary(&self) -> Vec<GlossaryTerm>;
135 fn noun(&self) -> Option<crate::locale::LocalizedText> {
137 None
138 }
139 fn record_operations(&self) -> Vec<OperationSpec> {
141 Vec::new()
142 }
143
144 fn briefing(
148 &self,
149 case_ref: CaseRef,
150 state: Option<&serde_json::Value>,
151 ) -> Result<Option<String>, ErasureError>;
152
153 fn start_preconditions(&self) -> Vec<StartPrecondition>;
159
160 fn start_behaviour(&self) -> StartBehaviour;
166
167 fn confirmation_subject(
173 &self,
174 case_ref: CaseRef,
175 state: Option<&serde_json::Value>,
176 act: &ResolvedAct,
177 ) -> Result<Option<ConfirmationSubject>, ErasureError>;
178
179 fn artifacts(&self, view: &ErasedWorkflowView) -> Result<Vec<ArtifactRef>, ErasureError>;
183
184 fn may_open_beside(&self, open: &[ErasedWorkflowView]) -> Result<(), ErasedCallError>;
192
193 fn narration_briefing(
200 &self,
201 stage: WritingStage,
202 view: &ErasedWorkflowView,
203 ) -> Result<Option<String>, ErasureError>;
204
205 fn enumerations(
211 &self,
212 case_ref: CaseRef,
213 state: Option<&serde_json::Value>,
214 ) -> Result<Vec<DomainEnumeration>, ErasureError>;
215
216 fn obligation_sentence(
219 &self,
220 obligation: &serde_json::Value,
221 ) -> Result<Option<crate::locale::LocalizedText>, ErasureError>;
222
223 fn next_steps(
226 &self,
227 case_ref: CaseRef,
228 state: Option<&serde_json::Value>,
229 ) -> Result<Vec<super::NextStep>, ErasureError> {
230 let _ = (case_ref, state);
231 Ok(Vec::new())
232 }
233
234 fn compile_act(
236 &self,
237 case_ref: CaseRef,
238 state: Option<&serde_json::Value>,
239 act: &ResolvedAct,
240 ) -> Result<Vec<serde_json::Value>, ErasedCallError>;
241
242 fn nothing_changed(
249 &self,
250 state: Option<&serde_json::Value>,
251 act: &ResolvedAct,
252 ) -> Option<crate::locale::LocalizedText> {
253 let _ = (state, act);
254 None
255 }
256
257 fn command_policy(
259 &self,
260 state: Option<&serde_json::Value>,
261 command: &serde_json::Value,
262 ) -> Result<CommandPolicy, ErasureError>;
263
264 fn validate_command(
266 &self,
267 state: Option<&serde_json::Value>,
268 command: &serde_json::Value,
269 ) -> Result<(), ErasedCallError>;
270
271 fn state_after(
273 &self,
274 state: Option<&serde_json::Value>,
275 command: &serde_json::Value,
276 ) -> Result<Option<serde_json::Value>, ErasureError> {
277 let _ = (state, command);
278 Ok(None)
279 }
280
281 fn receipts(
287 &self,
288 events: &[ReceiptEvent<serde_json::Value>],
289 locale: &Locale,
290 ) -> Result<Vec<OperationalReceipt>, ErasureError>;
291
292 fn build_interaction(
298 &self,
299 case_ref: CaseRef,
300 state: Option<&serde_json::Value>,
301 requirement: &InteractionRequirement,
302 ) -> Result<InteractionSpec, ErasedCallError>;
303}
304
305#[async_trait::async_trait]
318pub trait ErasedCaseLoader: Send + Sync {
319 async fn load_case(
325 &self,
326 account: &AccountId,
327 case_id: &CaseId,
328 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError>;
329}
330
331#[async_trait::async_trait]
333impl<E: ErasedExecutor + ?Sized> ErasedCaseLoader for E {
334 async fn load_case(
335 &self,
336 account: &AccountId,
337 case_id: &CaseId,
338 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
339 self.load(account, case_id).await
340 }
341}
342
343#[derive(Clone)]
351pub struct CaseLoaderHandle {
352 executor: Arc<dyn ErasedExecutor>,
353}
354
355impl CaseLoaderHandle {
356 #[must_use]
358 pub fn new(executor: Arc<dyn ErasedExecutor>) -> Self {
359 Self { executor }
360 }
361}
362
363impl fmt::Debug for CaseLoaderHandle {
364 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
365 f.debug_struct("CaseLoaderHandle").finish_non_exhaustive()
366 }
367}
368
369#[async_trait::async_trait]
370impl ErasedCaseLoader for CaseLoaderHandle {
371 async fn load_case(
372 &self,
373 account: &AccountId,
374 case_id: &CaseId,
375 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
376 self.executor.load(account, case_id).await
377 }
378}
379
380#[async_trait::async_trait]
382pub trait ErasedExecutor: Send + Sync {
383 async fn load(
385 &self,
386 account: &AccountId,
387 case_id: &CaseId,
388 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError>;
389
390 async fn execute(
392 &self,
393 batch: CommandBatch<serde_json::Value>,
394 ) -> Result<Commit<serde_json::Value, serde_json::Value>, ExecutionError>;
395}
396
397pub struct TypedWorkflowAdapter<W, E> {
399 definition: W,
400 executor: E,
401}
402
403impl<W, E> TypedWorkflowAdapter<W, E> {
404 #[must_use]
406 pub const fn new(definition: W, executor: E) -> Self {
407 Self {
408 definition,
409 executor,
410 }
411 }
412
413 #[must_use]
415 pub const fn definition(&self) -> &W {
416 &self.definition
417 }
418
419 #[must_use]
421 pub const fn executor(&self) -> &E {
422 &self.executor
423 }
424}
425
426impl<W: WorkflowDefinition, E> fmt::Debug for TypedWorkflowAdapter<W, E> {
427 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
428 f.debug_struct("TypedWorkflowAdapter")
429 .field("workflow", &self.definition.key())
430 .field("version", &self.definition.version())
431 .finish_non_exhaustive()
432 }
433}
434
435impl<W: WorkflowDefinition, E> TypedWorkflowAdapter<W, E> {
436 fn state(&self, value: Option<&serde_json::Value>) -> Result<Option<W::State>, ErasureError> {
437 value
438 .map(|v| {
439 W::State::deserialize(v).map_err(|_| ErasureError::StateDeserialization {
440 workflow: self.definition.key(),
441 })
442 })
443 .transpose()
444 }
445
446 fn command(&self, value: &serde_json::Value) -> Result<W::Command, ErasureError> {
447 W::Command::deserialize(value).map_err(|_| ErasureError::CommandDeserialization {
448 workflow: self.definition.key(),
449 })
450 }
451
452 fn serialize<T: Serialize>(&self, value: &T) -> Result<serde_json::Value, ErasureError> {
453 serde_json::to_value(value).map_err(|_| ErasureError::Serialization {
454 workflow: self.definition.key(),
455 })
456 }
457
458 fn typed_view(&self, case_ref: CaseRef, state: Option<&W::State>) -> ViewOf<W> {
459 self.definition.project(case_ref, state)
460 }
461
462 fn typed_view_from(&self, view: &ErasedWorkflowView) -> Result<ViewOf<W>, ErasureError> {
468 let mismatch = || ErasureError::StateDeserialization {
469 workflow: self.definition.key(),
470 };
471 let phase: W::Phase = serde_json::from_value(view.phase.clone()).map_err(|_| mismatch())?;
472 let mut obligations = Vec::with_capacity(view.obligations.len());
473 for obligation in &view.obligations {
474 obligations.push(
475 serde_json::from_value::<W::Obligation>(obligation.value.clone())
476 .map_err(|_| mismatch())?,
477 );
478 }
479 let outcome = view
480 .outcome
481 .clone()
482 .map(serde_json::from_value::<W::Outcome>)
483 .transpose()
484 .map_err(|_| mismatch())?;
485 Ok(WorkflowView {
486 case_ref: view.case_ref.clone(),
487 workflow_version: view.workflow_version.clone(),
488 phase,
489 obligations,
490 blocking_interaction: view.blocking_interaction.clone(),
491 notices: view.notices.clone(),
492 outcome,
493 })
494 }
495
496 fn check_same_case(&self, projected: &CaseRef, resolved: &CaseRef) -> Result<(), ErasureError> {
501 if projected.same_case(resolved)
502 && projected.expected_revision == resolved.expected_revision
503 {
504 Ok(())
505 } else {
506 Err(ErasureError::CaseMismatch {
507 workflow: self.definition.key(),
508 })
509 }
510 }
511}
512
513impl<W: WorkflowDefinition, E: Send + Sync> ErasedWorkflow for TypedWorkflowAdapter<W, E> {
514 fn key(&self) -> WorkflowKey {
515 self.definition.key()
516 }
517
518 fn version(&self) -> WorkflowVersion {
519 self.definition.version()
520 }
521
522 fn project(
523 &self,
524 case_ref: CaseRef,
525 state: Option<&serde_json::Value>,
526 ) -> Result<ErasedWorkflowView, ErasureError> {
527 let state = self.state(state)?;
528 let view = self.typed_view(case_ref, state.as_ref());
529 let ownership = self.definition.phase_ownership(&view.phase);
530 let mut erased = view
531 .erase(ownership)
532 .map_err(|_| ErasureError::Serialization {
533 workflow: self.definition.key(),
534 })?;
535 erased.state = self.definition.narratable_state(state.as_ref());
539 for obligation in &mut erased.obligations {
543 if let Ok(typed) = serde_json::from_value::<W::Obligation>(obligation.value.clone()) {
544 obligation.sentence = self.definition.obligation_sentence(&typed);
545 obligation.act = self.definition.obligation_act(state.as_ref(), &typed);
546 }
547 }
548 Ok(erased)
549 }
550
551 fn operations(
552 &self,
553 case_ref: CaseRef,
554 state: Option<&serde_json::Value>,
555 ) -> Result<Vec<OperationSpec>, ErasureError> {
556 let state = self.state(state)?;
557 let view = self.typed_view(case_ref, state.as_ref());
558 let workflow = self.definition.key();
559 let mut operations = self.definition.operations(&view);
560 for operation in &mut operations {
561 operation.workflow = workflow.clone();
562 operation
563 .validate()
564 .map_err(|error| ErasureError::InvalidOperation {
565 workflow: workflow.clone(),
566 reason: error.to_string(),
567 })?;
568 }
569 Ok(operations)
570 }
571
572 fn summary(&self) -> Option<String> {
573 self.definition.summary()
574 }
575
576 fn glossary(&self) -> Vec<GlossaryTerm> {
577 self.definition.glossary()
578 }
579
580 fn noun(&self) -> Option<crate::locale::LocalizedText> {
581 self.definition.noun()
582 }
583
584 fn record_operations(&self) -> Vec<OperationSpec> {
585 self.definition.record_operations()
586 }
587
588 fn briefing(
589 &self,
590 case_ref: CaseRef,
591 state: Option<&serde_json::Value>,
592 ) -> Result<Option<String>, ErasureError> {
593 let state = self.state(state)?;
594 let view = self.typed_view(case_ref, state.as_ref());
595 Ok(self.definition.briefing(&view))
596 }
597
598 fn start_preconditions(&self) -> Vec<StartPrecondition> {
599 self.definition.start_preconditions()
600 }
601
602 fn start_behaviour(&self) -> StartBehaviour {
603 self.definition.start_behaviour()
604 }
605
606 fn confirmation_subject(
607 &self,
608 case_ref: CaseRef,
609 state: Option<&serde_json::Value>,
610 act: &ResolvedAct,
611 ) -> Result<Option<ConfirmationSubject>, ErasureError> {
612 let state = self.state(state)?;
613 let view = self.typed_view(case_ref, state.as_ref());
614 Ok(self
615 .definition
616 .confirmation_subject(state.as_ref(), &view, act))
617 }
618
619 fn artifacts(&self, view: &ErasedWorkflowView) -> Result<Vec<ArtifactRef>, ErasureError> {
620 let typed = self.typed_view_from(view)?;
621 Ok(self.definition.artifacts(&typed))
622 }
623
624 fn may_open_beside(&self, open: &[ErasedWorkflowView]) -> Result<(), ErasedCallError> {
625 let typed: Vec<ViewOf<W>> = open
626 .iter()
627 .map(|view| self.typed_view_from(view))
628 .collect::<Result<_, _>>()?;
629 self.definition
630 .may_open_beside(&typed)
631 .map_err(ErasedCallError::from)
632 }
633
634 fn narration_briefing(
635 &self,
636 stage: WritingStage,
637 view: &ErasedWorkflowView,
638 ) -> Result<Option<String>, ErasureError> {
639 let typed = self.typed_view_from(view)?;
640 Ok(match stage {
641 WritingStage::Transition => self.definition.transition_briefing(&typed),
642 WritingStage::Answer => self.definition.answer_briefing(&typed),
643 })
644 }
645
646 fn enumerations(
647 &self,
648 case_ref: CaseRef,
649 state: Option<&serde_json::Value>,
650 ) -> Result<Vec<DomainEnumeration>, ErasureError> {
651 let state = self.state(state)?;
652 let view = self.typed_view(case_ref, state.as_ref());
653 Ok(self.definition.enumerations(&view))
654 }
655
656 fn next_steps(
657 &self,
658 case_ref: CaseRef,
659 state: Option<&serde_json::Value>,
660 ) -> Result<Vec<super::NextStep>, ErasureError> {
661 let state = self.state(state)?;
662 let view = self.typed_view(case_ref, state.as_ref());
663 Ok(self.definition.next_steps(state.as_ref(), &view))
664 }
665
666 fn obligation_sentence(
667 &self,
668 obligation: &serde_json::Value,
669 ) -> Result<Option<crate::locale::LocalizedText>, ErasureError> {
670 let typed: W::Obligation = serde_json::from_value(obligation.clone()).map_err(|_| {
671 ErasureError::StateDeserialization {
672 workflow: self.definition.key(),
673 }
674 })?;
675 Ok(self.definition.obligation_sentence(&typed))
676 }
677
678 fn compile_act(
679 &self,
680 case_ref: CaseRef,
681 state: Option<&serde_json::Value>,
682 act: &ResolvedAct,
683 ) -> Result<Vec<serde_json::Value>, ErasedCallError> {
684 self.check_same_case(&case_ref, &act.case_ref)?;
685 let state = self.state(state)?;
686 let view = self.typed_view(case_ref, state.as_ref());
687 let commands = self.definition.compile_act(state.as_ref(), &view, act)?;
688 commands
689 .iter()
690 .map(|c| self.serialize(c).map_err(ErasedCallError::from))
691 .collect()
692 }
693
694 fn nothing_changed(
695 &self,
696 state: Option<&serde_json::Value>,
697 act: &ResolvedAct,
698 ) -> Option<crate::locale::LocalizedText> {
699 let state = self.state(state).ok()?;
705 self.definition.nothing_changed(state.as_ref(), act)
706 }
707
708 fn command_policy(
709 &self,
710 state: Option<&serde_json::Value>,
711 command: &serde_json::Value,
712 ) -> Result<CommandPolicy, ErasureError> {
713 let state = self.state(state)?;
714 let command = self.command(command)?;
715 Ok(self.definition.command_policy(state.as_ref(), &command))
716 }
717
718 fn validate_command(
719 &self,
720 state: Option<&serde_json::Value>,
721 command: &serde_json::Value,
722 ) -> Result<(), ErasedCallError> {
723 let state = self.state(state)?;
724 let command = self.command(command)?;
725 self.definition
726 .validate_command(state.as_ref(), &command)
727 .map_err(ErasedCallError::from)
728 }
729
730 fn state_after(
731 &self,
732 state: Option<&serde_json::Value>,
733 command: &serde_json::Value,
734 ) -> Result<Option<serde_json::Value>, ErasureError> {
735 let state = self.state(state)?;
736 let command = self.command(command)?;
737 self.definition
738 .state_after(state.as_ref(), &command)
739 .map(|next| self.serialize(&next))
740 .transpose()
741 }
742
743 fn receipts(
744 &self,
745 events: &[ReceiptEvent<serde_json::Value>],
746 locale: &Locale,
747 ) -> Result<Vec<OperationalReceipt>, ErasureError> {
748 let events = events
753 .iter()
754 .map(|event| {
755 event.try_map_payload_ref(|payload| {
756 W::Event::deserialize(payload).map_err(|_| ErasureError::EventDeserialization {
757 workflow: self.definition.key(),
758 })
759 })
760 })
761 .collect::<Result<Vec<_>, _>>()?;
762 Ok(self.definition.receipts(&events, locale))
763 }
764
765 fn build_interaction(
766 &self,
767 case_ref: CaseRef,
768 state: Option<&serde_json::Value>,
769 requirement: &InteractionRequirement,
770 ) -> Result<InteractionSpec, ErasedCallError> {
771 let state = self.state(state)?;
772 let view = self.typed_view(case_ref.clone(), state.as_ref());
773 let spec = self
774 .definition
775 .build_interaction(state.as_ref(), &view, requirement)?;
776 self.check_same_case(&case_ref, &spec.case_ref)?;
779 spec.validate()?;
780 Ok(spec)
781 }
782}
783
784#[async_trait::async_trait]
785impl<W, E> ErasedExecutor for TypedWorkflowAdapter<W, E>
786where
787 W: WorkflowDefinition,
788 E: WorkflowExecutor<W>,
789{
790 async fn load(
791 &self,
792 account: &AccountId,
793 case_id: &CaseId,
794 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
795 let loaded = self.executor.load(account, case_id).await?;
796 let revision = loaded.revision;
797 let value = loaded
798 .value
799 .map(|s| serde_json::to_value(&s).map_err(|_| StoreError::Serialization))
800 .transpose()?;
801 Ok(Versioned::new(value, revision))
802 }
803
804 async fn execute(
805 &self,
806 batch: CommandBatch<serde_json::Value>,
807 ) -> Result<Commit<serde_json::Value, serde_json::Value>, ExecutionError> {
808 let typed = batch.try_map(|c| self.command(&c))?;
809 let commit = self.executor.execute(typed).await?;
810 let state = commit
811 .state
812 .as_ref()
813 .map(|s| self.serialize(s))
814 .transpose()?;
815 let events = commit
816 .events
817 .into_iter()
818 .map(|e: CommittedEvent<W::Event>| e.try_map_payload(|p| self.serialize(&p)))
819 .collect::<Result<Vec<_>, _>>()?;
820 Ok(Commit {
821 state,
822 new_revision: commit.new_revision,
823 events,
824 idempotency_replay: commit.idempotency_replay,
825 })
826 }
827}
828
829#[derive(Clone)]
831pub struct RegisteredWorkflow {
832 pub key: WorkflowKey,
834 pub version: WorkflowVersion,
836 pub definition: Arc<dyn ErasedWorkflow>,
838 pub executor: Arc<dyn ErasedExecutor>,
840}
841
842impl fmt::Debug for RegisteredWorkflow {
843 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
844 f.debug_struct("RegisteredWorkflow")
845 .field("key", &self.key)
846 .field("version", &self.version)
847 .finish_non_exhaustive()
848 }
849}
850
851impl RegisteredWorkflow {
852 #[must_use]
854 pub fn loader(&self) -> CaseLoaderHandle {
855 CaseLoaderHandle::new(Arc::clone(&self.executor))
856 }
857
858 pub fn check_version(&self, found: &WorkflowVersion) -> Result<(), ErasureError> {
864 if &self.version == found {
865 Ok(())
866 } else {
867 Err(ErasureError::VersionMismatch {
868 workflow: self.key.clone(),
869 registered: self.version.clone(),
870 found: found.clone(),
871 })
872 }
873 }
874}
875
876#[derive(Debug, Clone, Default)]
878pub struct WorkflowRegistry {
879 entries: IndexMap<WorkflowKey, RegisteredWorkflow>,
880}
881
882impl WorkflowRegistry {
883 #[must_use]
885 pub fn builder() -> WorkflowRegistryBuilder {
886 WorkflowRegistryBuilder::default()
887 }
888
889 #[must_use]
891 pub fn get(&self, key: &WorkflowKey) -> Option<&RegisteredWorkflow> {
892 self.entries.get(key)
893 }
894
895 pub fn require(&self, key: &WorkflowKey) -> Result<&RegisteredWorkflow, ErasureError> {
897 self.entries
898 .get(key)
899 .ok_or_else(|| ErasureError::UnknownWorkflow {
900 workflow: key.clone(),
901 })
902 }
903
904 pub fn check_version(
906 &self,
907 key: &WorkflowKey,
908 found: &WorkflowVersion,
909 ) -> Result<(), ErasureError> {
910 self.require(key)?.check_version(found)
911 }
912
913 #[must_use]
915 pub fn contains(&self, key: &WorkflowKey) -> bool {
916 self.entries.contains_key(key)
917 }
918
919 pub fn iter(&self) -> impl Iterator<Item = &RegisteredWorkflow> {
921 self.entries.values()
922 }
923
924 pub fn keys(&self) -> impl Iterator<Item = &WorkflowKey> {
926 self.entries.keys()
927 }
928
929 #[must_use]
931 pub fn len(&self) -> usize {
932 self.entries.len()
933 }
934
935 #[must_use]
937 pub fn is_empty(&self) -> bool {
938 self.entries.is_empty()
939 }
940
941 #[must_use]
948 pub fn definitions(&self) -> WorkflowDefinitions {
949 WorkflowDefinitions {
950 entries: self
951 .entries
952 .iter()
953 .map(|(key, registered)| (key.clone(), Arc::clone(®istered.definition)))
954 .collect(),
955 }
956 }
957
958 #[must_use]
965 pub fn read_only(&self) -> WorkflowReadRegistry {
966 WorkflowReadRegistry {
967 definitions: self.definitions(),
968 loaders: self
969 .entries
970 .iter()
971 .map(|(key, registered)| {
972 let loader: Arc<dyn ErasedCaseLoader> = Arc::new(registered.loader());
973 (key.clone(), loader)
974 })
975 .collect(),
976 }
977 }
978}
979
980#[derive(Clone, Default)]
986pub struct WorkflowDefinitions {
987 entries: IndexMap<WorkflowKey, Arc<dyn ErasedWorkflow>>,
988}
989
990impl fmt::Debug for WorkflowDefinitions {
991 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
992 f.debug_struct("WorkflowDefinitions")
993 .field("keys", &self.entries.keys().collect::<Vec<_>>())
994 .finish()
995 }
996}
997
998impl WorkflowDefinitions {
999 #[must_use]
1001 pub fn new() -> Self {
1002 Self::default()
1003 }
1004
1005 #[must_use]
1008 pub fn with(mut self, definition: Arc<dyn ErasedWorkflow>) -> Self {
1009 self.entries.insert(definition.key(), definition);
1010 self
1011 }
1012
1013 #[must_use]
1015 pub fn get(&self, key: &WorkflowKey) -> Option<&Arc<dyn ErasedWorkflow>> {
1016 self.entries.get(key)
1017 }
1018
1019 pub fn require(&self, key: &WorkflowKey) -> Result<&Arc<dyn ErasedWorkflow>, ErasureError> {
1025 self.entries
1026 .get(key)
1027 .ok_or_else(|| ErasureError::UnknownWorkflow {
1028 workflow: key.clone(),
1029 })
1030 }
1031
1032 #[must_use]
1034 pub fn contains(&self, key: &WorkflowKey) -> bool {
1035 self.entries.contains_key(key)
1036 }
1037
1038 pub fn keys(&self) -> impl Iterator<Item = &WorkflowKey> {
1040 self.entries.keys()
1041 }
1042
1043 pub fn iter(&self) -> impl Iterator<Item = &Arc<dyn ErasedWorkflow>> {
1045 self.entries.values()
1046 }
1047
1048 #[must_use]
1050 pub fn len(&self) -> usize {
1051 self.entries.len()
1052 }
1053
1054 #[must_use]
1056 pub fn is_empty(&self) -> bool {
1057 self.entries.is_empty()
1058 }
1059}
1060
1061impl From<&WorkflowRegistry> for WorkflowDefinitions {
1062 fn from(registry: &WorkflowRegistry) -> Self {
1063 registry.definitions()
1064 }
1065}
1066
1067impl From<&Arc<WorkflowRegistry>> for WorkflowDefinitions {
1068 fn from(registry: &Arc<WorkflowRegistry>) -> Self {
1069 registry.definitions()
1070 }
1071}
1072
1073impl From<Arc<WorkflowRegistry>> for WorkflowDefinitions {
1074 fn from(registry: Arc<WorkflowRegistry>) -> Self {
1075 registry.definitions()
1076 }
1077}
1078
1079#[derive(Clone, Default)]
1087pub struct WorkflowReadRegistry {
1088 definitions: WorkflowDefinitions,
1089 loaders: IndexMap<WorkflowKey, Arc<dyn ErasedCaseLoader>>,
1090}
1091
1092impl fmt::Debug for WorkflowReadRegistry {
1093 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1094 f.debug_struct("WorkflowReadRegistry")
1095 .field("definitions", &self.definitions)
1096 .field("loaders", &self.loaders.len())
1097 .finish()
1098 }
1099}
1100
1101impl WorkflowReadRegistry {
1102 #[must_use]
1104 pub fn from_definitions(definitions: WorkflowDefinitions) -> Self {
1105 Self {
1106 definitions,
1107 loaders: IndexMap::new(),
1108 }
1109 }
1110
1111 #[must_use]
1113 pub const fn definitions(&self) -> &WorkflowDefinitions {
1114 &self.definitions
1115 }
1116
1117 #[must_use]
1119 pub fn loader(&self, key: &WorkflowKey) -> Option<&Arc<dyn ErasedCaseLoader>> {
1120 self.loaders.get(key)
1121 }
1122
1123 pub fn require_loader(
1130 &self,
1131 key: &WorkflowKey,
1132 ) -> Result<&Arc<dyn ErasedCaseLoader>, ErasureError> {
1133 self.loaders
1134 .get(key)
1135 .ok_or_else(|| ErasureError::UnknownWorkflow {
1136 workflow: key.clone(),
1137 })
1138 }
1139
1140 #[must_use]
1142 pub fn has_no_loaders(&self) -> bool {
1143 self.loaders.is_empty()
1144 }
1145}
1146
1147impl From<&WorkflowRegistry> for WorkflowReadRegistry {
1148 fn from(registry: &WorkflowRegistry) -> Self {
1149 registry.read_only()
1150 }
1151}
1152
1153#[derive(Debug, Default)]
1155pub struct WorkflowRegistryBuilder {
1156 entries: Vec<RegisteredWorkflow>,
1157}
1158
1159impl WorkflowRegistryBuilder {
1160 #[must_use]
1162 pub fn register<W, E>(self, definition: W, executor: E) -> Self
1163 where
1164 W: WorkflowDefinition,
1165 E: WorkflowExecutor<W> + 'static,
1166 {
1167 let adapter = Arc::new(TypedWorkflowAdapter::new(definition, executor));
1168 let erased_definition: Arc<dyn ErasedWorkflow> = adapter.clone();
1169 let erased_executor: Arc<dyn ErasedExecutor> = adapter;
1170 self.register_erased(erased_definition, erased_executor)
1171 }
1172
1173 #[must_use]
1175 pub fn register_erased(
1176 mut self,
1177 definition: Arc<dyn ErasedWorkflow>,
1178 executor: Arc<dyn ErasedExecutor>,
1179 ) -> Self {
1180 self.entries.push(RegisteredWorkflow {
1181 key: definition.key(),
1182 version: definition.version(),
1183 definition,
1184 executor,
1185 });
1186 self
1187 }
1188
1189 pub fn build(self) -> Result<WorkflowRegistry, ErasureError> {
1191 let mut entries = IndexMap::with_capacity(self.entries.len());
1192 for entry in self.entries {
1193 if entries.contains_key(&entry.key) {
1194 return Err(ErasureError::DuplicateWorkflow {
1195 workflow: entry.key,
1196 });
1197 }
1198 entries.insert(entry.key.clone(), entry);
1199 }
1200 Ok(WorkflowRegistry { entries })
1201 }
1202}
1203
1204#[cfg(test)]
1205mod tests {
1206 use super::*;
1207 use crate::case::Versioned;
1208 use crate::flow::{PhaseOwnership, ViewOf, WorkflowDefinition, WorkflowView};
1209 use crate::ids::CaseRevision;
1210 use crate::locale::Locale;
1211 use crate::target::ResolvedAct;
1212 use serde::{Deserialize, Serialize};
1213
1214 #[derive(Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
1215 struct Nothing;
1216
1217 struct Toy;
1218
1219 impl WorkflowDefinition for Toy {
1220 type State = Nothing;
1221 type Phase = Nothing;
1222 type Obligation = Nothing;
1223 type Command = Nothing;
1224 type Event = Nothing;
1225 type Outcome = Nothing;
1226
1227 fn key(&self) -> WorkflowKey {
1228 WorkflowKey::from("toy")
1229 }
1230
1231 fn version(&self) -> WorkflowVersion {
1232 WorkflowVersion::from("1")
1233 }
1234
1235 fn phase_ownership(&self, _phase: &Self::Phase) -> PhaseOwnership {
1236 PhaseOwnership::System
1237 }
1238
1239 fn project(&self, case_ref: CaseRef, _state: Option<&Self::State>) -> ViewOf<Self> {
1240 WorkflowView::new(case_ref, self.version(), Nothing)
1241 }
1242
1243 fn operations(&self, _view: &ViewOf<Self>) -> Vec<crate::operation::OperationSpec> {
1244 Vec::new()
1245 }
1246
1247 fn compile_act(
1248 &self,
1249 _state: Option<&Self::State>,
1250 _view: &ViewOf<Self>,
1251 _act: &ResolvedAct,
1252 ) -> Result<Vec<Self::Command>, crate::error::DomainRejection> {
1253 Ok(Vec::new())
1254 }
1255
1256 fn command_policy(
1257 &self,
1258 _state: Option<&Self::State>,
1259 _command: &Self::Command,
1260 ) -> CommandPolicy {
1261 CommandPolicy::conservative()
1262 }
1263
1264 fn validate_command(
1265 &self,
1266 _state: Option<&Self::State>,
1267 _command: &Self::Command,
1268 ) -> Result<(), crate::error::DomainRejection> {
1269 Ok(())
1270 }
1271
1272 fn receipts(
1273 &self,
1274 _events: &[ReceiptEvent<Self::Event>],
1275 _locale: &Locale,
1276 ) -> Vec<OperationalReceipt> {
1277 Vec::new()
1278 }
1279 }
1280
1281 struct Executor;
1282
1283 #[async_trait::async_trait]
1284 impl crate::flow::WorkflowExecutor<Toy> for Executor {
1285 async fn load(
1286 &self,
1287 _account: &AccountId,
1288 _case_id: &CaseId,
1289 ) -> Result<Versioned<Option<Nothing>>, StoreError> {
1290 Ok(Versioned::new(Some(Nothing), CaseRevision(7)))
1291 }
1292
1293 async fn execute(
1294 &self,
1295 _batch: CommandBatch<Nothing>,
1296 ) -> Result<Commit<Nothing, Nothing>, ExecutionError> {
1297 panic!("the read-only projection must never reach execution");
1298 }
1299 }
1300
1301 fn registry() -> WorkflowRegistry {
1302 WorkflowRegistry::builder()
1303 .register(Toy, Executor)
1304 .build()
1305 .expect("one workflow, one key")
1306 }
1307
1308 #[tokio::test]
1309 async fn the_read_only_projection_loads_and_cannot_execute() {
1310 let registry = registry();
1311 let reading = registry.read_only();
1312 let key = WorkflowKey::from("toy");
1313
1314 assert!(reading.definitions().contains(&key));
1315 let loaded = reading
1316 .require_loader(&key)
1317 .expect("the projection carries a loader")
1318 .load_case(&AccountId::from("a"), &CaseId::from("c1"))
1319 .await
1320 .expect("the loader reads");
1321 assert_eq!(loaded.revision, CaseRevision(7));
1322 }
1325
1326 #[test]
1327 fn definitions_alone_carry_no_loader() {
1328 let definitions = registry().definitions();
1329 let reading = WorkflowReadRegistry::from_definitions(definitions);
1330 assert!(reading.has_no_loaders());
1331 assert!(reading.require_loader(&WorkflowKey::from("toy")).is_err());
1332 assert!(
1333 reading
1334 .definitions()
1335 .require(&WorkflowKey::from("toy"))
1336 .is_ok(),
1337 "the pure half is still there"
1338 );
1339 }
1340
1341 #[test]
1342 fn definitions_can_be_built_without_any_executor_at_all() {
1343 let definition: Arc<dyn ErasedWorkflow> = Arc::new(TypedWorkflowAdapter::new(Toy, ()));
1346 let definitions = WorkflowDefinitions::new().with(definition);
1347 assert_eq!(definitions.len(), 1);
1348 assert!(definitions.contains(&WorkflowKey::from("toy")));
1349 }
1350}