pg tx producer paths wake commit-atomically: publish_with_key_tx/enqueue_tx/outbox_enqueue_tx issue pg_notify on their mechanism-named channel (stream name / queue name / derived backing queue) inside the caller's tx — empty payload, best-effort log-and-swallow, no double-wake on the auto-commit paths; shared tx::wake_tx owner + module-doc section pinning the semantics; new tx_tests harness test pinning pre-commit silence, commit delivery on all three channels, and rollback silence; F-1 engine arm retired in review-wave-4 + implementation.md records (tx_publishes_compose_with_the_handle deterministic: 20 solo runs green; pg suite 116/116 x2 vs harness) (task pg-fix-tx-wake, review 002 Finding 2)
This commit is contained in:
1 parent
7ae426a01d
commit
9a5d1f4705
5 files changed
+261
-15
No files matched your search
@@ -5,7 +5,10 @@
|
||||
//! dispositions (no-ghosts across job rows, events, notifications, and
|
||||
//! offset saves), drop panic-safety + the unknowable-state discard
|
||||
//! arms, ops-after-consume → `Closed`, `with_tx` end-to-end, the
|
||||
//! 8000-byte notify boundary, and pool accounting under concurrency.
|
||||
//! 8000-byte notify boundary, pool accounting under concurrency, and
|
||||
//! the producer tx paths' commit-atomic `pg_notify` wakes (delivered
|
||||
//! at the caller's commit on the mechanism-named channel; rollback
|
||||
//! drops the row with its wake).
|
||||
//!
|
||||
//! Harness convention (the schema/open tasks'): connection settings
|
||||
//! ride the environment (`ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB`),
|
||||
@@ -1261,3 +1264,107 @@ async fn many_txs_commit_exactly_once_through_the_pool() {
|
||||
drop_schema(&admin, &schema).await;
|
||||
Arc::into_inner(store).unwrap().close();
|
||||
}
|
||||
|
||||
/// Acceptance: the producer tx paths issue their `pg_notify` wakes
|
||||
/// commit-atomically (review 002 Finding 2's fix) — a registered
|
||||
/// listener receives the wake at/after the caller's commit,
|
||||
/// deterministically, for all three channels (the stream name for
|
||||
/// `publish_tx`, the queue name for `enqueue_tx`, the derived backing
|
||||
/// queue for `outbox_enqueue_tx`); the wake is not statement-atomic
|
||||
/// (nothing delivers before the COMMIT); and a rolled-back tx drops
|
||||
/// the row *with* its wake (the no-ghosts arm — dead silence after
|
||||
/// rollback).
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn tx_producer_wakes_are_commit_atomic_and_rollback_discards_them() {
|
||||
let Some(dsn) = harness_dsn() else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("tx-wake")();
|
||||
let store = open_store(&dsn, test_opts(&schema)).await.unwrap();
|
||||
let admin = harness_client().await.unwrap();
|
||||
|
||||
let stream = instance_namer("wake_st")();
|
||||
let queue = instance_namer("wake_q")();
|
||||
let outbox = instance_namer("wake_o")();
|
||||
let derived = crate::resolution::outbox_backing_queue_name(&outbox);
|
||||
|
||||
// Registered + subscribed BEFORE any write rides.
|
||||
for channel in [&stream, &queue, &derived] {
|
||||
store.forwarder().register(channel).await.unwrap();
|
||||
}
|
||||
let mut stream_rx = store.forwarder().subscribe().unwrap();
|
||||
let mut queue_rx = store.forwarder().subscribe().unwrap();
|
||||
let mut outbox_rx = store.forwarder().subscribe().unwrap();
|
||||
|
||||
// One business tx producing on all three channels.
|
||||
let mut tx = store.begin_tx().await.unwrap();
|
||||
tx.publish_tx(&stream, serde_json::json!({"via": "tx"}))
|
||||
.await
|
||||
.unwrap();
|
||||
tx.enqueue_tx(&queue, Default::default(), serde_json::json!({"n": 1}))
|
||||
.await
|
||||
.unwrap();
|
||||
tx.outbox_enqueue_tx(&outbox, Default::default(), serde_json::json!({"n": 2}))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Statement-issued, not yet delivered: pg_notify inside a tx does
|
||||
// not emit at statement time (commit-atomic, not statement-atomic).
|
||||
tokio::time::sleep(Duration::from_millis(300)).await;
|
||||
for (rx, channel, tag) in [
|
||||
(&mut stream_rx, &stream, "stream"),
|
||||
(&mut queue_rx, &queue, "queue"),
|
||||
(&mut outbox_rx, &derived, "outbox"),
|
||||
] {
|
||||
assert!(
|
||||
!listener_hears(rx, channel, Duration::from_millis(100)).await,
|
||||
"{tag}: no wake may deliver before the caller's commit"
|
||||
);
|
||||
}
|
||||
|
||||
tx.commit().await.unwrap();
|
||||
|
||||
// The wakes deliver at/after the commit, deterministically.
|
||||
assert!(
|
||||
listener_hears(&mut stream_rx, &stream, Duration::from_secs(5)).await,
|
||||
"a committed publish_tx wakes the stream's channel at commit"
|
||||
);
|
||||
assert!(
|
||||
listener_hears(&mut queue_rx, &queue, Duration::from_secs(5)).await,
|
||||
"a committed enqueue_tx wakes the queue's channel at commit"
|
||||
);
|
||||
assert!(
|
||||
listener_hears(&mut outbox_rx, &derived, Duration::from_secs(5)).await,
|
||||
"a committed outbox_enqueue_tx wakes the derived backing channel at commit"
|
||||
);
|
||||
|
||||
// The no-ghosts arm: the rollback drops the event/job rows *with*
|
||||
// their wakes — the registered channels stay silent.
|
||||
let mut tx = store.begin_tx().await.unwrap();
|
||||
tx.publish_tx(&stream, serde_json::json!({"via": "ghost"}))
|
||||
.await
|
||||
.unwrap();
|
||||
tx.enqueue_tx(&queue, Default::default(), serde_json::json!({"n": 3}))
|
||||
.await
|
||||
.unwrap();
|
||||
tx.outbox_enqueue_tx(&outbox, Default::default(), serde_json::json!({"n": 4}))
|
||||
.await
|
||||
.unwrap();
|
||||
drop(tx);
|
||||
tokio::time::sleep(Duration::from_millis(400)).await;
|
||||
for (rx, channel, tag) in [
|
||||
(&mut stream_rx, &stream, "stream"),
|
||||
(&mut queue_rx, &queue, "queue"),
|
||||
(&mut outbox_rx, &derived, "outbox"),
|
||||
] {
|
||||
assert!(
|
||||
!listener_hears(rx, channel, Duration::from_millis(300)).await,
|
||||
"{tag}: a rolled-back tx must not wake (commit-atomic, not statement-atomic)"
|
||||
);
|
||||
}
|
||||
assert_eq!(raw_event_count(&admin, &schema, &stream).await, 1);
|
||||
|
||||
drop_schema(&admin, &schema).await;
|
||||
drop(store);
|
||||
}
|
||||
@@ -54,6 +54,21 @@
|
||||
//! the same byte string ADR-020 §4 pins as the stored row bytes —
|
||||
//! which is valid UTF-8 by construction.
|
||||
//!
|
||||
//! # The tx wakes
|
||||
//!
|
||||
//! The producer paths (`publish_with_key_tx` — which `publish_tx`
|
||||
//! rides, `enqueue_tx`, `outbox_enqueue_tx`) issue `pg_notify` on the
|
||||
//! mechanism's wake channel (the stream name / the queue name / the
|
||||
//! outbox's derived backing-queue name — the mechanism-name-is-the-
|
||||
//! channel realization) **inside the caller's transaction** ([`wake_tx`]).
|
||||
//! NOTIFY's native transactional delivery makes the wake commit-atomic:
|
||||
//! it delivers at the caller's COMMIT, and a rollback discards the row
|
||||
//! *with* its wake (the tx/no-ghosts rows' pin). Wakes are best-effort
|
||||
//! (ADR-006; the engine-wide posture — a wake must never fail a write):
|
||||
//! a failed statement is logged and swallowed, and since a failed
|
||||
//! statement aborts the tx server-side the truth still surfaces typed
|
||||
//! at COMMIT. `run_once` is a pull op, consumer-driven — no wake there.
|
||||
//!
|
||||
//! # Engine-side arithmetic
|
||||
//!
|
||||
//! The stamp/derivation arithmetic lives in [`crate::resolution`]
|
||||
@@ -260,6 +275,26 @@ fn enqueue_row_sql(schema: &str) -> String {
|
||||
)
|
||||
}
|
||||
|
||||
/// The in-tx wake statement (empty payload — `Wake { channel }` only,
|
||||
/// ADR-008 §3): issued inside the caller's transaction, so NOTIFY's
|
||||
/// native transactional delivery makes it commit-atomic — delivers at
|
||||
/// the caller's COMMIT, discarded with the row by a rollback. Best-
|
||||
/// effort at the op (the engine-wide posture — a wake must never fail
|
||||
/// a write): a failed statement is logged and swallowed; because a
|
||||
/// failed statement aborts the tx server-side, the truth still
|
||||
/// surfaces typed at COMMIT.
|
||||
async fn wake_tx(client: &tokio_postgres::Client, channel: &str) {
|
||||
if let Err(wake_err) = client
|
||||
.query_one("SELECT pg_notify($1, '')", &[&channel])
|
||||
.await
|
||||
{
|
||||
eprintln!(
|
||||
"alkstore-postgres: tx wake failed for {channel:?} \
|
||||
(best-effort: the durable row lands at the caller's commit): {wake_err}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Encode a trait-crossing payload `Value` into the exact stored bytes
|
||||
/// (the serde_json serialization — ADR-020 §4). The notify carriage
|
||||
/// reuses it: the JSON serialization is valid UTF-8 text by
|
||||
@@ -282,8 +317,19 @@ impl TxHandle for PgTxHandle {
|
||||
Box::pin(async move {
|
||||
let client = self.client()?;
|
||||
let bytes = encode_payload_bytes(&payload)?;
|
||||
self.enqueue_row(client, &queue, bytes, &opts, plain_queue_default_stamps())
|
||||
.await
|
||||
let id = self
|
||||
.enqueue_row(client, &queue, bytes, &opts, plain_queue_default_stamps())
|
||||
.await?;
|
||||
// Posture-parity wake (the contract obligation is the
|
||||
// stream path's): pg_notify on the queue's wake channel
|
||||
// inside the caller's tx — commit-atomic like every
|
||||
// in-tx notify, matching SQLite's watcher firing on the
|
||||
// committed write. Best-effort (the queues row's pinned
|
||||
// re-poll safety net covers a gap; a failed statement
|
||||
// aborts the tx server-side, so the truth surfaces typed
|
||||
// at COMMIT).
|
||||
wake_tx(client, &queue).await;
|
||||
Ok(id)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -324,6 +370,14 @@ impl TxHandle for PgTxHandle {
|
||||
)
|
||||
.await
|
||||
.map_err(pg_error)?;
|
||||
// Then wake, inside the caller's tx (the contract
|
||||
// obligation — engine-postgres.md pins `pg_notify` as the
|
||||
// streams wake trigger, and core-contract.md pins "events
|
||||
// never require polling to become visible" with rollback
|
||||
// dropping the event row *with* its wake): the stream's
|
||||
// wake channel, the mechanism-name-is-the-channel
|
||||
// realization — the auto-commit publish's twin.
|
||||
wake_tx(client, &stream).await;
|
||||
Ok(row.get(0))
|
||||
})
|
||||
}
|
||||
@@ -556,8 +610,14 @@ impl TxHandle for PgTxHandle {
|
||||
let client = self.client()?;
|
||||
let bytes = encode_payload_bytes(&payload)?;
|
||||
let queue = outbox_backing_queue_name(&outbox);
|
||||
self.enqueue_row(client, &queue, bytes, &opts, outbox_default_stamps())
|
||||
.await
|
||||
let id = self
|
||||
.enqueue_row(client, &queue, bytes, &opts, outbox_default_stamps())
|
||||
.await?;
|
||||
// The outbox's posture-parity twin of `enqueue_tx`'s wake:
|
||||
// the derived backing-queue name is the channel (the
|
||||
// scheduler's fire-wake targets the same name).
|
||||
wake_tx(client, &queue).await;
|
||||
Ok(id)
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
---
|
||||
status: draft
|
||||
last_updated: 2026-10-09 (wave-4 fix batch decomposed from the general review — 7 fix tasks + review-wave-4-fixes gate; wave 5 decomposition gates on the two HIGH fixes landing)
|
||||
last_updated: 2026-10-09 (pg-fix-tx-wake landed — the second HIGH fix, retiring F-1's engine arm; wave 5's remaining gate is the fix batch's review gate)
|
||||
---
|
||||
|
||||
# alkstore — Implementation plan
|
||||
@@ -244,8 +244,10 @@ decomposes. Specific gates:
|
||||
conformance code-read clean (0 findings); the flagged
|
||||
`tx_publishes_compose_with_the_handle` flake reproduced twice and
|
||||
recorded as F-1 (test-infra, wake-subscription under load — root
|
||||
cause unresolved; candidate dispositions in the task's Notes for
|
||||
wave 5's suite hardening); three no-action notes (F-2..F-4).
|
||||
cause unresolved at gate time; candidate dispositions in the task's
|
||||
Notes for wave 5's suite hardening); three no-action notes
|
||||
(F-2..F-4). F-1's root cause was later found engine-side by the
|
||||
general review below and retired by `pg-fix-tx-wake`.
|
||||
- **General review, wave 4** (2026-10-09,
|
||||
`docs/reviews/002-wave-4-general-review.md`) — two live-proven
|
||||
consumer-facing bugs (the forwarder's permanent death after one
|
||||
@@ -253,7 +255,12 @@ decomposes. Specific gates:
|
||||
covered; the tx enqueue/publish paths' missing `pg_notify` wake —
|
||||
F-1's root cause, engine-side), a `max_size: 0` open-hang, two
|
||||
narrow robustness gaps, doc mismatches, decode-duplication smells.
|
||||
Fixes recommended before wave 5 decomposes; not yet landed.
|
||||
Fixes recommended before wave 5 decomposes. The tx-wake fix landed
|
||||
(`pg-fix-tx-wake`, 2026-10-09): the tx producer paths now wake
|
||||
commit-atomically, retiring F-1's engine arm (record updated in
|
||||
`tasks/review-wave-4.md`) and making `tx_publishes_compose_with_the_handle`
|
||||
deterministic (20 consecutive solo runs green); the remaining fixes
|
||||
ride the wave-4 fix-batch decomposition below.
|
||||
- **Wave 4 decomposition** (2026-10-08) — shaped as wave 3's
|
||||
structural twin; wave-3 outcomes absorbed (constructors exist;
|
||||
arithmetic re-owned per engine; per-engine `@every` parser;
|
||||
|
||||
+57
-3
@@ -1,7 +1,7 @@
|
||||
---
|
||||
id: pg-fix-tx-wake
|
||||
name: Fix tx enqueue/publish paths issuing no pg_notify wake (review 002 Finding 2 — F-1's root cause)
|
||||
status: pending
|
||||
status: completed
|
||||
depends_on: []
|
||||
scope: narrow
|
||||
risk: low
|
||||
@@ -98,8 +98,62 @@ review-rounds line for the general review.
|
||||
|
||||
## Notes
|
||||
|
||||
> To be filled by implementation agent
|
||||
Fix-shape decisions of record (what the description left open):
|
||||
|
||||
- **One-owner wake statement**: all three producer paths call a
|
||||
shared `tx::wake_tx(client, channel)` — `SELECT pg_notify($1, '')`
|
||||
(empty payload, `Wake { channel }` only per ADR-008 §3) issued on
|
||||
the caller's tx client after the INSERT. Channel = the stream name
|
||||
(`publish_tx` rides `publish_with_key_tx`, so one site), the queue
|
||||
name (`enqueue_tx`), and the outbox's derived `__alkstore_outbox:{name}`
|
||||
(`outbox_enqueue_tx`).
|
||||
- **Wake error posture**: best-effort log-and-swallow (the engine-wide
|
||||
stance — a wake must never fail a write), matching the auto-commit
|
||||
publish/enqueue twins. Documented hazard: a failed statement aborts
|
||||
the pg tx server-side, so a swallowed wake failure still surfaces
|
||||
typed at COMMIT — no silent half-state; a rolled-back tx's queued
|
||||
notifies are discarded server-side with the row (the no-ghosts arm
|
||||
is native).
|
||||
- **No double-wake**: the wake is applied at the tx call sites, NOT
|
||||
inside the shared `enqueue_row`/INSERT — the auto-commit
|
||||
`Queue::enqueue` and the scheduler's fire-wake keep their own single
|
||||
statements on their pool connections, untouched.
|
||||
- **`run_once` untouched** (pull op, consumer-driven — no wake, per
|
||||
the task).
|
||||
- **Test shape**: one harness test in `tx_tests.rs`
|
||||
(`tx_producer_wakes_are_commit_atomic_and_rollback_discards_them`)
|
||||
pins all four arms — pre-commit silence (commit-atomic, not
|
||||
statement-atomic: 300 ms statement window + 100 ms listen, then
|
||||
commit), deterministic wake delivery at commit on all three
|
||||
registered channels, rollback silence, committed row-count truth.
|
||||
Uses `forwarder().register` + the raw fanout (`listener_hears`);
|
||||
`register` does no name validation, so the reserved derived outbox
|
||||
channel registers directly.
|
||||
- **Record updates went slightly beyond the strict line**: besides
|
||||
the review-rounds line, `docs/plans/implementation.md`'s frontmatter
|
||||
`last_updated` and the wave-4 gate bullet's "root cause unresolved"
|
||||
text got dated supersede/correction pointers (the stale-reference
|
||||
discipline — those conclusions were factually wrong, not just
|
||||
unresolved).
|
||||
|
||||
## Summary
|
||||
|
||||
> To be filled on completion
|
||||
**Landed**: the tx producer paths' commit-atomic `pg_notify` wakes
|
||||
(`tx.rs`: `wake_tx` helper + the three call sites + the "The tx wakes"
|
||||
module-doc section pinning the semantics) — `publish_with_key_tx`/
|
||||
`publish_tx` (the contract obligation: streams never require polling;
|
||||
rollback drops the row with its wake) and the posture-parity
|
||||
`enqueue_tx`/`outbox_enqueue_tx` twins (channel = queue name / derived
|
||||
backing queue); the new pinned tx-wake test; the F-1 retirement
|
||||
recorded in `tasks/review-wave-4.md` (including the correction that
|
||||
the gate's "the tx INSERT's commit *is* the NOTIFY carrier" conclusion
|
||||
was wrong) and `docs/plans/implementation.md`.
|
||||
|
||||
**Verified**: full pg lib suite 116/116 green against the harness
|
||||
server (two consecutive full runs, ~70 s each); the new wake test
|
||||
green; `tx_publishes_compose_with_the_handle` deterministic — 20
|
||||
consecutive solo runs green with zero deadline misses; workspace
|
||||
`cargo test` server-less green (harness tests skip cleanly);
|
||||
`cargo clippy --all-targets -- -D warnings` clean; `cargo fmt --check`
|
||||
clean. Auto-commit publish/enqueue unchanged (their wake tests ride
|
||||
the suite untouched).
|
||||
+21
-3
@@ -236,7 +236,8 @@ report — see Summary below for the disposition): reproduced once in 53
|
||||
solo runs and once in the first full-suite run of this session; wake
|
||||
path code-read clean; the stand-alone binary probe that mirrors the
|
||||
test's beats (64 runs) never missed (delivery latency ~20–210 ms).
|
||||
Root cause unresolved — recorded for wave 5 (see below).
|
||||
Root cause unresolved at gate time — resolved engine-side by the
|
||||
general review (see F-1's updates below).
|
||||
|
||||
Deferred (documented, not fixed — each needs its own small task or a
|
||||
wave-5/6 decision; none is a consumer-visible defect):
|
||||
@@ -252,7 +253,21 @@ wave-5/6 decision; none is a consumer-visible defect):
|
||||
subscriber's pre-commit drain parks on a wake that never comes. Fix
|
||||
is engine-side (one `pg_notify` in `publish_with_key_tx`), not the
|
||||
suite-hardening candidates below, which become
|
||||
defense-in-depth.** Failure signature: `must_recv_event`'s 15 s
|
||||
defense-in-depth.** **Update 2 (2026-10-09): retired —
|
||||
`pg-fix-tx-wake` landed the engine-side fix.** The tx producer
|
||||
paths (`publish_with_key_tx`, `enqueue_tx`, `outbox_enqueue_tx`)
|
||||
now issue their `pg_notify` wakes commit-atomically inside the
|
||||
caller's tx (stream name / queue name / derived backing queue as
|
||||
the channel; rollback drops the row with its wake), pinned by
|
||||
tests; `tx_publishes_compose_with_the_handle` is deterministic
|
||||
(20 consecutive solo runs green). The suite-hardening candidates
|
||||
below remain valid as defense-in-depth but are **no longer
|
||||
load-bearing** for this flake. Correction to the facts below: the
|
||||
"engine's wake path is honest — the tx INSERT's commit *is* the
|
||||
NOTIFY carrier" conclusion was wrong (nothing on that path ever
|
||||
sent a NOTIFY); the 64-run probe simply sampled the lucky
|
||||
attach-read/drain-vs-commit ordering every time.
|
||||
Failure signature: `must_recv_event`'s 15 s
|
||||
deadline exhausts — the bridge's re-drain never fires after a
|
||||
committed `publish_tx`. Facts established: the engine's wake path is
|
||||
honest (the tx INSERT's commit *is* the NOTIFY carrier — verified by
|
||||
@@ -339,6 +354,9 @@ stand-alone probe, whole-module and 30× full-suite runs clean); the
|
||||
miss is a still-unexplained wake-subscription interaction under load.
|
||||
Recorded as F-1 for wave 5's suite hardening (candidate dispositions
|
||||
in Notes); the task agents' "wait for wave 5 if it recurs" posture is
|
||||
adopted — it recurred, wave 5 gets the record.
|
||||
adopted — it recurred, wave 5 gets the record. (Superseded
|
||||
2026-10-09: the wake-path-honesty conclusion was wrong — the general
|
||||
review root-caused the flake engine-side — and `pg-fix-tx-wake`
|
||||
retired F-1; see F-1's updates in Notes.)
|
||||
|
||||
Wave 5 decomposition may proceed.
|
||||
Reference in new issue
Block a user