turnframe_store/outbox.rs
1//! The outbox: external side effects awaiting dispatch (spec §16.4, ADR-007).
2//!
3//! # Contract
4//!
5//! * `(destination, idempotency_key)` is unique: a given external action exists
6//! at most once per destination. A duplicate enqueue is `Conflict`.
7//! * [`OutboxWriter::claim_due`] has *skip-locked* semantics: it selects
8//! `Pending` entries whose `next_attempt_at` is absent or not after `now`,
9//! moves them to `Dispatching`, increments `attempt_count`, records the
10//! claiming worker and returns them. An entry claimed by one worker is not
11//! returned to another until it is rescheduled, released or marked.
12//! * Status changes follow `OutboxStatus::can_transition` from `turnframe-core`:
13//! `Dispatching` ends in `Completed`, `Failed`, `OutcomeUnknown` or goes back
14//! to `Pending` (retry later); `OutcomeUnknown` is settled by reconciliation
15//! to `Completed` or `Failed`, or rescheduled to `Pending` when the remote
16//! guarantees idempotency. Re-marking a terminal status with the same status
17//! is accepted; anything else illegal is `Conflict`.
18//! * The outbox is a system-owned queue read by dispatchers, not by end users;
19//! its rows are addressed by `outbox_id` and carry the `command_id` that
20//! links them to the account-scoped journal.
21
22use async_trait::async_trait;
23use chrono::{DateTime, Utc};
24use serde::{Deserialize, Serialize};
25use turnframe_core::event::OutboxEntry;
26use turnframe_core::ids::{CommandId, OutboxId};
27
28use crate::error::StoreError;
29
30/// Who holds an entry in `Dispatching` and since when.
31#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
32pub struct OutboxClaim {
33 /// Dispatcher identifier.
34 pub worker_id: String,
35 /// When the claim was taken.
36 pub claimed_at: DateTime<Utc>,
37}
38
39/// An outbox entry with the store-owned dispatch metadata.
40#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
41pub struct OutboxRecord {
42 /// The entry as the core sees it.
43 pub entry: OutboxEntry,
44 /// Current claim, while `Dispatching`.
45 #[serde(default, skip_serializing_if = "Option::is_none")]
46 pub claim: Option<OutboxClaim>,
47 /// Stable code of the last failure, if any.
48 #[serde(default, skip_serializing_if = "Option::is_none")]
49 pub last_failure: Option<String>,
50 /// Remote reference recorded with an unknown outcome, for reconciliation.
51 #[serde(default, skip_serializing_if = "Option::is_none")]
52 pub remote_ref: Option<String>,
53}
54
55impl OutboxRecord {
56 /// Wraps a freshly enqueued entry.
57 #[must_use]
58 pub fn new(entry: OutboxEntry) -> Self {
59 Self {
60 entry,
61 claim: None,
62 last_failure: None,
63 remote_ref: None,
64 }
65 }
66}
67
68/// The read half of the outbox (spec §22.1).
69///
70/// It says what is queued and what happened to it. Every state change —
71/// enqueueing, claiming, completing, failing, rescheduling, releasing — is the
72/// write half.
73#[async_trait]
74pub trait OutboxReader: Send + Sync {
75 /// Loads one record.
76 ///
77 /// # Errors
78 /// * `NotFound` when it does not exist.
79 async fn get(&self, outbox_id: &OutboxId) -> Result<OutboxRecord, StoreError>;
80
81 /// Every record produced by a command, ordered by `created_at` then id.
82 ///
83 /// # Errors
84 /// * [`StoreError`] when the listing could not be read.
85 async fn list_for_command(
86 &self,
87 command_id: &CommandId,
88 ) -> Result<Vec<OutboxRecord>, StoreError>;
89}
90
91/// The write half of the outbox (spec §22.1).
92///
93/// [`claim_due`](OutboxWriter::claim_due) is here and not in
94/// [`OutboxReader`] on purpose: it looks like a read and is not one. It takes
95/// the entries it returns, moving them to `Dispatching` under the caller's
96/// worker id, so a second caller cannot get them.
97#[async_trait]
98pub trait OutboxWriter: Send + Sync {
99 /// Enqueues a `Pending` entry.
100 ///
101 /// # Errors
102 /// * `Other(INVALID_RECORD)` when `entry.status` is not `Pending`.
103 /// * `Conflict` when `outbox_id` or `(destination, idempotency_key)` exists.
104 async fn enqueue(&self, entry: OutboxEntry) -> Result<(), StoreError>;
105
106 /// Claims up to `limit` due `Pending` entries for `worker_id` (see the
107 /// module documentation), oldest first by `created_at` then id, and
108 /// returns them already in `Dispatching` with `attempt_count` incremented.
109 ///
110 /// # Errors
111 /// * [`StoreError`] when the claim could not be written.
112 async fn claim_due(
113 &self,
114 now: DateTime<Utc>,
115 limit: usize,
116 worker_id: &str,
117 ) -> Result<Vec<OutboxEntry>, StoreError>;
118
119 /// `Dispatching | OutcomeUnknown → Completed`; clears the claim and stamps
120 /// `completed_at`. Accepted without change when already `Completed`.
121 ///
122 /// # Errors
123 /// * `NotFound` / `Conflict` on an illegal transition.
124 async fn mark_completed(&self, outbox_id: &OutboxId) -> Result<(), StoreError>;
125
126 /// Records a failure. `reason` is a stable code, never free text, and is
127 /// stored on the record either way.
128 ///
129 /// With `retry_at`, the entry goes back to `Pending` with
130 /// `next_attempt_at = retry_at` and no claim; that is legal from
131 /// `Dispatching`, `OutcomeUnknown` and `Pending` (which only moves the
132 /// time), and refused on a terminal status. Without `retry_at` the entry
133 /// becomes `Failed` and stamps `completed_at`; repeating that on an already
134 /// `Failed` entry is accepted without change.
135 ///
136 /// # Errors
137 /// * `NotFound` / `Conflict` on an illegal transition.
138 async fn mark_failed(
139 &self,
140 outbox_id: &OutboxId,
141 reason: String,
142 retry_at: Option<DateTime<Utc>>,
143 ) -> Result<(), StoreError>;
144
145 /// `Dispatching → OutcomeUnknown`, releasing the claim and recording the
146 /// remote reference when the remote returned one (I15).
147 ///
148 /// Repeating it on an already `OutcomeUnknown` entry is accepted and
149 /// changes nothing, except that a `remote_ref` is filled in when the entry
150 /// had none: reconciliation may learn the reference after the fact, and
151 /// losing it would leave the row unreconcilable.
152 ///
153 /// # Errors
154 /// * `NotFound` / `Conflict` on an illegal transition.
155 async fn mark_outcome_unknown(
156 &self,
157 outbox_id: &OutboxId,
158 remote_ref: Option<String>,
159 ) -> Result<(), StoreError>;
160
161 /// Returns the entry to `Pending` with `next_attempt_at`, releasing any
162 /// claim. Legal from `Dispatching`, `OutcomeUnknown` and `Pending` (which
163 /// only moves the time).
164 ///
165 /// # Errors
166 /// * `NotFound` / `Conflict` on a terminal status.
167 async fn reschedule(
168 &self,
169 outbox_id: &OutboxId,
170 next_attempt_at: DateTime<Utc>,
171 ) -> Result<(), StoreError>;
172
173 /// Releases every `Dispatching` entry whose claim was taken strictly
174 /// before `claimed_before` back to `Pending`, due immediately. Returns the
175 /// released identifiers. Meant for a reaper that recovers entries a crashed
176 /// dispatcher left behind.
177 ///
178 /// # Errors
179 /// * [`StoreError`] when the sweep could not be written.
180 async fn release_expired_claims(
181 &self,
182 claimed_before: DateTime<Utc>,
183 ) -> Result<Vec<OutboxId>, StoreError>;
184}
185
186/// The outbox (spec §22.1): both halves.
187///
188/// There is nothing to implement here: write [`OutboxReader`] and
189/// [`OutboxWriter`] and the blanket implementation below supplies this trait.
190pub trait OutboxStore: OutboxReader + OutboxWriter {}
191
192impl<T: OutboxReader + OutboxWriter + ?Sized> OutboxStore for T {}
193
194#[cfg(test)]
195mod tests {
196 use super::*;
197 use turnframe_core::command::IdempotencyKey;
198 use turnframe_core::event::OutboxStatus;
199
200 #[test]
201 fn record_round_trips() {
202 let record = OutboxRecord::new(OutboxEntry {
203 outbox_id: OutboxId::nil(),
204 command_id: CommandId::nil(),
205 destination: "airline".into(),
206 payload: serde_json::Value::Null,
207 idempotency_key: IdempotencyKey::new("k"),
208 status: OutboxStatus::Pending,
209 attempt_count: 0,
210 next_attempt_at: None,
211 created_at: DateTime::<Utc>::UNIX_EPOCH,
212 completed_at: None,
213 });
214 let json = serde_json::to_string(&record).unwrap();
215 assert_eq!(serde_json::from_str::<OutboxRecord>(&json).unwrap(), record);
216 }
217}