pg queue/tx/notify dedupe + doc truth-telling: job_from_row has one owner (queue.rs, tx.rs imports it) and get_job_tx reuses live_columns()/dead_columns() instead of inlining the 18-column lists — the contract-pinned Job shape now has exactly one decode owner (ADR-012 §2); sweep_expired's doc states the moved + retention-deleted sum in the in_tx multi-statement frame (the sum the SQLite twin returns — the doc was the only liar); notify.rs's closed-store listen error routes through the shared database_error helper; bridge_capacity() is a const with its rationale doc carried. No behavior change — identical SQL strings, error shapes, and capacity. (task pg-fix-dedupe-cleanup, review 002 minor notes + Finding 6 queue bullet)
This commit is contained in:
1 parent
12d0497b1c
commit
40625090f6
4 files changed
+61
-77
No files matched your search
@@ -105,9 +105,7 @@ pub(crate) fn listen(
|
||||
}
|
||||
Box::pin(async move {
|
||||
if closed.load(Ordering::Acquire) {
|
||||
return Err(Error::database(std::io::Error::other(
|
||||
"the store is closed",
|
||||
)));
|
||||
return Err(database_error("the store is closed"));
|
||||
}
|
||||
forwarder.register(&channel).await?;
|
||||
let broadcast_rx = forwarder.subscribe();
|
||||
@@ -119,7 +117,7 @@ pub(crate) fn listen(
|
||||
forwarder.unregister(&channel);
|
||||
return Err(database_error("the store is closed"));
|
||||
};
|
||||
let (wake_tx, wake_rx) = tokio::sync::mpsc::channel(bridge_capacity());
|
||||
let (wake_tx, wake_rx) = tokio::sync::mpsc::channel(BRIDGE_CAPACITY);
|
||||
let bridge = tokio::spawn(bridge_loop(broadcast_rx, wake_tx, channel.clone()));
|
||||
Ok(Box::new(PgWakeReceiver {
|
||||
wake_rx,
|
||||
@@ -134,9 +132,7 @@ pub(crate) fn listen(
|
||||
/// broadcast → receiver). The wake rate is the notify rate (bursty at
|
||||
/// worst); 1024 matches the fanout's bounded shape, lag surfaced at
|
||||
/// the broadcast layer (not here).
|
||||
fn bridge_capacity() -> usize {
|
||||
1024
|
||||
}
|
||||
const BRIDGE_CAPACITY: usize = 1024;
|
||||
|
||||
/// The broadcast bridge: loops the raw fanout, maps notifications on
|
||||
/// this receiver's channel (plus the reserved reconnect-wake — every
|
||||
|
||||
@@ -68,8 +68,9 @@
|
||||
//! ADR-020).
|
||||
//! - `sweep_expired()` — both-states no-stranded-rows move (any
|
||||
//! past-expiry row → dead with `"expired"`) + retention enforcement
|
||||
//! (ADR-010 §5), single-statement atomic; returns the moved count.
|
||||
//! No leader lock required.
|
||||
//! (ADR-010 §5), in the [`in_tx`] multi-statement frame; returns the
|
||||
//! moved + retention-deleted sum (the SQLite twin's same count). No
|
||||
//! leader lock required.
|
||||
|
||||
use std::collections::hash_map::RandomState;
|
||||
use std::hash::BuildHasher;
|
||||
@@ -192,7 +193,7 @@ fn table(schema: &str, name: &str) -> String {
|
||||
/// The live-row select's column list (the uniform 18-column shape the
|
||||
/// dead-table select mirrors with `last_error`/`died_at` in place of
|
||||
/// the NULL casts) — one shape, shared with the tx reads.
|
||||
fn live_columns() -> &'static str {
|
||||
pub(crate) fn live_columns() -> &'static str {
|
||||
"id, queue, state, payload, priority, run_at, attempts,
|
||||
max_attempts, worker_id, claimed_at, claim_expires_at,
|
||||
created_at, expires_at, visibility_timeout_s,
|
||||
@@ -202,7 +203,7 @@ fn live_columns() -> &'static str {
|
||||
|
||||
/// The dead-row select's column list (the 18-column live shape with
|
||||
/// `last_error`/`died_at` riding the NULL casts' slots).
|
||||
fn dead_columns() -> &'static str {
|
||||
pub(crate) fn dead_columns() -> &'static str {
|
||||
"id, queue, 'dead' AS state, payload, priority, run_at, attempts,
|
||||
max_attempts, worker_id, claimed_at, claim_expires_at,
|
||||
created_at, expires_at, visibility_timeout_s,
|
||||
@@ -406,8 +407,9 @@ pub(crate) fn claim_stmt_sql(schema: &str) -> String {
|
||||
/// Decode an 18-column job row into the contract value
|
||||
/// ([`Job::from_row`] — the `#[doc(hidden)]` constructor's field list
|
||||
/// of record, ADR-021 §2). Malformed rows map to `Error::Codec` (the
|
||||
/// decode-side posture guard). One decode owner: the tx reads and this
|
||||
/// module's claim/get rows share the shape.
|
||||
/// decode-side posture guard). The one decode owner: the tx reads
|
||||
/// ([`crate::tx`]) and this module's claim/get rows import it, so a
|
||||
/// contract-pinned `Job` shape edit has exactly one place to land.
|
||||
pub(crate) fn job_from_row(row: &tokio_postgres::Row) -> Result<Job> {
|
||||
let state = match row.get::<_, Option<&str>>(2) {
|
||||
Some("pending") => JobState::Pending,
|
||||
|
||||
+14
-61
@@ -84,10 +84,11 @@ use serde_json::Value;
|
||||
use tokio::runtime::Handle;
|
||||
|
||||
use alkstore::{
|
||||
BoxedFuture, EnqueueOpts, Error, Job, JobState, Result, StreamEvent, TxHandle,
|
||||
validate_local_name, validate_shared_name,
|
||||
BoxedFuture, EnqueueOpts, Error, Job, Result, StreamEvent, TxHandle, validate_local_name,
|
||||
validate_shared_name,
|
||||
};
|
||||
|
||||
use crate::queue::{dead_columns, job_from_row, live_columns};
|
||||
use crate::resolution::{
|
||||
ResolvedEnqueue, outbox_backing_queue_name, outbox_default_stamps, plain_queue_default_stamps,
|
||||
resolve_enqueue_opts, stamps_with_override,
|
||||
@@ -471,13 +472,9 @@ impl TxHandle for PgTxHandle {
|
||||
let row = client
|
||||
.query_opt(
|
||||
&format!(
|
||||
"SELECT id, queue, 'dead', payload, priority, run_at, attempts,
|
||||
max_attempts, worker_id, claimed_at, claim_expires_at,
|
||||
created_at, expires_at, visibility_timeout_s,
|
||||
backoff_base_s, dead_letter_retention_s,
|
||||
last_error, died_at
|
||||
FROM {} WHERE id = $1 AND queue = $2",
|
||||
self.table(tables::DEAD)
|
||||
"SELECT {dead} FROM {deadt} WHERE id = $1 AND queue = $2",
|
||||
dead = dead_columns(),
|
||||
deadt = self.table(tables::DEAD)
|
||||
),
|
||||
&[&job_id, &queue],
|
||||
)
|
||||
@@ -489,13 +486,9 @@ impl TxHandle for PgTxHandle {
|
||||
let row = client
|
||||
.query_opt(
|
||||
&format!(
|
||||
"SELECT id, queue, state, payload, priority, run_at, attempts,
|
||||
max_attempts, worker_id, claimed_at, claim_expires_at,
|
||||
created_at, expires_at, visibility_timeout_s,
|
||||
backoff_base_s, dead_letter_retention_s,
|
||||
NULL::text, NULL::bigint
|
||||
FROM {} WHERE id = $1 AND queue = $2",
|
||||
self.table(tables::JOB)
|
||||
"SELECT {live} FROM {jobt} WHERE id = $1 AND queue = $2",
|
||||
live = live_columns(),
|
||||
jobt = self.table(tables::JOB)
|
||||
),
|
||||
&[&job_id, &queue],
|
||||
)
|
||||
@@ -678,48 +671,8 @@ impl Drop for PgTxHandle {
|
||||
}
|
||||
}
|
||||
|
||||
/// Decode a job row (the uniform 18-column shape — the dead-table
|
||||
/// select carries `last_error`/`died_at`; the live one binds NULLs)
|
||||
/// into the contract value (`Job::from_row` — the `#[doc(hidden)]`
|
||||
/// constructor, ADR-021 §2's field list of record). Malformed rows map
|
||||
/// to `Error::Codec` (the decode-side posture guard).
|
||||
fn job_from_row(row: &tokio_postgres::Row) -> Result<Job> {
|
||||
let state = match row.get::<_, Option<&str>>(2) {
|
||||
Some("pending") => JobState::Pending,
|
||||
Some("processing") => JobState::Processing,
|
||||
Some("dead") => JobState::Dead,
|
||||
other => {
|
||||
return Err(Error::Codec(format!(
|
||||
"job row state must be pending|processing|dead, got {other:?}"
|
||||
)));
|
||||
}
|
||||
};
|
||||
let payload: Option<Vec<u8>> = row
|
||||
.try_get(3)
|
||||
.map_err(|e| Error::Codec(format!("job row payload must be bytea: {e}")))?;
|
||||
Ok(Job::from_row(
|
||||
row.get(0),
|
||||
row.get(1),
|
||||
state,
|
||||
payload.unwrap_or_default(),
|
||||
row.get(4),
|
||||
row.get(5),
|
||||
row.get(6),
|
||||
row.get(7),
|
||||
row.get(8),
|
||||
row.get(9),
|
||||
row.get(10),
|
||||
row.get(11),
|
||||
row.get(12),
|
||||
row.get(13),
|
||||
row.get(14),
|
||||
row.get(15),
|
||||
row.get(16),
|
||||
row.get(17),
|
||||
))
|
||||
}
|
||||
|
||||
// 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).
|
||||
// The row decodes ride one owner each: the stream-page decode
|
||||
// (`stream_events_from_rows`) lives in `stream.rs`, the job-row
|
||||
// decode (`job_from_row` + the select column lists) in `queue.rs` —
|
||||
// the one decode owners the tx reads and the auto-commit reads share
|
||||
// (ADR-012 §2's one-owner rule).
|
||||
Reference in new issue
Block a user