diff --git a/Cargo.lock b/Cargo.lock index 4f29c10..4474959 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -28,6 +28,7 @@ dependencies = [ "alkstore", "alkstore-contract-suite", "deadpool-postgres", + "thiserror", "tokio", "tokio-postgres", ] diff --git a/alkstore-postgres/Cargo.toml b/alkstore-postgres/Cargo.toml index b21c47b..2518a8a 100644 --- a/alkstore-postgres/Cargo.toml +++ b/alkstore-postgres/Cargo.toml @@ -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" } \ No newline at end of file diff --git a/alkstore-postgres/src/lib.rs b/alkstore-postgres/src/lib.rs index 8b13789..f4e0ed1 100644 --- a/alkstore-postgres/src/lib.rs +++ b/alkstore-postgres/src/lib.rs @@ -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, +}; diff --git a/alkstore-postgres/src/schema.rs b/alkstore-postgres/src/schema.rs new file mode 100644 index 0000000..6e52d6b --- /dev/null +++ b/alkstore-postgres/src/schema.rs @@ -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, + ) +} diff --git a/alkstore-postgres/tests/schema_tests.rs b/alkstore-postgres/tests/schema_tests.rs new file mode 100644 index 0000000..9463a0f --- /dev/null +++ b/alkstore-postgres/tests/schema_tests.rs @@ -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 { + 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 { + 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 { + 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 { + 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 { + 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 = 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; +} diff --git a/tasks/pg-engine-schema.md b/tasks/pg-engine-schema.md index c77925d..453b5d6 100644 --- a/tasks/pg-engine-schema.md +++ b/tasks/pg-engine-schema.md @@ -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 \ No newline at end of file +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. \ No newline at end of file