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
140 fn briefing(
144 &self,
145 case_ref: CaseRef,
146 state: Option<&serde_json::Value>,
147 ) -> Result<Option<String>, ErasureError>;
148
149 fn start_preconditions(&self) -> Vec<StartPrecondition>;
155
156 fn start_behaviour(&self) -> StartBehaviour;
162
163 fn confirmation_subject(
169 &self,
170 case_ref: CaseRef,
171 state: Option<&serde_json::Value>,
172 act: &ResolvedAct,
173 ) -> Result<Option<ConfirmationSubject>, ErasureError>;
174
175 fn artifacts(&self, view: &ErasedWorkflowView) -> Result<Vec<ArtifactRef>, ErasureError>;
179
180 fn may_open_beside(&self, open: &[ErasedWorkflowView]) -> Result<(), ErasedCallError>;
188
189 fn narration_briefing(
196 &self,
197 stage: WritingStage,
198 view: &ErasedWorkflowView,
199 ) -> Result<Option<String>, ErasureError>;
200
201 fn enumerations(
207 &self,
208 case_ref: CaseRef,
209 state: Option<&serde_json::Value>,
210 ) -> Result<Vec<DomainEnumeration>, ErasureError>;
211
212 fn obligation_sentence(
215 &self,
216 obligation: &serde_json::Value,
217 ) -> Result<Option<crate::locale::LocalizedText>, ErasureError>;
218
219 fn next_steps(
222 &self,
223 view: &ErasedWorkflowView,
224 ) -> Result<Vec<crate::locale::LocalizedText>, ErasureError> {
225 let _ = view;
226 Ok(Vec::new())
227 }
228
229 fn compile_act(
231 &self,
232 case_ref: CaseRef,
233 state: Option<&serde_json::Value>,
234 act: &ResolvedAct,
235 ) -> Result<Vec<serde_json::Value>, ErasedCallError>;
236
237 fn nothing_changed(
244 &self,
245 state: Option<&serde_json::Value>,
246 act: &ResolvedAct,
247 ) -> Option<crate::locale::LocalizedText> {
248 let _ = (state, act);
249 None
250 }
251
252 fn command_policy(
254 &self,
255 state: Option<&serde_json::Value>,
256 command: &serde_json::Value,
257 ) -> Result<CommandPolicy, ErasureError>;
258
259 fn validate_command(
261 &self,
262 state: Option<&serde_json::Value>,
263 command: &serde_json::Value,
264 ) -> Result<(), ErasedCallError>;
265
266 fn receipts(
272 &self,
273 events: &[ReceiptEvent<serde_json::Value>],
274 locale: &Locale,
275 ) -> Result<Vec<OperationalReceipt>, ErasureError>;
276
277 fn build_interaction(
283 &self,
284 case_ref: CaseRef,
285 state: Option<&serde_json::Value>,
286 requirement: &InteractionRequirement,
287 ) -> Result<InteractionSpec, ErasedCallError>;
288}
289
290#[async_trait::async_trait]
303pub trait ErasedCaseLoader: Send + Sync {
304 async fn load_case(
310 &self,
311 account: &AccountId,
312 case_id: &CaseId,
313 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError>;
314}
315
316#[async_trait::async_trait]
318impl<E: ErasedExecutor + ?Sized> ErasedCaseLoader for E {
319 async fn load_case(
320 &self,
321 account: &AccountId,
322 case_id: &CaseId,
323 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
324 self.load(account, case_id).await
325 }
326}
327
328#[derive(Clone)]
336pub struct CaseLoaderHandle {
337 executor: Arc<dyn ErasedExecutor>,
338}
339
340impl CaseLoaderHandle {
341 #[must_use]
343 pub fn new(executor: Arc<dyn ErasedExecutor>) -> Self {
344 Self { executor }
345 }
346}
347
348impl fmt::Debug for CaseLoaderHandle {
349 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
350 f.debug_struct("CaseLoaderHandle").finish_non_exhaustive()
351 }
352}
353
354#[async_trait::async_trait]
355impl ErasedCaseLoader for CaseLoaderHandle {
356 async fn load_case(
357 &self,
358 account: &AccountId,
359 case_id: &CaseId,
360 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
361 self.executor.load(account, case_id).await
362 }
363}
364
365#[async_trait::async_trait]
367pub trait ErasedExecutor: Send + Sync {
368 async fn load(
370 &self,
371 account: &AccountId,
372 case_id: &CaseId,
373 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError>;
374
375 async fn execute(
377 &self,
378 batch: CommandBatch<serde_json::Value>,
379 ) -> Result<Commit<serde_json::Value, serde_json::Value>, ExecutionError>;
380}
381
382pub struct TypedWorkflowAdapter<W, E> {
384 definition: W,
385 executor: E,
386}
387
388impl<W, E> TypedWorkflowAdapter<W, E> {
389 #[must_use]
391 pub const fn new(definition: W, executor: E) -> Self {
392 Self {
393 definition,
394 executor,
395 }
396 }
397
398 #[must_use]
400 pub const fn definition(&self) -> &W {
401 &self.definition
402 }
403
404 #[must_use]
406 pub const fn executor(&self) -> &E {
407 &self.executor
408 }
409}
410
411impl<W: WorkflowDefinition, E> fmt::Debug for TypedWorkflowAdapter<W, E> {
412 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
413 f.debug_struct("TypedWorkflowAdapter")
414 .field("workflow", &self.definition.key())
415 .field("version", &self.definition.version())
416 .finish_non_exhaustive()
417 }
418}
419
420impl<W: WorkflowDefinition, E> TypedWorkflowAdapter<W, E> {
421 fn state(&self, value: Option<&serde_json::Value>) -> Result<Option<W::State>, ErasureError> {
422 value
423 .map(|v| {
424 W::State::deserialize(v).map_err(|_| ErasureError::StateDeserialization {
425 workflow: self.definition.key(),
426 })
427 })
428 .transpose()
429 }
430
431 fn command(&self, value: &serde_json::Value) -> Result<W::Command, ErasureError> {
432 W::Command::deserialize(value).map_err(|_| ErasureError::CommandDeserialization {
433 workflow: self.definition.key(),
434 })
435 }
436
437 fn serialize<T: Serialize>(&self, value: &T) -> Result<serde_json::Value, ErasureError> {
438 serde_json::to_value(value).map_err(|_| ErasureError::Serialization {
439 workflow: self.definition.key(),
440 })
441 }
442
443 fn typed_view(&self, case_ref: CaseRef, state: Option<&W::State>) -> ViewOf<W> {
444 self.definition.project(case_ref, state)
445 }
446
447 fn typed_view_from(&self, view: &ErasedWorkflowView) -> Result<ViewOf<W>, ErasureError> {
453 let mismatch = || ErasureError::StateDeserialization {
454 workflow: self.definition.key(),
455 };
456 let phase: W::Phase = serde_json::from_value(view.phase.clone()).map_err(|_| mismatch())?;
457 let mut obligations = Vec::with_capacity(view.obligations.len());
458 for obligation in &view.obligations {
459 obligations.push(
460 serde_json::from_value::<W::Obligation>(obligation.value.clone())
461 .map_err(|_| mismatch())?,
462 );
463 }
464 let outcome = view
465 .outcome
466 .clone()
467 .map(serde_json::from_value::<W::Outcome>)
468 .transpose()
469 .map_err(|_| mismatch())?;
470 Ok(WorkflowView {
471 case_ref: view.case_ref.clone(),
472 workflow_version: view.workflow_version.clone(),
473 phase,
474 obligations,
475 blocking_interaction: view.blocking_interaction.clone(),
476 notices: view.notices.clone(),
477 outcome,
478 })
479 }
480
481 fn check_same_case(&self, projected: &CaseRef, resolved: &CaseRef) -> Result<(), ErasureError> {
486 if projected.same_case(resolved)
487 && projected.expected_revision == resolved.expected_revision
488 {
489 Ok(())
490 } else {
491 Err(ErasureError::CaseMismatch {
492 workflow: self.definition.key(),
493 })
494 }
495 }
496}
497
498impl<W: WorkflowDefinition, E: Send + Sync> ErasedWorkflow for TypedWorkflowAdapter<W, E> {
499 fn key(&self) -> WorkflowKey {
500 self.definition.key()
501 }
502
503 fn version(&self) -> WorkflowVersion {
504 self.definition.version()
505 }
506
507 fn project(
508 &self,
509 case_ref: CaseRef,
510 state: Option<&serde_json::Value>,
511 ) -> Result<ErasedWorkflowView, ErasureError> {
512 let state = self.state(state)?;
513 let view = self.typed_view(case_ref, state.as_ref());
514 let ownership = self.definition.phase_ownership(&view.phase);
515 let mut erased = view
516 .erase(ownership)
517 .map_err(|_| ErasureError::Serialization {
518 workflow: self.definition.key(),
519 })?;
520 erased.state = self.definition.narratable_state(state.as_ref());
524 for obligation in &mut erased.obligations {
528 if let Ok(typed) = serde_json::from_value::<W::Obligation>(obligation.value.clone()) {
529 obligation.sentence = self.definition.obligation_sentence(&typed);
530 obligation.act = self.definition.obligation_act(state.as_ref(), &typed);
531 }
532 }
533 Ok(erased)
534 }
535
536 fn operations(
537 &self,
538 case_ref: CaseRef,
539 state: Option<&serde_json::Value>,
540 ) -> Result<Vec<OperationSpec>, ErasureError> {
541 let state = self.state(state)?;
542 let view = self.typed_view(case_ref, state.as_ref());
543 let workflow = self.definition.key();
544 let mut operations = self.definition.operations(&view);
545 for operation in &mut operations {
546 operation.workflow = workflow.clone();
547 operation
548 .validate()
549 .map_err(|error| ErasureError::InvalidOperation {
550 workflow: workflow.clone(),
551 reason: error.to_string(),
552 })?;
553 }
554 Ok(operations)
555 }
556
557 fn summary(&self) -> Option<String> {
558 self.definition.summary()
559 }
560
561 fn glossary(&self) -> Vec<GlossaryTerm> {
562 self.definition.glossary()
563 }
564
565 fn noun(&self) -> Option<crate::locale::LocalizedText> {
566 self.definition.noun()
567 }
568
569 fn briefing(
570 &self,
571 case_ref: CaseRef,
572 state: Option<&serde_json::Value>,
573 ) -> Result<Option<String>, ErasureError> {
574 let state = self.state(state)?;
575 let view = self.typed_view(case_ref, state.as_ref());
576 Ok(self.definition.briefing(&view))
577 }
578
579 fn start_preconditions(&self) -> Vec<StartPrecondition> {
580 self.definition.start_preconditions()
581 }
582
583 fn start_behaviour(&self) -> StartBehaviour {
584 self.definition.start_behaviour()
585 }
586
587 fn confirmation_subject(
588 &self,
589 case_ref: CaseRef,
590 state: Option<&serde_json::Value>,
591 act: &ResolvedAct,
592 ) -> Result<Option<ConfirmationSubject>, ErasureError> {
593 let state = self.state(state)?;
594 let view = self.typed_view(case_ref, state.as_ref());
595 Ok(self
596 .definition
597 .confirmation_subject(state.as_ref(), &view, act))
598 }
599
600 fn artifacts(&self, view: &ErasedWorkflowView) -> Result<Vec<ArtifactRef>, ErasureError> {
601 let typed = self.typed_view_from(view)?;
602 Ok(self.definition.artifacts(&typed))
603 }
604
605 fn may_open_beside(&self, open: &[ErasedWorkflowView]) -> Result<(), ErasedCallError> {
606 let typed: Vec<ViewOf<W>> = open
607 .iter()
608 .map(|view| self.typed_view_from(view))
609 .collect::<Result<_, _>>()?;
610 self.definition
611 .may_open_beside(&typed)
612 .map_err(ErasedCallError::from)
613 }
614
615 fn narration_briefing(
616 &self,
617 stage: WritingStage,
618 view: &ErasedWorkflowView,
619 ) -> Result<Option<String>, ErasureError> {
620 let typed = self.typed_view_from(view)?;
621 Ok(match stage {
622 WritingStage::Transition => self.definition.transition_briefing(&typed),
623 WritingStage::Answer => self.definition.answer_briefing(&typed),
624 })
625 }
626
627 fn enumerations(
628 &self,
629 case_ref: CaseRef,
630 state: Option<&serde_json::Value>,
631 ) -> Result<Vec<DomainEnumeration>, ErasureError> {
632 let state = self.state(state)?;
633 let view = self.typed_view(case_ref, state.as_ref());
634 Ok(self.definition.enumerations(&view))
635 }
636
637 fn next_steps(
638 &self,
639 view: &ErasedWorkflowView,
640 ) -> Result<Vec<crate::locale::LocalizedText>, ErasureError> {
641 let typed = self.typed_view_from(view)?;
642 Ok(self.definition.next_steps(&typed))
643 }
644
645 fn obligation_sentence(
646 &self,
647 obligation: &serde_json::Value,
648 ) -> Result<Option<crate::locale::LocalizedText>, ErasureError> {
649 let typed: W::Obligation = serde_json::from_value(obligation.clone()).map_err(|_| {
650 ErasureError::StateDeserialization {
651 workflow: self.definition.key(),
652 }
653 })?;
654 Ok(self.definition.obligation_sentence(&typed))
655 }
656
657 fn compile_act(
658 &self,
659 case_ref: CaseRef,
660 state: Option<&serde_json::Value>,
661 act: &ResolvedAct,
662 ) -> Result<Vec<serde_json::Value>, ErasedCallError> {
663 self.check_same_case(&case_ref, &act.case_ref)?;
664 let state = self.state(state)?;
665 let view = self.typed_view(case_ref, state.as_ref());
666 let commands = self.definition.compile_act(state.as_ref(), &view, act)?;
667 commands
668 .iter()
669 .map(|c| self.serialize(c).map_err(ErasedCallError::from))
670 .collect()
671 }
672
673 fn nothing_changed(
674 &self,
675 state: Option<&serde_json::Value>,
676 act: &ResolvedAct,
677 ) -> Option<crate::locale::LocalizedText> {
678 let state = self.state(state).ok()?;
684 self.definition.nothing_changed(state.as_ref(), act)
685 }
686
687 fn command_policy(
688 &self,
689 state: Option<&serde_json::Value>,
690 command: &serde_json::Value,
691 ) -> Result<CommandPolicy, ErasureError> {
692 let state = self.state(state)?;
693 let command = self.command(command)?;
694 Ok(self.definition.command_policy(state.as_ref(), &command))
695 }
696
697 fn validate_command(
698 &self,
699 state: Option<&serde_json::Value>,
700 command: &serde_json::Value,
701 ) -> Result<(), ErasedCallError> {
702 let state = self.state(state)?;
703 let command = self.command(command)?;
704 self.definition
705 .validate_command(state.as_ref(), &command)
706 .map_err(ErasedCallError::from)
707 }
708
709 fn receipts(
710 &self,
711 events: &[ReceiptEvent<serde_json::Value>],
712 locale: &Locale,
713 ) -> Result<Vec<OperationalReceipt>, ErasureError> {
714 let events = events
719 .iter()
720 .map(|event| {
721 event.try_map_payload_ref(|payload| {
722 W::Event::deserialize(payload).map_err(|_| ErasureError::EventDeserialization {
723 workflow: self.definition.key(),
724 })
725 })
726 })
727 .collect::<Result<Vec<_>, _>>()?;
728 Ok(self.definition.receipts(&events, locale))
729 }
730
731 fn build_interaction(
732 &self,
733 case_ref: CaseRef,
734 state: Option<&serde_json::Value>,
735 requirement: &InteractionRequirement,
736 ) -> Result<InteractionSpec, ErasedCallError> {
737 let state = self.state(state)?;
738 let view = self.typed_view(case_ref.clone(), state.as_ref());
739 let spec = self
740 .definition
741 .build_interaction(state.as_ref(), &view, requirement)?;
742 self.check_same_case(&case_ref, &spec.case_ref)?;
745 spec.validate()?;
746 Ok(spec)
747 }
748}
749
750#[async_trait::async_trait]
751impl<W, E> ErasedExecutor for TypedWorkflowAdapter<W, E>
752where
753 W: WorkflowDefinition,
754 E: WorkflowExecutor<W>,
755{
756 async fn load(
757 &self,
758 account: &AccountId,
759 case_id: &CaseId,
760 ) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
761 let loaded = self.executor.load(account, case_id).await?;
762 let revision = loaded.revision;
763 let value = loaded
764 .value
765 .map(|s| serde_json::to_value(&s).map_err(|_| StoreError::Serialization))
766 .transpose()?;
767 Ok(Versioned::new(value, revision))
768 }
769
770 async fn execute(
771 &self,
772 batch: CommandBatch<serde_json::Value>,
773 ) -> Result<Commit<serde_json::Value, serde_json::Value>, ExecutionError> {
774 let typed = batch.try_map(|c| self.command(&c))?;
775 let commit = self.executor.execute(typed).await?;
776 let state = commit
777 .state
778 .as_ref()
779 .map(|s| self.serialize(s))
780 .transpose()?;
781 let events = commit
782 .events
783 .into_iter()
784 .map(|e: CommittedEvent<W::Event>| e.try_map_payload(|p| self.serialize(&p)))
785 .collect::<Result<Vec<_>, _>>()?;
786 Ok(Commit {
787 state,
788 new_revision: commit.new_revision,
789 events,
790 idempotency_replay: commit.idempotency_replay,
791 })
792 }
793}
794
795#[derive(Clone)]
797pub struct RegisteredWorkflow {
798 pub key: WorkflowKey,
800 pub version: WorkflowVersion,
802 pub definition: Arc<dyn ErasedWorkflow>,
804 pub executor: Arc<dyn ErasedExecutor>,
806}
807
808impl fmt::Debug for RegisteredWorkflow {
809 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
810 f.debug_struct("RegisteredWorkflow")
811 .field("key", &self.key)
812 .field("version", &self.version)
813 .finish_non_exhaustive()
814 }
815}
816
817impl RegisteredWorkflow {
818 #[must_use]
820 pub fn loader(&self) -> CaseLoaderHandle {
821 CaseLoaderHandle::new(Arc::clone(&self.executor))
822 }
823
824 pub fn check_version(&self, found: &WorkflowVersion) -> Result<(), ErasureError> {
830 if &self.version == found {
831 Ok(())
832 } else {
833 Err(ErasureError::VersionMismatch {
834 workflow: self.key.clone(),
835 registered: self.version.clone(),
836 found: found.clone(),
837 })
838 }
839 }
840}
841
842#[derive(Debug, Clone, Default)]
844pub struct WorkflowRegistry {
845 entries: IndexMap<WorkflowKey, RegisteredWorkflow>,
846}
847
848impl WorkflowRegistry {
849 #[must_use]
851 pub fn builder() -> WorkflowRegistryBuilder {
852 WorkflowRegistryBuilder::default()
853 }
854
855 #[must_use]
857 pub fn get(&self, key: &WorkflowKey) -> Option<&RegisteredWorkflow> {
858 self.entries.get(key)
859 }
860
861 pub fn require(&self, key: &WorkflowKey) -> Result<&RegisteredWorkflow, ErasureError> {
863 self.entries
864 .get(key)
865 .ok_or_else(|| ErasureError::UnknownWorkflow {
866 workflow: key.clone(),
867 })
868 }
869
870 pub fn check_version(
872 &self,
873 key: &WorkflowKey,
874 found: &WorkflowVersion,
875 ) -> Result<(), ErasureError> {
876 self.require(key)?.check_version(found)
877 }
878
879 #[must_use]
881 pub fn contains(&self, key: &WorkflowKey) -> bool {
882 self.entries.contains_key(key)
883 }
884
885 pub fn iter(&self) -> impl Iterator<Item = &RegisteredWorkflow> {
887 self.entries.values()
888 }
889
890 pub fn keys(&self) -> impl Iterator<Item = &WorkflowKey> {
892 self.entries.keys()
893 }
894
895 #[must_use]
897 pub fn len(&self) -> usize {
898 self.entries.len()
899 }
900
901 #[must_use]
903 pub fn is_empty(&self) -> bool {
904 self.entries.is_empty()
905 }
906
907 #[must_use]
914 pub fn definitions(&self) -> WorkflowDefinitions {
915 WorkflowDefinitions {
916 entries: self
917 .entries
918 .iter()
919 .map(|(key, registered)| (key.clone(), Arc::clone(®istered.definition)))
920 .collect(),
921 }
922 }
923
924 #[must_use]
931 pub fn read_only(&self) -> WorkflowReadRegistry {
932 WorkflowReadRegistry {
933 definitions: self.definitions(),
934 loaders: self
935 .entries
936 .iter()
937 .map(|(key, registered)| {
938 let loader: Arc<dyn ErasedCaseLoader> = Arc::new(registered.loader());
939 (key.clone(), loader)
940 })
941 .collect(),
942 }
943 }
944}
945
946#[derive(Clone, Default)]
952pub struct WorkflowDefinitions {
953 entries: IndexMap<WorkflowKey, Arc<dyn ErasedWorkflow>>,
954}
955
956impl fmt::Debug for WorkflowDefinitions {
957 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
958 f.debug_struct("WorkflowDefinitions")
959 .field("keys", &self.entries.keys().collect::<Vec<_>>())
960 .finish()
961 }
962}
963
964impl WorkflowDefinitions {
965 #[must_use]
967 pub fn new() -> Self {
968 Self::default()
969 }
970
971 #[must_use]
974 pub fn with(mut self, definition: Arc<dyn ErasedWorkflow>) -> Self {
975 self.entries.insert(definition.key(), definition);
976 self
977 }
978
979 #[must_use]
981 pub fn get(&self, key: &WorkflowKey) -> Option<&Arc<dyn ErasedWorkflow>> {
982 self.entries.get(key)
983 }
984
985 pub fn require(&self, key: &WorkflowKey) -> Result<&Arc<dyn ErasedWorkflow>, ErasureError> {
991 self.entries
992 .get(key)
993 .ok_or_else(|| ErasureError::UnknownWorkflow {
994 workflow: key.clone(),
995 })
996 }
997
998 #[must_use]
1000 pub fn contains(&self, key: &WorkflowKey) -> bool {
1001 self.entries.contains_key(key)
1002 }
1003
1004 pub fn keys(&self) -> impl Iterator<Item = &WorkflowKey> {
1006 self.entries.keys()
1007 }
1008
1009 pub fn iter(&self) -> impl Iterator<Item = &Arc<dyn ErasedWorkflow>> {
1011 self.entries.values()
1012 }
1013
1014 #[must_use]
1016 pub fn len(&self) -> usize {
1017 self.entries.len()
1018 }
1019
1020 #[must_use]
1022 pub fn is_empty(&self) -> bool {
1023 self.entries.is_empty()
1024 }
1025}
1026
1027impl From<&WorkflowRegistry> for WorkflowDefinitions {
1028 fn from(registry: &WorkflowRegistry) -> Self {
1029 registry.definitions()
1030 }
1031}
1032
1033impl From<&Arc<WorkflowRegistry>> for WorkflowDefinitions {
1034 fn from(registry: &Arc<WorkflowRegistry>) -> Self {
1035 registry.definitions()
1036 }
1037}
1038
1039impl From<Arc<WorkflowRegistry>> for WorkflowDefinitions {
1040 fn from(registry: Arc<WorkflowRegistry>) -> Self {
1041 registry.definitions()
1042 }
1043}
1044
1045#[derive(Clone, Default)]
1053pub struct WorkflowReadRegistry {
1054 definitions: WorkflowDefinitions,
1055 loaders: IndexMap<WorkflowKey, Arc<dyn ErasedCaseLoader>>,
1056}
1057
1058impl fmt::Debug for WorkflowReadRegistry {
1059 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1060 f.debug_struct("WorkflowReadRegistry")
1061 .field("definitions", &self.definitions)
1062 .field("loaders", &self.loaders.len())
1063 .finish()
1064 }
1065}
1066
1067impl WorkflowReadRegistry {
1068 #[must_use]
1070 pub fn from_definitions(definitions: WorkflowDefinitions) -> Self {
1071 Self {
1072 definitions,
1073 loaders: IndexMap::new(),
1074 }
1075 }
1076
1077 #[must_use]
1079 pub const fn definitions(&self) -> &WorkflowDefinitions {
1080 &self.definitions
1081 }
1082
1083 #[must_use]
1085 pub fn loader(&self, key: &WorkflowKey) -> Option<&Arc<dyn ErasedCaseLoader>> {
1086 self.loaders.get(key)
1087 }
1088
1089 pub fn require_loader(
1096 &self,
1097 key: &WorkflowKey,
1098 ) -> Result<&Arc<dyn ErasedCaseLoader>, ErasureError> {
1099 self.loaders
1100 .get(key)
1101 .ok_or_else(|| ErasureError::UnknownWorkflow {
1102 workflow: key.clone(),
1103 })
1104 }
1105
1106 #[must_use]
1108 pub fn has_no_loaders(&self) -> bool {
1109 self.loaders.is_empty()
1110 }
1111}
1112
1113impl From<&WorkflowRegistry> for WorkflowReadRegistry {
1114 fn from(registry: &WorkflowRegistry) -> Self {
1115 registry.read_only()
1116 }
1117}
1118
1119#[derive(Debug, Default)]
1121pub struct WorkflowRegistryBuilder {
1122 entries: Vec<RegisteredWorkflow>,
1123}
1124
1125impl WorkflowRegistryBuilder {
1126 #[must_use]
1128 pub fn register<W, E>(self, definition: W, executor: E) -> Self
1129 where
1130 W: WorkflowDefinition,
1131 E: WorkflowExecutor<W> + 'static,
1132 {
1133 let adapter = Arc::new(TypedWorkflowAdapter::new(definition, executor));
1134 let erased_definition: Arc<dyn ErasedWorkflow> = adapter.clone();
1135 let erased_executor: Arc<dyn ErasedExecutor> = adapter;
1136 self.register_erased(erased_definition, erased_executor)
1137 }
1138
1139 #[must_use]
1141 pub fn register_erased(
1142 mut self,
1143 definition: Arc<dyn ErasedWorkflow>,
1144 executor: Arc<dyn ErasedExecutor>,
1145 ) -> Self {
1146 self.entries.push(RegisteredWorkflow {
1147 key: definition.key(),
1148 version: definition.version(),
1149 definition,
1150 executor,
1151 });
1152 self
1153 }
1154
1155 pub fn build(self) -> Result<WorkflowRegistry, ErasureError> {
1157 let mut entries = IndexMap::with_capacity(self.entries.len());
1158 for entry in self.entries {
1159 if entries.contains_key(&entry.key) {
1160 return Err(ErasureError::DuplicateWorkflow {
1161 workflow: entry.key,
1162 });
1163 }
1164 entries.insert(entry.key.clone(), entry);
1165 }
1166 Ok(WorkflowRegistry { entries })
1167 }
1168}
1169
1170#[cfg(test)]
1171mod tests {
1172 use super::*;
1173 use crate::case::Versioned;
1174 use crate::flow::{PhaseOwnership, ViewOf, WorkflowDefinition, WorkflowView};
1175 use crate::ids::CaseRevision;
1176 use crate::locale::Locale;
1177 use crate::target::ResolvedAct;
1178 use serde::{Deserialize, Serialize};
1179
1180 #[derive(Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
1181 struct Nothing;
1182
1183 struct Toy;
1184
1185 impl WorkflowDefinition for Toy {
1186 type State = Nothing;
1187 type Phase = Nothing;
1188 type Obligation = Nothing;
1189 type Command = Nothing;
1190 type Event = Nothing;
1191 type Outcome = Nothing;
1192
1193 fn key(&self) -> WorkflowKey {
1194 WorkflowKey::from("toy")
1195 }
1196
1197 fn version(&self) -> WorkflowVersion {
1198 WorkflowVersion::from("1")
1199 }
1200
1201 fn phase_ownership(&self, _phase: &Self::Phase) -> PhaseOwnership {
1202 PhaseOwnership::System
1203 }
1204
1205 fn project(&self, case_ref: CaseRef, _state: Option<&Self::State>) -> ViewOf<Self> {
1206 WorkflowView::new(case_ref, self.version(), Nothing)
1207 }
1208
1209 fn operations(&self, _view: &ViewOf<Self>) -> Vec<crate::operation::OperationSpec> {
1210 Vec::new()
1211 }
1212
1213 fn compile_act(
1214 &self,
1215 _state: Option<&Self::State>,
1216 _view: &ViewOf<Self>,
1217 _act: &ResolvedAct,
1218 ) -> Result<Vec<Self::Command>, crate::error::DomainRejection> {
1219 Ok(Vec::new())
1220 }
1221
1222 fn command_policy(
1223 &self,
1224 _state: Option<&Self::State>,
1225 _command: &Self::Command,
1226 ) -> CommandPolicy {
1227 CommandPolicy::conservative()
1228 }
1229
1230 fn validate_command(
1231 &self,
1232 _state: Option<&Self::State>,
1233 _command: &Self::Command,
1234 ) -> Result<(), crate::error::DomainRejection> {
1235 Ok(())
1236 }
1237
1238 fn receipts(
1239 &self,
1240 _events: &[ReceiptEvent<Self::Event>],
1241 _locale: &Locale,
1242 ) -> Vec<OperationalReceipt> {
1243 Vec::new()
1244 }
1245 }
1246
1247 struct Executor;
1248
1249 #[async_trait::async_trait]
1250 impl crate::flow::WorkflowExecutor<Toy> for Executor {
1251 async fn load(
1252 &self,
1253 _account: &AccountId,
1254 _case_id: &CaseId,
1255 ) -> Result<Versioned<Option<Nothing>>, StoreError> {
1256 Ok(Versioned::new(Some(Nothing), CaseRevision(7)))
1257 }
1258
1259 async fn execute(
1260 &self,
1261 _batch: CommandBatch<Nothing>,
1262 ) -> Result<Commit<Nothing, Nothing>, ExecutionError> {
1263 panic!("the read-only projection must never reach execution");
1264 }
1265 }
1266
1267 fn registry() -> WorkflowRegistry {
1268 WorkflowRegistry::builder()
1269 .register(Toy, Executor)
1270 .build()
1271 .expect("one workflow, one key")
1272 }
1273
1274 #[tokio::test]
1275 async fn the_read_only_projection_loads_and_cannot_execute() {
1276 let registry = registry();
1277 let reading = registry.read_only();
1278 let key = WorkflowKey::from("toy");
1279
1280 assert!(reading.definitions().contains(&key));
1281 let loaded = reading
1282 .require_loader(&key)
1283 .expect("the projection carries a loader")
1284 .load_case(&AccountId::from("a"), &CaseId::from("c1"))
1285 .await
1286 .expect("the loader reads");
1287 assert_eq!(loaded.revision, CaseRevision(7));
1288 }
1291
1292 #[test]
1293 fn definitions_alone_carry_no_loader() {
1294 let definitions = registry().definitions();
1295 let reading = WorkflowReadRegistry::from_definitions(definitions);
1296 assert!(reading.has_no_loaders());
1297 assert!(reading.require_loader(&WorkflowKey::from("toy")).is_err());
1298 assert!(
1299 reading
1300 .definitions()
1301 .require(&WorkflowKey::from("toy"))
1302 .is_ok(),
1303 "the pure half is still there"
1304 );
1305 }
1306
1307 #[test]
1308 fn definitions_can_be_built_without_any_executor_at_all() {
1309 let definition: Arc<dyn ErasedWorkflow> = Arc::new(TypedWorkflowAdapter::new(Toy, ()));
1312 let definitions = WorkflowDefinitions::new().with(definition);
1313 assert_eq!(definitions.len(), 1);
1314 assert!(definitions.contains(&WorkflowKey::from("toy")));
1315 }
1316}