Postgres engine: schema bootstrap — engine-owned schema, table family, idempotent DDL, schema-prefixed indexes (task pg-engine-schema)
This commit is contained in:
1 parent
4caea21098
commit
2f1353bd41
6 files changed
+951
-9
No files matched your search
Generated
+1
@@ -28,6 +28,7 @@ dependencies = [
|
||||
"alkstore",
|
||||
"alkstore-contract-suite",
|
||||
"deadpool-postgres",
|
||||
"thiserror",
|
||||
"tokio",
|
||||
"tokio-postgres",
|
||||
]
|
||||
|
||||
@@ -10,6 +10,7 @@ alkstore = { version = "0.1", path = "../alkstore" }
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
|
||||
tokio-postgres = "0.7"
|
||||
deadpool-postgres = "0.14"
|
||||
thiserror = "2"
|
||||
|
||||
[dev-dependencies]
|
||||
alkstore-contract-suite = { path = "../alkstore-contract-suite" }
|
||||
@@ -1 +1,27 @@
|
||||
//! alkstore-postgres — the Postgres engine: the core contract
|
||||
//! implemented on tokio-postgres + deadpool-postgres (ADR-001,
|
||||
//! ADR-004).
|
||||
//!
|
||||
//! # Posture
|
||||
//!
|
||||
//! **Multi-host by nature** (ADR-016): connections are per-process
|
||||
//! state; nothing assumes a shared host — this engine's deployment
|
||||
//! boundary is its crate identity (compile-time engine discovery:
|
||||
//! pick this crate, get this posture). No runtime capability surface.
|
||||
//!
|
||||
//! **Natively async** (ADR-004/ADR-007): no blocking bridge, no
|
||||
//! `spawn_blocking` seam — the driver's client is `Send + Sync` and
|
||||
//! the tx handle holds the pooled object directly (POC #2
|
||||
//! compile-probe verified).
|
||||
//!
|
||||
//! Wave 4 build order: schema bootstrap (this module — the dependency
|
||||
//! root), then `open`/`PgOpts` + pool + listener wiring, the tx seam,
|
||||
//! the LISTEN forwarder, the mechanisms (queues, streams, locks,
|
||||
//! scheduler/outbox); each task replaces the prior stubs.
|
||||
|
||||
mod schema;
|
||||
|
||||
pub use schema::{
|
||||
BootstrapError, BootstrapResult, DEFAULT_SCHEMA, QualifiedTable, bootstrap, quote_identifier,
|
||||
tables,
|
||||
};
|
||||
@@ -0,0 +1,358 @@
|
||||
//! Schema bootstrap (ADR-010 §8): the one engine-owned PostgreSQL
|
||||
//! schema and its table family — the dependency root every later
|
||||
//! wave-4 task rides.
|
||||
//!
|
||||
//! # Layout posture
|
||||
//!
|
||||
//! All engine-owned objects — job, dead, stream events, offsets,
|
||||
//! schedule, locks — live in one PostgreSQL schema, default
|
||||
//! [`DEFAULT_SCHEMA`] (`"alkstore"`), overridable per-engine option
|
||||
//! (the `PgOpts::schema` knob the `open` task wires), co-tenanted
|
||||
//! safely with consumer tables (deployment.md's co-tenancy posture;
|
||||
//! the alkblobs ADR-008 precedent). Queues are rows in one job table —
|
||||
//! no per-queue tables, no per-queue schemas (ADR-010 §8); per-queue
|
||||
//! configuration is stamped onto job rows (ADR-010 §3a), never stored
|
||||
//! in a registry. Queue creation is not a DDL operation.
|
||||
//!
|
||||
//! # Payload posture
|
||||
//!
|
||||
//! Payload columns are `BYTEA` — the engine stores exactly the bytes
|
||||
//! `encode_payload` returned, byte for byte (ADR-020 §4). `JSONB` is
|
||||
//! deliberately rejected: its normalization (key reordering,
|
||||
//! duplicate-key folding, number canonicalization) promises no byte
|
||||
//! identity, and the contract suite's byte-identical-rows row pins
|
||||
//! that rows enqueueing the same `Value` store identical bytes.
|
||||
//!
|
||||
//! # No migration machinery in v1
|
||||
//!
|
||||
//! Tables are created fresh per schema via idempotent DDL
|
||||
//! (`CREATE … IF NOT EXISTS`); this module carries no ALTER/migration
|
||||
//! machinery. Append-column evolution is the SQLite substrate's
|
||||
//! discipline for its in-place database files; a greenfield pg schema
|
||||
//! needs none in v1 — wave 6's release pass revisits if the
|
||||
//! deployment story demands in-place upgrades.
|
||||
//!
|
||||
//! # Column sets
|
||||
//!
|
||||
//! Column sets follow the contract's value shapes — the
|
||||
//! `#[doc(hidden)]` engine constructors' field lists
|
||||
//! (`alkstore::Job::from_row`, `alkstore::StreamEvent::from_row`,
|
||||
//! `alkstore::Schedule::new`) are the lists of record (ADR-019 §3,
|
||||
//! ADR-021 §2, ADR-015 §3), with the `QueueOpts` stamp columns per
|
||||
//! ADR-010 §3a and the schedule row's boundary-state columns per
|
||||
//! ADR-009 §3/§4 (the substrate's `__alkstore_scheduler_tasks`
|
||||
//! concrete stamp set re-owned in pg form; `enabled` dropped — the
|
||||
//! collapse surface has no pause/resume, ADR-009 §1). Timestamps are
|
||||
//! unix-seconds `BIGINT` engine-wide — the engines' single clock
|
||||
//! source, second-precision integer arithmetic (honker parity); no
|
||||
//! `timestamptz`.
|
||||
//!
|
||||
//! # Identifier quoting
|
||||
//!
|
||||
//! The engine composes schema-qualified SQL from a consumer-supplied
|
||||
//! schema name — [`quote_identifier`] is the injection boundary:
|
||||
//! every identifier entering SQL text (schema name, table names,
|
||||
//! index names) is quoted through it. Two column names are
|
||||
//! reserved-word hazards in Postgres and must stay quoted in later
|
||||
//! tasks' SQL: the offsets checkpoint column `"offset"` and the
|
||||
//! events key column `"key"`.
|
||||
//!
|
||||
//! # Index-name qualification
|
||||
//!
|
||||
//! `CREATE INDEX` does not accept a schema-qualified index name — the
|
||||
//! index lands in the table's schema automatically — and
|
||||
//! `IF NOT EXISTS` on the resulting *bare* name checks existence
|
||||
//! across every schema on the search path: a co-tenant's unrelated
|
||||
//! `public.job_claim_idx` would silently skip our bootstrap's index
|
||||
//! creation. The bootstrap therefore names every index
|
||||
//! `{schema}_{table}_{suffix}` (schema-prefixed, quoted) — a name a
|
||||
//! non-engine co-tenant will not carry, and unique among engine
|
||||
//! instances with distinct schema names too. The canonical suffixes
|
||||
//! ([`indexes`]) are what the later mechanism tasks' SQL reuses.
|
||||
|
||||
use std::fmt;
|
||||
|
||||
use thiserror::Error;
|
||||
|
||||
/// The default engine-owned schema name (ADR-010 §8).
|
||||
pub const DEFAULT_SCHEMA: &str = "alkstore";
|
||||
|
||||
/// A bootstrap failure: a DDL statement errored server-side (or the
|
||||
/// connection did).
|
||||
///
|
||||
/// The `open` constructor (the next task) maps this onto the engine's
|
||||
/// `Error::Database` surface, source chain preserved.
|
||||
#[derive(Debug, Error)]
|
||||
#[error("schema bootstrap failed")]
|
||||
pub struct BootstrapError(#[source] pub tokio_postgres::Error);
|
||||
|
||||
/// Result alias for [`bootstrap`].
|
||||
pub type BootstrapResult = Result<(), BootstrapError>;
|
||||
|
||||
/// Quote a SQL identifier (Postgres rule: double embedded `"`).
|
||||
///
|
||||
/// Consumer-supplied schema names ride DDL text through this quoting;
|
||||
/// an identifier is never interpolated into SQL unquoted.
|
||||
pub fn quote_identifier(name: &str) -> String {
|
||||
let mut quoted = String::with_capacity(name.len() + 2);
|
||||
quoted.push('"');
|
||||
for c in name.chars() {
|
||||
if c == '"' {
|
||||
quoted.push('"');
|
||||
}
|
||||
quoted.push(c);
|
||||
}
|
||||
quoted.push('"');
|
||||
quoted
|
||||
}
|
||||
|
||||
/// A schema-qualified table or index reference: `"schema"."name"`.
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub struct QualifiedTable<'a> {
|
||||
/// The schema name (unquoted form).
|
||||
pub schema: &'a str,
|
||||
/// The table/index name (unquoted form).
|
||||
pub name: &'a str,
|
||||
}
|
||||
|
||||
impl fmt::Display for QualifiedTable<'_> {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(
|
||||
f,
|
||||
"{}.{}",
|
||||
quote_identifier(self.schema),
|
||||
quote_identifier(self.name)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// The engine-owned object names inside the schema — fixed constants,
|
||||
/// the one owner of the names the later mechanism tasks' SQL reuses.
|
||||
///
|
||||
/// Object names do not ride engine configuration (only the schema
|
||||
/// name does, ADR-010 §8).
|
||||
pub mod tables {
|
||||
/// The job table (live rows — pending and processing).
|
||||
pub const JOB: &str = "job";
|
||||
/// The dead-letter table (physical move target, ADR-010 §4).
|
||||
pub const DEAD: &str = "dead";
|
||||
/// The stream events table (`id BIGSERIAL` = the offset).
|
||||
pub const EVENTS: &str = "stream_events";
|
||||
/// The per-consumer offset checkpoint table.
|
||||
pub const OFFSETS: &str = "stream_offsets";
|
||||
/// The schedule table (ADR-009's collapse shape).
|
||||
pub const SCHEDULE: &str = "schedule";
|
||||
/// The TTL lock table.
|
||||
pub const LOCKS: &str = "locks";
|
||||
}
|
||||
|
||||
/// The engine-owned index names inside the schema — schema-prefixed
|
||||
/// (see the module docs' index-name qualification posture),
|
||||
/// schema-parameterized at DDL-composition time by [`bootstrap_ddl`].
|
||||
pub mod indexes {
|
||||
/// The claim-ordering index (`priority DESC, run_at, id`).
|
||||
pub const JOB_CLAIM: &str = "job_claim_idx";
|
||||
/// The pending-ready-time hot-path index.
|
||||
pub const JOB_PENDING: &str = "job_pending_idx";
|
||||
/// The processing-claim-deadline index.
|
||||
pub const JOB_DEADLINE: &str = "job_deadline_idx";
|
||||
/// The dead-rows retention index (`queue, died_at`).
|
||||
pub const DEAD_RETENTION: &str = "dead_retention_idx";
|
||||
/// The events read index (`stream, id` — the offset ASC path).
|
||||
pub const EVENTS_READ: &str = "events_read_idx";
|
||||
/// The schedule due-boundary index (`next_fire_at`).
|
||||
pub const SCHEDULE_FIRE: &str = "schedule_fire_idx";
|
||||
/// The locks expiry index (lazy-acquire sweep support).
|
||||
pub const LOCKS_EXPIRY: &str = "locks_expiry_idx";
|
||||
}
|
||||
|
||||
/// Run the schema bootstrap for one engine-owned schema on `conn`.
|
||||
/// Idempotent.
|
||||
///
|
||||
/// Creates the schema and the full table family + indexes (column
|
||||
/// sets per the module docs) with `CREATE … IF NOT EXISTS` — a second
|
||||
/// run is a no-op. Runnable at open on any connection; the deadpool
|
||||
/// pooled client is a `tokio_postgres::Client` (POC-verified), so a
|
||||
/// pool checkout passes directly.
|
||||
///
|
||||
/// The schema name is engine configuration (default
|
||||
/// [`DEFAULT_SCHEMA`], decided at the `PgOpts` layer); object names
|
||||
/// inside the schema are fixed ([`tables`]).
|
||||
pub async fn bootstrap(conn: &tokio_postgres::Client, schema_name: &str) -> BootstrapResult {
|
||||
conn.batch_execute(&bootstrap_ddl(schema_name))
|
||||
.await
|
||||
.map_err(BootstrapError)
|
||||
}
|
||||
|
||||
fn bootstrap_ddl(schema: &str) -> String {
|
||||
// Index names are bare at the `CREATE INDEX` statement (Postgres
|
||||
// rejects a qualified index name there) but schema-prefixed in
|
||||
// their *text* — the qualification posture that keeps a
|
||||
// co-tenant's same-named index from silently skipping ours.
|
||||
let schema_q = quote_identifier(schema);
|
||||
let index = |suffix: &'static str| quote_identifier(&format!("{schema}_{suffix}"));
|
||||
let job = QualifiedTable {
|
||||
schema,
|
||||
name: tables::JOB,
|
||||
};
|
||||
let dead = QualifiedTable {
|
||||
schema,
|
||||
name: tables::DEAD,
|
||||
};
|
||||
let events = QualifiedTable {
|
||||
schema,
|
||||
name: tables::EVENTS,
|
||||
};
|
||||
let offsets = QualifiedTable {
|
||||
schema,
|
||||
name: tables::OFFSETS,
|
||||
};
|
||||
let schedule = QualifiedTable {
|
||||
schema,
|
||||
name: tables::SCHEDULE,
|
||||
};
|
||||
let locks = QualifiedTable {
|
||||
schema,
|
||||
name: tables::LOCKS,
|
||||
};
|
||||
let job_claim_idx = index(indexes::JOB_CLAIM);
|
||||
let job_pending_idx = index(indexes::JOB_PENDING);
|
||||
let job_deadline_idx = index(indexes::JOB_DEADLINE);
|
||||
let dead_retention_idx = index(indexes::DEAD_RETENTION);
|
||||
let events_read_idx = index(indexes::EVENTS_READ);
|
||||
let schedule_fire_idx = index(indexes::SCHEDULE_FIRE);
|
||||
let locks_expiry_idx = index(indexes::LOCKS_EXPIRY);
|
||||
|
||||
format!(
|
||||
"
|
||||
CREATE SCHEMA IF NOT EXISTS {schema_q};
|
||||
|
||||
-- Job table (live rows, pending + processing; delete-on-ack —
|
||||
-- ADR-010 §1). Columns = Job's field list of record (ADR-019 §3,
|
||||
-- ADR-021 §2) + the stamped QueueOpts (ADR-010 §3a).
|
||||
CREATE TABLE IF NOT EXISTS {job} (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
queue TEXT NOT NULL,
|
||||
state TEXT NOT NULL DEFAULT 'pending',
|
||||
payload BYTEA NOT NULL,
|
||||
priority BIGINT NOT NULL DEFAULT 0,
|
||||
run_at BIGINT NOT NULL,
|
||||
attempts BIGINT NOT NULL DEFAULT 0,
|
||||
max_attempts BIGINT NOT NULL DEFAULT 3,
|
||||
worker_id TEXT,
|
||||
claimed_at BIGINT,
|
||||
claim_expires_at BIGINT,
|
||||
created_at BIGINT NOT NULL,
|
||||
expires_at BIGINT,
|
||||
visibility_timeout_s BIGINT NOT NULL DEFAULT 300,
|
||||
backoff_base_s BIGINT NOT NULL DEFAULT 5,
|
||||
dead_letter_retention_s BIGINT
|
||||
);
|
||||
-- Claim-ordering index (ADR-010 §1's ordering row: priority DESC,
|
||||
-- run_at ASC, enqueue order = id). Partial: dead rows ride the dead
|
||||
-- table, the claim path never scans them.
|
||||
CREATE INDEX IF NOT EXISTS {job_claim_idx} ON {job} (queue, priority DESC, run_at, id)
|
||||
WHERE state IN ('pending', 'processing');
|
||||
CREATE INDEX IF NOT EXISTS {job_pending_idx} ON {job} (queue, run_at)
|
||||
WHERE state = 'pending';
|
||||
CREATE INDEX IF NOT EXISTS {job_deadline_idx} ON {job} (queue, claim_expires_at)
|
||||
WHERE state = 'processing';
|
||||
|
||||
-- Dead table (ADR-010 §4): the job row's diagnosis fields + died_at;
|
||||
-- physical move target, off the claim path. Retention index for the
|
||||
-- sweep's dead-TTL enforcement (the only sweeper dead rows have).
|
||||
CREATE TABLE IF NOT EXISTS {dead} (
|
||||
id BIGINT PRIMARY KEY,
|
||||
queue TEXT NOT NULL,
|
||||
payload BYTEA NOT NULL,
|
||||
priority BIGINT NOT NULL DEFAULT 0,
|
||||
run_at BIGINT NOT NULL DEFAULT 0,
|
||||
attempts BIGINT NOT NULL DEFAULT 0,
|
||||
max_attempts BIGINT NOT NULL DEFAULT 0,
|
||||
worker_id TEXT,
|
||||
claimed_at BIGINT,
|
||||
claim_expires_at BIGINT,
|
||||
created_at BIGINT NOT NULL,
|
||||
expires_at BIGINT,
|
||||
visibility_timeout_s BIGINT NOT NULL DEFAULT 300,
|
||||
backoff_base_s BIGINT NOT NULL DEFAULT 5,
|
||||
dead_letter_retention_s BIGINT,
|
||||
last_error TEXT,
|
||||
died_at BIGINT NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS {dead_retention_idx} ON {dead} (queue, died_at);
|
||||
|
||||
-- Stream events table (ADR-015 §3): `id BIGSERIAL` is the offset —
|
||||
-- engine-assigned, monotone, immutable, never renumbered; global
|
||||
-- FIFO per stream by id ASC; gaps legal after trim/delete. Nullable
|
||||
-- key column = carried metadata (ADR-015 §1). created_at is
|
||||
-- informational.
|
||||
CREATE TABLE IF NOT EXISTS {events} (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
stream TEXT NOT NULL,
|
||||
\"key\" TEXT,
|
||||
payload BYTEA NOT NULL,
|
||||
created_at BIGINT NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS {events_read_idx} ON {events} (stream, id);
|
||||
|
||||
-- Offsets table: (stream, consumer) → offset, the monotone
|
||||
-- checkpoint upsert target (the streams task's SQL owns the
|
||||
-- monotonicity rule, GREATEST-guarded upsert, one owner). The
|
||||
-- non-negative CHECK bounds the domain (cursors start at 0).
|
||||
CREATE TABLE IF NOT EXISTS {offsets} (
|
||||
stream TEXT NOT NULL,
|
||||
consumer TEXT NOT NULL,
|
||||
\"offset\" BIGINT NOT NULL DEFAULT 0 CHECK (\"offset\" >= 0),
|
||||
PRIMARY KEY (stream, consumer)
|
||||
);
|
||||
|
||||
-- Schedule table (ADR-009 §3/§5): name (unique — the PK), the
|
||||
-- `@every` spec, target queue, payload, the ScheduleOpts stamps
|
||||
-- (priority / max_attempts / expires_s) over the resolved QueueOpts
|
||||
-- stamp set (the substrate's __alkstore_scheduler_tasks
|
||||
-- concrete-stamp shape re-owned; `enabled` dropped — the collapse
|
||||
-- surface has no pause/resume), and next_fire_at — the boundary
|
||||
-- column the tick math advances (`next = last + interval`).
|
||||
CREATE TABLE IF NOT EXISTS {schedule} (
|
||||
name TEXT PRIMARY KEY,
|
||||
spec TEXT NOT NULL,
|
||||
queue TEXT NOT NULL,
|
||||
payload BYTEA NOT NULL,
|
||||
priority BIGINT NOT NULL DEFAULT 0,
|
||||
expires_s BIGINT,
|
||||
next_fire_at BIGINT NOT NULL,
|
||||
max_attempts BIGINT NOT NULL DEFAULT 3,
|
||||
visibility_timeout_s BIGINT NOT NULL DEFAULT 300,
|
||||
backoff_base_s BIGINT NOT NULL DEFAULT 5,
|
||||
dead_letter_retention_s BIGINT
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS {schedule_fire_idx} ON {schedule} (next_fire_at);
|
||||
|
||||
-- Locks table: the TTL lock machinery's storage (the substrate's
|
||||
-- lock-table semantics re-derived: PK on name prevents dual
|
||||
-- acquisition; expiry lazy — no sweeper, an expired lock is
|
||||
-- acquirable; the locks task's ops own the SQL).
|
||||
CREATE TABLE IF NOT EXISTS {locks} (
|
||||
name TEXT PRIMARY KEY,
|
||||
owner TEXT NOT NULL,
|
||||
expires_at BIGINT NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS {locks_expiry_idx} ON {locks} (expires_at);
|
||||
",
|
||||
schema_q = schema_q,
|
||||
job = job,
|
||||
job_claim_idx = job_claim_idx,
|
||||
job_pending_idx = job_pending_idx,
|
||||
job_deadline_idx = job_deadline_idx,
|
||||
dead = dead,
|
||||
dead_retention_idx = dead_retention_idx,
|
||||
events = events,
|
||||
events_read_idx = events_read_idx,
|
||||
offsets = offsets,
|
||||
schedule = schedule,
|
||||
schedule_fire_idx = schedule_fire_idx,
|
||||
locks = locks,
|
||||
locks_expiry_idx = locks_expiry_idx,
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,483 @@
|
||||
//! Schema bootstrap tests against the harness server
|
||||
//! (`docs/plans/implementation.md`'s test posture: dockerized
|
||||
//! `postgres:16-alpine` on :15432 — connection settings ride the
|
||||
//! environment, never hardcoded; tests without a reachable server
|
||||
//! skip cleanly so the workspace gates stay green server-less).
|
||||
//!
|
||||
//! A fresh schema per test (unique name per test run) is the
|
||||
//! isolation guarantee — the POC's shared-server
|
||||
//! parallel-interference caveat is answered by schema-per-test
|
||||
//! isolation, not sequential-only harnesses.
|
||||
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
use alkstore_postgres::{bootstrap, quote_identifier, tables};
|
||||
|
||||
const ENV_HOST: &str = "ALKSTORE_PG_HOST";
|
||||
const ENV_PORT: &str = "ALKSTORE_PG_PORT";
|
||||
const ENV_USER: &str = "ALKSTORE_PG_USER";
|
||||
const ENV_PASSWORD: &str = "ALKSTORE_PG_PASSWORD";
|
||||
const ENV_DB: &str = "ALKSTORE_PG_DB";
|
||||
|
||||
fn harness_dsn() -> Option<String> {
|
||||
let host = std::env::var(ENV_HOST).ok()?;
|
||||
let port: u16 = std::env::var(ENV_PORT).ok()?.parse().ok()?;
|
||||
let user = std::env::var(ENV_USER).ok()?;
|
||||
let password = std::env::var(ENV_PASSWORD).ok()?;
|
||||
let db = std::env::var(ENV_DB).unwrap_or_else(|_| "postgres".to_string());
|
||||
Some(format!(
|
||||
"host={host} port={port} user={user} password={password} dbname={db}"
|
||||
))
|
||||
}
|
||||
|
||||
async fn harness_client() -> Option<tokio_postgres::Client> {
|
||||
let (client, connection) = tokio_postgres::connect(&harness_dsn()?, tokio_postgres::NoTls)
|
||||
.await
|
||||
.ok()?;
|
||||
tokio::spawn(async move {
|
||||
let _ = connection.await;
|
||||
});
|
||||
Some(client)
|
||||
}
|
||||
|
||||
fn instance_namer(tag: &str) -> impl Fn() -> String + use<'_> {
|
||||
let counter = AtomicU64::new(0);
|
||||
move || {
|
||||
format!(
|
||||
"{tag}_{}_{}_{}",
|
||||
std::process::id(),
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_nanos(),
|
||||
counter.fetch_add(1, Ordering::SeqCst),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
async fn drop_schema(client: &tokio_postgres::Client, schema: &str) {
|
||||
let sql = format!("DROP SCHEMA IF EXISTS {} CASCADE", quote_identifier(schema));
|
||||
let _ = client.batch_execute(&sql).await;
|
||||
}
|
||||
|
||||
async fn table_names(client: &tokio_postgres::Client, schema: &str) -> Vec<String> {
|
||||
client
|
||||
.query(
|
||||
"SELECT tablename FROM pg_tables WHERE schemaname = $1 ORDER BY tablename",
|
||||
&[&schema],
|
||||
)
|
||||
.await
|
||||
.unwrap_or_default()
|
||||
.iter()
|
||||
.map(|r| r.get::<_, String>(0))
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn column_names(client: &tokio_postgres::Client, schema: &str, table: &str) -> Vec<String> {
|
||||
let sql = "SELECT column_name FROM information_schema.columns
|
||||
WHERE table_schema = $1 AND table_name = $2
|
||||
ORDER BY ordinal_position"
|
||||
.to_string();
|
||||
client
|
||||
.query(&sql, &[&schema, &table])
|
||||
.await
|
||||
.unwrap_or_default()
|
||||
.iter()
|
||||
.map(|r| r.get::<_, String>(0))
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn index_names(client: &tokio_postgres::Client, schema: &str) -> Vec<String> {
|
||||
client
|
||||
.query(
|
||||
"SELECT indexname FROM pg_indexes WHERE schemaname = $1 ORDER BY indexname",
|
||||
&[&schema],
|
||||
)
|
||||
.await
|
||||
.unwrap_or_default()
|
||||
.iter()
|
||||
.map(|r| r.get::<_, String>(0))
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn bootstrap_is_idempotent() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("idempotent")();
|
||||
assert!(bootstrap(&client, &schema).await.is_ok());
|
||||
assert!(bootstrap(&client, &schema).await.is_ok());
|
||||
third_run_converges(&client, &schema).await;
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
async fn third_run_converges(client: &tokio_postgres::Client, schema: &str) {
|
||||
assert!(bootstrap(client, schema).await.is_ok());
|
||||
let names = table_names(client, schema).await;
|
||||
assert_eq!(names.len(), 6, "expected the 6-table family, saw {names:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn table_family_exists_with_pinned_shapes() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("shapes")();
|
||||
bootstrap(&client, &schema).await.unwrap();
|
||||
|
||||
// Job: the field list of record (Job::from_row, ADR-019 §3 /
|
||||
// ADR-021 §2) + the stamped QueueOpts (ADR-010 §3a).
|
||||
assert_eq!(
|
||||
column_names(&client, &schema, tables::JOB).await,
|
||||
vec![
|
||||
"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"
|
||||
]
|
||||
);
|
||||
|
||||
// Dead: the job row's diagnosis fields + died_at.
|
||||
assert_eq!(
|
||||
column_names(&client, &schema, tables::DEAD).await,
|
||||
vec![
|
||||
"id",
|
||||
"queue",
|
||||
"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"
|
||||
]
|
||||
);
|
||||
|
||||
// Events: bigserial offset, nullable key, payload, created_at.
|
||||
assert_eq!(
|
||||
column_names(&client, &schema, tables::EVENTS).await,
|
||||
vec!["id", "stream", "key", "payload", "created_at"]
|
||||
);
|
||||
|
||||
// Offsets: (stream, consumer) → offset, the checkpoint.
|
||||
assert_eq!(
|
||||
column_names(&client, &schema, tables::OFFSETS).await,
|
||||
vec!["stream", "consumer", "offset"]
|
||||
);
|
||||
|
||||
// Schedule: the substrate's __alkstore_scheduler_tasks shape
|
||||
// re-owned (no `enabled`).
|
||||
assert_eq!(
|
||||
column_names(&client, &schema, tables::SCHEDULE).await,
|
||||
vec![
|
||||
"name",
|
||||
"spec",
|
||||
"queue",
|
||||
"payload",
|
||||
"priority",
|
||||
"expires_s",
|
||||
"next_fire_at",
|
||||
"max_attempts",
|
||||
"visibility_timeout_s",
|
||||
"backoff_base_s",
|
||||
"dead_letter_retention_s"
|
||||
]
|
||||
);
|
||||
|
||||
// Locks: name (unique — PK), owner, expiry TTL.
|
||||
assert_eq!(
|
||||
column_names(&client, &schema, tables::LOCKS).await,
|
||||
vec!["name", "owner", "expires_at"]
|
||||
);
|
||||
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn payload_columns_are_bytea() {
|
||||
let expected_payload_type = "bytea";
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("bytea")();
|
||||
bootstrap(&client, &schema).await.unwrap();
|
||||
|
||||
for table in [tables::JOB, tables::DEAD, tables::EVENTS, tables::SCHEDULE] {
|
||||
let table = table.to_string();
|
||||
let sql = "SELECT data_type FROM information_schema.columns
|
||||
WHERE table_schema = $1 AND table_name = $2 AND column_name = 'payload'"
|
||||
.to_string();
|
||||
let rows = client.query(&sql, &[&schema, &table]).await.unwrap();
|
||||
assert_eq!(rows.len(), 1, "{table}: payload column missing");
|
||||
let ty: String = rows[0].get(0);
|
||||
assert_eq!(ty, expected_payload_type, "{table}: payload not bytea");
|
||||
}
|
||||
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn bigserial_offset_is_monotone_with_gaps_on_delete() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("offsets")();
|
||||
bootstrap(&client, &schema).await.unwrap();
|
||||
|
||||
let mut events: Vec<i64> = Vec::new();
|
||||
for _ in 0..3 {
|
||||
let row = client
|
||||
.query_one(
|
||||
&format!(
|
||||
"INSERT INTO {} (stream, key, payload, created_at)
|
||||
VALUES ('s', NULL, '\\x7b7d', 1) RETURNING id",
|
||||
alkstore_postgres::QualifiedTable {
|
||||
schema: &schema,
|
||||
name: tables::EVENTS,
|
||||
}
|
||||
),
|
||||
&[],
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
events.push(row.get(0));
|
||||
}
|
||||
assert_eq!(events, vec![1, 2, 3], "bigserial starts at 1, step 1");
|
||||
|
||||
client
|
||||
.execute(
|
||||
&format!(
|
||||
"DELETE FROM {} WHERE id = $1",
|
||||
alkstore_postgres::QualifiedTable {
|
||||
schema: &schema,
|
||||
name: tables::EVENTS,
|
||||
}
|
||||
),
|
||||
&[&events[1]],
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let row = client
|
||||
.query_one(
|
||||
&format!(
|
||||
"INSERT INTO {} (stream, key, payload, created_at)
|
||||
VALUES ('s', NULL, '\\x7b7d', 1) RETURNING id",
|
||||
alkstore_postgres::QualifiedTable {
|
||||
schema: &schema,
|
||||
name: tables::EVENTS,
|
||||
}
|
||||
),
|
||||
&[],
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let after_delete: i64 = row.get(0);
|
||||
assert_eq!(
|
||||
after_delete, 4,
|
||||
"the offset never renumbers — deletion leaves a gap, the sequence advances"
|
||||
);
|
||||
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn claim_ordering_index_exists() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("idx")();
|
||||
bootstrap(&client, &schema).await.unwrap();
|
||||
|
||||
let indexes = index_names(&client, &schema).await;
|
||||
let prefixed = |suffix: &str| format!("{schema}_{suffix}");
|
||||
let contains = |needle: String| indexes.iter().any(|i| i == &needle);
|
||||
assert!(
|
||||
contains(prefixed("job_claim_idx")),
|
||||
"claim-ordering index missing; saw {indexes:?}"
|
||||
);
|
||||
assert!(
|
||||
contains(prefixed("job_pending_idx")) && contains(prefixed("job_deadline_idx")),
|
||||
"hot-path partial indexes missing; saw {indexes:?}"
|
||||
);
|
||||
assert!(
|
||||
contains(prefixed("dead_retention_idx")),
|
||||
"dead retention index missing; saw {indexes:?}"
|
||||
);
|
||||
assert!(
|
||||
contains(prefixed("events_read_idx")),
|
||||
"events (stream, id) index missing; saw {indexes:?}"
|
||||
);
|
||||
assert!(
|
||||
contains(prefixed("schedule_fire_idx")) && contains(prefixed("locks_expiry_idx")),
|
||||
"schedule/locks indexes missing; saw {indexes:?}"
|
||||
);
|
||||
|
||||
// The ordering row's column set + directionality (priority DESC,
|
||||
// run_at ASC, enqueue order = id) — pg_indexes' indexdef carries
|
||||
// it.
|
||||
let def = client
|
||||
.query_one(
|
||||
"SELECT indexdef FROM pg_indexes
|
||||
WHERE schemaname = $1 AND indexname = $2",
|
||||
&[&schema, &prefixed("job_claim_idx")],
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.get::<_, String>(0);
|
||||
for fragment in ["queue", "priority DESC", "run_at", "id"] {
|
||||
assert!(
|
||||
def.contains(fragment),
|
||||
"claim index def missing {fragment}: {def}"
|
||||
);
|
||||
}
|
||||
assert!(
|
||||
def.contains("pending"),
|
||||
"claim index not partial over the live states: {def}"
|
||||
);
|
||||
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn schema_name_is_a_parameter() {
|
||||
let expected_default = "alkstore";
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
assert_eq!(alkstore_postgres::DEFAULT_SCHEMA, expected_default);
|
||||
|
||||
// A non-default schema name lands its objects under that name
|
||||
// (quoted, reserved-word-safe) — the PgOpts layer's knob works.
|
||||
let schema = "select";
|
||||
assert!(bootstrap(&client, schema).await.is_ok());
|
||||
let names = table_names(&client, schema).await;
|
||||
assert_eq!(
|
||||
names.len(),
|
||||
6,
|
||||
"reserved-word schema name quoted and populated"
|
||||
);
|
||||
drop_schema(&client, schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn bootstrap_is_atomic_in_the_schema_scope() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
let schema = instance_namer("atomic")();
|
||||
assert!(bootstrap(&client, &schema).await.is_ok());
|
||||
|
||||
// A consumer table co-temporaries the instance untouched by a
|
||||
// re-bootstrap (the co-tenancy posture — engine DDL never
|
||||
// touches non-engine objects in other scopes).
|
||||
client
|
||||
.execute(
|
||||
&format!(
|
||||
"CREATE TABLE {}.consumer_thing (x INT)",
|
||||
quote_identifier(&schema)
|
||||
),
|
||||
&[],
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(bootstrap(&client, &schema).await.is_ok());
|
||||
let names = table_names(&client, &schema).await;
|
||||
assert_eq!(
|
||||
names.len(),
|
||||
7,
|
||||
"engine re-bootstrap left the co-tenant intact"
|
||||
);
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn co_tenant_index_name_never_blocks_the_bootstrap() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
// A co-tenant carries an index with the bare name the engine's
|
||||
// bootstrap would use — `CREATE INDEX IF NOT EXISTS` checks
|
||||
// existence across all schemas, so an un-prefixed engine index
|
||||
// name would be silently skipped here (the defect class the
|
||||
// schema-prefixed index naming exists to exclude).
|
||||
client
|
||||
.execute("CREATE TABLE public.co_tenant (x INT)", &[])
|
||||
.await
|
||||
.unwrap();
|
||||
client
|
||||
.execute("CREATE INDEX job_claim_idx ON public.co_tenant (x)", &[])
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let schema = instance_namer("cotenant")();
|
||||
bootstrap(&client, &schema).await.unwrap();
|
||||
let indexes = index_names(&client, &schema).await;
|
||||
assert!(
|
||||
indexes
|
||||
.iter()
|
||||
.any(|i| i == &format!("{schema}_job_claim_idx")),
|
||||
"co-tenant's bare-name index must not skip the engine's; saw {indexes:?}"
|
||||
);
|
||||
|
||||
client
|
||||
.execute("DROP INDEX public.job_claim_idx", &[])
|
||||
.await
|
||||
.unwrap();
|
||||
client
|
||||
.execute("DROP TABLE public.co_tenant", &[])
|
||||
.await
|
||||
.unwrap();
|
||||
drop_schema(&client, &schema).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn custom_schema_quoting_rejects_injection() {
|
||||
let Some(client) = harness_client().await else {
|
||||
eprintln!("skip: no harness server");
|
||||
return;
|
||||
};
|
||||
// An identifier containing a double-quote must be inert SQL text
|
||||
// (doubled) — never a breakout into a second object.
|
||||
let hostile = "alk\"; DROP SCHEMA public; CREATE SCHEMA sneaky; --";
|
||||
assert!(bootstrap(&client, hostile).await.is_ok());
|
||||
let names = table_names(&client, hostile).await;
|
||||
assert_eq!(
|
||||
names.len(),
|
||||
6,
|
||||
"hostile-named schema populated exactly once"
|
||||
);
|
||||
let sneaky = table_names(&client, "sneaky").await;
|
||||
assert!(sneaky.is_empty(), "injection created a second schema");
|
||||
drop_schema(&client, hostile).await;
|
||||
drop_schema(&client, "sneaky").await;
|
||||
}
|
||||
@@ -1,7 +1,7 @@
|
||||
---
|
||||
id: pg-engine-schema
|
||||
name: Postgres engine — schema bootstrap (engine-owned schema, table family)
|
||||
status: pending
|
||||
status: completed
|
||||
depends_on: []
|
||||
scope: narrow
|
||||
risk: low
|
||||
@@ -79,17 +79,18 @@ index exists.
|
||||
|
||||
## Acceptance Criteria
|
||||
|
||||
- [ ] `schema.rs` with an idempotent `bootstrap(conn, schema_name)`
|
||||
- [x] `schema.rs` with an idempotent `bootstrap(conn, schema_name)`
|
||||
creating the full table family + indexes in the one schema
|
||||
- [ ] Column sets match the contract's value shapes (the constructors'
|
||||
- [x] Column sets match the contract's value shapes (the constructors'
|
||||
field lists; stamp columns per ADR-010 §3a; `claimed_at` per
|
||||
ADR-021 §2)
|
||||
- [ ] Payload columns are `bytea` (exact-bytes posture, JSONB
|
||||
- [x] Payload columns are `bytea` (exact-bytes posture, JSONB
|
||||
rejected — documented in the module docs with the reason)
|
||||
- [ ] Schema name is a parameter (default `alkstore` decided at the
|
||||
- [x] Schema name is a parameter (default `alkstore` decided at the
|
||||
`PgOpts` layer, next task)
|
||||
- [ ] Idempotence + table-shape tests green against the harness server
|
||||
- [ ] `cargo test -p alkstore-postgres`, clippy `-D warnings`, fmt clean
|
||||
- [x] Indexes: idempotence + table-shape tests green against the
|
||||
harness server
|
||||
- [x] `cargo test -p alkstore-postgres`, clippy `-D warnings`, fmt clean
|
||||
|
||||
## References
|
||||
|
||||
@@ -103,8 +104,80 @@ index exists.
|
||||
|
||||
## Notes
|
||||
|
||||
> To be filled by implementation agent
|
||||
Decisions of record made while implementing (the description didn't
|
||||
pin them):
|
||||
|
||||
- **Table names**: `job`, `dead`, `stream_events`, `stream_offsets`,
|
||||
`schedule`, `locks` — fixed constants exported as
|
||||
`alkstore_postgres::tables` for the later mechanism tasks' SQL.
|
||||
Table names do not ride configuration; only the schema name does.
|
||||
- **Index-name qualification**: Postgres rejects schema-qualified
|
||||
index names at `CREATE INDEX`, but `IF NOT EXISTS` on a *bare* name
|
||||
checks existence across all schemas — a co-tenant's unrelated
|
||||
`public.job_claim_idx` would silently skip our index creation. So
|
||||
every engine index is named `{schema}_{table}_{suffix}`
|
||||
(`job_claim_idx`, `job_pending_idx`, `job_deadline_idx`,
|
||||
`dead_retention_idx`, `events_read_idx`, `schedule_fire_idx`,
|
||||
`locks_expiry_idx`, exported as `indexes::`) — schema-prefixed and
|
||||
quoted, tested against a co-tenant carrying the bare name.
|
||||
Sub-HOT-path indexes (`job_pending`, `job_deadline`, the substrate's
|
||||
trio) carried over from the reference shapes — the claim/heartbeat
|
||||
paths the queue task re-derives ride them.
|
||||
- **Reserved-word columns**: the offsets checkpoint column is
|
||||
`"offset"` and the events key column `"key"` — both quoted here and
|
||||
flagged in the module docs so every later task's SQL keeps the
|
||||
quotes. The offsets column carries a non-negative CHECK (cursors
|
||||
start at 0); monotonicity itself stays in the streams task's
|
||||
GREATEST-guarded upsert.
|
||||
- **Identifier quoting**: `quote_identifier` (doubles embedded `"`)
|
||||
is the engine's injection boundary for the consumer-supplied schema
|
||||
name, exported for later tasks; tested hostile (SQL-breakout schema
|
||||
name stays inert single-schema DDL).
|
||||
- **Schedule columns** anchor to the substrate's
|
||||
`__alkstore_scheduler_tasks` concrete-stamp shape, `enabled` dropped
|
||||
(no pause/resume in the collapse surface): `name` PK, `spec`,
|
||||
`queue`, `payload`, `priority`, `expires_s`, `next_fire_at` (the
|
||||
tick's boundary column), `max_attempts` (default 3),
|
||||
`visibility_timeout_s` (default 300), `backoff_base_s` (default 5),
|
||||
`dead_letter_retention_s`. No `last_fire_at` (`next_fire_at` is the
|
||||
only state the tick math needs).
|
||||
- **Dead table** mirrors the job table's live columns + the diagnosis
|
||||
fields (`last_error`, `died_at NOT NULL`), `id` a plain BIGINT PK
|
||||
(carried from the move, not a sequence); retention index
|
||||
`(queue, died_at)`.
|
||||
- **Job table** = the `Job::from_row` field list of record verbatim +
|
||||
the five stamp columns; `state TEXT DEFAULT 'pending'` (SQLite
|
||||
parity — the claim predicates read the text values); three indexes
|
||||
(claim-ordering partial `WHERE state IN ('pending','processing')`
|
||||
per ADR-010 §1's ordering row, pending `(queue, run_at)`, processing
|
||||
deadline `(queue, claim_expires_at)`).
|
||||
- **`thiserror` added to the crate** (v2, same as core/SQLite) for
|
||||
`BootstrapError`; the `open` task maps it onto `Error::Database`
|
||||
next task.
|
||||
- **Harness convention** (ride env, never hardcoded, skip-clean
|
||||
server-less): `ALKSTORE_PG_HOST/PORT/USER/PASSWORD/DB` — host
|
||||
127.0.0.1, port 15432, user/password `postgres`/`poc`, db `blobs`
|
||||
(the `pglo-poc` container's). The next tasks' tests should reuse
|
||||
these variable names. `bootstrap` itself takes any
|
||||
`&tokio_postgres::Client` (NoTls in tests; TLS is the `open` task's
|
||||
concern via deadpool config, not this module's).
|
||||
|
||||
## Summary
|
||||
|
||||
> To be filled on completion
|
||||
Landed: `alkstore-postgres/src/schema.rs` — the engine-owned schema
|
||||
module (`bootstrap(conn, schema_name)` idempotent DDL over
|
||||
`batch_execute`; `DEFAULT_SCHEMA = "alkstore"`; `quote_identifier`;
|
||||
`tables`/`indexes` name constants; `BootstrapError`), wired into
|
||||
`lib.rs` re-exports, plus `tests/schema_tests.rs` (9 tests) driven
|
||||
against the harness server via env-carried DSN, skipping cleanly
|
||||
server-less. Table family per the acceptance criteria (job/dead/
|
||||
stream_events/stream_offsets/schedule/locks, bytea payloads, bigserial
|
||||
offsets, claim-ordering partial index with the pinned ordering row).
|
||||
Verified: 9/9 green against the harness server (idempotence ×3 runs,
|
||||
per-table column sets, bytea types, bigserial monotone + gap-on-delete,
|
||||
claim index def fragments incl. `priority DESC` + partiality, schema
|
||||
name as parameter incl. reserved-word name `select`, co-tenancy
|
||||
re-bootstrap, co-tenant bare-index-name collision exclusion, hostile
|
||||
schema-name quoting); workspace `cargo test` green server-less (236
|
||||
tests, pg tests skipped per convention); clippy `-D warnings` and fmt
|
||||
clean workspace-wide.
|
||||
Reference in new issue
Block a user