use std::fmt;
use std::sync::Arc;
use indexmap::IndexMap;
use serde::{Deserialize, Serialize};
use crate::case::{CaseRef, Versioned};
use crate::command::{CommandBatch, CommandPolicy};
use crate::error::{ErasedCallError, ErasureError, ExecutionError, StoreError};
use crate::event::{ArtifactRef, Commit, CommittedEvent, OperationalReceipt, ReceiptEvent};
use crate::flow::{
ConfirmationSubject, DomainEnumeration, InteractionRequirement, ObligationId, PhaseOwnership,
StartBehaviour, StartPrecondition, ViewOf, WorkflowDefinition, WorkflowExecutor,
WorkflowNotice, WorkflowView, WritingStage,
};
use crate::ids::{AccountId, CaseId, WorkflowKey, WorkflowVersion};
use crate::interaction::InteractionSpec;
use crate::locale::Locale;
use crate::operation::{GlossaryTerm, OperationSpec};
use crate::target::ResolvedAct;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ErasedObligation {
pub id: ObligationId,
pub value: serde_json::Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sentence: Option<crate::locale::LocalizedText>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub act: Option<super::ObligationAct>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ErasedWorkflowView {
pub case_ref: CaseRef,
pub workflow_version: WorkflowVersion,
pub phase: serde_json::Value,
pub phase_ownership: PhaseOwnership,
pub obligations: Vec<ErasedObligation>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub blocking_interaction: Option<InteractionRequirement>,
#[serde(default)]
pub notices: Vec<WorkflowNotice>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub state: Vec<crate::flow::StateField>,
}
impl ErasedWorkflowView {
#[must_use]
pub fn is_complete(&self) -> bool {
self.outcome.is_some()
}
#[must_use]
pub fn is_user_owned(&self) -> bool {
self.phase_ownership == PhaseOwnership::User
}
#[must_use]
pub fn obligation_ids(&self) -> Vec<&ObligationId> {
self.obligations.iter().map(|o| &o.id).collect()
}
}
pub trait ErasedWorkflow: Send + Sync {
fn key(&self) -> WorkflowKey;
fn version(&self) -> WorkflowVersion;
fn project(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<ErasedWorkflowView, ErasureError>;
fn operations(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<Vec<OperationSpec>, ErasureError>;
fn summary(&self) -> Option<String>;
fn glossary(&self) -> Vec<GlossaryTerm>;
fn noun(&self) -> Option<crate::locale::LocalizedText> {
None
}
fn briefing(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<Option<String>, ErasureError>;
fn start_preconditions(&self) -> Vec<StartPrecondition>;
fn start_behaviour(&self) -> StartBehaviour;
fn confirmation_subject(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
act: &ResolvedAct,
) -> Result<Option<ConfirmationSubject>, ErasureError>;
fn artifacts(&self, view: &ErasedWorkflowView) -> Result<Vec<ArtifactRef>, ErasureError>;
fn may_open_beside(&self, open: &[ErasedWorkflowView]) -> Result<(), ErasedCallError>;
fn narration_briefing(
&self,
stage: WritingStage,
view: &ErasedWorkflowView,
) -> Result<Option<String>, ErasureError>;
fn enumerations(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<Vec<DomainEnumeration>, ErasureError>;
fn obligation_sentence(
&self,
obligation: &serde_json::Value,
) -> Result<Option<crate::locale::LocalizedText>, ErasureError>;
fn next_steps(
&self,
view: &ErasedWorkflowView,
) -> Result<Vec<crate::locale::LocalizedText>, ErasureError> {
let _ = view;
Ok(Vec::new())
}
fn compile_act(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
act: &ResolvedAct,
) -> Result<Vec<serde_json::Value>, ErasedCallError>;
fn nothing_changed(
&self,
state: Option<&serde_json::Value>,
act: &ResolvedAct,
) -> Option<crate::locale::LocalizedText> {
let _ = (state, act);
None
}
fn command_policy(
&self,
state: Option<&serde_json::Value>,
command: &serde_json::Value,
) -> Result<CommandPolicy, ErasureError>;
fn validate_command(
&self,
state: Option<&serde_json::Value>,
command: &serde_json::Value,
) -> Result<(), ErasedCallError>;
fn receipts(
&self,
events: &[ReceiptEvent<serde_json::Value>],
locale: &Locale,
) -> Result<Vec<OperationalReceipt>, ErasureError>;
fn build_interaction(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
requirement: &InteractionRequirement,
) -> Result<InteractionSpec, ErasedCallError>;
}
#[async_trait::async_trait]
pub trait ErasedCaseLoader: Send + Sync {
async fn load_case(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<serde_json::Value>>, StoreError>;
}
#[async_trait::async_trait]
impl<E: ErasedExecutor + ?Sized> ErasedCaseLoader for E {
async fn load_case(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
self.load(account, case_id).await
}
}
#[derive(Clone)]
pub struct CaseLoaderHandle {
executor: Arc<dyn ErasedExecutor>,
}
impl CaseLoaderHandle {
#[must_use]
pub fn new(executor: Arc<dyn ErasedExecutor>) -> Self {
Self { executor }
}
}
impl fmt::Debug for CaseLoaderHandle {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("CaseLoaderHandle").finish_non_exhaustive()
}
}
#[async_trait::async_trait]
impl ErasedCaseLoader for CaseLoaderHandle {
async fn load_case(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
self.executor.load(account, case_id).await
}
}
#[async_trait::async_trait]
pub trait ErasedExecutor: Send + Sync {
async fn load(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<serde_json::Value>>, StoreError>;
async fn execute(
&self,
batch: CommandBatch<serde_json::Value>,
) -> Result<Commit<serde_json::Value, serde_json::Value>, ExecutionError>;
}
pub struct TypedWorkflowAdapter<W, E> {
definition: W,
executor: E,
}
impl<W, E> TypedWorkflowAdapter<W, E> {
#[must_use]
pub const fn new(definition: W, executor: E) -> Self {
Self {
definition,
executor,
}
}
#[must_use]
pub const fn definition(&self) -> &W {
&self.definition
}
#[must_use]
pub const fn executor(&self) -> &E {
&self.executor
}
}
impl<W: WorkflowDefinition, E> fmt::Debug for TypedWorkflowAdapter<W, E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("TypedWorkflowAdapter")
.field("workflow", &self.definition.key())
.field("version", &self.definition.version())
.finish_non_exhaustive()
}
}
impl<W: WorkflowDefinition, E> TypedWorkflowAdapter<W, E> {
fn state(&self, value: Option<&serde_json::Value>) -> Result<Option<W::State>, ErasureError> {
value
.map(|v| {
W::State::deserialize(v).map_err(|_| ErasureError::StateDeserialization {
workflow: self.definition.key(),
})
})
.transpose()
}
fn command(&self, value: &serde_json::Value) -> Result<W::Command, ErasureError> {
W::Command::deserialize(value).map_err(|_| ErasureError::CommandDeserialization {
workflow: self.definition.key(),
})
}
fn serialize<T: Serialize>(&self, value: &T) -> Result<serde_json::Value, ErasureError> {
serde_json::to_value(value).map_err(|_| ErasureError::Serialization {
workflow: self.definition.key(),
})
}
fn typed_view(&self, case_ref: CaseRef, state: Option<&W::State>) -> ViewOf<W> {
self.definition.project(case_ref, state)
}
fn typed_view_from(&self, view: &ErasedWorkflowView) -> Result<ViewOf<W>, ErasureError> {
let mismatch = || ErasureError::StateDeserialization {
workflow: self.definition.key(),
};
let phase: W::Phase = serde_json::from_value(view.phase.clone()).map_err(|_| mismatch())?;
let mut obligations = Vec::with_capacity(view.obligations.len());
for obligation in &view.obligations {
obligations.push(
serde_json::from_value::<W::Obligation>(obligation.value.clone())
.map_err(|_| mismatch())?,
);
}
let outcome = view
.outcome
.clone()
.map(serde_json::from_value::<W::Outcome>)
.transpose()
.map_err(|_| mismatch())?;
Ok(WorkflowView {
case_ref: view.case_ref.clone(),
workflow_version: view.workflow_version.clone(),
phase,
obligations,
blocking_interaction: view.blocking_interaction.clone(),
notices: view.notices.clone(),
outcome,
})
}
fn check_same_case(&self, projected: &CaseRef, resolved: &CaseRef) -> Result<(), ErasureError> {
if projected.same_case(resolved)
&& projected.expected_revision == resolved.expected_revision
{
Ok(())
} else {
Err(ErasureError::CaseMismatch {
workflow: self.definition.key(),
})
}
}
}
impl<W: WorkflowDefinition, E: Send + Sync> ErasedWorkflow for TypedWorkflowAdapter<W, E> {
fn key(&self) -> WorkflowKey {
self.definition.key()
}
fn version(&self) -> WorkflowVersion {
self.definition.version()
}
fn project(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<ErasedWorkflowView, ErasureError> {
let state = self.state(state)?;
let view = self.typed_view(case_ref, state.as_ref());
let ownership = self.definition.phase_ownership(&view.phase);
let mut erased = view
.erase(ownership)
.map_err(|_| ErasureError::Serialization {
workflow: self.definition.key(),
})?;
erased.state = self.definition.narratable_state(state.as_ref());
for obligation in &mut erased.obligations {
if let Ok(typed) = serde_json::from_value::<W::Obligation>(obligation.value.clone()) {
obligation.sentence = self.definition.obligation_sentence(&typed);
obligation.act = self.definition.obligation_act(state.as_ref(), &typed);
}
}
Ok(erased)
}
fn operations(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<Vec<OperationSpec>, ErasureError> {
let state = self.state(state)?;
let view = self.typed_view(case_ref, state.as_ref());
let workflow = self.definition.key();
let mut operations = self.definition.operations(&view);
for operation in &mut operations {
operation.workflow = workflow.clone();
operation
.validate()
.map_err(|error| ErasureError::InvalidOperation {
workflow: workflow.clone(),
reason: error.to_string(),
})?;
}
Ok(operations)
}
fn summary(&self) -> Option<String> {
self.definition.summary()
}
fn glossary(&self) -> Vec<GlossaryTerm> {
self.definition.glossary()
}
fn noun(&self) -> Option<crate::locale::LocalizedText> {
self.definition.noun()
}
fn briefing(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<Option<String>, ErasureError> {
let state = self.state(state)?;
let view = self.typed_view(case_ref, state.as_ref());
Ok(self.definition.briefing(&view))
}
fn start_preconditions(&self) -> Vec<StartPrecondition> {
self.definition.start_preconditions()
}
fn start_behaviour(&self) -> StartBehaviour {
self.definition.start_behaviour()
}
fn confirmation_subject(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
act: &ResolvedAct,
) -> Result<Option<ConfirmationSubject>, ErasureError> {
let state = self.state(state)?;
let view = self.typed_view(case_ref, state.as_ref());
Ok(self
.definition
.confirmation_subject(state.as_ref(), &view, act))
}
fn artifacts(&self, view: &ErasedWorkflowView) -> Result<Vec<ArtifactRef>, ErasureError> {
let typed = self.typed_view_from(view)?;
Ok(self.definition.artifacts(&typed))
}
fn may_open_beside(&self, open: &[ErasedWorkflowView]) -> Result<(), ErasedCallError> {
let typed: Vec<ViewOf<W>> = open
.iter()
.map(|view| self.typed_view_from(view))
.collect::<Result<_, _>>()?;
self.definition
.may_open_beside(&typed)
.map_err(ErasedCallError::from)
}
fn narration_briefing(
&self,
stage: WritingStage,
view: &ErasedWorkflowView,
) -> Result<Option<String>, ErasureError> {
let typed = self.typed_view_from(view)?;
Ok(match stage {
WritingStage::Transition => self.definition.transition_briefing(&typed),
WritingStage::Answer => self.definition.answer_briefing(&typed),
})
}
fn enumerations(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
) -> Result<Vec<DomainEnumeration>, ErasureError> {
let state = self.state(state)?;
let view = self.typed_view(case_ref, state.as_ref());
Ok(self.definition.enumerations(&view))
}
fn next_steps(
&self,
view: &ErasedWorkflowView,
) -> Result<Vec<crate::locale::LocalizedText>, ErasureError> {
let typed = self.typed_view_from(view)?;
Ok(self.definition.next_steps(&typed))
}
fn obligation_sentence(
&self,
obligation: &serde_json::Value,
) -> Result<Option<crate::locale::LocalizedText>, ErasureError> {
let typed: W::Obligation = serde_json::from_value(obligation.clone()).map_err(|_| {
ErasureError::StateDeserialization {
workflow: self.definition.key(),
}
})?;
Ok(self.definition.obligation_sentence(&typed))
}
fn compile_act(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
act: &ResolvedAct,
) -> Result<Vec<serde_json::Value>, ErasedCallError> {
self.check_same_case(&case_ref, &act.case_ref)?;
let state = self.state(state)?;
let view = self.typed_view(case_ref, state.as_ref());
let commands = self.definition.compile_act(state.as_ref(), &view, act)?;
commands
.iter()
.map(|c| self.serialize(c).map_err(ErasedCallError::from))
.collect()
}
fn nothing_changed(
&self,
state: Option<&serde_json::Value>,
act: &ResolvedAct,
) -> Option<crate::locale::LocalizedText> {
let state = self.state(state).ok()?;
self.definition.nothing_changed(state.as_ref(), act)
}
fn command_policy(
&self,
state: Option<&serde_json::Value>,
command: &serde_json::Value,
) -> Result<CommandPolicy, ErasureError> {
let state = self.state(state)?;
let command = self.command(command)?;
Ok(self.definition.command_policy(state.as_ref(), &command))
}
fn validate_command(
&self,
state: Option<&serde_json::Value>,
command: &serde_json::Value,
) -> Result<(), ErasedCallError> {
let state = self.state(state)?;
let command = self.command(command)?;
self.definition
.validate_command(state.as_ref(), &command)
.map_err(ErasedCallError::from)
}
fn receipts(
&self,
events: &[ReceiptEvent<serde_json::Value>],
locale: &Locale,
) -> Result<Vec<OperationalReceipt>, ErasureError> {
let events = events
.iter()
.map(|event| {
event.try_map_payload_ref(|payload| {
W::Event::deserialize(payload).map_err(|_| ErasureError::EventDeserialization {
workflow: self.definition.key(),
})
})
})
.collect::<Result<Vec<_>, _>>()?;
Ok(self.definition.receipts(&events, locale))
}
fn build_interaction(
&self,
case_ref: CaseRef,
state: Option<&serde_json::Value>,
requirement: &InteractionRequirement,
) -> Result<InteractionSpec, ErasedCallError> {
let state = self.state(state)?;
let view = self.typed_view(case_ref.clone(), state.as_ref());
let spec = self
.definition
.build_interaction(state.as_ref(), &view, requirement)?;
self.check_same_case(&case_ref, &spec.case_ref)?;
spec.validate()?;
Ok(spec)
}
}
#[async_trait::async_trait]
impl<W, E> ErasedExecutor for TypedWorkflowAdapter<W, E>
where
W: WorkflowDefinition,
E: WorkflowExecutor<W>,
{
async fn load(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<serde_json::Value>>, StoreError> {
let loaded = self.executor.load(account, case_id).await?;
let revision = loaded.revision;
let value = loaded
.value
.map(|s| serde_json::to_value(&s).map_err(|_| StoreError::Serialization))
.transpose()?;
Ok(Versioned::new(value, revision))
}
async fn execute(
&self,
batch: CommandBatch<serde_json::Value>,
) -> Result<Commit<serde_json::Value, serde_json::Value>, ExecutionError> {
let typed = batch.try_map(|c| self.command(&c))?;
let commit = self.executor.execute(typed).await?;
let state = commit
.state
.as_ref()
.map(|s| self.serialize(s))
.transpose()?;
let events = commit
.events
.into_iter()
.map(|e: CommittedEvent<W::Event>| e.try_map_payload(|p| self.serialize(&p)))
.collect::<Result<Vec<_>, _>>()?;
Ok(Commit {
state,
new_revision: commit.new_revision,
events,
idempotency_replay: commit.idempotency_replay,
})
}
}
#[derive(Clone)]
pub struct RegisteredWorkflow {
pub key: WorkflowKey,
pub version: WorkflowVersion,
pub definition: Arc<dyn ErasedWorkflow>,
pub executor: Arc<dyn ErasedExecutor>,
}
impl fmt::Debug for RegisteredWorkflow {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("RegisteredWorkflow")
.field("key", &self.key)
.field("version", &self.version)
.finish_non_exhaustive()
}
}
impl RegisteredWorkflow {
#[must_use]
pub fn loader(&self) -> CaseLoaderHandle {
CaseLoaderHandle::new(Arc::clone(&self.executor))
}
pub fn check_version(&self, found: &WorkflowVersion) -> Result<(), ErasureError> {
if &self.version == found {
Ok(())
} else {
Err(ErasureError::VersionMismatch {
workflow: self.key.clone(),
registered: self.version.clone(),
found: found.clone(),
})
}
}
}
#[derive(Debug, Clone, Default)]
pub struct WorkflowRegistry {
entries: IndexMap<WorkflowKey, RegisteredWorkflow>,
}
impl WorkflowRegistry {
#[must_use]
pub fn builder() -> WorkflowRegistryBuilder {
WorkflowRegistryBuilder::default()
}
#[must_use]
pub fn get(&self, key: &WorkflowKey) -> Option<&RegisteredWorkflow> {
self.entries.get(key)
}
pub fn require(&self, key: &WorkflowKey) -> Result<&RegisteredWorkflow, ErasureError> {
self.entries
.get(key)
.ok_or_else(|| ErasureError::UnknownWorkflow {
workflow: key.clone(),
})
}
pub fn check_version(
&self,
key: &WorkflowKey,
found: &WorkflowVersion,
) -> Result<(), ErasureError> {
self.require(key)?.check_version(found)
}
#[must_use]
pub fn contains(&self, key: &WorkflowKey) -> bool {
self.entries.contains_key(key)
}
pub fn iter(&self) -> impl Iterator<Item = &RegisteredWorkflow> {
self.entries.values()
}
pub fn keys(&self) -> impl Iterator<Item = &WorkflowKey> {
self.entries.keys()
}
#[must_use]
pub fn len(&self) -> usize {
self.entries.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
#[must_use]
pub fn definitions(&self) -> WorkflowDefinitions {
WorkflowDefinitions {
entries: self
.entries
.iter()
.map(|(key, registered)| (key.clone(), Arc::clone(®istered.definition)))
.collect(),
}
}
#[must_use]
pub fn read_only(&self) -> WorkflowReadRegistry {
WorkflowReadRegistry {
definitions: self.definitions(),
loaders: self
.entries
.iter()
.map(|(key, registered)| {
let loader: Arc<dyn ErasedCaseLoader> = Arc::new(registered.loader());
(key.clone(), loader)
})
.collect(),
}
}
}
#[derive(Clone, Default)]
pub struct WorkflowDefinitions {
entries: IndexMap<WorkflowKey, Arc<dyn ErasedWorkflow>>,
}
impl fmt::Debug for WorkflowDefinitions {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("WorkflowDefinitions")
.field("keys", &self.entries.keys().collect::<Vec<_>>())
.finish()
}
}
impl WorkflowDefinitions {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn with(mut self, definition: Arc<dyn ErasedWorkflow>) -> Self {
self.entries.insert(definition.key(), definition);
self
}
#[must_use]
pub fn get(&self, key: &WorkflowKey) -> Option<&Arc<dyn ErasedWorkflow>> {
self.entries.get(key)
}
pub fn require(&self, key: &WorkflowKey) -> Result<&Arc<dyn ErasedWorkflow>, ErasureError> {
self.entries
.get(key)
.ok_or_else(|| ErasureError::UnknownWorkflow {
workflow: key.clone(),
})
}
#[must_use]
pub fn contains(&self, key: &WorkflowKey) -> bool {
self.entries.contains_key(key)
}
pub fn keys(&self) -> impl Iterator<Item = &WorkflowKey> {
self.entries.keys()
}
pub fn iter(&self) -> impl Iterator<Item = &Arc<dyn ErasedWorkflow>> {
self.entries.values()
}
#[must_use]
pub fn len(&self) -> usize {
self.entries.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
}
impl From<&WorkflowRegistry> for WorkflowDefinitions {
fn from(registry: &WorkflowRegistry) -> Self {
registry.definitions()
}
}
impl From<&Arc<WorkflowRegistry>> for WorkflowDefinitions {
fn from(registry: &Arc<WorkflowRegistry>) -> Self {
registry.definitions()
}
}
impl From<Arc<WorkflowRegistry>> for WorkflowDefinitions {
fn from(registry: Arc<WorkflowRegistry>) -> Self {
registry.definitions()
}
}
#[derive(Clone, Default)]
pub struct WorkflowReadRegistry {
definitions: WorkflowDefinitions,
loaders: IndexMap<WorkflowKey, Arc<dyn ErasedCaseLoader>>,
}
impl fmt::Debug for WorkflowReadRegistry {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("WorkflowReadRegistry")
.field("definitions", &self.definitions)
.field("loaders", &self.loaders.len())
.finish()
}
}
impl WorkflowReadRegistry {
#[must_use]
pub fn from_definitions(definitions: WorkflowDefinitions) -> Self {
Self {
definitions,
loaders: IndexMap::new(),
}
}
#[must_use]
pub const fn definitions(&self) -> &WorkflowDefinitions {
&self.definitions
}
#[must_use]
pub fn loader(&self, key: &WorkflowKey) -> Option<&Arc<dyn ErasedCaseLoader>> {
self.loaders.get(key)
}
pub fn require_loader(
&self,
key: &WorkflowKey,
) -> Result<&Arc<dyn ErasedCaseLoader>, ErasureError> {
self.loaders
.get(key)
.ok_or_else(|| ErasureError::UnknownWorkflow {
workflow: key.clone(),
})
}
#[must_use]
pub fn has_no_loaders(&self) -> bool {
self.loaders.is_empty()
}
}
impl From<&WorkflowRegistry> for WorkflowReadRegistry {
fn from(registry: &WorkflowRegistry) -> Self {
registry.read_only()
}
}
#[derive(Debug, Default)]
pub struct WorkflowRegistryBuilder {
entries: Vec<RegisteredWorkflow>,
}
impl WorkflowRegistryBuilder {
#[must_use]
pub fn register<W, E>(self, definition: W, executor: E) -> Self
where
W: WorkflowDefinition,
E: WorkflowExecutor<W> + 'static,
{
let adapter = Arc::new(TypedWorkflowAdapter::new(definition, executor));
let erased_definition: Arc<dyn ErasedWorkflow> = adapter.clone();
let erased_executor: Arc<dyn ErasedExecutor> = adapter;
self.register_erased(erased_definition, erased_executor)
}
#[must_use]
pub fn register_erased(
mut self,
definition: Arc<dyn ErasedWorkflow>,
executor: Arc<dyn ErasedExecutor>,
) -> Self {
self.entries.push(RegisteredWorkflow {
key: definition.key(),
version: definition.version(),
definition,
executor,
});
self
}
pub fn build(self) -> Result<WorkflowRegistry, ErasureError> {
let mut entries = IndexMap::with_capacity(self.entries.len());
for entry in self.entries {
if entries.contains_key(&entry.key) {
return Err(ErasureError::DuplicateWorkflow {
workflow: entry.key,
});
}
entries.insert(entry.key.clone(), entry);
}
Ok(WorkflowRegistry { entries })
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::case::Versioned;
use crate::flow::{PhaseOwnership, ViewOf, WorkflowDefinition, WorkflowView};
use crate::ids::CaseRevision;
use crate::locale::Locale;
use crate::target::ResolvedAct;
use serde::{Deserialize, Serialize};
#[derive(Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
struct Nothing;
struct Toy;
impl WorkflowDefinition for Toy {
type State = Nothing;
type Phase = Nothing;
type Obligation = Nothing;
type Command = Nothing;
type Event = Nothing;
type Outcome = Nothing;
fn key(&self) -> WorkflowKey {
WorkflowKey::from("toy")
}
fn version(&self) -> WorkflowVersion {
WorkflowVersion::from("1")
}
fn phase_ownership(&self, _phase: &Self::Phase) -> PhaseOwnership {
PhaseOwnership::System
}
fn project(&self, case_ref: CaseRef, _state: Option<&Self::State>) -> ViewOf<Self> {
WorkflowView::new(case_ref, self.version(), Nothing)
}
fn operations(&self, _view: &ViewOf<Self>) -> Vec<crate::operation::OperationSpec> {
Vec::new()
}
fn compile_act(
&self,
_state: Option<&Self::State>,
_view: &ViewOf<Self>,
_act: &ResolvedAct,
) -> Result<Vec<Self::Command>, crate::error::DomainRejection> {
Ok(Vec::new())
}
fn command_policy(
&self,
_state: Option<&Self::State>,
_command: &Self::Command,
) -> CommandPolicy {
CommandPolicy::conservative()
}
fn validate_command(
&self,
_state: Option<&Self::State>,
_command: &Self::Command,
) -> Result<(), crate::error::DomainRejection> {
Ok(())
}
fn receipts(
&self,
_events: &[ReceiptEvent<Self::Event>],
_locale: &Locale,
) -> Vec<OperationalReceipt> {
Vec::new()
}
}
struct Executor;
#[async_trait::async_trait]
impl crate::flow::WorkflowExecutor<Toy> for Executor {
async fn load(
&self,
_account: &AccountId,
_case_id: &CaseId,
) -> Result<Versioned<Option<Nothing>>, StoreError> {
Ok(Versioned::new(Some(Nothing), CaseRevision(7)))
}
async fn execute(
&self,
_batch: CommandBatch<Nothing>,
) -> Result<Commit<Nothing, Nothing>, ExecutionError> {
panic!("the read-only projection must never reach execution");
}
}
fn registry() -> WorkflowRegistry {
WorkflowRegistry::builder()
.register(Toy, Executor)
.build()
.expect("one workflow, one key")
}
#[tokio::test]
async fn the_read_only_projection_loads_and_cannot_execute() {
let registry = registry();
let reading = registry.read_only();
let key = WorkflowKey::from("toy");
assert!(reading.definitions().contains(&key));
let loaded = reading
.require_loader(&key)
.expect("the projection carries a loader")
.load_case(&AccountId::from("a"), &CaseId::from("c1"))
.await
.expect("the loader reads");
assert_eq!(loaded.revision, CaseRevision(7));
}
#[test]
fn definitions_alone_carry_no_loader() {
let definitions = registry().definitions();
let reading = WorkflowReadRegistry::from_definitions(definitions);
assert!(reading.has_no_loaders());
assert!(reading.require_loader(&WorkflowKey::from("toy")).is_err());
assert!(
reading
.definitions()
.require(&WorkflowKey::from("toy"))
.is_ok(),
"the pure half is still there"
);
}
#[test]
fn definitions_can_be_built_without_any_executor_at_all() {
let definition: Arc<dyn ErasedWorkflow> = Arc::new(TypedWorkflowAdapter::new(Toy, ()));
let definitions = WorkflowDefinitions::new().with(definition);
assert_eq!(definitions.len(), 1);
assert!(definitions.contains(&WorkflowKey::from("toy")));
}
}