turnframe_core/flow/mod.rs
1//! Flow Map V2: the deterministic workflow projector (spec §8).
2//!
3//! A [`WorkflowDefinition`] projects persisted state into a [`WorkflowView`]:
4//! exactly one lifecycle phase, zero or more parameterized obligations, at most
5//! one blocking [`InteractionRequirement`], informational notices and an
6//! outcome that is present only when the workflow is complete. Projection is
7//! pure (I2): same version + same state ⇒ same view, no I/O.
8//!
9//! Execution lives in a separate [`WorkflowExecutor`] because applications
10//! mutate SQL rows, call services or fold event-sourced aggregates.
11//!
12//! [`registry`] erases the generics at the boundary so the runtime can host
13//! several workflows without knowing their concrete types; [`invariants`]
14//! checks the §8.4 rules on any view.
15
16pub mod invariants;
17pub mod registry;
18
19use std::time::Duration;
20
21use serde::{Deserialize, Serialize};
22
23use crate::case::{CaseRef, Versioned};
24use crate::command::RiskClass;
25use crate::error::{DomainRejection, ExecutionError, HashError, StoreError};
26use crate::ids::{AccountId, CaseId, OperationKey, WorkflowKey, WorkflowVersion};
27use crate::interaction::{
28 InteractionKind, InteractionPayload, InteractionSpec, TextResolutionPolicy,
29};
30use crate::locale::{Locale, LocalizedText};
31use crate::operation::{GlossaryTerm, OperationSpec};
32use crate::response::NoticeSeverity;
33use crate::target::ResolvedAct;
34
35pub use crate::command::{CommandBatch, CommandPolicy};
36pub use crate::event::{
37 Commit, CommittedEvent, EventRedaction, OperationalReceipt, ReceiptEvent, RedactedEvent,
38};
39pub use invariants::{check_erased_view, check_view};
40pub use registry::{
41 CaseLoaderHandle, ErasedCaseLoader, ErasedExecutor, ErasedObligation, ErasedWorkflow,
42 ErasedWorkflowView, RegisteredWorkflow, TypedWorkflowAdapter, WorkflowDefinitions,
43 WorkflowReadRegistry, WorkflowRegistry, WorkflowRegistryBuilder,
44};
45
46/// Who must act for the case to leave its current phase.
47#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
48#[serde(rename_all = "snake_case")]
49pub enum PhaseOwnership {
50 /// The user must answer a blocking interaction (I6).
51 User,
52 /// The system acts (a job, a policy).
53 System,
54 /// An external party acts (an airline, an intermediary, a recipient).
55 External,
56 /// Nothing more happens; the outcome is present.
57 Terminal,
58}
59
60/// Stable identifier of an obligation: the canonical JSON of its value.
61///
62/// Two obligations with the same identifier in one view are a map defect
63/// (spec §8.4). Parameterized obligations therefore carry their entity ids.
64#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
65#[serde(transparent)]
66pub struct ObligationId(pub String);
67
68impl ObligationId {
69 /// Derives the identifier of an obligation.
70 pub fn of<O: Serialize + ?Sized>(obligation: &O) -> Result<Self, HashError> {
71 crate::hash::canonical_json(obligation).map(Self)
72 }
73
74 /// Borrows the canonical JSON.
75 #[must_use]
76 pub fn as_str(&self) -> &str {
77 &self.0
78 }
79}
80
81impl std::fmt::Display for ObligationId {
82 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
83 f.write_str(&self.0)
84 }
85}
86
87/// A step the user may take next, offered by a case that owes nothing: an operation, the
88/// arguments already known, and the words the reply offers it in.
89#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
90#[non_exhaustive]
91pub struct NextStep {
92 /// The operation it runs.
93 pub operation: OperationKey,
94 /// The words the reply offers it in.
95 pub words: LocalizedText,
96 /// The arguments already known, by name; the rest are asked when it is taken up.
97 pub arguments: serde_json::Map<String, serde_json::Value>,
98}
99
100impl NextStep {
101 /// A step running `operation`, offered in `words`, with no argument known yet.
102 #[must_use]
103 pub fn new(operation: impl Into<OperationKey>, words: LocalizedText) -> Self {
104 Self {
105 operation: operation.into(),
106 words,
107 arguments: serde_json::Map::new(),
108 }
109 }
110
111 /// With the arguments already known, an object by name: any other value gives none.
112 #[must_use]
113 pub fn with_arguments(mut self, arguments: serde_json::Value) -> Self {
114 self.arguments = match arguments {
115 serde_json::Value::Object(named) => named,
116 _ => serde_json::Map::new(),
117 };
118 self
119 }
120}
121
122/// An informational, non-blocking element of a view.
123#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
124pub struct WorkflowNotice {
125 /// Stable code (e.g. `"trip.traveler_missing_email"`).
126 pub code: String,
127 /// Severity.
128 pub severity: NoticeSeverity,
129 /// Copy.
130 pub text: LocalizedText,
131}
132
133/// What a user-owned phase requires from the user (I6).
134///
135/// The requirement is data; the workflow turns it into a full
136/// [`InteractionSpec`] through [`WorkflowDefinition::build_interaction`], where
137/// it can consult state. When the projection can already describe the whole
138/// card, it may set `payload` and rely on the default implementation.
139#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
140pub struct InteractionRequirement {
141 /// Stable key per phase (e.g. `"send_confirmation"`).
142 pub key: String,
143 /// Shape of the interaction.
144 pub kind: InteractionKind,
145 /// Whether it owns unqualified answers for the case (I5).
146 pub blocking: bool,
147 /// Whether a revision change leaves it valid.
148 pub revision_independent: bool,
149 /// Whether typed text may resolve it.
150 pub text_resolution: TextResolutionPolicy,
151 /// Highest risk class an answer to this card authorizes. Conservative by
152 /// default, so a requirement that forgets it cannot be resolved from text.
153 #[serde(default = "RiskClass::conservative")]
154 pub confirms_risk: RiskClass,
155 /// Time to live.
156 #[serde(default, skip_serializing_if = "Option::is_none")]
157 pub expires_in: Option<Duration>,
158 /// Full payload when the projection can describe it.
159 #[serde(default, skip_serializing_if = "Option::is_none")]
160 pub payload: Option<InteractionPayload>,
161}
162
163impl InteractionRequirement {
164 /// A blocking, revision-bound requirement with `Never` text resolution.
165 #[must_use]
166 pub fn blocking(key: impl Into<String>, kind: InteractionKind) -> Self {
167 Self {
168 key: key.into(),
169 kind,
170 blocking: true,
171 revision_independent: false,
172 text_resolution: TextResolutionPolicy::Never,
173 confirms_risk: RiskClass::conservative(),
174 expires_in: None,
175 payload: None,
176 }
177 }
178
179 /// A non-blocking requirement: a card the case offers without owning the
180 /// answers to everything else.
181 ///
182 /// The flag existed and could not be set, which made every declared card
183 /// blocking whether or not the question wanted one. A blocking card owns
184 /// unqualified answers for the whole case (I5), so putting one beside a
185 /// question whose ordinary answer is a value leaves the user reading
186 /// buttons that cannot say what they came to say.
187 ///
188 /// The other half of that problem is answered by
189 /// the claim the act declares, which lets the refusal be
190 /// spoken instead of pressed. This is the smaller lever, kept because a
191 /// field nobody can write is not a field.
192 #[must_use]
193 pub fn non_blocking(key: impl Into<String>, kind: InteractionKind) -> Self {
194 Self {
195 blocking: false,
196 ..Self::blocking(key, kind)
197 }
198 }
199
200 /// Declares the highest risk class an answer authorizes. Lowering it is
201 /// what makes a card resolvable from typed text (spec §15.7).
202 #[must_use]
203 pub fn with_confirms_risk(mut self, risk: RiskClass) -> Self {
204 self.confirms_risk = risk;
205 self
206 }
207
208 /// Attaches a full payload.
209 #[must_use]
210 pub fn with_payload(mut self, payload: InteractionPayload) -> Self {
211 self.payload = Some(payload);
212 self
213 }
214
215 /// Sets the text resolution policy.
216 #[must_use]
217 pub fn with_text_resolution(mut self, policy: TextResolutionPolicy) -> Self {
218 self.text_resolution = policy;
219 self
220 }
221
222 /// Builds a spec for `case_ref`, using `payload` or a title-only payload
223 /// named after the key.
224 ///
225 /// A title-only payload is answerable for no kind, so a requirement without
226 /// a payload must be completed by
227 /// [`WorkflowDefinition::build_interaction`]; the spec is validated there.
228 #[must_use]
229 pub fn to_spec(&self, case_ref: CaseRef) -> InteractionSpec {
230 let payload = self
231 .payload
232 .clone()
233 .unwrap_or_else(|| InteractionPayload::new(self.key.clone()));
234 InteractionSpec {
235 key: self.key.clone(),
236 case_ref,
237 kind: self.kind,
238 blocking: self.blocking,
239 payload,
240 expires_in: self.expires_in,
241 text_resolution: self.text_resolution.clone(),
242 confirms_risk: self.confirms_risk,
243 binds_to_revision: !self.revision_independent,
244 }
245 }
246}
247
248/// The pure projection of a case (spec §8.1).
249#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
250pub struct WorkflowView<P, O, T> {
251 /// The case and the revision it was projected at.
252 pub case_ref: CaseRef,
253 /// Version of the definition that produced the view.
254 pub workflow_version: WorkflowVersion,
255 /// Exactly one lifecycle phase (I3).
256 pub phase: P,
257 /// Every currently open obligation, possibly parameterized (I4), **in the
258 /// order the workflow wants them asked for**.
259 ///
260 /// A set is the truth about what is open and is what an index wants; it
261 /// cannot say which one is being asked. Told to pick "the most useful one"
262 /// a writer picks, and a collection flow was asked for its bank account
263 /// first and its beneficiary second — the last question of the flow — and
264 /// then read all six out as a menu. So the list is ordered by declaration
265 /// and the writer is told to ask for the first.
266 ///
267 /// Ordering rather than a single "current" obligation, because some
268 /// domains genuinely have several open at once — a proposal reviewed in one
269 /// answer, three extras with no payer that one message can close two
270 /// of — and a shape that forced exactly one would make those unsayable.
271 pub obligations: Vec<O>,
272 /// Zero or one blocking interaction requirement (I5, I6).
273 #[serde(default = "Option::default", skip_serializing_if = "Option::is_none")]
274 pub blocking_interaction: Option<InteractionRequirement>,
275 /// Informational, non-blocking elements.
276 #[serde(default = "Vec::new")]
277 pub notices: Vec<WorkflowNotice>,
278 /// Present only when the workflow is complete.
279 #[serde(default = "Option::default", skip_serializing_if = "Option::is_none")]
280 pub outcome: Option<T>,
281}
282
283impl<P, O, T> WorkflowView<P, O, T> {
284 /// A view with a phase and nothing else.
285 #[must_use]
286 pub fn new(case_ref: CaseRef, workflow_version: WorkflowVersion, phase: P) -> Self {
287 Self {
288 case_ref,
289 workflow_version,
290 phase,
291 obligations: Vec::new(),
292 blocking_interaction: None,
293 notices: Vec::new(),
294 outcome: None,
295 }
296 }
297
298 /// Adds obligations.
299 #[must_use]
300 pub fn with_obligations(mut self, obligations: impl IntoIterator<Item = O>) -> Self {
301 self.obligations.extend(obligations);
302 self
303 }
304
305 /// Sets the blocking requirement.
306 #[must_use]
307 pub fn with_blocking_interaction(mut self, requirement: InteractionRequirement) -> Self {
308 self.blocking_interaction = Some(requirement);
309 self
310 }
311
312 /// Adds a notice.
313 #[must_use]
314 pub fn with_notice(mut self, notice: WorkflowNotice) -> Self {
315 self.notices.push(notice);
316 self
317 }
318
319 /// Sets the outcome.
320 #[must_use]
321 pub fn with_outcome(mut self, outcome: T) -> Self {
322 self.outcome = Some(outcome);
323 self
324 }
325
326 /// Returns `true` when an outcome is present.
327 #[must_use]
328 pub fn is_complete(&self) -> bool {
329 self.outcome.is_some()
330 }
331
332 /// Returns `true` when obligations remain.
333 #[must_use]
334 pub fn has_obligations(&self) -> bool {
335 !self.obligations.is_empty()
336 }
337}
338
339impl<P: Serialize, O: Serialize, T: Serialize> WorkflowView<P, O, T> {
340 /// Erases the generics into canonical JSON, attaching the phase ownership
341 /// the definition declares.
342 pub fn erase(&self, ownership: PhaseOwnership) -> Result<ErasedWorkflowView, HashError> {
343 let mut obligations = Vec::with_capacity(self.obligations.len());
344 for obligation in &self.obligations {
345 obligations.push(ErasedObligation {
346 id: ObligationId::of(obligation)?,
347 value: crate::hash::canonical_value(obligation)?,
348 // Filled by the registry, which holds the definition: see
349 // `ErasedObligation::sentence`.
350 sentence: None,
351 act: None,
352 });
353 }
354 Ok(ErasedWorkflowView {
355 case_ref: self.case_ref.clone(),
356 workflow_version: self.workflow_version.clone(),
357 phase: crate::hash::canonical_value(&self.phase)?,
358 phase_ownership: ownership,
359 obligations,
360 blocking_interaction: self.blocking_interaction.clone(),
361 notices: self.notices.clone(),
362 outcome: self
363 .outcome
364 .as_ref()
365 .map(crate::hash::canonical_value)
366 .transpose()?,
367 // Filled by the registry, which is the layer that still holds the
368 // state: erasing a view has already dropped it.
369 state: Vec::new(),
370 })
371 }
372}
373
374/// One value a case holds, as the stage that ANSWERS may state it.
375///
376/// # Why a workflow declares this and the runtime does not read it
377///
378/// The composer knows what every case still NEEDS — obligations travel on the
379/// view and reach the writer as facts — and nothing at all about what it
380/// already HAS. So a question the user asks about their own record («what did
381/// you save as the company name?») arrives at the answering stage with the list
382/// of missing fields, the sources, and no answer in it. What comes back is
383/// «nothing is set», in good faith, over a record that holds the value.
384///
385/// The state itself cannot be handed over wholesale: it is the workflow's own
386/// type, it may carry values a person must not be told back, and a JSON dump is
387/// not a fact. So the workflow says which of its values are sayable and under
388/// which name, and the runtime turns each into a
389/// [`NarratableFact::StateValue`](crate::response::NarratableFact::StateValue)
390/// — the variant that has existed for exactly this and that nothing filled.
391///
392/// Returning nothing, the default, keeps the previous behaviour.
393#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
394pub struct StateField {
395 /// The path the workflow names it with.
396 pub field: String,
397 /// The value, as the writer may state it.
398 pub value: serde_json::Value,
399 /// Whether it tells this record apart from others of its workflow, such as a
400 /// number or a counterpart's name. Identifying fields are what a record is shown
401 /// by when the model has to choose among several.
402 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
403 pub identifying: bool,
404}
405
406impl StateField {
407 /// A field named `field` holding `value`.
408 #[must_use]
409 pub fn new(field: impl Into<String>, value: serde_json::Value) -> Self {
410 Self {
411 field: field.into(),
412 value,
413 identifying: false,
414 }
415 }
416
417 /// Marks the field as telling this record apart from others.
418 #[must_use]
419 pub const fn identifying(mut self) -> Self {
420 self.identifying = true;
421 self
422 }
423}
424
425/// Shorthand for the view type of a definition.
426pub type ViewOf<W> = WorkflowView<
427 <W as WorkflowDefinition>::Phase,
428 <W as WorkflowDefinition>::Obligation,
429 <W as WorkflowDefinition>::Outcome,
430>;
431
432/// What must already be true of **another** case before this workflow may be
433/// started.
434///
435/// # Why a declaration
436///
437/// A projector is pure and cannot read another case, and the executor is the wrong place
438/// for a domain rule about when a case may exist: «a traveler is registered only while a
439/// trip is open» belongs to neither. The workflow declares what it needs; the runtime,
440/// which has the other cases in hand already, decides whether it is there.
441///
442/// # What it costs when it is not met
443///
444/// The start act is not offered at all, so it cannot be proposed and cannot be
445/// refused later. The [`reason`](Self::reason) travels in the catalogue instead,
446/// which is the honest limit of this shape: with no act there is no refusal to
447/// attach a notice to, so whether the user hears *why* depends on the sentence
448/// the writing stage produces. A deployment that needs the reason guaranteed
449/// should keep an operation that refuses it rather than one that is absent.
450#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
451pub struct StartPrecondition {
452 /// The workflow the required case belongs to.
453 pub workflow: WorkflowKey,
454 /// The phases that satisfy it, as that workflow's phase serializes.
455 pub phases: Vec<serde_json::Value>,
456 /// Why, in words a user can read, when it is not met.
457 pub reason: LocalizedText,
458}
459
460impl StartPrecondition {
461 /// Requires a case of `workflow` to be in one of `phases`, already
462 /// serialized.
463 ///
464 /// The infallible primitive. Prefer [`Self::requires`], which takes the
465 /// other workflow's own phase type so the compiler checks the spelling.
466 #[must_use]
467 pub fn new(
468 workflow: impl Into<WorkflowKey>,
469 phases: impl IntoIterator<Item = serde_json::Value>,
470 reason: LocalizedText,
471 ) -> Self {
472 Self {
473 workflow: workflow.into(),
474 phases: phases.into_iter().collect(),
475 reason,
476 }
477 }
478
479 /// Requires a case of `workflow` to be in one of `phases`.
480 ///
481 /// The phases are serialized here, so an adopter names them with the other
482 /// workflow's own type and the compiler checks the spelling.
483 ///
484 /// # Errors
485 ///
486 /// [`HashError`] when a phase does not serialize, which is the workflow's
487 /// own type failing to round-trip.
488 pub fn requires<P: Serialize>(
489 workflow: impl Into<WorkflowKey>,
490 phases: &[P],
491 reason: LocalizedText,
492 ) -> Result<Self, HashError> {
493 let mut serialized = Vec::with_capacity(phases.len());
494 for phase in phases {
495 serialized.push(crate::hash::canonical_value(phase)?);
496 }
497 Ok(Self {
498 workflow: workflow.into(),
499 phases: serialized,
500 reason,
501 })
502 }
503
504 /// Whether `view` is a case that satisfies this precondition.
505 #[must_use]
506 pub fn satisfied_by(&self, view: &ErasedWorkflowView) -> bool {
507 view.case_ref.workflow == self.workflow && self.phases.contains(&view.phase)
508 }
509}
510
511/// What starting this workflow means when the account already has a case of it.
512///
513/// `StartWorkflow` is the runtime's own door: no catalogue entry, no target. What
514/// *start* means differs by workflow: a second trip is ordinary, a second profile of the
515/// same traveler is not, and only the domain knows which.
516///
517/// Under [`ResumesOpenCase`](Self::ResumesOpenCase) the act reaches the case
518/// the turn can already see. Seeing one it resolves to it, so `compile_act`
519/// will be asked to start an open case and the honest answers are no commands
520/// or a rejection; seeing several the act is refused, because the door carries
521/// no target and a selection card would come back with the same question;
522/// seeing none it mints as before.
523///
524/// The limit is the directory's: a case the turn did not load cannot be seen. An
525/// operation aimed at a new record is governed by
526/// [`WorkflowDefinition::may_open_beside`] instead.
527#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
528#[serde(rename_all = "snake_case")]
529pub enum StartBehaviour {
530 /// Every start opens a new record.
531 ///
532 /// The default, and what every workflow did before this existed.
533 #[default]
534 OpensNewCase,
535 /// A start reaches the case that is already open, when the turn sees one.
536 ResumesOpenCase,
537}
538
539impl StartBehaviour {
540 /// Whether a start should reach a case the turn already has.
541 #[must_use]
542 pub const fn resumes(self) -> bool {
543 matches!(self, Self::ResumesOpenCase)
544 }
545}
546
547/// What a confirmation the **policy engine** raised is about, in the domain's
548/// own words.
549///
550/// # The card nobody could write
551///
552/// An operation whose confirmation policy demands a click gets a card from the
553/// policy engine, built from `ConfirmationCopy`: a title and two labels, all
554/// per-*kind*. So every act in a deployment that needs a click and declares no
555/// card of its own draws the same box. Asked to delete a draft, a user got
556/// "Confermi?" over Cancel and Confirm, with nothing anywhere naming what was
557/// about to be deleted.
558///
559/// It is the right card for the policy — a destructive act must be a click, and
560/// it was — and the wrong card for a person, who is being asked to confirm
561/// something the card does not name.
562///
563/// A workflow can already describe a card its **projection** declares, through
564/// `build_interaction`. That door is shut for a confirmation the engine raises,
565/// because there is no requirement to hang it on. This is the same door on that
566/// path: the engine is the only layer that knows a confirmation is *needed*, and
567/// the domain is the only one that knows what it is *about*.
568///
569/// Returning `None` keeps the generic box, so a workflow that says nothing sees
570/// no change.
571#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
572pub struct ConfirmationSubject {
573 /// Replaces the per-kind title, when the domain has a better question.
574 #[serde(default, skip_serializing_if = "Option::is_none")]
575 pub title: Option<LocalizedText>,
576 /// Sets the body, which the generic card has none of.
577 #[serde(default, skip_serializing_if = "Option::is_none")]
578 pub body: Option<LocalizedText>,
579}
580
581impl ConfirmationSubject {
582 /// A confirmation that asks its own question.
583 #[must_use]
584 pub fn asking(title: LocalizedText) -> Self {
585 Self {
586 title: Some(title),
587 body: None,
588 }
589 }
590
591 /// A confirmation that keeps the generic question and explains underneath.
592 #[must_use]
593 pub fn describing(body: LocalizedText) -> Self {
594 Self {
595 title: None,
596 body: Some(body),
597 }
598 }
599
600 /// Adds a body to a question.
601 #[must_use]
602 pub fn with_body(mut self, body: LocalizedText) -> Self {
603 self.body = Some(body);
604 self
605 }
606}
607
608/// Which of the two writing stages a briefing is addressed to.
609///
610/// They are told apart because they must be. The transition acknowledges what
611/// happened and asks for what is still open; the answer stage answers what the
612/// user asked. A string written for one of them and delivered to both is how
613/// the runtime's own care — keeping every question away from the transition —
614/// was undone by an adopter sentence it had no way to notice.
615#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
616#[serde(rename_all = "snake_case")]
617pub enum WritingStage {
618 /// The acknowledgement, which is also the stage that asks.
619 Transition,
620 /// The block that answers one question the user asked.
621 Answer,
622}
623
624/// One value a field accepts, as the workflow names it to a user.
625#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
626pub struct EnumeratedValue {
627 /// Stable identifier, as the domain stores it.
628 pub id: String,
629 /// What a person calls it, in the languages the workflow answers in.
630 pub label: LocalizedText,
631}
632
633impl EnumeratedValue {
634 /// A value with an id and a label.
635 #[must_use]
636 pub fn new(id: impl Into<String>, label: LocalizedText) -> Self {
637 Self {
638 id: id.into(),
639 label,
640 }
641 }
642}
643
644crate::ids::string_id! {
645 /// A field or concept a question may be about, as a workflow names it:
646 /// `proposed.company.registered_address`.
647 QuestionReference
648}
649
650/// Every value one subject accepts, and nothing else.
651///
652/// Declared per view, so a workflow whose accepted values depend on the phase
653/// or on the record declares what is true now rather than what is true in
654/// general. See
655/// [`WorkflowDefinition::enumerations`] for what the runtime does with it.
656#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
657pub struct DomainEnumeration {
658 /// The field or concept, in the vocabulary a plan's questions use.
659 pub subject: QuestionReference,
660 /// The complete set. A partial one would be worse than none: it would make
661 /// the runtime state as exhaustive a list that is not.
662 pub values: Vec<EnumeratedValue>,
663 /// A sentence to put before the values, in the workflow's own words.
664 ///
665 /// Optional, and server-authored like a receipt's body. The runtime writes
666 /// no sentence of its own here, because a sentence about a domain is the
667 /// domain's to write.
668 #[serde(default, skip_serializing_if = "Option::is_none")]
669 pub preamble: Option<LocalizedText>,
670}
671
672impl DomainEnumeration {
673 /// A subject and the values it accepts.
674 #[must_use]
675 pub fn new(subject: impl Into<QuestionReference>, values: Vec<EnumeratedValue>) -> Self {
676 Self {
677 subject: subject.into(),
678 values,
679 preamble: None,
680 }
681 }
682
683 /// Adds the sentence that introduces the values.
684 #[must_use]
685 pub fn with_preamble(mut self, preamble: LocalizedText) -> Self {
686 self.preamble = Some(preamble);
687 self
688 }
689}
690
691/// How much of a workflow's guidance reaches the model, per case.
692///
693/// # Why nothing is bounded by default
694///
695/// The right number depends on the context window of the model a deployment
696/// runs and on what else that deployment puts in a turn. A workflow knows
697/// neither, and neither does the library, so a shipped number would be a guess
698/// made once on behalf of everybody — and it would be a guess that silently
699/// deleted the end of a workflow's guidance for anyone whose situation it did
700/// not fit. Nothing is bounded until a deployment says so, and there is no
701/// ceiling on what it may say: a limit an adopter cannot raise is not a
702/// configurable limit.
703///
704/// # Why cutting is right here and wrong elsewhere
705///
706/// This is the adopter's own text on its way *into* a prompt, not a model's
707/// answer on its way out to a user. Shortening what a deployment wrote costs
708/// the model some guidance and is visible in what was sent; shortening what a
709/// model wrote hands a person half a sentence. And the cut is marked, so a
710/// briefing that arrives as a prefix says so rather than reading as a complete
711/// rule that happens to be shorter.
712#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
713#[serde(transparent)]
714pub struct BriefingBudget {
715 max_bytes: Option<usize>,
716}
717
718impl BriefingBudget {
719 /// The default: no bound at all.
720 #[must_use]
721 pub const fn conservative() -> Self {
722 Self { max_bytes: None }
723 }
724
725 /// A budget of `max_bytes`, taken as given.
726 #[must_use]
727 pub const fn new(max_bytes: usize) -> Self {
728 Self {
729 max_bytes: Some(max_bytes),
730 }
731 }
732
733 /// The bound in force, or `None` when the briefing is not bounded.
734 #[must_use]
735 pub const fn max_bytes(self) -> Option<usize> {
736 self.max_bytes
737 }
738
739 /// Cuts `text` to the budget, marking the cut. Returns it unchanged when
740 /// no budget is set, which is what ships.
741 ///
742 /// The marker matters more than the limit. A model shown a prefix with no
743 /// sign that it is one will follow half a rule as though it were the whole
744 /// rule, which is worse than never having been told the rule.
745 #[must_use]
746 pub fn apply(self, text: &str) -> String {
747 let Some(max_bytes) = self.max_bytes else {
748 return text.to_owned();
749 };
750 if text.len() <= max_bytes {
751 return text.to_owned();
752 }
753 let mut end = max_bytes;
754 while end > 0 && !text.is_char_boundary(end) {
755 end -= 1;
756 }
757 format!("{}… [briefing truncated]", &text[..end])
758 }
759}
760
761/// A workflow: pure projection plus deterministic compilation and policy
762/// (spec §8.2).
763pub trait WorkflowDefinition: Send + Sync + 'static {
764 /// Persisted case state.
765 type State: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + 'static;
766 /// Lifecycle phase.
767 type Phase: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + Eq + 'static;
768 /// Open obligation, possibly parameterized.
769 type Obligation: Clone
770 + Send
771 + Sync
772 + Serialize
773 + serde::de::DeserializeOwned
774 + Eq
775 + std::hash::Hash
776 + 'static;
777 /// Typed command.
778 type Command: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + 'static;
779 /// Typed domain event.
780 type Event: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + 'static;
781 /// Terminal outcome.
782 type Outcome: Clone + Send + Sync + Serialize + serde::de::DeserializeOwned + Eq + 'static;
783
784 /// Stable key.
785 fn key(&self) -> WorkflowKey;
786
787 /// Version; must change when projection semantics change (spec §8.4).
788 fn version(&self) -> WorkflowVersion;
789
790 /// Who must act in a phase. Drives the §8.4 invariants.
791 fn phase_ownership(&self, phase: &Self::Phase) -> PhaseOwnership;
792
793 /// Pure projection (I2).
794 ///
795 /// An absent state means the case does not exist **yet**, and never that it
796 /// no longer does: a case's identity outlives its content, so removal is a
797 /// status the state carries and never an absence. Its one use is a case the
798 /// application's case directory has offered but that has not been created,
799 /// which is projected to a pre-draft phase and never to a terminal one. A
800 /// projector that gives an absent state a terminal phase or an outcome is
801 /// reported by the state explorer in the test kit.
802 fn project(&self, case_ref: CaseRef, state: Option<&Self::State>) -> ViewOf<Self>;
803
804 /// What this case HOLDS, for the stage that answers questions.
805 ///
806 /// The obligations on the view say what a case still needs; this says what
807 /// it already has. Without it the answering stage is handed a list of
808 /// missing fields and no values, and «what did you record as the company
809 /// name?» comes back as «nothing», truthfully as far as the brief goes,
810 /// over a record that holds one.
811 ///
812 /// Takes the state rather than the view for the same reason
813 /// [`Self::compile_act`] does: the view is a projection and a projection
814 /// drops the values.
815 ///
816 /// Declare only what a person may be told back. A value under an
817 /// obligation is fine — it is theirs, they gave it — and a value the domain
818 /// holds for its own bookkeeping is not.
819 ///
820 /// The default is empty, which is the previous behaviour.
821 fn narratable_state(&self, state: Option<&Self::State>) -> Vec<StateField> {
822 let _ = state;
823 Vec::new()
824 }
825
826 /// One line saying what the workflow is for, shown when a message is split into
827 /// requests. `None`, the default, shows the workflow's key alone.
828 fn summary(&self) -> Option<String> {
829 None
830 }
831
832 /// Terms this workflow's users say, and what they mean here.
833 fn glossary(&self) -> Vec<GlossaryTerm> {
834 Vec::new()
835 }
836
837 /// What one record of this workflow is called, in each language its users speak
838 /// («traveler», «viaggiatore»), for the sentences the runtime writes about one. `None`,
839 /// the default, uses the workflow's key.
840 fn noun(&self) -> Option<crate::locale::LocalizedText> {
841 None
842 }
843
844 /// The operations offered in this view, with their arguments, labels and examples.
845 fn operations(&self, view: &ViewOf<Self>) -> Vec<OperationSpec>;
846
847 /// Every operation a record of this workflow may offer, in any phase. Understanding is
848 /// shown them while none of its records is in view, so a request for one is told there
849 /// is none yet, and offered to open one, instead of going unread. The default is empty.
850 fn record_operations(&self) -> Vec<OperationSpec> {
851 Vec::new()
852 }
853
854 /// Guidance for understanding a turn about a record in *this* view. `None`, the
855 /// default, is ordinary.
856 ///
857 /// It is instructions, not data: never interpolate text a user wrote. It varies by
858 /// view, not by record contents, which keeps it reviewable.
859 fn briefing(&self, view: &ViewOf<Self>) -> Option<String> {
860 let _ = view;
861 None
862 }
863
864 /// One obligation, in words a person would recognise.
865 ///
866 /// The stage that writes is handed obligations as the domain's own serialized
867 /// values, which a `match` needs and a sentence does not. Returning `None`, the
868 /// default, lets the value travel as it is, and a structured one such as
869 /// `{"fill":{"row":1}}` is then guessed at. Say it as the question it is, «row 1
870 /// has no A: what is it?»: localized copy for a reader, saying what is missing
871 /// rather than what to do about it.
872 fn obligation_sentence(&self, obligation: &Self::Obligation) -> Option<LocalizedText> {
873 let _ = obligation;
874 None
875 }
876
877 /// The act that answers `obligation`, and the values it already knows.
878 ///
879 /// Asked «row 1 has no A: what is it?», the user answers «X»: the answer is A, and
880 /// the row is the obligation's. With an act named here the reply's question carries
881 /// it, so a bare answer completes it. `None`, the default, leaves the answer to be
882 /// routed as any other message.
883 fn obligation_act(
884 &self,
885 state: Option<&Self::State>,
886 obligation: &Self::Obligation,
887 ) -> Option<ObligationAct> {
888 let _ = (state, obligation);
889 None
890 }
891
892 /// What must already be true of another case before this workflow may be
893 /// started.
894 ///
895 /// Every precondition must hold, or the workflow's "no case yet" operations
896 /// are absent from the catalogue and the model cannot propose starting it.
897 /// A declaration and not a read: see [`StartPrecondition`], which also says
898 /// what this shape cannot guarantee.
899 ///
900 /// The default is empty, which requires nothing and changes nothing.
901 fn start_preconditions(&self) -> Vec<StartPrecondition> {
902 Vec::new()
903 }
904
905 /// What a confirmation the policy engine raises is about.
906 ///
907 /// Called when an act compiles to commands that policy says need a click,
908 /// with the state and the act that produced them — which is everything
909 /// needed to write "Delete the record for Mario Rossi?" where the engine
910 /// would otherwise draw its per-kind box. See [`ConfirmationSubject`].
911 ///
912 /// The default is `None`, which keeps that box exactly as it was.
913 fn confirmation_subject(
914 &self,
915 state: Option<&Self::State>,
916 view: &ViewOf<Self>,
917 act: &ResolvedAct,
918 ) -> Option<ConfirmationSubject> {
919 let _ = (state, view, act);
920 None
921 }
922
923 /// What starting this workflow means when a case of it is already open.
924 ///
925 /// See [`StartBehaviour`], which carries the whole argument. The default
926 /// mints a new case every time, which is what this door did before the
927 /// declaration existed.
928 fn start_behaviour(&self) -> StartBehaviour {
929 StartBehaviour::OpensNewCase
930 }
931
932 /// Whether a new case may be opened while `open` ones are: an operation aimed at a
933 /// new record, the door [`StartBehaviour`] does not govern.
934 ///
935 /// `open` holds every case of this workflow the turn can address, plus any it
936 /// already minted, projected with no state. Only the domain knows whether a second
937 /// one is reasonable (an unfinished draft beside another is not; one waiting to be
938 /// sent may be), and neither `compile_act` nor `validate_command` sees the others.
939 /// A refusal rejects the act with the domain's own sentence and the rest of the
940 /// turn stands. The default admits everything.
941 fn may_open_beside(&self, open: &[ViewOf<Self>]) -> Result<(), DomainRejection> {
942 let _ = open;
943 Ok(())
944 }
945
946 /// What this workflow wants said while it **acknowledges** the turn and
947 /// asks for what is still open.
948 ///
949 /// # Why there are two of these and not one
950 ///
951 /// There was one, and it went to both writing stages. The runtime goes to
952 /// real lengths to keep a question away from this stage — the questions are
953 /// not in its brief and the words they are made of are cut out of the
954 /// message it is shown — because a model handed a question answers it, and
955 /// the user reads the same explanation twice. Then one adopter string
956 /// reached both stages and put the instruction straight back, and the
957 /// duplicate came out again.
958 ///
959 /// The guarantee has to be whole or it is not one. So a workflow addresses
960 /// each stage by name: what to say while acknowledging is not what to say
961 /// while answering, and a workflow with something for one and nothing for
962 /// the other says exactly that by leaving the other at `None`.
963 ///
964 /// This is the stage that asks, so guidance about *what to ask for and in
965 /// what order* belongs here.
966 ///
967 /// Like [`Self::briefing`] it takes the view and not the state, so the
968 /// guidance varies exactly as much as the projection does and can never
969 /// carry a user's words into a prompt. The composer sees the view projected
970 /// **after** the turn committed, so what it is briefed about is the case as
971 /// it now stands.
972 fn transition_briefing(&self, view: &ViewOf<Self>) -> Option<String> {
973 let _ = view;
974 None
975 }
976
977 /// What this workflow wants said while it **answers** a question the user
978 /// asked.
979 ///
980 /// The other half of [`Self::transition_briefing`], and deliberately a
981 /// different string: guidance about how to explain a domain concept has no
982 /// business reaching the stage whose job is to acknowledge and ask.
983 fn answer_briefing(&self, view: &ViewOf<Self>) -> Option<String> {
984 let _ = view;
985 None
986 }
987
988 /// The complete sets of values this workflow accepts, for the fields where
989 /// there is one.
990 ///
991 /// # The claim nothing was checking
992 ///
993 /// The claim guard verifies that prose does not say an action happened
994 /// when it did not. "These are the values this field accepts" is not a
995 /// claim about an action; it is a claim about the domain, it is exactly as
996 /// harmful when false, and it went out unchecked. A workflow accepting two
997 /// legal forms was asked which forms exist and answered with three, the
998 /// third being a plausible name for nothing. A user who takes that advice
999 /// types a value the domain will refuse.
1000 ///
1001 /// A firmer prompt does not fix it. A sentence naming three plausible
1002 /// things is what a language model produces when nothing decides how many
1003 /// there are, and the values already existed in the workflow — as prose in
1004 /// a briefing, which is to say as a suggestion.
1005 ///
1006 /// # What the runtime does with it
1007 ///
1008 /// A question whose references name an enumerated subject, and whose basis
1009 /// is general domain knowledge — "which forms are there", not "which one
1010 /// does this record have" — is answered from this declaration and no model
1011 /// is asked. The values reach the user as structured data carrying the
1012 /// workflow's own labels, so there is no sentence for a third value to
1013 /// appear in. Guarding prose after the fact was the alternative, and it
1014 /// cannot be done: reading an answer cannot tell an invented value from a
1015 /// real one, which is why the declaration answers instead of checking.
1016 ///
1017 /// It also tells the runtime that such a question *is* answerable, which
1018 /// is what stops it being dropped as the assistant's own next step: the
1019 /// field a flow is collecting is the field a user asks about, and without
1020 /// a declared answer the two are indistinguishable.
1021 ///
1022 /// The default is empty, which declares nothing and changes nothing.
1023 fn enumerations(&self, view: &ViewOf<Self>) -> Vec<DomainEnumeration> {
1024 let _ = view;
1025 Vec::new()
1026 }
1027
1028 /// What the user may do next once the case owes nothing, each an operation with its
1029 /// words: the reply offers them when the case needs nothing more, and only those the
1030 /// runtime's dry run of the act on `state` accepts (I22).
1031 ///
1032 /// The default is empty, which offers nothing.
1033 fn next_steps(&self, state: Option<&Self::State>, view: &ViewOf<Self>) -> Vec<NextStep> {
1034 let _ = (state, view);
1035 Vec::new()
1036 }
1037
1038 /// The documents this case has, for the turn to put in front of the user.
1039 ///
1040 /// # The type that nothing could fill
1041 ///
1042 /// [`ArtifactRef`](crate::event::ArtifactRef),
1043 /// [`ArtifactView`](crate::response::ArtifactView) and
1044 /// [`OperationalReceipt::artifact_refs`](crate::event::OperationalReceipt::artifact_refs)
1045 /// all existed, and no workflow could fill any of them. The only hook that
1046 /// came close is [`Self::receipts`], which is handed the events and nothing
1047 /// else — and whether a case has a document is a question about the state
1048 /// those events folded into, not about the events. A list of "line added"
1049 /// cannot answer it, so every implementation ended at an empty vector.
1050 ///
1051 /// What that cost is concrete: a user was asked to authorise the
1052 /// irreversible transmission of a document they had never seen, because the
1053 /// preview that used to sit in the conversation had nowhere to come from.
1054 ///
1055 /// # Why the view and not the receipt
1056 ///
1057 /// A receipt exists only where something committed. A document does not stop
1058 /// existing on a turn that writes nothing — the user asks a question, or
1059 /// refuses a card and the requirement stays down — and on those turns there
1060 /// is no receipt to hang it on. An artifact outlives the turn that produced
1061 /// it, so it is declared from the projection, like everything else that is
1062 /// true of a case rather than of a moment.
1063 ///
1064 /// # What the runtime does with it
1065 ///
1066 /// One [`ResponseBlock::Artifact`](crate::response::ResponseBlock::Artifact)
1067 /// per declaration, for the cases **the turn was about**. A case the turn
1068 /// merely loaded does not put a document in the reply, for the same reason
1069 /// its briefing does not: a conversation about one record is not an occasion
1070 /// to show another.
1071 ///
1072 /// Re-declaring the same artifact after an edit is not a duplicate and is
1073 /// not suppressed. It is the same document at a later revision, the turn
1074 /// order says which is which, and what a surface does with the earlier ones
1075 /// is a rendering decision the runtime has no business taking.
1076 ///
1077 /// The default is empty, which declares nothing and changes nothing.
1078 fn artifacts(&self, view: &ViewOf<Self>) -> Vec<crate::event::ArtifactRef> {
1079 let _ = view;
1080 Vec::new()
1081 }
1082
1083 /// Compiles a resolved act into typed commands (spec §21.2).
1084 fn compile_act(
1085 &self,
1086 state: Option<&Self::State>,
1087 view: &ViewOf<Self>,
1088 act: &ResolvedAct,
1089 ) -> Result<Vec<Self::Command>, DomainRejection>;
1090
1091 /// Why an act that compiled nothing changed nothing, in the reader's own
1092 /// words.
1093 ///
1094 /// Called only when [`compile_act`](Self::compile_act) returned no commands
1095 /// at all. That is a legitimate answer — the state the act asks for is the
1096 /// state the case is already in — but the runtime cannot say which of a
1097 /// workflow's several reasons it was, and the writing stage is handed the
1098 /// operation's name and nothing else.
1099 ///
1100 /// # What that costs when nobody implements it
1101 ///
1102 /// One workflow compiles nothing for two quite different reasons: the
1103 /// singleton whose start was proposed a second time, and the write that
1104 /// tells a field what it already says. Given only «this act changed
1105 /// nothing», a writer invents a reason, and the one it reaches for is a
1106 /// refusal: asked to correct a value and told nothing changed, it answered
1107 /// «I cannot do that here» — which was not true, and left the user with no
1108 /// idea what to say next. The sentence a person needed was «that is already
1109 /// the value; tell me what you want instead», and only the workflow knows
1110 /// it.
1111 ///
1112 /// This is the same channel [`DomainRejection`] gives a refusal, for the
1113 /// same reason: a workflow that wrote a sentence knows more than the
1114 /// runtime does.
1115 ///
1116 /// # The default
1117 ///
1118 /// `None`, which is exactly what every workflow said before this existed —
1119 /// the fact reaches the writing stage naming the operation and no more.
1120 fn nothing_changed(
1121 &self,
1122 state: Option<&Self::State>,
1123 act: &ResolvedAct,
1124 ) -> Option<LocalizedText> {
1125 let _ = (state, act);
1126 None
1127 }
1128
1129 /// Policy of a command. Unknown commands must default to
1130 /// [`CommandPolicy::conservative`].
1131 fn command_policy(&self, state: Option<&Self::State>, command: &Self::Command)
1132 -> CommandPolicy;
1133
1134 /// Deterministic validation before execution.
1135 fn validate_command(
1136 &self,
1137 state: Option<&Self::State>,
1138 command: &Self::Command,
1139 ) -> Result<(), DomainRejection>;
1140
1141 /// The state a valid `command` leaves `state` in, when the workflow can tell without
1142 /// executing it. An act on a record an earlier act of the same message opens is
1143 /// compiled and validated against the state the opening leaves; with `None`, the
1144 /// default, it is checked against no state, and a domain that needs the record refuses.
1145 fn state_after(
1146 &self,
1147 state: Option<&Self::State>,
1148 command: &Self::Command,
1149 ) -> Option<Self::State> {
1150 let _ = (state, command);
1151 None
1152 }
1153
1154 /// Renders receipts from committed events (spec §17.3).
1155 ///
1156 /// The events arrive as [`ReceiptEvent`]s, not bare payloads, because a
1157 /// receipt must cite the [`EventId`](crate::ids::EventId)s that authorize
1158 /// its claim (I16): a `Success` receipt with no event ids is refused by
1159 /// [`claim_guard::verify`](crate::response::claim_guard::verify). Derive the
1160 /// receipt id with
1161 /// [`ReceiptId::derive`](crate::ids::ReceiptId::derive) so a replayed turn
1162 /// renders the same receipts.
1163 ///
1164 /// # The redacted case
1165 ///
1166 /// [`ReceiptEvent::Redacted`] means the event happened and its payload was
1167 /// erased (see the [`event`](crate::event) module). The event is still in
1168 /// the ledger, at its position, with its identity and its type, so a
1169 /// receipt rendered over it is still backed and still passes the claim
1170 /// guard — but it cannot say what changed, and it must not read as though
1171 /// it could. Write copy that is true of an erased event: that this step is
1172 /// on record and its detail is gone. Rendering nothing at all is worse than
1173 /// it looks, because a turn that quietly drops a receipt reads as a turn in
1174 /// which nothing happened.
1175 fn receipts(
1176 &self,
1177 events: &[ReceiptEvent<Self::Event>],
1178 locale: &Locale,
1179 ) -> Vec<OperationalReceipt>;
1180
1181 /// Turns a requirement of the view into a full interaction spec.
1182 ///
1183 /// The default uses the requirement's payload, so a requirement that
1184 /// carries none must be completed here: the engine validates the result
1185 /// ([`InteractionSpec::validate`]) and refuses a card nobody could answer.
1186 fn build_interaction(
1187 &self,
1188 state: Option<&Self::State>,
1189 view: &ViewOf<Self>,
1190 requirement: &InteractionRequirement,
1191 ) -> Result<InteractionSpec, DomainRejection> {
1192 let _ = state;
1193 Ok(requirement.to_spec(view.case_ref.clone()))
1194 }
1195}
1196
1197/// Loads and mutates cases of one workflow (spec §8.2).
1198#[async_trait::async_trait]
1199pub trait WorkflowExecutor<W: WorkflowDefinition>: Send + Sync {
1200 /// Loads a case for an account. `None` value means the case does not exist;
1201 /// its revision is then [`crate::ids::CaseRevision::ZERO`]. An executor
1202 /// must never delete a case to express completion, because a completed case
1203 /// keeps its row and moves to a terminal status, so an absent state always
1204 /// means the case has not been created yet.
1205 async fn load(
1206 &self,
1207 account: &AccountId,
1208 case_id: &CaseId,
1209 ) -> Result<Versioned<Option<W::State>>, StoreError>;
1210
1211 /// Executes a batch under its atomicity scope with revision and
1212 /// idempotency checks (I13, I14).
1213 async fn execute(
1214 &self,
1215 batch: CommandBatch<W::Command>,
1216 ) -> Result<Commit<W::State, W::Event>, ExecutionError>;
1217}
1218
1219/// The read-only half of [`WorkflowExecutor`]: loading a case, and nothing
1220/// else (spec §8.2).
1221///
1222/// It exists for callers that must be unable to mutate a case — the plan-only
1223/// turn path of `turnframe-runtime` above all, which needs the state a turn is
1224/// planned against and must not be able to execute a batch. A guarantee that
1225/// says "this code simply never calls `execute`" is not a guarantee; being
1226/// handed a value that has no `execute` is.
1227///
1228/// **There is nothing to implement.** Every [`WorkflowExecutor`] is a
1229/// `CaseLoader` through the blanket implementation below, so an adopter writes
1230/// exactly what they write today. The method is called `load_case` rather than
1231/// `load` so that a type which is both never makes a call site ambiguous.
1232///
1233/// ```rust
1234/// # use turnframe_core::flow::{CaseLoader, WorkflowDefinition, WorkflowExecutor};
1235/// /// Accepts any executor, but can only read through it.
1236/// fn planning_only<W: WorkflowDefinition, E: WorkflowExecutor<W>>(executor: E) -> impl CaseLoader<W> {
1237/// executor
1238/// }
1239/// ```
1240#[async_trait::async_trait]
1241pub trait CaseLoader<W: WorkflowDefinition>: Send + Sync {
1242 /// Loads a case for an account. `None` value means the case does not exist;
1243 /// its revision is then [`crate::ids::CaseRevision::ZERO`]. See
1244 /// [`WorkflowExecutor::load`] for what an absent state does and does not
1245 /// mean.
1246 ///
1247 /// # Errors
1248 ///
1249 /// [`StoreError`] when the case could not be read.
1250 async fn load_case(
1251 &self,
1252 account: &AccountId,
1253 case_id: &CaseId,
1254 ) -> Result<Versioned<Option<W::State>>, StoreError>;
1255}
1256
1257/// Every executor loads.
1258#[async_trait::async_trait]
1259impl<W, E> CaseLoader<W> for E
1260where
1261 W: WorkflowDefinition,
1262 E: WorkflowExecutor<W> + ?Sized,
1263{
1264 async fn load_case(
1265 &self,
1266 account: &AccountId,
1267 case_id: &CaseId,
1268 ) -> Result<Versioned<Option<W::State>>, StoreError> {
1269 self.load(account, case_id).await
1270 }
1271}
1272
1273/// A shared executor is an executor.
1274///
1275/// [`WorkflowRegistryBuilder::register`] takes the executor by value, so an
1276/// application that also holds its own handle on it — to seed a case, to read a
1277/// revision back, to share one connection pool between two workflows — would
1278/// otherwise have to wrap the `Arc` in a newtype just to re-implement two
1279/// forwarding methods. This impl is that newtype, written once.
1280#[async_trait::async_trait]
1281impl<W, E> WorkflowExecutor<W> for std::sync::Arc<E>
1282where
1283 W: WorkflowDefinition,
1284 E: WorkflowExecutor<W> + ?Sized,
1285{
1286 async fn load(
1287 &self,
1288 account: &AccountId,
1289 case_id: &CaseId,
1290 ) -> Result<Versioned<Option<W::State>>, StoreError> {
1291 (**self).load(account, case_id).await
1292 }
1293
1294 async fn execute(
1295 &self,
1296 batch: CommandBatch<W::Command>,
1297 ) -> Result<Commit<W::State, W::Event>, ExecutionError> {
1298 (**self).execute(batch).await
1299 }
1300}
1301
1302#[cfg(test)]
1303mod tests {
1304 use super::*;
1305 use crate::ids::CaseRevision;
1306
1307 #[test]
1308 fn requirement_to_spec_defaults() {
1309 let req = InteractionRequirement::blocking("send", InteractionKind::ConfirmCommand);
1310 let spec = req.to_spec(CaseRef::new("trip", "i1", CaseRevision(1)));
1311 assert_eq!(spec.key, "send");
1312 assert!(spec.blocking);
1313 assert!(spec.binds_to_revision);
1314 assert_eq!(spec.payload.title.default, "send");
1315 }
1316
1317 #[test]
1318 fn erase_produces_stable_obligation_ids() {
1319 #[derive(Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
1320 enum Ob {
1321 Line { id: u32 },
1322 }
1323 let view: WorkflowView<&str, Ob, ()> = WorkflowView::new(
1324 CaseRef::new("w", "c", CaseRevision(1)),
1325 WorkflowVersion::from("1"),
1326 "collecting",
1327 )
1328 .with_obligations([Ob::Line { id: 2 }, Ob::Line { id: 1 }]);
1329 let erased = view.erase(PhaseOwnership::System).unwrap();
1330 assert_eq!(erased.obligations[0].id.as_str(), r#"{"Line":{"id":2}}"#);
1331 assert_eq!(erased.phase, serde_json::json!("collecting"));
1332 assert!(erased.outcome.is_none());
1333 }
1334}
1335
1336/// The act that answers an obligation: its operation, the values the obligation fixes,
1337/// and the values the answer gives.
1338#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1339#[non_exhaustive]
1340pub struct ObligationAct {
1341 /// The operation.
1342 pub operation: crate::ids::OperationKey,
1343 /// The arguments the obligation fixes, by name.
1344 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
1345 pub given: std::collections::BTreeMap<String, serde_json::Value>,
1346 /// The arguments the answer gives.
1347 pub asks: Vec<String>,
1348}
1349
1350impl ObligationAct {
1351 /// An act of `operation` asking for `asks`, with nothing fixed yet.
1352 #[must_use]
1353 pub fn new(
1354 operation: impl Into<crate::ids::OperationKey>,
1355 asks: impl IntoIterator<Item = impl Into<String>>,
1356 ) -> Self {
1357 Self {
1358 operation: operation.into(),
1359 given: std::collections::BTreeMap::new(),
1360 asks: asks.into_iter().map(Into::into).collect(),
1361 }
1362 }
1363
1364 /// Fixes argument `name` to `value`.
1365 #[must_use]
1366 pub fn given(mut self, name: impl Into<String>, value: serde_json::Value) -> Self {
1367 self.given.insert(name.into(), value);
1368 self
1369 }
1370}