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}