1use async_trait::async_trait;
21use chrono::{DateTime, Utc};
22use serde::{Deserialize, Serialize};
23use turnframe_core::case::CaseKey;
24use turnframe_core::event::OutboxEntry;
25use turnframe_core::ids::{AccountId, CaseRevision, CommandId, EventId, InteractionId, TurnId};
26use turnframe_core::interaction::Interaction;
27use turnframe_core::replay::{ReplayRecord, TurnPhase};
28
29use crate::error::{StoreError, bundle_account_mismatch, invalid_record};
30use crate::events::EventBatch;
31use crate::interaction::{InvalidationReason, ResolutionOutcome};
32use crate::journal::JournalOutcome;
33
34#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
36pub struct JournalCompletion {
37 pub command_id: CommandId,
39 pub outcome: JournalOutcome,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
45pub struct InteractionFinish {
46 pub interaction_id: InteractionId,
48 pub outcome: ResolutionOutcome,
50}
51
52#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
54pub struct InteractionInsert {
55 pub interaction: Interaction,
57 pub replace_blocking: bool,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65pub struct CaseInvalidation {
66 pub case_key: CaseKey,
68 pub new_revision: CaseRevision,
70 pub reason: InvalidationReason,
72}
73
74#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
76pub struct TurnPhaseUpdate {
77 pub turn_id: TurnId,
79 pub phase: TurnPhase,
81}
82
83#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
88#[non_exhaustive]
89pub struct CommitBundle {
90 #[serde(default)]
92 pub journal_completions: Vec<JournalCompletion>,
93 #[serde(default)]
95 pub events: Vec<EventBatch>,
96 #[serde(default)]
98 pub interaction_finishes: Vec<InteractionFinish>,
99 #[serde(default)]
101 pub interaction_invalidations: Vec<CaseInvalidation>,
102 #[serde(default)]
104 pub interaction_inserts: Vec<InteractionInsert>,
105 #[serde(default)]
107 pub outbox_entries: Vec<OutboxEntry>,
108 #[serde(default, skip_serializing_if = "Option::is_none")]
110 pub replay_record: Option<ReplayRecord>,
111 #[serde(default, skip_serializing_if = "Option::is_none")]
113 pub turn_phase: Option<TurnPhaseUpdate>,
114}
115
116impl CommitBundle {
117 #[must_use]
119 pub fn new() -> Self {
120 Self::default()
121 }
122
123 #[must_use]
125 pub fn with_journal_completion(
126 mut self,
127 command_id: CommandId,
128 outcome: JournalOutcome,
129 ) -> Self {
130 self.journal_completions.push(JournalCompletion {
131 command_id,
132 outcome,
133 });
134 self
135 }
136
137 #[must_use]
139 pub fn with_events(mut self, batch: EventBatch) -> Self {
140 self.events.push(batch);
141 self
142 }
143
144 #[must_use]
146 pub fn with_interaction_finish(
147 mut self,
148 interaction_id: InteractionId,
149 outcome: ResolutionOutcome,
150 ) -> Self {
151 self.interaction_finishes.push(InteractionFinish {
152 interaction_id,
153 outcome,
154 });
155 self
156 }
157
158 #[must_use]
160 pub fn with_invalidation(
161 mut self,
162 case_key: CaseKey,
163 new_revision: CaseRevision,
164 reason: InvalidationReason,
165 ) -> Self {
166 self.interaction_invalidations.push(CaseInvalidation {
167 case_key,
168 new_revision,
169 reason,
170 });
171 self
172 }
173
174 #[must_use]
176 pub fn with_interaction_insert(
177 mut self,
178 interaction: Interaction,
179 replace_blocking: bool,
180 ) -> Self {
181 self.interaction_inserts.push(InteractionInsert {
182 interaction,
183 replace_blocking,
184 });
185 self
186 }
187
188 #[must_use]
190 pub fn with_outbox_entry(mut self, entry: OutboxEntry) -> Self {
191 self.outbox_entries.push(entry);
192 self
193 }
194
195 #[must_use]
197 pub fn with_replay_record(mut self, record: ReplayRecord) -> Self {
198 self.replay_record = Some(record);
199 self
200 }
201
202 #[must_use]
204 pub fn with_turn_phase(mut self, turn_id: TurnId, phase: TurnPhase) -> Self {
205 self.turn_phase = Some(TurnPhaseUpdate { turn_id, phase });
206 self
207 }
208
209 #[must_use]
211 pub fn is_empty(&self) -> bool {
212 self.journal_completions.is_empty()
213 && self.events.is_empty()
214 && self.interaction_finishes.is_empty()
215 && self.interaction_invalidations.is_empty()
216 && self.interaction_inserts.is_empty()
217 && self.outbox_entries.is_empty()
218 && self.replay_record.is_none()
219 && self.turn_phase.is_none()
220 }
221
222 pub fn validate(&self, account: &AccountId) -> Result<(), StoreError> {
232 for batch in &self.events {
233 if &batch.account_id != account {
234 return Err(bundle_account_mismatch());
235 }
236 if batch.is_empty() {
237 return Err(invalid_record());
238 }
239 }
240 if self
241 .interaction_inserts
242 .iter()
243 .any(|insert| &insert.interaction.account_id != account)
244 {
245 return Err(bundle_account_mismatch());
246 }
247 if self
248 .replay_record
249 .as_ref()
250 .is_some_and(|record| &record.account_id != account)
251 {
252 return Err(bundle_account_mismatch());
253 }
254 Ok(())
255 }
256}
257
258#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
260pub struct CommitReceipt {
261 pub event_ids: Vec<EventId>,
263 pub inserted_interactions: Vec<InteractionId>,
265 pub invalidated_interactions: Vec<InteractionId>,
268 pub committed_at: DateTime<Utc>,
270}
271
272#[async_trait]
285pub trait CommitStore: Send + Sync {
286 async fn commit(
293 &self,
294 account: &AccountId,
295 bundle: CommitBundle,
296 ) -> Result<CommitReceipt, StoreError>;
297}
298
299#[cfg(test)]
300mod tests {
301 use super::*;
302 use crate::error::{codes, has_code};
303 use turnframe_core::ids::ConversationId;
304
305 #[test]
306 fn empty_bundle_is_empty_and_valid() {
307 let bundle = CommitBundle::new();
308 assert!(bundle.is_empty());
309 assert!(bundle.validate(&AccountId::from("a")).is_ok());
310 let with_phase = bundle.with_turn_phase(TurnId::nil(), TurnPhase::Committed);
311 assert!(!with_phase.is_empty());
312 }
313
314 #[test]
315 fn validate_refuses_foreign_and_empty_items() {
316 let account = AccountId::from("a");
317 let foreign = CommitBundle::new().with_events(EventBatch::new(
318 AccountId::from("b"),
319 CaseKey::new("w", "c"),
320 CommandId::nil(),
321 CaseRevision(1),
322 vec![],
323 ));
324 assert!(has_code(
325 &foreign.validate(&account).unwrap_err(),
326 codes::BUNDLE_ACCOUNT_MISMATCH
327 ));
328 let empty = CommitBundle::new().with_events(EventBatch::new(
329 account.clone(),
330 CaseKey::new("w", "c"),
331 CommandId::nil(),
332 CaseRevision(1),
333 vec![],
334 ));
335 assert!(has_code(
336 &empty.validate(&account).unwrap_err(),
337 codes::INVALID_RECORD
338 ));
339 let replay = CommitBundle::new().with_replay_record(ReplayRecord::received(
340 TurnId::nil(),
341 ConversationId::nil(),
342 AccountId::from("b"),
343 DateTime::<Utc>::UNIX_EPOCH,
344 ));
345 assert!(has_code(
346 &replay.validate(&account).unwrap_err(),
347 codes::BUNDLE_ACCOUNT_MISMATCH
348 ));
349 }
350}