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
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
//! What the in-memory store cannot prove.
//!
//! The conformance suite says this adapter obeys the persistence contract; it
//! runs one operation at a time, which is all a single-process store can be
//! asked to do. These tests cover the other half — the claims that are only
//! true because PostgreSQL is underneath, and that would still pass if the
//! adapter had quietly replaced a compare-and-swap with a read followed by a
//! write.
//!
//! Every race here uses two independent pools, so the two sides are two real
//! connections in two real transactions, and repeats itself enough times that a
//! lost update would have to be lucky every round to stay hidden.

#![allow(
    clippy::unwrap_used,
    clippy::expect_used,
    clippy::panic,
    clippy::print_stdout
)]

mod support;

use std::sync::Arc;

use tokio::sync::Barrier;
use turnframe_core::ids::{
    CaseRevision, CommandId, EventId, InteractionId, OptionId, OutboxId, RedactionAuthority, TurnId,
};
use turnframe_core::interaction::InteractionStatus;
use turnframe_core::replay::{ReplayRecord, TurnPhase};
use turnframe_store::commit::{CommitBundle, CommitStore};
use turnframe_store::error::StoreError;
use turnframe_store::events::{EventJournalReader, EventJournalWriter};
use turnframe_store::interaction::{InteractionReader, InteractionWriter, InvalidationReason};
use turnframe_store::outbox::{OutboxReader, OutboxWriter};
use turnframe_store::replay::ReplayReader;

/// The schema the account-scoped tests of this binary share.
const SCHEMA: &str = "tf_test_concurrency";
/// The outbox is a queue with no tenant, so the test that sweeps it needs a
/// schema nobody else writes to.
const OUTBOX_SCHEMA: &str = "tf_test_outbox_claim";
/// How many times a race is repeated before it is believed.
const ROUNDS: usize = 8;

/// Two transactions invalidate the same case at the same expected revision.
///
/// One statement does the selecting and the writing, so the loser does not
/// re-read a stale snapshot and overwrite the winner: it waits on the row lock,
/// re-evaluates `status = 'active'` against the row the winner committed, and
/// updates nothing. Exactly one caller is told it invalidated the card, which is
/// what stops two commits both believing they retired it.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn two_transactions_racing_the_same_expected_revision_only_one_wins() {
    let Some(url) = support::database_url() else {
        support::skipped("two_transactions_racing_the_same_expected_revision_only_one_wins");
        return;
    };
    let first = support::co_tenant_of(&url, SCHEMA).await;
    let second = support::second_pool(&url, SCHEMA).await;
    let account = support::unique_account("revision-race");

    for round in 0..ROUNDS {
        let case_id = format!("case-{round}");
        let bound = support::card(&account, support::case(&case_id, 1), support::at(0));
        first.insert(bound.clone()).await.unwrap();

        let gate = Arc::new(Barrier::new(2));
        let key = support::case(&case_id, 0).key();
        let left = {
            let (store, account, key, gate) =
                (first.clone(), account.clone(), key.clone(), gate.clone());
            tokio::spawn(async move {
                gate.wait().await;
                store
                    .invalidate_for_case(
                        &account,
                        &key,
                        CaseRevision(2),
                        InvalidationReason::RevisionChanged,
                    )
                    .await
            })
        };
        let right = {
            let (store, account, key, gate) = (second.clone(), account.clone(), key, gate);
            tokio::spawn(async move {
                gate.wait().await;
                store
                    .invalidate_for_case(
                        &account,
                        &key,
                        CaseRevision(2),
                        InvalidationReason::RevisionChanged,
                    )
                    .await
            })
        };
        let (left, right) = (left.await.unwrap().unwrap(), right.await.unwrap().unwrap());

        assert_eq!(
            left.len() + right.len(),
            1,
            "round {round}: exactly one caller may invalidate the card, got {left:?} and {right:?}"
        );
        let winner = if left.is_empty() { &right } else { &left };
        assert_eq!(winner, &vec![bound.id], "round {round}");

        // And the card really moved, once, with the revision that retired it.
        let stale = InteractionReader::get(&first, &account, &bound.id)
            .await
            .unwrap();
        assert_eq!(stale.status(), InteractionStatus::Invalidated);
        assert_eq!(
            stale.invalidation.and_then(|record| record.new_revision),
            Some(CaseRevision(2)),
            "round {round}"
        );
    }
}

/// Two clicks on the same card, from two connections, at the same instant.
///
/// The compare-and-swap and the write are one statement, so the second click
/// loses on the row rather than silently re-resolving a card that is already
/// executing its commands.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_second_answer_to_the_same_card_loses_the_compare_and_swap() {
    let Some(url) = support::database_url() else {
        support::skipped("a_second_answer_to_the_same_card_loses_the_compare_and_swap");
        return;
    };
    let first = support::co_tenant_of(&url, SCHEMA).await;
    let second = support::second_pool(&url, SCHEMA).await;
    let account = support::unique_account("resolution-race");

    for round in 0..ROUNDS {
        let card = support::card(
            &account,
            support::case(&format!("case-{round}"), 1),
            support::at(0),
        );
        first.insert(card.clone()).await.unwrap();
        let turn = TurnId::new();

        let gate = Arc::new(Barrier::new(2));
        let mut answers = Vec::new();
        for store in [first.clone(), second.clone()] {
            let (account, id, gate) = (account.clone(), card.id, gate.clone());
            answers.push(tokio::spawn(async move {
                gate.wait().await;
                store
                    .begin_resolution(
                        &account,
                        &id,
                        InteractionStatus::Active,
                        OptionId::from("ack"),
                        turn,
                    )
                    .await
            }));
        }
        let mut outcomes = Vec::new();
        for answer in answers {
            outcomes.push(answer.await.unwrap());
        }
        let accepted = outcomes.iter().filter(|outcome| outcome.is_ok()).count();
        assert_eq!(accepted, 1, "round {round}: {outcomes:?}");
        assert!(
            outcomes
                .iter()
                .any(|outcome| matches!(outcome, Err(StoreError::Conflict))),
            "round {round}: the loser must be told it lost, got {outcomes:?}"
        );
    }
}

/// Two dispatchers sweep the queue at once and divide it between them.
///
/// `FOR UPDATE SKIP LOCKED` is the whole claim: a row another worker is already
/// taking is passed over rather than waited on, so no external request is ever
/// sent twice, and no row is left behind either.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn two_workers_claiming_the_outbox_never_get_the_same_entry() {
    let Some(url) = support::database_url() else {
        support::skipped("two_workers_claiming_the_outbox_never_get_the_same_entry");
        return;
    };
    let first = support::sole_owner_of(&url, OUTBOX_SCHEMA).await;
    let second = support::second_pool(&url, OUTBOX_SCHEMA).await;
    let command = CommandId::new();

    let mut enqueued = Vec::new();
    for index in 0..40 {
        let id = OutboxId::new();
        enqueued.push(id);
        first
            .enqueue(support::outbox_entry(
                id,
                command,
                "claim-race",
                &format!("key-{index}"),
                support::at(index),
            ))
            .await
            .unwrap();
    }

    let gate = Arc::new(Barrier::new(2));
    let mut workers = Vec::new();
    for (name, store) in [("worker-a", first.clone()), ("worker-b", second.clone())] {
        let gate = gate.clone();
        workers.push(tokio::spawn(async move {
            gate.wait().await;
            let mut claimed = Vec::new();
            loop {
                let batch = store
                    .claim_due(support::at(100), 3, name)
                    .await
                    .expect("claiming is not an error");
                if batch.is_empty() {
                    break;
                }
                claimed.extend(batch.into_iter().map(|entry| entry.outbox_id));
            }
            claimed
        }));
    }
    let mut claimed = Vec::new();
    for worker in workers {
        claimed.extend(worker.await.unwrap());
    }

    let mut sorted = claimed.clone();
    sorted.sort_unstable();
    let mut unique = sorted.clone();
    unique.dedup();
    assert_eq!(
        sorted.len(),
        unique.len(),
        "no row may be handed to two workers"
    );
    let mut expected = enqueued;
    expected.sort_unstable();
    assert_eq!(sorted, expected, "every row must be claimed exactly once");
}

/// The blocking slot is held by an index, not by a check in the adapter.
///
/// The adapter refuses a second open blocking card, and so does a raw insert
/// that goes around the adapter entirely: the rule survives a bug in this crate,
/// a migration script, or a `psql` session.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn the_partial_unique_index_refuses_a_second_blocking_card() {
    let Some(url) = support::database_url() else {
        support::skipped("the_partial_unique_index_refuses_a_second_blocking_card");
        return;
    };
    let store = support::co_tenant_of(&url, SCHEMA).await;
    let account = support::unique_account("blocking-slot");
    let case_ref = support::case("case-1", 1);

    let occupant = support::card(&account, case_ref.clone(), support::at(0));
    store.insert(occupant.clone()).await.unwrap();
    assert_eq!(
        store
            .insert(support::card(&account, case_ref.clone(), support::at(1)))
            .await,
        Err(StoreError::Conflict),
        "the adapter refuses a second blocking card"
    );

    // Now around the adapter: the index itself must refuse the row.
    let error = sqlx::query(
        "INSERT INTO tf_interaction (
             account_id, interaction_id, conversation_id, workflow_key, case_id, case_revision,
             kind, blocking, revision_independent, payload_hash, interaction, status, created_at
         ) VALUES ($1, $2, $3, $4, $5, 1, 'single_select', true, false, 'x', '{}'::jsonb,
                   'active', now())",
    )
    .bind(account.as_str())
    .bind(InteractionId::new().as_uuid())
    .bind(occupant.conversation_id.as_uuid())
    .bind(case_ref.workflow.as_str())
    .bind(case_ref.case_id.as_str())
    .execute(store.pool())
    .await
    .expect_err("the index must refuse a second open blocking card");
    let database = error.as_database_error().expect("a server-side refusal");
    assert_eq!(
        database.constraint(),
        Some("tf_one_open_blocking_interaction_per_case"),
        "the refusal must come from the partial index, not from something else"
    );

    // The index is partial on purpose: it covers open blocking cards and
    // nothing else.
    let aside = support::card(&account, case_ref.clone(), support::at(2));
    let non_blocking = turnframe_core::interaction::Interaction {
        blocking: false,
        ..aside
    };
    store
        .insert(non_blocking)
        .await
        .expect("a non-blocking card does not take the slot");

    let other_case = support::card(&account, support::case("case-2", 1), support::at(3));
    store
        .insert(other_case)
        .await
        .expect("another case has a slot of its own");

    // And once the occupant leaves the slot, the next card may take it.
    store
        .begin_resolution(
            &account,
            &occupant.id,
            InteractionStatus::Active,
            OptionId::from("ack"),
            TurnId::new(),
        )
        .await
        .unwrap();
    assert_eq!(
        store
            .insert(support::card(&account, case_ref.clone(), support::at(4)))
            .await,
        Err(StoreError::Conflict),
        "a card that is resolving still holds the slot"
    );
    store
        .finish_resolution(
            &account,
            &occupant.id,
            turnframe_store::interaction::ResolutionOutcome::Resolved {
                event_ids: Vec::new(),
            },
        )
        .await
        .unwrap();
    store
        .insert(support::card(&account, case_ref, support::at(5)))
        .await
        .expect("a settled card has left the slot");
}

/// A bundle in flight is invisible, and a bundle rolled back never happened.
///
/// The first half holds a real transaction open and looks at the database
/// through a second connection: nothing the bundle wrote is there. The second
/// half lets a bundle fail on its last item the way the contract requires and
/// checks, again from the other connection, that the items before it are gone
/// too — and then that the same bundle without the bad item lands in full, so
/// the rollback was the item's doing and not a store that writes nothing.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_bundle_rolled_back_mid_flight_is_never_visible_to_another_connection() {
    let Some(url) = support::database_url() else {
        support::skipped("a_bundle_rolled_back_mid_flight_is_never_visible_to_another_connection");
        return;
    };
    let writer = support::co_tenant_of(&url, SCHEMA).await;
    let reader = support::second_pool(&url, SCHEMA).await;
    let account = support::unique_account("rollback");
    let command = CommandId::new();

    let inserted = support::card(&account, support::case("case-1", 1), support::at(0));
    let outbox_id = OutboxId::new();
    let turn = TurnId::new();
    let event_ids = vec![EventId::new()];
    let bundle = || {
        CommitBundle::new()
            .with_events(support::event_batch(
                &account, "case-1", command, 1, &event_ids,
            ))
            .with_interaction_insert(inserted.clone(), false)
            .with_outbox_entry(support::outbox_entry(
                outbox_id,
                command,
                "rollback-test",
                &format!("key-{outbox_id}"),
                support::at(0),
            ))
            .with_replay_record(ReplayRecord::received(
                turn,
                inserted.conversation_id,
                account.clone(),
                support::at(0),
            ))
    };

    // Held open, not committed: another connection sees none of it.
    let mut transaction = writer.pool().begin().await.unwrap();
    writer
        .commit_in(&mut transaction, &account, bundle())
        .await
        .expect("every item of the bundle applies");
    assert_eq!(
        InteractionReader::get(&reader, &account, &inserted.id).await,
        Err(StoreError::NotFound),
        "an uncommitted card must not be visible to another connection"
    );
    assert!(
        reader
            .get_by_ids(&account, &event_ids)
            .await
            .unwrap()
            .is_empty()
    );
    transaction.rollback().await.unwrap();

    assert_eq!(
        InteractionReader::get(&reader, &account, &inserted.id).await,
        Err(StoreError::NotFound),
        "a rolled-back card must never appear"
    );
    assert_eq!(
        OutboxReader::get(&reader, &outbox_id).await,
        Err(StoreError::NotFound)
    );
    assert_eq!(
        ReplayReader::get(&reader, &account, &turn).await,
        Err(StoreError::NotFound)
    );

    // The contract's own failure mode: the last item names a turn that does not
    // exist, so everything before it must roll back with it.
    assert_eq!(
        writer
            .commit(
                &account,
                bundle().with_turn_phase(TurnId::new(), TurnPhase::Committed)
            )
            .await
            .map(|receipt| receipt.event_ids),
        Err(StoreError::NotFound)
    );
    assert_eq!(
        InteractionReader::get(&reader, &account, &inserted.id).await,
        Err(StoreError::NotFound)
    );
    assert_eq!(
        reader
            .count(&account, &support::case("case-1", 0).key())
            .await
            .unwrap(),
        0,
        "no event of a failed bundle may be visible"
    );

    // The same bundle without the illegal item lands in full.
    let receipt = writer.commit(&account, bundle()).await.unwrap();
    assert_eq!(receipt.event_ids, event_ids);
    assert_eq!(
        InteractionReader::get(&reader, &account, &inserted.id)
            .await
            .unwrap()
            .status(),
        InteractionStatus::Active
    );
    assert_eq!(
        reader.get_by_ids(&account, &event_ids).await.unwrap().len(),
        1
    );
    OutboxReader::get(&reader, &outbox_id).await.unwrap();
    ReplayReader::get(&reader, &account, &turn).await.unwrap();
}

/// Two erasure requests for the same event, at the same moment.
///
/// An erasure arrives from a support tool, a queue and a retry, so two of them
/// racing is the normal case rather than the exotic one. Both must be told the
/// same thing, and what they are told has to be the record that is actually in
/// the row: a second caller that saw its own authority stamped while the row
/// kept the first one would report an erasure that never happened under that
/// name. `COALESCE` over a locked row is what settles it — the loser waits, re-
/// reads the row the winner committed, and keeps that record.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn two_erasure_requests_for_the_same_event_agree_on_one_record() {
    let Some(url) = support::database_url() else {
        support::skipped("two_erasure_requests_for_the_same_event_agree_on_one_record");
        return;
    };
    let first = support::co_tenant_of(&url, SCHEMA).await;
    let second = support::second_pool(&url, SCHEMA).await;
    let account = support::unique_account("erasure-race");

    for round in 0..ROUNDS {
        let case_id = format!("erasure-{round}");
        let event_id = EventId::new();
        EventJournalWriter::append(
            &first,
            support::event_batch(&account, &case_id, CommandId::new(), 1, &[event_id]),
        )
        .await
        .unwrap();

        let gate = Arc::new(Barrier::new(2));
        let left = {
            let (store, account, gate) = (first.clone(), account.clone(), gate.clone());
            tokio::spawn(async move {
                gate.wait().await;
                store
                    .redact_payload(
                        &account,
                        &event_id,
                        &RedactionAuthority::from("erasure-request-left"),
                    )
                    .await
            })
        };
        let right = {
            let (store, account, gate) = (second.clone(), account.clone(), gate.clone());
            tokio::spawn(async move {
                gate.wait().await;
                store
                    .redact_payload(
                        &account,
                        &event_id,
                        &RedactionAuthority::from("erasure-request-right"),
                    )
                    .await
            })
        };
        let (left, right) = (left.await.unwrap(), right.await.unwrap());
        let left = left.expect("the first erasure request");
        let right = right.expect("the second erasure request");
        assert_eq!(
            left, right,
            "both callers must be told about the same erasure"
        );

        let stored = EventJournalReader::get_by_ids(&first, &account, &[event_id])
            .await
            .unwrap();
        let stored = stored.first().expect("the event is still in the ledger");
        assert_eq!(
            stored.redaction.as_ref(),
            Some(&left),
            "the row must carry the record both callers were given"
        );
        assert!(stored.payload.is_null(), "the payload must be gone");
    }
}