pub mod invariants;
pub mod registry;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use crate::case::{CaseRef, Versioned};
use crate::command::RiskClass;
use crate::error::{DomainRejection, ExecutionError, HashError, StoreError};
use crate::ids::{AccountId, CaseId, WorkflowKey, WorkflowVersion};
use crate::interaction::{
InteractionKind, InteractionPayload, InteractionSpec, TextResolutionPolicy,
};
use crate::locale::{Locale, LocalizedText};
use crate::operation::{GlossaryTerm, OperationSpec};
use crate::response::NoticeSeverity;
use crate::target::ResolvedAct;
pub use crate::command::{CommandBatch, CommandPolicy};
pub use crate::event::{
Commit, CommittedEvent, EventRedaction, OperationalReceipt, ReceiptEvent, RedactedEvent,
};
pub use invariants::{check_erased_view, check_view};
pub use registry::{
CaseLoaderHandle, ErasedCaseLoader, ErasedExecutor, ErasedObligation, ErasedWorkflow,
ErasedWorkflowView, RegisteredWorkflow, TypedWorkflowAdapter, WorkflowDefinitions,
WorkflowReadRegistry, WorkflowRegistry, WorkflowRegistryBuilder,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PhaseOwnership {
User,
System,
External,
Terminal,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(transparent)]
pub struct ObligationId(pub String);
impl ObligationId {
pub fn of<O: Serialize + ?Sized>(obligation: &O) -> Result<Self, HashError> {
crate::hash::canonical_json(obligation).map(Self)
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
impl std::fmt::Display for ObligationId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkflowNotice {
pub code: String,
pub severity: NoticeSeverity,
pub text: LocalizedText,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct InteractionRequirement {
pub key: String,
pub kind: InteractionKind,
pub blocking: bool,
pub revision_independent: bool,
pub text_resolution: TextResolutionPolicy,
#[serde(default = "RiskClass::conservative")]
pub confirms_risk: RiskClass,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expires_in: Option<Duration>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub payload: Option<InteractionPayload>,
}
impl InteractionRequirement {
#[must_use]
pub fn blocking(key: impl Into<String>, kind: InteractionKind) -> Self {
Self {
key: key.into(),
kind,
blocking: true,
revision_independent: false,
text_resolution: TextResolutionPolicy::Never,
confirms_risk: RiskClass::conservative(),
expires_in: None,
payload: None,
}
}
#[must_use]
pub fn non_blocking(key: impl Into<String>, kind: InteractionKind) -> Self {
Self {
blocking: false,
..Self::blocking(key, kind)
}
}
#[must_use]
pub fn with_confirms_risk(mut self, risk: RiskClass) -> Self {
self.confirms_risk = risk;
self
}
#[must_use]
pub fn with_payload(mut self, payload: InteractionPayload) -> Self {
self.payload = Some(payload);
self
}
#[must_use]
pub fn with_text_resolution(mut self, policy: TextResolutionPolicy) -> Self {
self.text_resolution = policy;
self
}
#[must_use]
pub fn to_spec(&self, case_ref: CaseRef) -> InteractionSpec {
let payload = self
.payload
.clone()
.unwrap_or_else(|| InteractionPayload::new(self.key.clone()));
InteractionSpec {
key: self.key.clone(),
case_ref,
kind: self.kind,
blocking: self.blocking,
payload,
expires_in: self.expires_in,
text_resolution: self.text_resolution.clone(),
confirms_risk: self.confirms_risk,
binds_to_revision: !self.revision_independent,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkflowView<P, O, T> {
pub case_ref: CaseRef,
pub workflow_version: WorkflowVersion,
pub phase: P,
pub obligations: Vec<O>,
#[serde(default = "Option::default", skip_serializing_if = "Option::is_none")]
pub blocking_interaction: Option<InteractionRequirement>,
#[serde(default = "Vec::new")]
pub notices: Vec<WorkflowNotice>,
#[serde(default = "Option::default", skip_serializing_if = "Option::is_none")]
pub outcome: Option<T>,
}
impl<P, O, T> WorkflowView<P, O, T> {
#[must_use]
pub fn new(case_ref: CaseRef, workflow_version: WorkflowVersion, phase: P) -> Self {
Self {
case_ref,
workflow_version,
phase,
obligations: Vec::new(),
blocking_interaction: None,
notices: Vec::new(),
outcome: None,
}
}
#[must_use]
pub fn with_obligations(mut self, obligations: impl IntoIterator<Item = O>) -> Self {
self.obligations.extend(obligations);
self
}
#[must_use]
pub fn with_blocking_interaction(mut self, requirement: InteractionRequirement) -> Self {
self.blocking_interaction = Some(requirement);
self
}
#[must_use]
pub fn with_notice(mut self, notice: WorkflowNotice) -> Self {
self.notices.push(notice);
self
}
#[must_use]
pub fn with_outcome(mut self, outcome: T) -> Self {
self.outcome = Some(outcome);
self
}
#[must_use]
pub fn is_complete(&self) -> bool {
self.outcome.is_some()
}
#[must_use]
pub fn has_obligations(&self) -> bool {
!self.obligations.is_empty()
}
}
impl<P: Serialize, O: Serialize, T: Serialize> WorkflowView<P, O, T> {
pub fn erase(&self, ownership: PhaseOwnership) -> Result<ErasedWorkflowView, HashError> {
let mut obligations = Vec::with_capacity(self.obligations.len());
for obligation in &self.obligations {
obligations.push(ErasedObligation {
id: ObligationId::of(obligation)?,
value: crate::hash::canonical_value(obligation)?,
sentence: None,
act: None,
});
}
Ok(ErasedWorkflowView {
case_ref: self.case_ref.clone(),
workflow_version: self.workflow_version.clone(),
phase: crate::hash::canonical_value(&self.phase)?,
phase_ownership: ownership,
obligations,
blocking_interaction: self.blocking_interaction.clone(),
notices: self.notices.clone(),
outcome: self
.outcome
.as_ref()
.map(crate::hash::canonical_value)
.transpose()?,
state: Vec::new(),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct StateField {
pub field: String,
pub value: serde_json::Value,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub identifying: bool,
}
impl StateField {
#[must_use]
pub fn new(field: impl Into<String>, value: serde_json::Value) -> Self {
Self {
field: field.into(),
value,
identifying: false,
}
}
#[must_use]
pub const fn identifying(mut self) -> Self {
self.identifying = true;
self
}
}
pub type ViewOf<W> = WorkflowView<
<W as WorkflowDefinition>::Phase,
<W as WorkflowDefinition>::Obligation,
<W as WorkflowDefinition>::Outcome,
>;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct StartPrecondition {
pub workflow: WorkflowKey,
pub phases: Vec<serde_json::Value>,
pub reason: LocalizedText,
}
impl StartPrecondition {
#[must_use]
pub fn new(
workflow: impl Into<WorkflowKey>,
phases: impl IntoIterator<Item = serde_json::Value>,
reason: LocalizedText,
) -> Self {
Self {
workflow: workflow.into(),
phases: phases.into_iter().collect(),
reason,
}
}
pub fn requires<P: Serialize>(
workflow: impl Into<WorkflowKey>,
phases: &[P],
reason: LocalizedText,
) -> Result<Self, HashError> {
let mut serialized = Vec::with_capacity(phases.len());
for phase in phases {
serialized.push(crate::hash::canonical_value(phase)?);
}
Ok(Self {
workflow: workflow.into(),
phases: serialized,
reason,
})
}
#[must_use]
pub fn satisfied_by(&self, view: &ErasedWorkflowView) -> bool {
view.case_ref.workflow == self.workflow && self.phases.contains(&view.phase)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum StartBehaviour {
#[default]
OpensNewCase,
ResumesOpenCase,
}
impl StartBehaviour {
#[must_use]
pub const fn resumes(self) -> bool {
matches!(self, Self::ResumesOpenCase)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
pub struct ConfirmationSubject {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<LocalizedText>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub body: Option<LocalizedText>,
}
impl ConfirmationSubject {
#[must_use]
pub fn asking(title: LocalizedText) -> Self {
Self {
title: Some(title),
body: None,
}
}
#[must_use]
pub fn describing(body: LocalizedText) -> Self {
Self {
title: None,
body: Some(body),
}
}
#[must_use]
pub fn with_body(mut self, body: LocalizedText) -> Self {
self.body = Some(body);
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WritingStage {
Transition,
Answer,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct EnumeratedValue {
pub id: String,
pub label: LocalizedText,
}
impl EnumeratedValue {
#[must_use]
pub fn new(id: impl Into<String>, label: LocalizedText) -> Self {
Self {
id: id.into(),
label,
}
}
}
crate::ids::string_id! {
QuestionReference
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DomainEnumeration {
pub subject: QuestionReference,
pub values: Vec<EnumeratedValue>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub preamble: Option<LocalizedText>,
}
impl DomainEnumeration {
#[must_use]
pub fn new(subject: impl Into<QuestionReference>, values: Vec<EnumeratedValue>) -> Self {
Self {
subject: subject.into(),
values,
preamble: None,
}
}
#[must_use]
pub fn with_preamble(mut self, preamble: LocalizedText) -> Self {
self.preamble = Some(preamble);
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(transparent)]
pub struct BriefingBudget {
max_bytes: Option<usize>,
}
impl BriefingBudget {
#[must_use]
pub const fn conservative() -> Self {
Self { max_bytes: None }
}
#[must_use]
pub const fn new(max_bytes: usize) -> Self {
Self {
max_bytes: Some(max_bytes),
}
}
#[must_use]
pub const fn max_bytes(self) -> Option<usize> {
self.max_bytes
}
#[must_use]
pub fn apply(self, text: &str) -> String {
let Some(max_bytes) = self.max_bytes else {
return text.to_owned();
};
if text.len() <= max_bytes {
return text.to_owned();
}
let mut end = max_bytes;
while end > 0 && !text.is_char_boundary(end) {
end -= 1;
}
format!("{}… [briefing truncated]", &text[..end])
}
}
pub trait WorkflowDefinition: Send + Sync + 'static {
type State: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + 'static;
type Phase: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + Eq + 'static;
type Obligation: Clone
+ Send
+ Sync
+ Serialize
+ serde::de::DeserializeOwned
+ Eq
+ std::hash::Hash
+ 'static;
type Command: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + 'static;
type Event: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + 'static;
type Outcome: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + Eq + 'static;
fn key(&self) -> WorkflowKey;
fn version(&self) -> WorkflowVersion;
fn phase_ownership(&self, phase: &Self::Phase) -> PhaseOwnership;
fn project(&self, case_ref: CaseRef, state: Option<&Self::State>) -> ViewOf<Self>;
fn narratable_state(&self, state: Option<&Self::State>) -> Vec<StateField> {
let _ = state;
Vec::new()
}
fn summary(&self) -> Option<String> {
None
}
fn glossary(&self) -> Vec<GlossaryTerm> {
Vec::new()
}
fn noun(&self) -> Option<crate::locale::LocalizedText> {
None
}
fn operations(&self, view: &ViewOf<Self>) -> Vec<OperationSpec>;
fn briefing(&self, view: &ViewOf<Self>) -> Option<String> {
let _ = view;
None
}
fn obligation_sentence(&self, obligation: &Self::Obligation) -> Option<LocalizedText> {
let _ = obligation;
None
}
fn obligation_act(
&self,
state: Option<&Self::State>,
obligation: &Self::Obligation,
) -> Option<ObligationAct> {
let _ = (state, obligation);
None
}
fn start_preconditions(&self) -> Vec<StartPrecondition> {
Vec::new()
}
fn confirmation_subject(
&self,
state: Option<&Self::State>,
view: &ViewOf<Self>,
act: &ResolvedAct,
) -> Option<ConfirmationSubject> {
let _ = (state, view, act);
None
}
fn start_behaviour(&self) -> StartBehaviour {
StartBehaviour::OpensNewCase
}
fn may_open_beside(&self, open: &[ViewOf<Self>]) -> Result<(), DomainRejection> {
let _ = open;
Ok(())
}
fn transition_briefing(&self, view: &ViewOf<Self>) -> Option<String> {
let _ = view;
None
}
fn answer_briefing(&self, view: &ViewOf<Self>) -> Option<String> {
let _ = view;
None
}
fn enumerations(&self, view: &ViewOf<Self>) -> Vec<DomainEnumeration> {
let _ = view;
Vec::new()
}
fn next_steps(&self, view: &ViewOf<Self>) -> Vec<crate::locale::LocalizedText> {
let _ = view;
Vec::new()
}
fn artifacts(&self, view: &ViewOf<Self>) -> Vec<crate::event::ArtifactRef> {
let _ = view;
Vec::new()
}
fn compile_act(
&self,
state: Option<&Self::State>,
view: &ViewOf<Self>,
act: &ResolvedAct,
) -> Result<Vec<Self::Command>, DomainRejection>;
fn nothing_changed(
&self,
state: Option<&Self::State>,
act: &ResolvedAct,
) -> Option<LocalizedText> {
let _ = (state, act);
None
}
fn command_policy(&self, state: Option<&Self::State>, command: &Self::Command)
-> CommandPolicy;
fn validate_command(
&self,
state: Option<&Self::State>,
command: &Self::Command,
) -> Result<(), DomainRejection>;
fn receipts(
&self,
events: &[ReceiptEvent<Self::Event>],
locale: &Locale,
) -> Vec<OperationalReceipt>;
fn build_interaction(
&self,
state: Option<&Self::State>,
view: &ViewOf<Self>,
requirement: &InteractionRequirement,
) -> Result<InteractionSpec, DomainRejection> {
let _ = state;
Ok(requirement.to_spec(view.case_ref.clone()))
}
}
#[async_trait::async_trait]
pub trait WorkflowExecutor<W: WorkflowDefinition>: Send + Sync {
async fn load(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<W::State>>, StoreError>;
async fn execute(
&self,
batch: CommandBatch<W::Command>,
) -> Result<Commit<W::State, W::Event>, ExecutionError>;
}
#[async_trait::async_trait]
pub trait CaseLoader<W: WorkflowDefinition>: Send + Sync {
async fn load_case(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<W::State>>, StoreError>;
}
#[async_trait::async_trait]
impl<W, E> CaseLoader<W> for E
where
W: WorkflowDefinition,
E: WorkflowExecutor<W> + ?Sized,
{
async fn load_case(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<W::State>>, StoreError> {
self.load(account, case_id).await
}
}
#[async_trait::async_trait]
impl<W, E> WorkflowExecutor<W> for std::sync::Arc<E>
where
W: WorkflowDefinition,
E: WorkflowExecutor<W> + ?Sized,
{
async fn load(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<W::State>>, StoreError> {
(**self).load(account, case_id).await
}
async fn execute(
&self,
batch: CommandBatch<W::Command>,
) -> Result<Commit<W::State, W::Event>, ExecutionError> {
(**self).execute(batch).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ids::CaseRevision;
#[test]
fn requirement_to_spec_defaults() {
let req = InteractionRequirement::blocking("send", InteractionKind::ConfirmCommand);
let spec = req.to_spec(CaseRef::new("trip", "i1", CaseRevision(1)));
assert_eq!(spec.key, "send");
assert!(spec.blocking);
assert!(spec.binds_to_revision);
assert_eq!(spec.payload.title.default, "send");
}
#[test]
fn erase_produces_stable_obligation_ids() {
#[derive(Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
enum Ob {
Line { id: u32 },
}
let view: WorkflowView<&str, Ob, ()> = WorkflowView::new(
CaseRef::new("w", "c", CaseRevision(1)),
WorkflowVersion::from("1"),
"collecting",
)
.with_obligations([Ob::Line { id: 2 }, Ob::Line { id: 1 }]);
let erased = view.erase(PhaseOwnership::System).unwrap();
assert_eq!(erased.obligations[0].id.as_str(), r#"{"Line":{"id":2}}"#);
assert_eq!(erased.phase, serde_json::json!("collecting"));
assert!(erased.outcome.is_none());
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct ObligationAct {
pub operation: crate::ids::OperationKey,
#[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
pub given: std::collections::BTreeMap<String, serde_json::Value>,
pub asks: Vec<String>,
}
impl ObligationAct {
#[must_use]
pub fn new(
operation: impl Into<crate::ids::OperationKey>,
asks: impl IntoIterator<Item = impl Into<String>>,
) -> Self {
Self {
operation: operation.into(),
given: std::collections::BTreeMap::new(),
asks: asks.into_iter().map(Into::into).collect(),
}
}
#[must_use]
pub fn given(mut self, name: impl Into<String>, value: serde_json::Value) -> Self {
self.given.insert(name.into(), value);
self
}
}