Skip to main content

turnframe_store_postgres/
outbox.rs

1//! The outbox over `tf_outbox` (spec §16.4, ADR-007).
2//!
3//! # Why a claim is safe
4//!
5//! [`claim_due`] selects the due rows `FOR UPDATE SKIP LOCKED` and marks them in
6//! the same statement. A row another dispatcher already holds is skipped rather
7//! than waited on, so two workers sweeping at the same instant divide the queue
8//! between them and never both send the same external request. The claim only
9//! becomes visible when the claiming transaction commits, which is why every
10//! call here runs inside one.
11//!
12//! # Why it is not account-scoped
13//!
14//! Every other table in this schema leads with `account_id`. The outbox is the
15//! deliberate exception the persistence contract names: it is a system-owned
16//! dispatch queue addressed by [`OutboxId`], never by user input, and it carries
17//! `command_id` so a row can be traced back to the account-scoped journal entry
18//! that produced it. Nothing here takes an account, so nothing here can leak one
19//! tenant's rows to another's request — a dispatcher is not serving a request.
20
21use async_trait::async_trait;
22use chrono::{DateTime, Utc};
23use sqlx::PgConnection;
24use sqlx::postgres::PgRow;
25use turnframe_core::command::IdempotencyKey;
26use turnframe_core::event::{OutboxEntry, OutboxStatus};
27use turnframe_core::ids::{CommandId, OutboxId};
28use turnframe_store::error::{StoreError, invalid_record};
29use turnframe_store::outbox::{OutboxClaim, OutboxReader, OutboxRecord, OutboxWriter};
30use uuid::Uuid;
31
32use crate::codec::{
33    attempts_from_sql, attempts_to_sql, column, from_label, label, limit_to_sql, now, to_json,
34};
35use crate::error::store_error;
36use crate::store::{PgStores, commit};
37
38/// Everything a read needs to rebuild an [`OutboxRecord`].
39const RECORD_COLUMNS: &str = "outbox_id, command_id, destination, payload, idempotency_key, \
40     status, attempt_count, next_attempt_at, created_at, completed_at, claim_worker_id, \
41     claim_taken_at, last_failure, remote_ref";
42
43/// The same columns qualified for the claim statement, which joins the rows it
44/// locked and would otherwise leave `outbox_id` ambiguous.
45const CLAIMED_COLUMNS: &str = "o.outbox_id, o.command_id, o.destination, o.payload, \
46     o.idempotency_key, o.status, o.attempt_count, o.next_attempt_at, o.created_at, \
47     o.completed_at, o.claim_worker_id, o.claim_taken_at, o.last_failure, o.remote_ref";
48
49/// The statuses a row never leaves.
50const TERMINAL_STATUSES: &str = "('completed', 'failed')";
51
52/// Enqueues a `Pending` row.
53pub(crate) async fn enqueue(conn: &mut PgConnection, entry: OutboxEntry) -> Result<(), StoreError> {
54    if entry.status != OutboxStatus::Pending {
55        return Err(invalid_record());
56    }
57    sqlx::query(
58        "INSERT INTO tf_outbox (
59             outbox_id, command_id, destination, payload, idempotency_key, status,
60             attempt_count, next_attempt_at, created_at, completed_at
61         ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)",
62    )
63    .bind(entry.outbox_id.as_uuid())
64    .bind(entry.command_id.as_uuid())
65    .bind(entry.destination.as_str())
66    .bind(to_json(&entry.payload)?)
67    .bind(entry.idempotency_key.as_str())
68    .bind(label(&entry.status)?)
69    .bind(attempts_to_sql(entry.attempt_count)?)
70    .bind(entry.next_attempt_at)
71    .bind(entry.created_at)
72    .bind(entry.completed_at)
73    .execute(conn)
74    .await
75    .map_err(|error| store_error(&error))?;
76    Ok(())
77}
78
79/// Loads one row.
80pub(crate) async fn get(
81    conn: &mut PgConnection,
82    outbox_id: &OutboxId,
83) -> Result<OutboxRecord, StoreError> {
84    let statement = format!("SELECT {RECORD_COLUMNS} FROM tf_outbox WHERE outbox_id = $1");
85    let row = sqlx::query(&statement)
86        .bind(outbox_id.as_uuid())
87        .fetch_optional(conn)
88        .await
89        .map_err(|error| store_error(&error))?
90        .ok_or(StoreError::NotFound)?;
91    decode_record(&row)
92}
93
94/// The rows one command produced, oldest first.
95pub(crate) async fn list_for_command(
96    conn: &mut PgConnection,
97    command_id: &CommandId,
98) -> Result<Vec<OutboxRecord>, StoreError> {
99    let statement = format!(
100        "SELECT {RECORD_COLUMNS} FROM tf_outbox
101         WHERE command_id = $1
102         ORDER BY created_at, outbox_id"
103    );
104    let rows = sqlx::query(&statement)
105        .bind(command_id.as_uuid())
106        .fetch_all(conn)
107        .await
108        .map_err(|error| store_error(&error))?;
109    rows.iter().map(decode_record).collect()
110}
111
112/// Claims up to `limit` due rows for `worker_id`.
113///
114/// `SKIP LOCKED` is what makes two dispatchers safe to run at once: a row
115/// another transaction is already claiming is passed over instead of waited on.
116pub(crate) async fn claim_due(
117    conn: &mut PgConnection,
118    at: DateTime<Utc>,
119    limit: usize,
120    worker_id: &str,
121) -> Result<Vec<OutboxEntry>, StoreError> {
122    let statement = format!(
123        "WITH due AS (
124             SELECT outbox_id FROM tf_outbox
125             WHERE status = 'pending' AND (next_attempt_at IS NULL OR next_attempt_at <= $1)
126             ORDER BY created_at, outbox_id
127             LIMIT $2
128             FOR UPDATE SKIP LOCKED
129         )
130         UPDATE tf_outbox o
131         SET status = 'dispatching',
132             attempt_count = o.attempt_count + 1,
133             claim_worker_id = $3,
134             claim_taken_at = $1
135         FROM due
136         WHERE o.outbox_id = due.outbox_id
137         RETURNING {CLAIMED_COLUMNS}"
138    );
139    let rows = sqlx::query(&statement)
140        .bind(at)
141        .bind(limit_to_sql(limit))
142        .bind(worker_id)
143        .fetch_all(conn)
144        .await
145        .map_err(|error| store_error(&error))?;
146    let mut claimed = rows
147        .iter()
148        .map(decode_record)
149        .collect::<Result<Vec<OutboxRecord>, StoreError>>()?;
150    // `RETURNING` has no order of its own, and the contract names one.
151    claimed.sort_by(|a, b| {
152        a.entry
153            .created_at
154            .cmp(&b.entry.created_at)
155            .then_with(|| a.entry.outbox_id.cmp(&b.entry.outbox_id))
156    });
157    Ok(claimed.into_iter().map(|record| record.entry).collect())
158}
159
160/// Settles a row as delivered.
161pub(crate) async fn mark_completed(
162    conn: &mut PgConnection,
163    outbox_id: &OutboxId,
164    at: DateTime<Utc>,
165) -> Result<(), StoreError> {
166    let updated = sqlx::query(
167        "UPDATE tf_outbox
168         SET status = 'completed', completed_at = $2, claim_worker_id = NULL, claim_taken_at = NULL
169         WHERE outbox_id = $1 AND status IN ('dispatching', 'outcome_unknown')",
170    )
171    .bind(outbox_id.as_uuid())
172    .bind(at)
173    .execute(&mut *conn)
174    .await
175    .map_err(|error| store_error(&error))?;
176    if updated.rows_affected() == 1 {
177        return Ok(());
178    }
179    accept_if(conn, outbox_id, OutboxStatus::Completed).await
180}
181
182/// Records a failure, either requeueing the row or settling it.
183pub(crate) async fn mark_failed(
184    conn: &mut PgConnection,
185    outbox_id: &OutboxId,
186    reason: String,
187    retry_at: Option<DateTime<Utc>>,
188    at: DateTime<Utc>,
189) -> Result<(), StoreError> {
190    match retry_at {
191        Some(retry_at) => requeue(conn, outbox_id, Some(reason), retry_at).await,
192        None => fail(conn, outbox_id, reason, at).await,
193    }
194}
195
196/// Puts a row back in the queue, due at `retry_at`.
197async fn requeue(
198    conn: &mut PgConnection,
199    outbox_id: &OutboxId,
200    reason: Option<String>,
201    retry_at: DateTime<Utc>,
202) -> Result<(), StoreError> {
203    let statement = format!(
204        "UPDATE tf_outbox
205         SET status = 'pending', next_attempt_at = $2,
206             last_failure = COALESCE($3, last_failure),
207             claim_worker_id = NULL, claim_taken_at = NULL
208         WHERE outbox_id = $1 AND status NOT IN {TERMINAL_STATUSES}"
209    );
210    let updated = sqlx::query(&statement)
211        .bind(outbox_id.as_uuid())
212        .bind(retry_at)
213        .bind(reason)
214        .execute(&mut *conn)
215        .await
216        .map_err(|error| store_error(&error))?;
217    if updated.rows_affected() == 1 {
218        return Ok(());
219    }
220    Err(missing_or_conflict(conn, outbox_id).await?)
221}
222
223/// Settles a row as definitively refused.
224async fn fail(
225    conn: &mut PgConnection,
226    outbox_id: &OutboxId,
227    reason: String,
228    at: DateTime<Utc>,
229) -> Result<(), StoreError> {
230    let updated = sqlx::query(
231        "UPDATE tf_outbox
232         SET status = 'failed', completed_at = $2, last_failure = $3,
233             claim_worker_id = NULL, claim_taken_at = NULL
234         WHERE outbox_id = $1 AND status IN ('dispatching', 'outcome_unknown')",
235    )
236    .bind(outbox_id.as_uuid())
237    .bind(at)
238    .bind(reason.as_str())
239    .execute(&mut *conn)
240    .await
241    .map_err(|error| store_error(&error))?;
242    if updated.rows_affected() == 1 {
243        return Ok(());
244    }
245    // Repeating a definite failure is accepted, and records the latest reason.
246    let repeated = sqlx::query(
247        "UPDATE tf_outbox SET last_failure = $2 WHERE outbox_id = $1 AND status = 'failed'",
248    )
249    .bind(outbox_id.as_uuid())
250    .bind(reason.as_str())
251    .execute(&mut *conn)
252    .await
253    .map_err(|error| store_error(&error))?;
254    if repeated.rows_affected() == 1 {
255        return Ok(());
256    }
257    Err(missing_or_conflict(conn, outbox_id).await?)
258}
259
260/// Records that the remote was called and the result is unknown (I15).
261pub(crate) async fn mark_outcome_unknown(
262    conn: &mut PgConnection,
263    outbox_id: &OutboxId,
264    remote_ref: Option<String>,
265) -> Result<(), StoreError> {
266    let updated = sqlx::query(
267        "UPDATE tf_outbox
268         SET status = 'outcome_unknown', remote_ref = $2,
269             claim_worker_id = NULL, claim_taken_at = NULL
270         WHERE outbox_id = $1 AND status = 'dispatching'",
271    )
272    .bind(outbox_id.as_uuid())
273    .bind(remote_ref.as_deref())
274    .execute(&mut *conn)
275    .await
276    .map_err(|error| store_error(&error))?;
277    if updated.rows_affected() == 1 {
278        return Ok(());
279    }
280    // Reconciliation may learn the remote reference after the fact, and losing
281    // it would leave the row impossible to reconcile.
282    let repeated = sqlx::query(
283        "UPDATE tf_outbox SET remote_ref = COALESCE(remote_ref, $2)
284         WHERE outbox_id = $1 AND status = 'outcome_unknown'",
285    )
286    .bind(outbox_id.as_uuid())
287    .bind(remote_ref.as_deref())
288    .execute(&mut *conn)
289    .await
290    .map_err(|error| store_error(&error))?;
291    if repeated.rows_affected() == 1 {
292        return Ok(());
293    }
294    Err(missing_or_conflict(conn, outbox_id).await?)
295}
296
297/// Returns a row to the queue at a chosen time.
298pub(crate) async fn reschedule(
299    conn: &mut PgConnection,
300    outbox_id: &OutboxId,
301    next_attempt_at: DateTime<Utc>,
302) -> Result<(), StoreError> {
303    requeue(conn, outbox_id, None, next_attempt_at).await
304}
305
306/// Frees the rows a dispatcher died holding.
307pub(crate) async fn release_expired_claims(
308    conn: &mut PgConnection,
309    claimed_before: DateTime<Utc>,
310) -> Result<Vec<OutboxId>, StoreError> {
311    let rows = sqlx::query(
312        "UPDATE tf_outbox
313         SET status = 'pending', next_attempt_at = NULL,
314             claim_worker_id = NULL, claim_taken_at = NULL
315         WHERE status = 'dispatching' AND claim_taken_at < $1
316         RETURNING outbox_id",
317    )
318    .bind(claimed_before)
319    .fetch_all(conn)
320    .await
321    .map_err(|error| store_error(&error))?;
322    let mut released = rows
323        .iter()
324        .map(|row| column::<Uuid>(row, "outbox_id").map(OutboxId::from))
325        .collect::<Result<Vec<OutboxId>, StoreError>>()?;
326    released.sort_unstable();
327    Ok(released)
328}
329
330/// Accepts a transition that has already happened, and refuses anything else.
331async fn accept_if(
332    conn: &mut PgConnection,
333    outbox_id: &OutboxId,
334    settled: OutboxStatus,
335) -> Result<(), StoreError> {
336    match read_status(conn, outbox_id).await? {
337        None => Err(StoreError::NotFound),
338        Some(status) if status == settled => Ok(()),
339        Some(_) => Err(StoreError::Conflict),
340    }
341}
342
343/// Whether a compare-and-swap found no row because it does not exist or because
344/// it had moved on.
345async fn missing_or_conflict(
346    conn: &mut PgConnection,
347    outbox_id: &OutboxId,
348) -> Result<StoreError, StoreError> {
349    Ok(match read_status(conn, outbox_id).await? {
350        None => StoreError::NotFound,
351        Some(_) => StoreError::Conflict,
352    })
353}
354
355/// The current status of a row, if it exists.
356async fn read_status(
357    conn: &mut PgConnection,
358    outbox_id: &OutboxId,
359) -> Result<Option<OutboxStatus>, StoreError> {
360    let row = sqlx::query("SELECT status FROM tf_outbox WHERE outbox_id = $1")
361        .bind(outbox_id.as_uuid())
362        .fetch_optional(conn)
363        .await
364        .map_err(|error| store_error(&error))?;
365    match row {
366        None => Ok(None),
367        Some(row) => {
368            let status: String = column(&row, "status")?;
369            from_label(&status).map(Some)
370        }
371    }
372}
373
374/// Rebuilds a record from its row.
375fn decode_record(row: &PgRow) -> Result<OutboxRecord, StoreError> {
376    let outbox_id: Uuid = column(row, "outbox_id")?;
377    let command_id: Uuid = column(row, "command_id")?;
378    let idempotency_key: String = column(row, "idempotency_key")?;
379    let status: String = column(row, "status")?;
380    let attempts: i32 = column(row, "attempt_count")?;
381    let worker_id: Option<String> = column(row, "claim_worker_id")?;
382    let claimed_at: Option<DateTime<Utc>> = column(row, "claim_taken_at")?;
383    Ok(OutboxRecord {
384        entry: OutboxEntry {
385            outbox_id: OutboxId::from(outbox_id),
386            command_id: CommandId::from(command_id),
387            destination: column(row, "destination")?,
388            payload: column(row, "payload")?,
389            idempotency_key: IdempotencyKey::new(idempotency_key),
390            status: from_label(&status)?,
391            attempt_count: attempts_from_sql(attempts)?,
392            next_attempt_at: column(row, "next_attempt_at")?,
393            created_at: column(row, "created_at")?,
394            completed_at: column(row, "completed_at")?,
395        },
396        claim: worker_id
397            .zip(claimed_at)
398            .map(|(worker_id, claimed_at)| OutboxClaim {
399                worker_id,
400                claimed_at,
401            }),
402        last_failure: column(row, "last_failure")?,
403        remote_ref: column(row, "remote_ref")?,
404    })
405}
406
407#[async_trait]
408impl OutboxReader for PgStores {
409    async fn get(&self, outbox_id: &OutboxId) -> Result<OutboxRecord, StoreError> {
410        let mut conn = self.connection().await?;
411        get(&mut conn, outbox_id).await
412    }
413
414    async fn list_for_command(
415        &self,
416        command_id: &CommandId,
417    ) -> Result<Vec<OutboxRecord>, StoreError> {
418        let mut conn = self.connection().await?;
419        list_for_command(&mut conn, command_id).await
420    }
421}
422
423#[async_trait]
424impl OutboxWriter for PgStores {
425    async fn enqueue(&self, entry: OutboxEntry) -> Result<(), StoreError> {
426        let mut transaction = self.transaction().await?;
427        enqueue(&mut transaction, entry).await?;
428        commit(transaction).await
429    }
430
431    async fn claim_due(
432        &self,
433        at: DateTime<Utc>,
434        limit: usize,
435        worker_id: &str,
436    ) -> Result<Vec<OutboxEntry>, StoreError> {
437        let mut transaction = self.transaction().await?;
438        let claimed = claim_due(&mut transaction, at, limit, worker_id).await?;
439        commit(transaction).await?;
440        Ok(claimed)
441    }
442
443    async fn mark_completed(&self, outbox_id: &OutboxId) -> Result<(), StoreError> {
444        let mut transaction = self.transaction().await?;
445        mark_completed(&mut transaction, outbox_id, now()).await?;
446        commit(transaction).await
447    }
448
449    async fn mark_failed(
450        &self,
451        outbox_id: &OutboxId,
452        reason: String,
453        retry_at: Option<DateTime<Utc>>,
454    ) -> Result<(), StoreError> {
455        let mut transaction = self.transaction().await?;
456        mark_failed(&mut transaction, outbox_id, reason, retry_at, now()).await?;
457        commit(transaction).await
458    }
459
460    async fn mark_outcome_unknown(
461        &self,
462        outbox_id: &OutboxId,
463        remote_ref: Option<String>,
464    ) -> Result<(), StoreError> {
465        let mut transaction = self.transaction().await?;
466        mark_outcome_unknown(&mut transaction, outbox_id, remote_ref).await?;
467        commit(transaction).await
468    }
469
470    async fn reschedule(
471        &self,
472        outbox_id: &OutboxId,
473        next_attempt_at: DateTime<Utc>,
474    ) -> Result<(), StoreError> {
475        let mut transaction = self.transaction().await?;
476        reschedule(&mut transaction, outbox_id, next_attempt_at).await?;
477        commit(transaction).await
478    }
479
480    async fn release_expired_claims(
481        &self,
482        claimed_before: DateTime<Utc>,
483    ) -> Result<Vec<OutboxId>, StoreError> {
484        let mut transaction = self.transaction().await?;
485        let released = release_expired_claims(&mut transaction, claimed_before).await?;
486        commit(transaction).await?;
487        Ok(released)
488    }
489}