Skip to main content

turnframe_store/
conversation.rs

1//! Conversations, persisted turns and the crash-recovery phase marker
2//! (spec §22.3, §23.1).
3//!
4//! # Contract
5//!
6//! * A conversation belongs to exactly one account. Every lookup is scoped by
7//!   `(account, id)`; a conversation of another tenant is `NotFound`.
8//! * A user turn is persisted **as received** ([`StoredUserTurn`]) and an
9//!   assistant turn is persisted **exactly as returned** to the client: the same
10//!   [`AssistantTurn`] value, with the same ordered blocks, must come back from
11//!   [`ConversationReader::load_turn`]. Reload never reconstructs cards or receipts
12//!   from free text (ADR-012 point 8).
13//! * Appending a user turn also creates its phase marker at
14//!   [`TurnPhase::Received`]. The marker is the crash-recovery anchor of
15//!   spec §23.1: a recovery sweep lists turns whose phase is not terminal and
16//!   decides, per turn, whether to re-interpret, resume by idempotency key or
17//!   regenerate the response.
18//! * A terminal phase ([`TurnPhase::Delivered`], [`TurnPhase::Failed`]) is
19//!   final: setting any other phase afterwards is a `Conflict`. Setting the same
20//!   phase again is accepted (idempotent).
21
22use async_trait::async_trait;
23use chrono::{DateTime, Utc};
24use serde::{Deserialize, Serialize};
25use turnframe_core::ids::{AccountId, ConversationId, TurnId};
26use turnframe_core::replay::TurnPhase;
27use turnframe_core::response::AssistantTurn;
28use turnframe_core::turn::TurnInput;
29
30use crate::error::StoreError;
31
32/// A conversation (a chat thread) owned by an account.
33#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
34pub struct ConversationRecord {
35    /// Identifier.
36    pub id: ConversationId,
37    /// Owning tenant.
38    pub account_id: AccountId,
39    /// Creation time, supplied by the caller.
40    pub created_at: DateTime<Utc>,
41    /// Application-defined metadata (title, channel, external ids...). Never
42    /// interpreted by the library.
43    #[serde(default)]
44    pub metadata: serde_json::Value,
45}
46
47impl ConversationRecord {
48    /// Builds a record without metadata.
49    #[must_use]
50    pub fn new(id: ConversationId, account_id: AccountId, created_at: DateTime<Utc>) -> Self {
51        Self {
52            id,
53            account_id,
54            created_at,
55            metadata: serde_json::Value::Null,
56        }
57    }
58
59    /// Attaches metadata.
60    #[must_use]
61    pub fn with_metadata(mut self, metadata: serde_json::Value) -> Self {
62        self.metadata = metadata;
63        self
64    }
65}
66
67/// A user turn as accepted by the runtime, with the time it was received.
68#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
69pub struct StoredUserTurn {
70    /// The input, unchanged.
71    pub input: TurnInput,
72    /// When the runtime accepted it.
73    pub received_at: DateTime<Utc>,
74}
75
76impl StoredUserTurn {
77    /// Pairs an input with its reception time.
78    #[must_use]
79    pub fn new(input: TurnInput, received_at: DateTime<Utc>) -> Self {
80        Self { input, received_at }
81    }
82
83    /// The owning account (the actor's account).
84    #[must_use]
85    pub fn account_id(&self) -> &AccountId {
86        &self.input.actor.account_id
87    }
88
89    /// The turn identifier.
90    #[must_use]
91    pub fn turn_id(&self) -> TurnId {
92        self.input.turn_id
93    }
94
95    /// The conversation the turn belongs to.
96    #[must_use]
97    pub fn conversation_id(&self) -> ConversationId {
98        self.input.conversation_id
99    }
100}
101
102/// A persisted turn: the user input, the assistant turn once composed, and the
103/// current phase marker.
104#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
105pub struct StoredTurn {
106    /// The user side.
107    pub user: StoredUserTurn,
108    /// The assistant side, present once persisted.
109    #[serde(default, skip_serializing_if = "Option::is_none")]
110    pub assistant: Option<AssistantTurn>,
111    /// Current phase.
112    pub phase: TurnPhase,
113}
114
115/// The crash-recovery marker of one turn (spec §23.1).
116#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
117pub struct TurnPhaseMarker {
118    /// Owning tenant.
119    pub account_id: AccountId,
120    /// The conversation.
121    pub conversation_id: ConversationId,
122    /// The turn.
123    pub turn_id: TurnId,
124    /// Last persisted phase.
125    pub phase: TurnPhase,
126    /// When the phase was last written, stamped by the store's clock.
127    pub updated_at: DateTime<Utc>,
128}
129
130/// Which turns a recovery sweep looks at.
131///
132/// Recovery is a system operation, not a user lookup, so it may legitimately
133/// span every tenant; every returned marker still carries its `account_id`.
134#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
135pub enum RecoveryScope {
136    /// Only turns of one account.
137    Account(AccountId),
138    /// Turns of every account.
139    AllAccounts,
140}
141
142impl RecoveryScope {
143    /// Returns `true` when `account` falls inside the scope.
144    #[must_use]
145    pub fn includes(&self, account: &AccountId) -> bool {
146        match self {
147            Self::Account(scoped) => scoped == account,
148            Self::AllAccounts => true,
149        }
150    }
151}
152
153/// The read half of the conversation contract.
154///
155/// A holder of this trait can answer questions about conversations, turns and
156/// phase markers and cannot change any of them. It is what the plan-only path
157/// of [`turnframe-runtime`](https://docs.rs/turnframe-runtime) is handed, so
158/// "this path does not write" is a fact about its types rather than a promise
159/// about its code.
160#[async_trait]
161pub trait ConversationReader: Send + Sync {
162    /// Loads a conversation of `account`.
163    ///
164    /// # Errors
165    /// * `NotFound` when it does not exist for this account.
166    async fn load_conversation(
167        &self,
168        account: &AccountId,
169        id: &ConversationId,
170    ) -> Result<ConversationRecord, StoreError>;
171
172    /// Loads the most recent `limit` turns of a conversation in chronological
173    /// order (oldest of the selected first). Ordering is by `received_at`, then
174    /// `turn_id`.
175    ///
176    /// # Errors
177    /// * `NotFound` when the conversation does not exist for `account`.
178    async fn load_recent_turns(
179        &self,
180        account: &AccountId,
181        conversation: &ConversationId,
182        limit: usize,
183    ) -> Result<Vec<StoredTurn>, StoreError>;
184
185    /// Loads one turn.
186    ///
187    /// # Errors
188    /// * `NotFound` when it does not exist for `account`.
189    async fn load_turn(
190        &self,
191        account: &AccountId,
192        turn_id: &TurnId,
193    ) -> Result<StoredTurn, StoreError>;
194
195    /// Reads the phase marker of a turn.
196    ///
197    /// # Errors
198    /// * `NotFound` when the turn does not exist for `account`.
199    async fn turn_phase(
200        &self,
201        account: &AccountId,
202        turn_id: &TurnId,
203    ) -> Result<TurnPhaseMarker, StoreError>;
204
205    /// Lists up to `limit` markers of turns whose phase is not terminal, oldest
206    /// first (by `received_at`, then `turn_id`).
207    ///
208    /// # Errors
209    /// * [`StoreError`] when the sweep could not be read.
210    async fn list_unfinished_turns(
211        &self,
212        scope: RecoveryScope,
213        limit: usize,
214    ) -> Result<Vec<TurnPhaseMarker>, StoreError>;
215}
216
217/// The write half of the conversation contract.
218#[async_trait]
219pub trait ConversationWriter: Send + Sync {
220    /// Creates a conversation.
221    ///
222    /// # Errors
223    /// * `Conflict` when a conversation with the same `(account_id, id)` exists.
224    async fn create_conversation(&self, record: ConversationRecord) -> Result<(), StoreError>;
225
226    /// Appends a user turn to its conversation and creates its phase marker at
227    /// [`TurnPhase::Received`].
228    ///
229    /// The account is the actor's account inside the input.
230    ///
231    /// # Errors
232    /// * `NotFound` when the conversation does not exist for the actor's account.
233    /// * `Conflict` when a turn with the same `(account, turn_id)` exists.
234    async fn append_user_turn(&self, turn: StoredUserTurn) -> Result<(), StoreError>;
235
236    /// Persists the assistant turn that answers a user turn, exactly as returned
237    /// to the client.
238    ///
239    /// # Errors
240    /// * `NotFound` when the user turn does not exist for `account`.
241    /// * `Conflict` when the turn already has an assistant turn.
242    /// * `Other(IDENTITY_MISMATCH)` when `turn.conversation_id` differs from the
243    ///   user turn's conversation.
244    async fn append_assistant_turn(
245        &self,
246        account: &AccountId,
247        turn: AssistantTurn,
248    ) -> Result<(), StoreError>;
249
250    /// Writes the phase marker of a turn and returns the updated marker.
251    ///
252    /// # Errors
253    /// * `NotFound` when the turn does not exist for `account`.
254    /// * `Conflict` when the current phase is terminal and `phase` differs.
255    async fn set_turn_phase(
256        &self,
257        account: &AccountId,
258        turn_id: &TurnId,
259        phase: TurnPhase,
260    ) -> Result<TurnPhaseMarker, StoreError>;
261}
262
263/// Persistence of conversations, turns and phase markers: both halves.
264///
265/// Every method is account-scoped; see the module documentation for the rules
266/// implementations must honour. The conformance suite proves them
267/// ([`crate::conformance::check_conversation_turn_persistence`]).
268///
269/// There is nothing to implement here. Write
270/// [`ConversationReader`] and [`ConversationWriter`] and this trait follows
271/// from the blanket implementation below, so `Arc<dyn ConversationStore>` keeps
272/// naming one value that does everything.
273pub trait ConversationStore: ConversationReader + ConversationWriter {}
274
275impl<T: ConversationReader + ConversationWriter + ?Sized> ConversationStore for T {}
276
277#[cfg(test)]
278mod tests {
279    use super::*;
280    use turnframe_core::locale::Locale;
281    use turnframe_core::turn::ActorContext;
282
283    #[test]
284    fn stored_user_turn_accessors() {
285        let input = TurnInput {
286            turn_id: TurnId::nil(),
287            conversation_id: ConversationId::nil(),
288            actor: ActorContext::new("acct", "u"),
289            text: Some("hi".into()),
290            interaction_response: None,
291            attachments: Vec::new(),
292            origin: None,
293            locale: Locale::from("en"),
294            effort: None,
295        };
296        let turn = StoredUserTurn::new(input, DateTime::<Utc>::UNIX_EPOCH);
297        assert_eq!(turn.account_id(), &AccountId::from("acct"));
298        assert_eq!(turn.turn_id(), TurnId::nil());
299        assert_eq!(turn.conversation_id(), ConversationId::nil());
300    }
301
302    #[test]
303    fn recovery_scope_membership() {
304        let a = AccountId::from("a");
305        assert!(RecoveryScope::AllAccounts.includes(&a));
306        assert!(RecoveryScope::Account(a.clone()).includes(&a));
307        assert!(!RecoveryScope::Account(AccountId::from("b")).includes(&a));
308    }
309
310    #[test]
311    fn conversation_record_round_trips() {
312        let record = ConversationRecord::new(
313            ConversationId::nil(),
314            AccountId::from("a"),
315            DateTime::<Utc>::UNIX_EPOCH,
316        )
317        .with_metadata(serde_json::json!({"title": "t"}));
318        let json = serde_json::to_string(&record).unwrap();
319        assert_eq!(
320            serde_json::from_str::<ConversationRecord>(&json).unwrap(),
321            record
322        );
323    }
324}