Skip to main content

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}