turnframe-store-postgres 0.1.0

PostgreSQL reference store implementation for Turnframe
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
//! The outbox over `tf_outbox` (spec §16.4, ADR-007).
//!
//! # Why a claim is safe
//!
//! [`claim_due`] selects the due rows `FOR UPDATE SKIP LOCKED` and marks them in
//! the same statement. A row another dispatcher already holds is skipped rather
//! than waited on, so two workers sweeping at the same instant divide the queue
//! between them and never both send the same external request. The claim only
//! becomes visible when the claiming transaction commits, which is why every
//! call here runs inside one.
//!
//! # Why it is not account-scoped
//!
//! Every other table in this schema leads with `account_id`. The outbox is the
//! deliberate exception the persistence contract names: it is a system-owned
//! dispatch queue addressed by [`OutboxId`], never by user input, and it carries
//! `command_id` so a row can be traced back to the account-scoped journal entry
//! that produced it. Nothing here takes an account, so nothing here can leak one
//! tenant's rows to another's request — a dispatcher is not serving a request.

use async_trait::async_trait;
use chrono::{DateTime, Utc};
use sqlx::PgConnection;
use sqlx::postgres::PgRow;
use turnframe_core::command::IdempotencyKey;
use turnframe_core::event::{OutboxEntry, OutboxStatus};
use turnframe_core::ids::{CommandId, OutboxId};
use turnframe_store::error::{StoreError, invalid_record};
use turnframe_store::outbox::{OutboxClaim, OutboxReader, OutboxRecord, OutboxWriter};
use uuid::Uuid;

use crate::codec::{
    attempts_from_sql, attempts_to_sql, column, from_label, label, limit_to_sql, now, to_json,
};
use crate::error::store_error;
use crate::store::{PgStores, commit};

/// Everything a read needs to rebuild an [`OutboxRecord`].
const RECORD_COLUMNS: &str = "outbox_id, command_id, destination, payload, idempotency_key, \
     status, attempt_count, next_attempt_at, created_at, completed_at, claim_worker_id, \
     claim_taken_at, last_failure, remote_ref";

/// The same columns qualified for the claim statement, which joins the rows it
/// locked and would otherwise leave `outbox_id` ambiguous.
const CLAIMED_COLUMNS: &str = "o.outbox_id, o.command_id, o.destination, o.payload, \
     o.idempotency_key, o.status, o.attempt_count, o.next_attempt_at, o.created_at, \
     o.completed_at, o.claim_worker_id, o.claim_taken_at, o.last_failure, o.remote_ref";

/// The statuses a row never leaves.
const TERMINAL_STATUSES: &str = "('completed', 'failed')";

/// Enqueues a `Pending` row.
pub(crate) async fn enqueue(conn: &mut PgConnection, entry: OutboxEntry) -> Result<(), StoreError> {
    if entry.status != OutboxStatus::Pending {
        return Err(invalid_record());
    }
    sqlx::query(
        "INSERT INTO tf_outbox (
             outbox_id, command_id, destination, payload, idempotency_key, status,
             attempt_count, next_attempt_at, created_at, completed_at
         ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)",
    )
    .bind(entry.outbox_id.as_uuid())
    .bind(entry.command_id.as_uuid())
    .bind(entry.destination.as_str())
    .bind(to_json(&entry.payload)?)
    .bind(entry.idempotency_key.as_str())
    .bind(label(&entry.status)?)
    .bind(attempts_to_sql(entry.attempt_count)?)
    .bind(entry.next_attempt_at)
    .bind(entry.created_at)
    .bind(entry.completed_at)
    .execute(conn)
    .await
    .map_err(|error| store_error(&error))?;
    Ok(())
}

/// Loads one row.
pub(crate) async fn get(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
) -> Result<OutboxRecord, StoreError> {
    let statement = format!("SELECT {RECORD_COLUMNS} FROM tf_outbox WHERE outbox_id = $1");
    let row = sqlx::query(&statement)
        .bind(outbox_id.as_uuid())
        .fetch_optional(conn)
        .await
        .map_err(|error| store_error(&error))?
        .ok_or(StoreError::NotFound)?;
    decode_record(&row)
}

/// The rows one command produced, oldest first.
pub(crate) async fn list_for_command(
    conn: &mut PgConnection,
    command_id: &CommandId,
) -> Result<Vec<OutboxRecord>, StoreError> {
    let statement = format!(
        "SELECT {RECORD_COLUMNS} FROM tf_outbox
         WHERE command_id = $1
         ORDER BY created_at, outbox_id"
    );
    let rows = sqlx::query(&statement)
        .bind(command_id.as_uuid())
        .fetch_all(conn)
        .await
        .map_err(|error| store_error(&error))?;
    rows.iter().map(decode_record).collect()
}

/// Claims up to `limit` due rows for `worker_id`.
///
/// `SKIP LOCKED` is what makes two dispatchers safe to run at once: a row
/// another transaction is already claiming is passed over instead of waited on.
pub(crate) async fn claim_due(
    conn: &mut PgConnection,
    at: DateTime<Utc>,
    limit: usize,
    worker_id: &str,
) -> Result<Vec<OutboxEntry>, StoreError> {
    let statement = format!(
        "WITH due AS (
             SELECT outbox_id FROM tf_outbox
             WHERE status = 'pending' AND (next_attempt_at IS NULL OR next_attempt_at <= $1)
             ORDER BY created_at, outbox_id
             LIMIT $2
             FOR UPDATE SKIP LOCKED
         )
         UPDATE tf_outbox o
         SET status = 'dispatching',
             attempt_count = o.attempt_count + 1,
             claim_worker_id = $3,
             claim_taken_at = $1
         FROM due
         WHERE o.outbox_id = due.outbox_id
         RETURNING {CLAIMED_COLUMNS}"
    );
    let rows = sqlx::query(&statement)
        .bind(at)
        .bind(limit_to_sql(limit))
        .bind(worker_id)
        .fetch_all(conn)
        .await
        .map_err(|error| store_error(&error))?;
    let mut claimed = rows
        .iter()
        .map(decode_record)
        .collect::<Result<Vec<OutboxRecord>, StoreError>>()?;
    // `RETURNING` has no order of its own, and the contract names one.
    claimed.sort_by(|a, b| {
        a.entry
            .created_at
            .cmp(&b.entry.created_at)
            .then_with(|| a.entry.outbox_id.cmp(&b.entry.outbox_id))
    });
    Ok(claimed.into_iter().map(|record| record.entry).collect())
}

/// Settles a row as delivered.
pub(crate) async fn mark_completed(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    at: DateTime<Utc>,
) -> Result<(), StoreError> {
    let updated = sqlx::query(
        "UPDATE tf_outbox
         SET status = 'completed', completed_at = $2, claim_worker_id = NULL, claim_taken_at = NULL
         WHERE outbox_id = $1 AND status IN ('dispatching', 'outcome_unknown')",
    )
    .bind(outbox_id.as_uuid())
    .bind(at)
    .execute(&mut *conn)
    .await
    .map_err(|error| store_error(&error))?;
    if updated.rows_affected() == 1 {
        return Ok(());
    }
    accept_if(conn, outbox_id, OutboxStatus::Completed).await
}

/// Records a failure, either requeueing the row or settling it.
pub(crate) async fn mark_failed(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    reason: String,
    retry_at: Option<DateTime<Utc>>,
    at: DateTime<Utc>,
) -> Result<(), StoreError> {
    match retry_at {
        Some(retry_at) => requeue(conn, outbox_id, Some(reason), retry_at).await,
        None => fail(conn, outbox_id, reason, at).await,
    }
}

/// Puts a row back in the queue, due at `retry_at`.
async fn requeue(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    reason: Option<String>,
    retry_at: DateTime<Utc>,
) -> Result<(), StoreError> {
    let statement = format!(
        "UPDATE tf_outbox
         SET status = 'pending', next_attempt_at = $2,
             last_failure = COALESCE($3, last_failure),
             claim_worker_id = NULL, claim_taken_at = NULL
         WHERE outbox_id = $1 AND status NOT IN {TERMINAL_STATUSES}"
    );
    let updated = sqlx::query(&statement)
        .bind(outbox_id.as_uuid())
        .bind(retry_at)
        .bind(reason)
        .execute(&mut *conn)
        .await
        .map_err(|error| store_error(&error))?;
    if updated.rows_affected() == 1 {
        return Ok(());
    }
    Err(missing_or_conflict(conn, outbox_id).await?)
}

/// Settles a row as definitively refused.
async fn fail(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    reason: String,
    at: DateTime<Utc>,
) -> Result<(), StoreError> {
    let updated = sqlx::query(
        "UPDATE tf_outbox
         SET status = 'failed', completed_at = $2, last_failure = $3,
             claim_worker_id = NULL, claim_taken_at = NULL
         WHERE outbox_id = $1 AND status IN ('dispatching', 'outcome_unknown')",
    )
    .bind(outbox_id.as_uuid())
    .bind(at)
    .bind(reason.as_str())
    .execute(&mut *conn)
    .await
    .map_err(|error| store_error(&error))?;
    if updated.rows_affected() == 1 {
        return Ok(());
    }
    // Repeating a definite failure is accepted, and records the latest reason.
    let repeated = sqlx::query(
        "UPDATE tf_outbox SET last_failure = $2 WHERE outbox_id = $1 AND status = 'failed'",
    )
    .bind(outbox_id.as_uuid())
    .bind(reason.as_str())
    .execute(&mut *conn)
    .await
    .map_err(|error| store_error(&error))?;
    if repeated.rows_affected() == 1 {
        return Ok(());
    }
    Err(missing_or_conflict(conn, outbox_id).await?)
}

/// Records that the remote was called and the result is unknown (I15).
pub(crate) async fn mark_outcome_unknown(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    remote_ref: Option<String>,
) -> Result<(), StoreError> {
    let updated = sqlx::query(
        "UPDATE tf_outbox
         SET status = 'outcome_unknown', remote_ref = $2,
             claim_worker_id = NULL, claim_taken_at = NULL
         WHERE outbox_id = $1 AND status = 'dispatching'",
    )
    .bind(outbox_id.as_uuid())
    .bind(remote_ref.as_deref())
    .execute(&mut *conn)
    .await
    .map_err(|error| store_error(&error))?;
    if updated.rows_affected() == 1 {
        return Ok(());
    }
    // Reconciliation may learn the remote reference after the fact, and losing
    // it would leave the row impossible to reconcile.
    let repeated = sqlx::query(
        "UPDATE tf_outbox SET remote_ref = COALESCE(remote_ref, $2)
         WHERE outbox_id = $1 AND status = 'outcome_unknown'",
    )
    .bind(outbox_id.as_uuid())
    .bind(remote_ref.as_deref())
    .execute(&mut *conn)
    .await
    .map_err(|error| store_error(&error))?;
    if repeated.rows_affected() == 1 {
        return Ok(());
    }
    Err(missing_or_conflict(conn, outbox_id).await?)
}

/// Returns a row to the queue at a chosen time.
pub(crate) async fn reschedule(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    next_attempt_at: DateTime<Utc>,
) -> Result<(), StoreError> {
    requeue(conn, outbox_id, None, next_attempt_at).await
}

/// Frees the rows a dispatcher died holding.
pub(crate) async fn release_expired_claims(
    conn: &mut PgConnection,
    claimed_before: DateTime<Utc>,
) -> Result<Vec<OutboxId>, StoreError> {
    let rows = sqlx::query(
        "UPDATE tf_outbox
         SET status = 'pending', next_attempt_at = NULL,
             claim_worker_id = NULL, claim_taken_at = NULL
         WHERE status = 'dispatching' AND claim_taken_at < $1
         RETURNING outbox_id",
    )
    .bind(claimed_before)
    .fetch_all(conn)
    .await
    .map_err(|error| store_error(&error))?;
    let mut released = rows
        .iter()
        .map(|row| column::<Uuid>(row, "outbox_id").map(OutboxId::from))
        .collect::<Result<Vec<OutboxId>, StoreError>>()?;
    released.sort_unstable();
    Ok(released)
}

/// Accepts a transition that has already happened, and refuses anything else.
async fn accept_if(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
    settled: OutboxStatus,
) -> Result<(), StoreError> {
    match read_status(conn, outbox_id).await? {
        None => Err(StoreError::NotFound),
        Some(status) if status == settled => Ok(()),
        Some(_) => Err(StoreError::Conflict),
    }
}

/// Whether a compare-and-swap found no row because it does not exist or because
/// it had moved on.
async fn missing_or_conflict(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
) -> Result<StoreError, StoreError> {
    Ok(match read_status(conn, outbox_id).await? {
        None => StoreError::NotFound,
        Some(_) => StoreError::Conflict,
    })
}

/// The current status of a row, if it exists.
async fn read_status(
    conn: &mut PgConnection,
    outbox_id: &OutboxId,
) -> Result<Option<OutboxStatus>, StoreError> {
    let row = sqlx::query("SELECT status FROM tf_outbox WHERE outbox_id = $1")
        .bind(outbox_id.as_uuid())
        .fetch_optional(conn)
        .await
        .map_err(|error| store_error(&error))?;
    match row {
        None => Ok(None),
        Some(row) => {
            let status: String = column(&row, "status")?;
            from_label(&status).map(Some)
        }
    }
}

/// Rebuilds a record from its row.
fn decode_record(row: &PgRow) -> Result<OutboxRecord, StoreError> {
    let outbox_id: Uuid = column(row, "outbox_id")?;
    let command_id: Uuid = column(row, "command_id")?;
    let idempotency_key: String = column(row, "idempotency_key")?;
    let status: String = column(row, "status")?;
    let attempts: i32 = column(row, "attempt_count")?;
    let worker_id: Option<String> = column(row, "claim_worker_id")?;
    let claimed_at: Option<DateTime<Utc>> = column(row, "claim_taken_at")?;
    Ok(OutboxRecord {
        entry: OutboxEntry {
            outbox_id: OutboxId::from(outbox_id),
            command_id: CommandId::from(command_id),
            destination: column(row, "destination")?,
            payload: column(row, "payload")?,
            idempotency_key: IdempotencyKey::new(idempotency_key),
            status: from_label(&status)?,
            attempt_count: attempts_from_sql(attempts)?,
            next_attempt_at: column(row, "next_attempt_at")?,
            created_at: column(row, "created_at")?,
            completed_at: column(row, "completed_at")?,
        },
        claim: worker_id
            .zip(claimed_at)
            .map(|(worker_id, claimed_at)| OutboxClaim {
                worker_id,
                claimed_at,
            }),
        last_failure: column(row, "last_failure")?,
        remote_ref: column(row, "remote_ref")?,
    })
}

#[async_trait]
impl OutboxReader for PgStores {
    async fn get(&self, outbox_id: &OutboxId) -> Result<OutboxRecord, StoreError> {
        let mut conn = self.connection().await?;
        get(&mut conn, outbox_id).await
    }

    async fn list_for_command(
        &self,
        command_id: &CommandId,
    ) -> Result<Vec<OutboxRecord>, StoreError> {
        let mut conn = self.connection().await?;
        list_for_command(&mut conn, command_id).await
    }
}

#[async_trait]
impl OutboxWriter for PgStores {
    async fn enqueue(&self, entry: OutboxEntry) -> Result<(), StoreError> {
        let mut transaction = self.transaction().await?;
        enqueue(&mut transaction, entry).await?;
        commit(transaction).await
    }

    async fn claim_due(
        &self,
        at: DateTime<Utc>,
        limit: usize,
        worker_id: &str,
    ) -> Result<Vec<OutboxEntry>, StoreError> {
        let mut transaction = self.transaction().await?;
        let claimed = claim_due(&mut transaction, at, limit, worker_id).await?;
        commit(transaction).await?;
        Ok(claimed)
    }

    async fn mark_completed(&self, outbox_id: &OutboxId) -> Result<(), StoreError> {
        let mut transaction = self.transaction().await?;
        mark_completed(&mut transaction, outbox_id, now()).await?;
        commit(transaction).await
    }

    async fn mark_failed(
        &self,
        outbox_id: &OutboxId,
        reason: String,
        retry_at: Option<DateTime<Utc>>,
    ) -> Result<(), StoreError> {
        let mut transaction = self.transaction().await?;
        mark_failed(&mut transaction, outbox_id, reason, retry_at, now()).await?;
        commit(transaction).await
    }

    async fn mark_outcome_unknown(
        &self,
        outbox_id: &OutboxId,
        remote_ref: Option<String>,
    ) -> Result<(), StoreError> {
        let mut transaction = self.transaction().await?;
        mark_outcome_unknown(&mut transaction, outbox_id, remote_ref).await?;
        commit(transaction).await
    }

    async fn reschedule(
        &self,
        outbox_id: &OutboxId,
        next_attempt_at: DateTime<Utc>,
    ) -> Result<(), StoreError> {
        let mut transaction = self.transaction().await?;
        reschedule(&mut transaction, outbox_id, next_attempt_at).await?;
        commit(transaction).await
    }

    async fn release_expired_claims(
        &self,
        claimed_before: DateTime<Utc>,
    ) -> Result<Vec<OutboxId>, StoreError> {
        let mut transaction = self.transaction().await?;
        let released = release_expired_claims(&mut transaction, claimed_before).await?;
        commit(transaction).await?;
        Ok(released)
    }
}