Skip to main content

turnframe_store/
interaction.rs

1//! Persistent interactions: immutable payloads, the one-blocking-per-case slot
2//! and compare-and-swap resolution (spec §15.5, §15.6).
3//!
4//! # Contract
5//!
6//! * **Payloads are immutable.** There is no method that changes a persisted
7//!   payload; the client only ever echoes identifiers back (I7).
8//! * **One blocking interaction per case.** At most one interaction with
9//!   `blocking == true` whose status is open (`Active` or `Resolving`) may
10//!   exist per `(account, workflow, case_id)`. [`InteractionWriter::insert`]
11//!   refuses a second one with `Conflict`;
12//!   [`InteractionWriter::insert_replacing_blocking`] invalidates the `Active`
13//!   occupant (reason [`InvalidationReason::Superseded`]) and inserts the new
14//!   card under a new identifier. A `Resolving` occupant is never replaced: its
15//!   commands are executing, so the call fails with `Conflict`.
16//! * **Tenant isolation.** Every lookup is scoped by account; an identifier of
17//!   another tenant is `NotFound`, indistinguishable from an unknown one.
18//! * **Compare-and-swap resolution.** [`InteractionWriter::begin_resolution`]
19//!   moves `Active → Resolving` only when the current status equals the
20//!   expected one; [`InteractionWriter::finish_resolution`] settles a `Resolving`
21//!   interaction to `Resolved`, `Failed` or back to `Active`. Repeating a finish
22//!   with the same outcome is accepted (idempotent); a different outcome after a
23//!   terminal one is `Conflict`.
24//! * **Revision invalidation.** A case revision change invalidates `Active`
25//!   interactions bound to another revision unless they declare revision
26//!   independence (`revision_independent == true`). `Resolving` interactions
27//!   are left alone: their own commit is usually what moved the revision, and
28//!   `finish_resolution` settles them.
29//! * **Expiry** moves `Active` interactions whose `expires_at` has passed to
30//!   `Expired`; it is a system sweep across tenants.
31
32use async_trait::async_trait;
33use chrono::{DateTime, Utc};
34use serde::{Deserialize, Serialize};
35use turnframe_core::case::CaseKey;
36use turnframe_core::ids::{
37    AccountId, CaseRevision, ConversationId, EventId, InteractionId, OptionId, TurnId,
38};
39use turnframe_core::interaction::{Interaction, InteractionStatus};
40
41use crate::error::StoreError;
42
43/// How a `Resolving` interaction is settled.
44#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
45#[serde(tag = "kind", rename_all = "snake_case")]
46pub enum ResolutionOutcome {
47    /// The associated commands committed; the events authorize the receipt.
48    Resolved {
49        /// Committed events backing the resolution.
50        event_ids: Vec<EventId>,
51    },
52    /// The associated commands failed and policy keeps the card closed.
53    Failed {
54        /// Stable failure code (never free text).
55        code: String,
56    },
57    /// The associated commands failed and policy restores the card so the user
58    /// may answer again.
59    RestoreActive,
60}
61
62impl ResolutionOutcome {
63    /// The status the outcome leads to.
64    #[must_use]
65    pub fn target_status(&self) -> InteractionStatus {
66        match self {
67            Self::Resolved { .. } => InteractionStatus::Resolved,
68            Self::Failed { .. } => InteractionStatus::Failed,
69            Self::RestoreActive => InteractionStatus::Active,
70        }
71    }
72}
73
74/// Why an interaction was invalidated.
75#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
76#[serde(tag = "kind", rename_all = "snake_case")]
77#[non_exhaustive]
78pub enum InvalidationReason {
79    /// The case moved to another revision.
80    RevisionChanged,
81    /// A new blocking interaction replaced it.
82    Superseded {
83        /// The replacement.
84        by: InteractionId,
85    },
86    /// The case reached a terminal phase.
87    CaseClosed,
88    /// An operator or policy decided, identified by a stable code.
89    Administrative {
90        /// Stable code.
91        code: String,
92    },
93}
94
95/// What is recorded when an interaction is invalidated.
96#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
97pub struct InvalidationRecord {
98    /// The reason.
99    pub reason: InvalidationReason,
100    /// The revision the case moved to, when the invalidation came from a
101    /// revision sweep ([`InteractionWriter::invalidate_for_case`]).
102    ///
103    /// `None` when a replacement card or a policy closed this one, where no
104    /// revision is involved.
105    #[serde(default, skip_serializing_if = "Option::is_none")]
106    pub new_revision: Option<CaseRevision>,
107    /// When it happened, stamped by the store's clock.
108    pub at: DateTime<Utc>,
109}
110
111/// A persisted interaction together with the store-owned resolution metadata
112/// that the core record does not carry.
113#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
114pub struct InteractionRecord {
115    /// The interaction as the core sees it (payload, status, resolved option).
116    pub interaction: Interaction,
117    /// Turn whose response began the resolution.
118    #[serde(default, skip_serializing_if = "Option::is_none")]
119    pub resolved_by_turn: Option<TurnId>,
120    /// Events that backed a `Resolved` outcome. A second click on a resolved
121    /// interaction is answered from these (spec §15.5).
122    #[serde(default)]
123    pub resolution_event_ids: Vec<EventId>,
124    /// Stable code of a `Failed` outcome.
125    #[serde(default, skip_serializing_if = "Option::is_none")]
126    pub failure_code: Option<String>,
127    /// Why and when it was invalidated, when status is `Invalidated`.
128    #[serde(default, skip_serializing_if = "Option::is_none")]
129    pub invalidation: Option<InvalidationRecord>,
130}
131
132impl InteractionRecord {
133    /// Wraps a freshly inserted interaction with empty metadata.
134    #[must_use]
135    pub fn new(interaction: Interaction) -> Self {
136        Self {
137            interaction,
138            resolved_by_turn: None,
139            resolution_event_ids: Vec::new(),
140            failure_code: None,
141            invalidation: None,
142        }
143    }
144
145    /// The interaction identifier.
146    #[must_use]
147    pub fn id(&self) -> InteractionId {
148        self.interaction.id
149    }
150
151    /// The current status.
152    #[must_use]
153    pub fn status(&self) -> InteractionStatus {
154        self.interaction.status
155    }
156}
157
158/// The read half of the interaction contract.
159///
160/// It answers what cards exist and what they say, and can create, resolve,
161/// invalidate or expire none of them. This is the half the plan-only path of
162/// [`turnframe-runtime`](https://docs.rs/turnframe-runtime) holds, which is why
163/// a shadow turn cannot leave a card a real user could click.
164#[async_trait]
165pub trait InteractionReader: Send + Sync {
166    /// Loads an interaction of `account`.
167    ///
168    /// # Errors
169    /// * `NotFound` when it does not exist for this account.
170    async fn get(
171        &self,
172        account: &AccountId,
173        id: &InteractionId,
174    ) -> Result<InteractionRecord, StoreError>;
175
176    /// Lists the open (`Active` or `Resolving`) interactions of a conversation,
177    /// ordered by `created_at` then id.
178    ///
179    /// # Errors
180    /// * [`StoreError`] when the listing could not be read.
181    async fn list_open_for_conversation(
182        &self,
183        account: &AccountId,
184        conversation: &ConversationId,
185    ) -> Result<Vec<Interaction>, StoreError>;
186
187    /// Lists the open (`Active` or `Resolving`) interactions of a case, ordered
188    /// by `created_at` then id.
189    ///
190    /// # Errors
191    /// * [`StoreError`] when the listing could not be read.
192    async fn list_open_for_case(
193        &self,
194        account: &AccountId,
195        case_key: &CaseKey,
196    ) -> Result<Vec<Interaction>, StoreError>;
197
198    /// Whether a blocking card of `case_key` bound to `revision` has already
199    /// been answered by the user.
200    ///
201    /// # The loop this ends
202    ///
203    /// A `ConfirmCommand` payload must offer a way to decline, and declining
204    /// ends the card **without effect**. Nothing is written, the case does not
205    /// move, and on the next turn the projection declares the same requirement
206    /// and the same card goes back up. The user presses "not now" and is asked
207    /// again, for ever.
208    ///
209    /// The domain cannot see this: a projector reads state, and a decline
210    /// leaves none. The runtime can, because a requirement belongs to a case at
211    /// a revision, and "not now" means "ask me when the document changes". A
212    /// card answered at the revision the case is still on has been answered.
213    ///
214    /// Answered is `Resolved`, `Declined` or `Dismissed` — the three the user
215    /// causes. A card the case outgrew (`Invalidated`) or that timed out
216    /// (`Expired`) was not answered by anybody, and a requirement still standing
217    /// after one is raised again as it always was.
218    ///
219    /// # Errors
220    /// * [`StoreError`] when the lookup could not be read.
221    async fn blocking_answered_at(
222        &self,
223        account: &AccountId,
224        case_key: &CaseKey,
225        revision: CaseRevision,
226    ) -> Result<bool, StoreError>;
227}
228
229/// The write half of the interaction contract.
230#[async_trait]
231pub trait InteractionWriter: Send + Sync {
232    /// Inserts a new `Active` interaction.
233    ///
234    /// # Errors
235    /// * `Other(INVALID_RECORD)` when `interaction.status` is not `Active`.
236    /// * `Conflict` when the identifier exists for the account, or when the
237    ///   interaction is blocking and its case already has an open blocking
238    ///   interaction (I5).
239    async fn insert(&self, interaction: Interaction) -> Result<(), StoreError>;
240
241    /// Inserts a new `Active` interaction, invalidating the `Active` blocking
242    /// occupant of the same case when there is one. Returns the identifiers it
243    /// invalidated (empty when the slot was free or the new card is not
244    /// blocking).
245    ///
246    /// # Errors
247    /// * `Other(INVALID_RECORD)` when `interaction.status` is not `Active`.
248    /// * `Conflict` when the identifier exists, or when the occupant is
249    ///   `Resolving` (its commands are executing; it cannot be replaced).
250    async fn insert_replacing_blocking(
251        &self,
252        interaction: Interaction,
253    ) -> Result<Vec<InteractionId>, StoreError>;
254
255    /// Compare-and-swap `expected_status → Resolving`, recording the chosen
256    /// option and the resolving turn. The only legal source status is
257    /// `Active`; the explicit parameter makes a lost race visible to the
258    /// caller instead of hiding it behind a re-read.
259    ///
260    /// # Errors
261    /// * `NotFound` when the interaction does not exist for `account`.
262    /// * `Conflict` when the current status differs from `expected_status`, or
263    ///   when `expected_status → Resolving` is not a legal transition.
264    async fn begin_resolution(
265        &self,
266        account: &AccountId,
267        id: &InteractionId,
268        expected_status: InteractionStatus,
269        option_id: OptionId,
270        resolved_by: TurnId,
271    ) -> Result<InteractionRecord, StoreError>;
272
273    /// Settles a `Resolving` interaction.
274    ///
275    /// Repeating the call with an outcome whose target status is already the
276    /// current status and whose data (event ids, failure code) matches is
277    /// accepted without change. `RestoreActive` clears the resolving turn, the
278    /// chosen option and `resolved_at` so the card is answerable again.
279    ///
280    /// `Resolved` and `Failed` leave `resolved_at` as
281    /// [`InteractionWriter::begin_resolution`] stamped it: it means "when the
282    /// answer was accepted". When the commands actually committed is recorded
283    /// by the command journal and by the events, and is not duplicated here.
284    ///
285    /// # Errors
286    /// * `NotFound` when the interaction does not exist for `account`.
287    /// * `Conflict` when the status is neither `Resolving` nor the outcome's
288    ///   target, or when the target status matches but the data differs.
289    async fn finish_resolution(
290        &self,
291        account: &AccountId,
292        id: &InteractionId,
293        outcome: ResolutionOutcome,
294    ) -> Result<InteractionRecord, StoreError>;
295
296    /// Invalidates every `Active` interaction of the case that is bound to a
297    /// revision other than `new_revision` and does not declare revision
298    /// independence. Returns the invalidated identifiers in list order.
299    ///
300    /// # Errors
301    /// * [`StoreError`] when the sweep could not be written.
302    async fn invalidate_for_case(
303        &self,
304        account: &AccountId,
305        case_key: &CaseKey,
306        new_revision: CaseRevision,
307        reason: InvalidationReason,
308    ) -> Result<Vec<InteractionId>, StoreError>;
309
310    /// Invalidates every `Active` interaction of the case, whatever revision it
311    /// is bound to, and whether or not it declares revision independence.
312    /// Returns the invalidated identifiers in list order.
313    ///
314    /// This is the operator's path, not the runtime's, and it exists because
315    /// [`invalidate_for_case`](Self::invalidate_for_case) cannot express it: a
316    /// card is invalidated there for having been bound to a revision the case
317    /// has left, so a decision that has nothing to do with the revision — a
318    /// workflow rolled back to a version that cannot compile the option the
319    /// card offers, an account suspended, a card withdrawn — has no way to run.
320    /// Without it the honest rollback step is to leave the user holding an
321    /// option nothing will honour.
322    ///
323    /// A `Resolving` interaction is deliberately left alone: a command it
324    /// authorized is in flight, and taking the card away underneath it would
325    /// settle nothing while making the outcome unattributable.
326    ///
327    /// # Errors
328    /// * [`StoreError`] when the sweep could not be written.
329    async fn invalidate_case_cards(
330        &self,
331        account: &AccountId,
332        case_key: &CaseKey,
333        reason: InvalidationReason,
334    ) -> Result<Vec<InteractionId>, StoreError>;
335
336    /// Moves every `Active` interaction whose `expires_at <= now` to `Expired`,
337    /// across all accounts. Returns the expired identifiers in list order.
338    ///
339    /// # Errors
340    /// * [`StoreError`] when the sweep could not be written.
341    async fn expire_due(&self, now: DateTime<Utc>) -> Result<Vec<InteractionId>, StoreError>;
342}
343
344/// Persistence of interactions with the rules of spec §15.5 and §15.6: both
345/// halves.
346///
347/// See the module documentation for the contract; the conformance suite proves
348/// it (`check_blocking_interaction_*`, `check_cross_tenant_*`,
349/// `check_begin_resolution_cas`, `check_finish_resolution_idempotent`,
350/// `check_revision_invalidation_respects_independence`,
351/// `check_interaction_expiry` in [`crate::conformance`]).
352///
353/// There is nothing to implement here: write [`InteractionReader`] and
354/// [`InteractionWriter`] and the blanket implementation below supplies this
355/// trait.
356pub trait InteractionStore: InteractionReader + InteractionWriter {}
357
358impl<T: InteractionReader + InteractionWriter + ?Sized> InteractionStore for T {}
359
360#[cfg(test)]
361mod tests {
362    use super::*;
363
364    #[test]
365    fn outcome_targets() {
366        assert_eq!(
367            ResolutionOutcome::Resolved { event_ids: vec![] }.target_status(),
368            InteractionStatus::Resolved
369        );
370        assert_eq!(
371            ResolutionOutcome::Failed { code: "x".into() }.target_status(),
372            InteractionStatus::Failed
373        );
374        assert_eq!(
375            ResolutionOutcome::RestoreActive.target_status(),
376            InteractionStatus::Active
377        );
378    }
379
380    #[test]
381    fn reason_serializes_tagged() {
382        let json = serde_json::to_value(InvalidationReason::Superseded {
383            by: InteractionId::nil(),
384        })
385        .unwrap();
386        assert_eq!(json["kind"], "superseded");
387    }
388}