Postgres engine: streams — StreamHandle (auto-commit publishes + best-effort pg_notify wake, ASC reads with the extent guard, monotone offsets, pool-connection trim) and the durable subscribe receiver (async bridge, wake-driven re-drains, reconnect gap-heal, shutdown-only terminal close, stateless idle wake runtime for the sync save) (task pg-engine-streams)
This commit is contained in:
1 parent
8f5c2add5e
commit
cb067bb4be
7 files changed
+2254
-45
No files matched your search
@@ -39,6 +39,7 @@ mod resolution;
|
||||
mod schema;
|
||||
mod seam;
|
||||
mod store;
|
||||
mod stream;
|
||||
mod tx;
|
||||
|
||||
pub use opts::{DEFAULT_MAX_SIZE, PgOpts};
|
||||
|
||||
@@ -261,12 +261,22 @@ impl Store for PgStore {
|
||||
|
||||
fn stream<'a>(
|
||||
&'a self,
|
||||
_name: &str,
|
||||
name: &str,
|
||||
) -> alkstore::BoxedFuture<'a, alkstore::Result<Box<dyn alkstore::StreamHandle>>> {
|
||||
if let Err(e) = alkstore::validate_shared_name(name) {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
if let Err(e) = self.closed_check() {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
Box::pin(async { Err(stub("stream wiring lands with the streams task")) })
|
||||
let handle = crate::stream::PgStreamHandle::new(
|
||||
name.to_string(),
|
||||
self.pool.clone(),
|
||||
self.schema.clone(),
|
||||
self.forwarder.clone(),
|
||||
self.closed.clone(),
|
||||
);
|
||||
Box::pin(async move { Ok(Box::new(handle) as Box<dyn alkstore::StreamHandle>) })
|
||||
}
|
||||
|
||||
fn queue<'a>(
|
||||
@@ -348,5 +358,8 @@ mod open_tests;
|
||||
#[cfg(test)]
|
||||
mod notify_tests;
|
||||
|
||||
#[cfg(test)]
|
||||
mod stream_tests;
|
||||
|
||||
#[cfg(test)]
|
||||
mod tx_tests;
|
||||
@@ -659,11 +659,10 @@ async fn store_trait_methods_are_wiring_stubs() {
|
||||
.notify("stub_ch_probe", serde_json::json!("x"))
|
||||
.await
|
||||
.unwrap();
|
||||
let err = match store.stream("s").await {
|
||||
Err(e) => e,
|
||||
Ok(_) => panic!("stream must be a stub"),
|
||||
};
|
||||
assert!(matches!(err, alkstore::Error::Database(_)));
|
||||
// The stream constructor is wired (the streams task's): returns a
|
||||
// named handle — full behavior the stream tests' scope.
|
||||
let handle = store.stream("s").await.unwrap();
|
||||
assert_eq!(handle.name(), "s");
|
||||
let err = stub_err!(store.queue("q", alkstore::QueueOpts::default()).await);
|
||||
assert!(matches!(err, alkstore::Error::Database(_)));
|
||||
let err = stub_err!(store.outbox("o").await);
|
||||
@@ -692,8 +691,8 @@ async fn store_trait_methods_are_wiring_stubs() {
|
||||
// The remaining stub messages name their landing task (the
|
||||
// wave-3 posture's "… wiring lands with the … task" shape) — in
|
||||
// the source chain (`Database`'s Display is opaque; ADR-008 §5's
|
||||
// detail-carriage). `stream` sampled as the representative.
|
||||
let err = stub_err!(store.stream("ch").await);
|
||||
// detail-carriage). `queue` sampled as the representative.
|
||||
let err = stub_err!(store.queue("ch", alkstore::QueueOpts::default()).await);
|
||||
match err {
|
||||
alkstore::Error::Database(source) => {
|
||||
let msg = source.to_string();
|
||||
@@ -702,7 +701,7 @@ async fn store_trait_methods_are_wiring_stubs() {
|
||||
"stub errors name their landing task, got: {msg}"
|
||||
);
|
||||
}
|
||||
other => panic!("stream stub must be Database, got {other:?}"),
|
||||
other => panic!("queue stub must be Database, got {other:?}"),
|
||||
}
|
||||
|
||||
drop(store);
|
||||
|
||||
File diff suppressed because it is too large.
Load diff
@@ -0,0 +1,856 @@
|
||||
//! The stream mechanism's Postgres arm (`pg-engine-streams`; ADR-015;
|
||||
//! ADR-019 §1/§6; ADR-021 §5; ADR-023 §2):
|
||||
//!
|
||||
//! - `publish` / `publish_with_key` — the auto-commit append paths:
|
||||
//! pool checkout, one parameterized `INSERT INTO … RETURNING id`
|
||||
//! (`id BIGSERIAL` is the assigned offset — engine-assigned,
|
||||
//! monotone per stream, immutable; the nullable `"key"` column is
|
||||
//! carried metadata with the non-empty-key rule: an empty `Some` key
|
||||
//! is `InvalidName`, matching the tx path), `encode_payload` at the
|
||||
//! seam with the typed `Codec` error `?`-ed, return to pool. Then
|
||||
//! **wake**: a `SELECT pg_notify($1, '')` on the stream's wake
|
||||
//! channel — the mechanism-name-is-the-channel realization (the
|
||||
//! plan's Wave 4 section) — as a *separate statement* on the same
|
||||
//! pooled connection (auto-commit: its own implicit transaction;
|
||||
//! the wake is best-effort, not commit-atomic with the insert —
|
||||
//! the no-replay hole covers any gap between the two statements;
|
||||
//! **the durable row is the truth** and every read path re-reads
|
||||
//! from it). A failed wake is logged and swallowed — the durable
|
||||
//! row already landed, and wake feeds must never fail a publish.
|
||||
//! - `read_since` / `read_from_consumer` — pool reads, `ORDER BY id
|
||||
//! ASC` (global FIFO per stream, ADR-015 §4), decoded into
|
||||
//! [`StreamEvent`] via the core constructor's `from_row` (the
|
||||
//! decode owner lives here, shared with the tx path). The extent
|
||||
//! guard (ADR-023 §2) short-circuits `limit <= 0` to the empty
|
||||
//! `Vec` at trait-impl entry, before any round trip — pg's `LIMIT`
|
||||
//! with a negative is a server error, which the guard makes
|
||||
//! unreachable. `read_from_consumer` anchors at the consumer's
|
||||
//! stored checkpoint (absent consumer = 0, ADR-019 §1).
|
||||
//! - `save_offset` / `get_offset` — the monotone upsert (a save below
|
||||
//! the stored checkpoint is a silent no-op — one owner of the
|
||||
//! monotonicity rule, the `GREATEST`-clamped + `WHERE`-guarded
|
||||
//! upsert `save_offset_tx` also rides) and the checkpoint read.
|
||||
//! Both save forms' `Err` arm carries `Database` only (ADR-021 §5).
|
||||
//! - `trim_to(horizon)` — `DELETE … WHERE id <= horizon` on a pool
|
||||
//! connection (wakes nothing — ADR-015 §5's posture; the boundary
|
||||
//! argument is total per ADR-023 §2: a negative horizon deletes
|
||||
//! nothing, idempotently). Surviving events keep their offsets
|
||||
//! (gaps legal, never renumbered).
|
||||
//! - `subscribe(consumer)` — the durable receiver, the structural
|
||||
//! twin of the SQLite engine's bridge with the pg deltas (the
|
||||
//! wake substrate is the forwarder's fanout, natively async — no
|
||||
//! blocking bridge thread):
|
||||
//! - **Attach-before-read**: the forwarder channel subscription is
|
||||
//! taken *before* the initial cursor read, so a commit cannot
|
||||
//! fall between the attach read and the wake wait (a wake for a
|
||||
//! commit already re-read sits queued; a re-drain that sees
|
||||
//! nothing new costs an empty pass, never a missed event).
|
||||
//! - Initial cursor read of the consumer's *stored* offset; page-
|
||||
//! drain to the tail (`id ASC`); park on the wake feed (this
|
||||
//! stream's wake channel plus the reserved reconnect-wake —
|
||||
//! both trigger re-drains); re-drain on wake. Wakes are
|
||||
//! best-effort hints (ADR-006): overtrigger on purpose, coalesc-
|
||||
//! ed by the re-read from the in-memory cursor.
|
||||
//! - **Close semantics (ADR-021 §5's pg arm)**: the receiver stays
|
||||
//! open across forwarder reconnects — the synthetic
|
||||
//! reconnect-wake arrives on the wake feed and simply triggers a
|
||||
//! re-drain (which *is* the gap healing: the no-replay hole's
|
||||
//! missed wakes are covered by the re-drain reading the durable
|
||||
//! rows); closes only at engine shutdown (`recv() -> None`,
|
||||
//! terminal).
|
||||
//! - **Position semantics**: `position` = the last-*yielded* event's
|
||||
//! offset (0 before any yield ⇒ `save_offset()` before any yield
|
||||
//! is a no-op, ADR-021 §5); the bridge's internal cursor advances
|
||||
//! on drained-from-storage events — separate from the position —
|
||||
//! so a stalled consumer's `save_offset()` never checkpoints
|
||||
//! unobserved events. Backpressure: a bounded tokio channel and
|
||||
//! backpressure-aware `send().await` — a stalled consumer parks
|
||||
//! the bridge task and delivery resumes exactly where it stopped.
|
||||
//! - The receiver's `read_since` is offset-anchored and extent-
|
||||
//! guarded like the direct form; its sync `save_offset()` sends
|
||||
//! the last-yielded offset through the same monotone op (ADR-019
|
||||
//! §6). The trait's sync signature on this natively-async engine
|
||||
//! is bridged by a dedicated wake runtime (see
|
||||
//! [`WAKE_RUNTIME`]) — the honest sync→async seam, never
|
||||
//! `spawn_blocking`-reentrancy tricks.
|
||||
//! - `recv()`'s `Err` carries `Database` only (ADR-021 §5 — non-
|
||||
//! `Database` read errors are remapped with the original as
|
||||
//! source, the SQLite `database_only` helper's shape); `try_recv`
|
||||
//! arms mirror the wake receiver: idle `Ok(None)`, `Err(Closed)`
|
||||
//! after the shutdown-close disconnect.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::sync::OnceLock;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
||||
use deadpool_postgres::Pool;
|
||||
use serde_json::Value;
|
||||
use tokio::runtime::{Handle, Runtime, RuntimeFlavor};
|
||||
|
||||
use alkstore::{
|
||||
BoxedFuture, Error, EventReceiver, Result, StreamEvent, StreamHandle, validate_local_name,
|
||||
validate_shared_name,
|
||||
};
|
||||
|
||||
use crate::forwarder::{Forwarder, RawNotification};
|
||||
use crate::schema::{QualifiedTable, tables};
|
||||
use crate::seam::{database_error, pg_error, pool_error};
|
||||
use crate::tx::encode_payload_bytes;
|
||||
|
||||
/// How many events a wake-driven re-read fetches per page while
|
||||
/// draining to the tail.
|
||||
const SUBSCRIBE_PAGE: i64 = 256;
|
||||
|
||||
/// The bridge's tokio-channel capacity — the smoothing buffer between
|
||||
/// the bridge task's re-read loop and the consumer.
|
||||
const EVENT_CHANNEL_CAPACITY: usize = 16;
|
||||
|
||||
/// The per-process idle wake runtime the receiver's **sync**
|
||||
/// `save_offset` (and any other sync receiver op needing async
|
||||
/// storage work) drives its round trip through.
|
||||
///
|
||||
/// Why this exists: the trait's `save_offset(&mut self)` is synchronous,
|
||||
/// but on this engine the checkpoint save is an async pool round trip —
|
||||
/// there is no blocking writer slot to lease (the pg posture is no
|
||||
/// `spawn_blocking` seam and pool clients are runtime-tied). The bridge
|
||||
/// options were probed exhaustively:
|
||||
///
|
||||
/// - `Handle::current().block_on` from within an async context — panics
|
||||
/// ("cannot start a runtime from within a runtime") on both flavors;
|
||||
/// - `Handle::block_on` from a *foreign* runtime's worker/blocking
|
||||
/// threads — also panics (it blocks the calling thread while that
|
||||
/// thread belongs to another runtime's driver, the same enter guard);
|
||||
/// - `Runtime::new().block_on` nested inside a worker thread — panics
|
||||
/// likewise, and builds a fresh runtime per call (unbounded cost).
|
||||
///
|
||||
/// The survivor is a dedicated, runtime-*independent* driver: a fresh
|
||||
/// current-thread runtime per blocked save fails the budget rule, so
|
||||
/// instead a single lazily-built multi-thread runtime is stored
|
||||
/// process-wide ([`WAKE_RUNTIME`]); the save calls
|
||||
/// `wake_runtime().block_on(save_future)` **from a spawned std
|
||||
/// thread** — the one context that is never inside any runtime (a
|
||||
/// task's stack runs *on* a worker thread; a fresh thread starts
|
||||
/// outside every enter guard). Concurrent saves park on their own
|
||||
/// threads and multiplex through one shared runtime; a save issued
|
||||
/// from a foreign runtime's `spawn_blocking` or from a plain thread
|
||||
/// skips the spawn and blocks directly.
|
||||
///
|
||||
/// The runtime carries **no engine state** — it is purely the sync
|
||||
/// call's exec context for its one `.await`; tokio-postgres clients
|
||||
/// bind to whatever runtime drives *them*, and each save checks out
|
||||
/// its own pool connection inside the driven future, on the runtime
|
||||
/// that will drive its connection's I/O. Idle cost: one parked
|
||||
/// runtime (one blocked driver thread, ~nothing else).
|
||||
fn wake_runtime() -> &'static Runtime {
|
||||
static WAKE_RUNTIME: OnceLock<Runtime> = OnceLock::new();
|
||||
WAKE_RUNTIME.get_or_init(|| {
|
||||
tokio::runtime::Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.enable_all()
|
||||
.build()
|
||||
.expect("the wake runtime builds (a fixed 1-worker multi-thread runtime)")
|
||||
})
|
||||
}
|
||||
|
||||
/// Drive `fut` to completion from a **synchronous** context, on the
|
||||
/// dedicated idle wake runtime
|
||||
/// (see [`wake_runtime`] for the full rationale).
|
||||
///
|
||||
/// The calling context decides the driving posture:
|
||||
///
|
||||
/// - plain thread / foreign-runtime blocking thread → `block_on`
|
||||
/// directly on the calling thread (already outside every runtime —
|
||||
/// no spawn, no extra thread);
|
||||
/// - inside a multi-thread runtime's worker (the async consumer calling
|
||||
/// `save_offset().unwrap()` in its task) → `block_in_place` frees
|
||||
/// this worker, then `block_on` on the dedicated runtime (the
|
||||
/// probed-safe shape — a runtime's `block_on` is legal from a
|
||||
/// different runtime's `block_in_place` context);
|
||||
/// - inside a current-thread runtime's worker → a spawned std thread
|
||||
/// does the blocking (the only escape from a current-thread worker —
|
||||
/// `block_in_place` panics there), the caller joins it.
|
||||
fn drive_sync<F>(fut: F) -> Result<F::Output>
|
||||
where
|
||||
F: std::future::Future + Send,
|
||||
F::Output: Send,
|
||||
{
|
||||
let in_runtime = Handle::try_current().is_ok();
|
||||
if !in_runtime {
|
||||
return Ok(wake_runtime().block_on(fut));
|
||||
}
|
||||
let handle = Handle::current();
|
||||
match handle.runtime_flavor() {
|
||||
RuntimeFlavor::MultiThread => {
|
||||
tokio::task::block_in_place(move || Ok(wake_runtime().block_on(fut)))
|
||||
}
|
||||
RuntimeFlavor::CurrentThread => {
|
||||
let boxed: std::pin::Pin<Box<dyn std::future::Future<Output = F::Output> + Send>> =
|
||||
Box::pin(fut);
|
||||
let joined = std::thread::scope(|scope| {
|
||||
let inner = scope.spawn(move || wake_runtime().block_on(boxed));
|
||||
inner
|
||||
.join()
|
||||
.map_err(|_| database_error("the receiver save's driver thread panicked"))
|
||||
})?;
|
||||
Ok(joined)
|
||||
}
|
||||
flavor => Err(database_error(format!(
|
||||
"the receiver save cannot bridge from runtime flavor {flavor:?}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
/// The stream-scoped handle (see the module docs).
|
||||
pub(crate) struct PgStreamHandle {
|
||||
name: String,
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
forwarder: Arc<Forwarder>,
|
||||
closed: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for PgStreamHandle {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("PgStreamHandle")
|
||||
.field("name", &self.name)
|
||||
.finish_non_exhaustive()
|
||||
}
|
||||
}
|
||||
|
||||
impl PgStreamHandle {
|
||||
pub(crate) fn new(
|
||||
name: String,
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
forwarder: Arc<Forwarder>,
|
||||
closed: Arc<AtomicBool>,
|
||||
) -> Self {
|
||||
Self {
|
||||
name,
|
||||
pool,
|
||||
schema,
|
||||
forwarder,
|
||||
closed,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Qualified table reference for this engine-owned schema.
|
||||
fn table(schema: &str, name: &str) -> String {
|
||||
QualifiedTable { schema, name }.to_string()
|
||||
}
|
||||
|
||||
/// The auto-commit publish path (see the module docs). Stream-name
|
||||
/// validation fires at the entry point, before any round trip and
|
||||
/// before the payload is touched (ADR-008 §4); a present key must be
|
||||
/// non-empty (the tx path's rule).
|
||||
fn publish_with_key(
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
closed: Arc<AtomicBool>,
|
||||
stream: String,
|
||||
key: Option<String>,
|
||||
payload: Value,
|
||||
) -> BoxedFuture<'static, Result<i64>> {
|
||||
if let Err(e) = validate_shared_name(&stream) {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
let key = match key {
|
||||
Some(k) if k.is_empty() || k.trim().is_empty() => {
|
||||
return Box::pin(async move { Err(Error::InvalidName { name: k }) });
|
||||
}
|
||||
other => other,
|
||||
};
|
||||
Box::pin(async move {
|
||||
if closed.load(Ordering::Acquire) {
|
||||
return Err(database_error("the store is closed"));
|
||||
}
|
||||
let bytes = encode_payload_bytes(&payload)?;
|
||||
let created_at = crate::resolution::now_unix();
|
||||
let conn = pool.get().await.map_err(pool_error)?;
|
||||
let events = table(&schema, tables::EVENTS);
|
||||
let row = conn
|
||||
.query_one(
|
||||
&format!(
|
||||
"INSERT INTO {events} (stream, \"key\", payload, created_at)
|
||||
VALUES ($1, $2, $3, $4)
|
||||
RETURNING id"
|
||||
),
|
||||
&[&stream, &key, &bytes, &created_at],
|
||||
)
|
||||
.await
|
||||
.map_err(pg_error)?;
|
||||
let offset: i64 = row.get(0);
|
||||
// Then wake: pg_notify on the stream's wake channel — a
|
||||
// separate statement (its own implicit tx on this
|
||||
// auto-commit connection; best-effort, not commit-atomic
|
||||
// with the insert — the no-replay hole covers a gap between
|
||||
// them; the durable row is the truth and every reader
|
||||
// re-reads from it). A wake failure is logged and swallowed:
|
||||
// the durable row landed; no wake fee may ever fail (or
|
||||
// stall) a publish.
|
||||
if let Err(wake_err) = conn.query_one("SELECT pg_notify($1, '')", &[&stream]).await {
|
||||
eprintln!(
|
||||
"alkstore-postgres: stream publish wake failed for {stream:?} \
|
||||
(the durable row landed; readers re-drain on their next hint): {wake_err}"
|
||||
);
|
||||
}
|
||||
Ok(offset)
|
||||
})
|
||||
}
|
||||
|
||||
/// One decoded page from the pool: `id > from`, ASC, up to `limit`
|
||||
/// (the extent guard is the caller's — this helper is contract-blind
|
||||
/// to it). The stream decode owner: shared by the auto-commit reads
|
||||
/// (this module) and the tx reads (`tx.rs`) — the `stream` field
|
||||
/// mapping and the `Codec` posture guard live once.
|
||||
pub(crate) async fn stream_read_page(
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
stream: String,
|
||||
from: i64,
|
||||
limit: i64,
|
||||
) -> Result<Vec<StreamEvent>> {
|
||||
let conn = pool.get().await.map_err(pool_error)?;
|
||||
let events = table(&schema, tables::EVENTS);
|
||||
let rows = conn
|
||||
.query(
|
||||
&format!(
|
||||
"SELECT id, \"key\", payload, created_at
|
||||
FROM {events}
|
||||
WHERE stream = $1 AND id > $2
|
||||
ORDER BY id ASC
|
||||
LIMIT $3"
|
||||
),
|
||||
&[&stream, &from, &limit],
|
||||
)
|
||||
.await
|
||||
.map_err(pg_error)?;
|
||||
stream_events_from_rows(&rows, stream)
|
||||
}
|
||||
|
||||
/// Decode a read page into the contract's [`StreamEvent`] values (the
|
||||
/// core constructor's `from_row`). Malformed rows map to `Error::Codec`
|
||||
/// (the decode-side posture guard, ADR-021 §2). One owner: the tx
|
||||
/// module imports this; no forked decoders.
|
||||
pub(crate) fn stream_events_from_rows(
|
||||
rows: &[tokio_postgres::Row],
|
||||
stream: String,
|
||||
) -> Result<Vec<StreamEvent>> {
|
||||
let mut out = Vec::with_capacity(rows.len());
|
||||
for row in rows {
|
||||
let key: Option<String> = row
|
||||
.try_get(1)
|
||||
.map_err(|e| Error::Codec(format!("stream row key must be text: {e}")))?;
|
||||
let payload: Option<Vec<u8>> = row
|
||||
.try_get(2)
|
||||
.map_err(|e| Error::Codec(format!("stream row payload must be bytea: {e}")))?;
|
||||
out.push(StreamEvent::from_row(
|
||||
row.get(0),
|
||||
stream.clone(),
|
||||
key,
|
||||
payload.unwrap_or_default(),
|
||||
row.get(3),
|
||||
));
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
/// The checkpoint read helper (absent consumer = 0, ADR-019 §1).
|
||||
async fn stream_get_offset(pool: &Pool, schema: &str, stream: &str, consumer: &str) -> Result<i64> {
|
||||
let conn = pool.get().await.map_err(pool_error)?;
|
||||
let offsets = table(schema, tables::OFFSETS);
|
||||
let row = conn
|
||||
.query_opt(
|
||||
&format!("SELECT \"offset\" FROM {offsets} WHERE stream = $1 AND consumer = $2"),
|
||||
&[&stream, &consumer],
|
||||
)
|
||||
.await
|
||||
.map_err(pg_error)?;
|
||||
Ok(row.map(|r| r.get(0)).unwrap_or(0))
|
||||
}
|
||||
|
||||
/// The monotone checkpoint save (one owner of the monotonicity rule —
|
||||
/// the same GREATEST-clamped, WHERE-guarded upsert `save_offset_tx`
|
||||
/// rides; a save at-or-below the stored checkpoint updates nothing,
|
||||
/// ADR-019 §6). First save seeds the checkpoint, clamped at 0 (the
|
||||
/// offsets table's non-negative domain CHECK would abort a negative
|
||||
/// seed; the clamp preserves both the domain and the monotone rule —
|
||||
/// a negative save at-or-below the stored checkpoint is the same
|
||||
/// silent no-op a 0-save is).
|
||||
async fn stream_save_offset(
|
||||
pool: &Pool,
|
||||
schema: &str,
|
||||
stream: &str,
|
||||
consumer: &str,
|
||||
offset: i64,
|
||||
) -> Result<()> {
|
||||
let conn = pool.get().await.map_err(pool_error)?;
|
||||
let offsets = table(schema, tables::OFFSETS);
|
||||
conn.execute(
|
||||
&format!(
|
||||
"INSERT INTO {offsets} AS o (stream, consumer, \"offset\")
|
||||
VALUES ($1, $2, GREATEST($3::bigint, 0))
|
||||
ON CONFLICT (stream, consumer) DO UPDATE
|
||||
SET \"offset\" = EXCLUDED.\"offset\"
|
||||
WHERE EXCLUDED.\"offset\" > o.\"offset\""
|
||||
),
|
||||
&[&stream, &consumer, &offset],
|
||||
)
|
||||
.await
|
||||
.map_err(pg_error)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
impl StreamHandle for PgStreamHandle {
|
||||
fn name(&self) -> &str {
|
||||
&self.name
|
||||
}
|
||||
|
||||
fn publish<'a>(&'a self, payload: Value) -> BoxedFuture<'a, Result<i64>> {
|
||||
publish_with_key(
|
||||
self.pool.clone(),
|
||||
self.schema.clone(),
|
||||
self.closed.clone(),
|
||||
self.name.clone(),
|
||||
None,
|
||||
payload,
|
||||
)
|
||||
}
|
||||
|
||||
fn publish_with_key<'a>(
|
||||
&'a self,
|
||||
key: Option<String>,
|
||||
payload: Value,
|
||||
) -> BoxedFuture<'a, Result<i64>> {
|
||||
publish_with_key(
|
||||
self.pool.clone(),
|
||||
self.schema.clone(),
|
||||
self.closed.clone(),
|
||||
self.name.clone(),
|
||||
key,
|
||||
payload,
|
||||
)
|
||||
}
|
||||
|
||||
fn read_since<'a>(
|
||||
&'a self,
|
||||
offset: i64,
|
||||
limit: i64,
|
||||
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
// The extent guard (ADR-023 §2): `limit <= 0` reads nothing —
|
||||
// a value, never an error — at trait-impl entry, before any
|
||||
// round trip (pg's LIMIT with a negative is a server error,
|
||||
// unreachable behind this guard).
|
||||
if limit <= 0 {
|
||||
return Box::pin(async { Ok(Vec::new()) });
|
||||
}
|
||||
Box::pin(stream_read_page(
|
||||
self.pool.clone(),
|
||||
self.schema.clone(),
|
||||
self.name.clone(),
|
||||
offset,
|
||||
limit,
|
||||
))
|
||||
}
|
||||
|
||||
fn read_from_consumer<'a>(
|
||||
&'a self,
|
||||
consumer: &str,
|
||||
limit: i64,
|
||||
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
||||
if let Err(e) = validate_local_name(consumer) {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
if limit <= 0 {
|
||||
return Box::pin(async { Ok(Vec::new()) });
|
||||
}
|
||||
let pool = self.pool.clone();
|
||||
let schema = self.schema.clone();
|
||||
let stream = self.name.clone();
|
||||
let consumer = consumer.to_string();
|
||||
Box::pin(async move {
|
||||
let from = stream_get_offset(&pool, &schema, &stream, &consumer).await?;
|
||||
stream_read_page(pool, schema, stream, from, limit).await
|
||||
})
|
||||
}
|
||||
|
||||
fn save_offset<'a>(&'a self, consumer: &str, offset: i64) -> BoxedFuture<'a, Result<()>> {
|
||||
if let Err(e) = validate_local_name(consumer) {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
let pool = self.pool.clone();
|
||||
let schema = self.schema.clone();
|
||||
let stream = self.name.clone();
|
||||
let consumer = consumer.to_string();
|
||||
Box::pin(
|
||||
async move { stream_save_offset(&pool, &schema, &stream, &consumer, offset).await },
|
||||
)
|
||||
}
|
||||
|
||||
fn get_offset<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result<i64>> {
|
||||
if let Err(e) = validate_local_name(consumer) {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
let pool = self.pool.clone();
|
||||
let schema = self.schema.clone();
|
||||
let stream = self.name.clone();
|
||||
let consumer = consumer.to_string();
|
||||
Box::pin(async move { stream_get_offset(&pool, &schema, &stream, &consumer).await })
|
||||
}
|
||||
|
||||
fn trim_to<'a>(&'a self, horizon: i64) -> BoxedFuture<'a, Result<i64>> {
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
let pool = self.pool.clone();
|
||||
let schema = self.schema.clone();
|
||||
let stream = self.name.clone();
|
||||
Box::pin(async move {
|
||||
// Wakes nothing (ADR-015 §5's posture — trim is a
|
||||
// bounded-growth op, not an event). The boundary argument
|
||||
// is total (ADR-023 §2): a negative horizon deletes
|
||||
// nothing idempotently (plain SQL, no guard — `id <= -7`
|
||||
// matches no rows).
|
||||
let conn = pool.get().await.map_err(pool_error)?;
|
||||
let events = table(&schema, tables::EVENTS);
|
||||
let deleted = conn
|
||||
.execute(
|
||||
&format!("DELETE FROM {events} WHERE stream = $1 AND id <= $2"),
|
||||
&[&stream, &horizon],
|
||||
)
|
||||
.await
|
||||
.map_err(pg_error)?;
|
||||
Ok(deleted as i64)
|
||||
})
|
||||
}
|
||||
|
||||
fn subscribe<'a>(&'a self, consumer: &str) -> BoxedFuture<'a, Result<Box<dyn EventReceiver>>> {
|
||||
if let Err(e) = validate_local_name(consumer) {
|
||||
return Box::pin(async move { Err(e) });
|
||||
}
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
let name = self.name.clone();
|
||||
let pool = self.pool.clone();
|
||||
let schema = self.schema.clone();
|
||||
let forwarder = self.forwarder.clone();
|
||||
let closed = self.closed.clone();
|
||||
let consumer = consumer.to_string();
|
||||
Box::pin(async move { spawn_bridge(name, consumer, pool, schema, forwarder, closed).await })
|
||||
}
|
||||
}
|
||||
|
||||
/// The `subscribe` bridge construction: the forwarder subscription is
|
||||
/// taken before any read (a commit cannot fall between the attach
|
||||
/// read and the wake wait), the bridge task then replays from the
|
||||
/// stored offset and parks on the wake feed. Subscribing a failed
|
||||
/// forwarder (`None` — shutdown raced) fails closed after
|
||||
/// unregistering the just-registered LISTEN (no registration leak —
|
||||
/// the notify path's same rule); a bridge-spawn failure (runtime
|
||||
/// gone) also fails `Database`.
|
||||
async fn spawn_bridge(
|
||||
stream: String,
|
||||
consumer: String,
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
forwarder: Arc<Forwarder>,
|
||||
store_closed: Arc<AtomicBool>,
|
||||
) -> Result<Box<dyn EventReceiver>> {
|
||||
// Registration first — a failed registration here (never: LISTEN
|
||||
// failures carry the registry entry's recovery, but a shutdown
|
||||
// race can kill the command channel) fails the whole subscribe.
|
||||
forwarder.register(&stream).await?;
|
||||
let Some(broadcast_rx) = forwarder.subscribe() else {
|
||||
// The subscribe raced (or arrived past) a shutdown: fail
|
||||
// closed — no receiver can ever deliver again. Unregister
|
||||
// first: the acked LISTEN must not leak past the failed
|
||||
// listen (a next-closest `subscribe` re-issues it).
|
||||
forwarder.unregister(&stream);
|
||||
return Err(database_error("the store is closed"));
|
||||
};
|
||||
let (event_tx, event_rx) =
|
||||
tokio::sync::mpsc::channel::<Result<StreamEvent>>(EVENT_CHANNEL_CAPACITY);
|
||||
let spawned = tokio::spawn(run_subscribe_loop(
|
||||
stream.clone(),
|
||||
consumer.clone(),
|
||||
pool.clone(),
|
||||
schema.clone(),
|
||||
store_closed.clone(),
|
||||
broadcast_rx,
|
||||
event_tx,
|
||||
));
|
||||
// tokio::spawn only fails on a shut-down runtime — and the boxed
|
||||
// receiver must never surface as a lie after its bridge never
|
||||
// started. (Also: no panics cross the seam, the W-2 posture.)
|
||||
if spawned.is_finished() {
|
||||
spawned.abort();
|
||||
return Err(database_error(
|
||||
"failed to spawn the stream bridge task (runtime shutting down)",
|
||||
));
|
||||
}
|
||||
Ok(Box::new(PgEventReceiver {
|
||||
stream,
|
||||
consumer,
|
||||
rx: event_rx,
|
||||
forwarder,
|
||||
position: 0,
|
||||
pool,
|
||||
schema,
|
||||
store_closed,
|
||||
_bridge: spawned,
|
||||
}))
|
||||
}
|
||||
|
||||
/// The bridge loop body (the subscribe task; the SQLite arm's
|
||||
/// std-thread shape, natively async here):
|
||||
///
|
||||
/// 1. read the consumer's stored offset (the attach read),
|
||||
/// 2. drain pages until the tail (`id ASC`),
|
||||
/// 3. park on the wake feed (this stream's channel + the reserved
|
||||
/// reconnect-wake), re-drain on wake,
|
||||
/// 4. a read failure while the store is open surfaces one
|
||||
/// `Err(Database)` and waits for the next wake; a read failure
|
||||
/// after store close exits the loop (the consumer sees the close).
|
||||
///
|
||||
/// Wake-feed disconnect (the forwarder's shutdown take drops the
|
||||
/// broadcast) is the driver arm: the loop exits, the mpsc sender
|
||||
/// drops, the consumer's `recv()` yields `None` — terminal, only at
|
||||
/// engine shutdown.
|
||||
enum WakeTick {
|
||||
Ok,
|
||||
Closed,
|
||||
}
|
||||
|
||||
async fn wait_wake(rx: &mut tokio::sync::broadcast::Receiver<RawNotification>) -> WakeTick {
|
||||
match rx.recv().await {
|
||||
Ok(_raw) => WakeTick::Ok,
|
||||
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {
|
||||
// Still surfaced — treat like a wake: re-drain.
|
||||
WakeTick::Ok
|
||||
}
|
||||
Err(tokio::sync::broadcast::error::RecvError::Closed) => WakeTick::Closed,
|
||||
}
|
||||
}
|
||||
|
||||
/// One drain-page arm: a decoded page (possibly empty), or the loop's
|
||||
/// exit state.
|
||||
enum PageArm {
|
||||
Page(Vec<StreamEvent>),
|
||||
Err(Error),
|
||||
Closed,
|
||||
}
|
||||
|
||||
async fn drain_page(
|
||||
pool: &Pool,
|
||||
schema: &str,
|
||||
stream: &str,
|
||||
cursor: i64,
|
||||
store_closed: &AtomicBool,
|
||||
) -> PageArm {
|
||||
if store_closed.load(Ordering::Acquire) {
|
||||
return PageArm::Closed;
|
||||
}
|
||||
match stream_read_page(
|
||||
pool.clone(),
|
||||
schema.to_string(),
|
||||
stream.to_string(),
|
||||
cursor,
|
||||
SUBSCRIBE_PAGE,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(events) => PageArm::Page(events),
|
||||
Err(_e) if store_closed.load(Ordering::Acquire) => PageArm::Closed,
|
||||
Err(e) => PageArm::Err(database_only(e)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_subscribe_loop(
|
||||
stream: String,
|
||||
consumer: String,
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
store_closed: Arc<AtomicBool>,
|
||||
mut wake_rx: tokio::sync::broadcast::Receiver<RawNotification>,
|
||||
event_tx: tokio::sync::mpsc::Sender<Result<StreamEvent>>,
|
||||
) {
|
||||
// 1. The attach read: the consumer's stored offset.
|
||||
if store_closed.load(Ordering::Acquire) {
|
||||
return;
|
||||
}
|
||||
let mut cursor = match stream_get_offset(&pool, &schema, &stream, &consumer).await {
|
||||
Ok(offset) => offset,
|
||||
Err(_e) if store_closed.load(Ordering::Acquire) => return,
|
||||
Err(e) => {
|
||||
let _ = event_tx.send(Err(database_only(e))).await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
loop {
|
||||
// 2. Drain to the tail, page by page.
|
||||
loop {
|
||||
match drain_page(&pool, &schema, &stream, cursor, &store_closed).await {
|
||||
PageArm::Closed => return,
|
||||
PageArm::Err(e) => {
|
||||
if event_tx.send(Err(e)).await.is_err() {
|
||||
return;
|
||||
}
|
||||
break;
|
||||
}
|
||||
PageArm::Page(events) => {
|
||||
if events.is_empty() {
|
||||
break;
|
||||
}
|
||||
for event in events {
|
||||
cursor = event.offset;
|
||||
if event_tx.send(Ok(event)).await.is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// 3. Park on the wake feed (this stream's channel — and the
|
||||
// reserved reconnect-wake: the same re-drain heals a
|
||||
// gap's missed wakes by reading the durable rows).
|
||||
match wait_wake(&mut wake_rx).await {
|
||||
WakeTick::Closed => return,
|
||||
WakeTick::Ok => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// `recv()`'s `Err` arm carries `Database` only (ADR-021 §5): a
|
||||
/// read error that mapped elsewhere (the decode-side `Codec` posture
|
||||
/// guard) is remapped into the opaque fallback with the original error
|
||||
/// preserved as the source.
|
||||
fn database_only(e: Error) -> Error {
|
||||
match e {
|
||||
already @ Error::Database(_) => already,
|
||||
other => Error::database(other),
|
||||
}
|
||||
}
|
||||
|
||||
/// The durable subscription receiver: the consumer side of the
|
||||
/// bridge. Dropping it drops the mpsc receiver — the bridge's next
|
||||
/// `send` fails and the bridge task exits (and its `_bridge` handle
|
||||
/// drop-detaches the task body), while the forwarder registration
|
||||
/// rides the receiver's `Drop` (unregister → UNLISTEN at the
|
||||
/// last-subscriber drop — one refcount owner per channel).
|
||||
pub(crate) struct PgEventReceiver {
|
||||
stream: String,
|
||||
consumer: String,
|
||||
rx: tokio::sync::mpsc::Receiver<Result<StreamEvent>>,
|
||||
forwarder: Arc<Forwarder>,
|
||||
position: i64,
|
||||
pool: Pool,
|
||||
schema: String,
|
||||
store_closed: Arc<AtomicBool>,
|
||||
_bridge: tokio::task::JoinHandle<()>,
|
||||
}
|
||||
|
||||
impl Drop for PgEventReceiver {
|
||||
fn drop(&mut self) {
|
||||
self.forwarder.unregister(&self.stream);
|
||||
}
|
||||
}
|
||||
|
||||
impl PgEventReceiver {
|
||||
/// The monotone save of this receiver's position through the
|
||||
/// shared op. The sync signature on the natively-async engine is
|
||||
/// bridged on the dedicated idle wake runtime (see
|
||||
/// [`drive_sync`]; the runtime carries no engine state — the
|
||||
/// save future checks out its own pool connection). `position ==
|
||||
/// 0` (nothing yielded yet) is the no-op arm — the stored
|
||||
/// checkpoint untouched (ADR-021 §5).
|
||||
fn save_now(&mut self) -> Result<()> {
|
||||
let position = self.position;
|
||||
if position == 0 {
|
||||
return Ok(());
|
||||
}
|
||||
if self.store_closed.load(Ordering::Acquire) {
|
||||
return Err(database_error("the store is closed"));
|
||||
}
|
||||
let pool = self.pool.clone();
|
||||
let schema = self.schema.clone();
|
||||
let stream = self.stream.clone();
|
||||
let consumer = self.consumer.clone();
|
||||
drive_sync(async move {
|
||||
stream_save_offset(&pool, &schema, &stream, &consumer, position).await
|
||||
})?
|
||||
}
|
||||
}
|
||||
|
||||
impl EventReceiver for PgEventReceiver {
|
||||
fn recv<'a>(&'a mut self) -> BoxedFuture<'a, Option<Result<StreamEvent>>> {
|
||||
Box::pin(async move {
|
||||
match self.rx.recv().await {
|
||||
Some(Ok(event)) => {
|
||||
self.position = event.offset;
|
||||
Some(Ok(event))
|
||||
}
|
||||
// A transient read failure — the bridge keeps running
|
||||
// and re-reads on the next wake.
|
||||
Some(Err(e)) => Some(Err(e)),
|
||||
// The bridge exited: shutdown disconnect — `None` =
|
||||
// closed, terminal (it never reopens; recovery is a
|
||||
// fresh `subscribe(consumer)`).
|
||||
None => None,
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
fn try_recv(&mut self) -> Result<Option<StreamEvent>> {
|
||||
match self.rx.try_recv() {
|
||||
Ok(Ok(event)) => {
|
||||
self.position = event.offset;
|
||||
Ok(Some(event))
|
||||
}
|
||||
Ok(Err(e)) => Err(e),
|
||||
Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None),
|
||||
Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => Err(Error::Closed),
|
||||
}
|
||||
}
|
||||
|
||||
fn read_since<'a>(
|
||||
&'a mut self,
|
||||
offset: i64,
|
||||
limit: i64,
|
||||
) -> BoxedFuture<'a, Result<Vec<StreamEvent>>> {
|
||||
if self.store_closed.load(Ordering::Acquire) {
|
||||
return Box::pin(async { Err(database_error("the store is closed")) });
|
||||
}
|
||||
if limit <= 0 {
|
||||
return Box::pin(async { Ok(Vec::new()) });
|
||||
}
|
||||
Box::pin(stream_read_page(
|
||||
self.pool.clone(),
|
||||
self.schema.clone(),
|
||||
self.stream.clone(),
|
||||
offset,
|
||||
limit,
|
||||
))
|
||||
}
|
||||
|
||||
fn save_offset(&mut self) -> Result<()> {
|
||||
self.save_now()
|
||||
}
|
||||
|
||||
fn offset(&self) -> i64 {
|
||||
self.position
|
||||
}
|
||||
}
|
||||
@@ -79,6 +79,7 @@ use crate::resolution::{
|
||||
};
|
||||
use crate::schema::{QualifiedTable, tables};
|
||||
use crate::seam::{pg_error, pool_error};
|
||||
use crate::stream::stream_events_from_rows;
|
||||
|
||||
/// The notification payload boundary (the contract's pinned limit —
|
||||
/// client-side checked, typed before any round trip; the predicate
|
||||
@@ -637,27 +638,7 @@ fn job_from_row(row: &tokio_postgres::Row) -> Result<Job> {
|
||||
))
|
||||
}
|
||||
|
||||
/// Decode a read page into the contract's [`StreamEvent`] values (the
|
||||
/// core constructor's `from_row`). Malformed rows map to `Error::Codec`.
|
||||
fn stream_events_from_rows(
|
||||
rows: &[tokio_postgres::Row],
|
||||
stream: String,
|
||||
) -> Result<Vec<StreamEvent>> {
|
||||
let mut out = Vec::with_capacity(rows.len());
|
||||
for row in rows {
|
||||
let key: Option<String> = row
|
||||
.try_get(1)
|
||||
.map_err(|e| Error::Codec(format!("stream row key must be text: {e}")))?;
|
||||
let payload: Option<Vec<u8>> = row
|
||||
.try_get(2)
|
||||
.map_err(|e| Error::Codec(format!("stream row payload must be bytea: {e}")))?;
|
||||
out.push(StreamEvent::from_row(
|
||||
row.get(0),
|
||||
stream.clone(),
|
||||
key,
|
||||
payload.unwrap_or_default(),
|
||||
row.get(3),
|
||||
));
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
// The stream-page decode (`stream_events_from_rows`) lives in
|
||||
// `stream.rs` — the one decode owner the tx reads and the auto-commit
|
||||
// reads share (the `stream` field mapping and the `Codec` posture
|
||||
// guard live once; ADR-012 §2's one-owner rule).
|
||||
Reference in new issue
Block a user