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}