v0.3.3-pre.003

This commit is contained in:
2026-08-30 09:22:53 +02:00
parent 9712c7e1f7
commit 61bf7ba468
17 changed files with 799 additions and 107 deletions

View File

@@ -0,0 +1,133 @@
CREATE TABLE ksp_store_identity (
singleton SMALLINT PRIMARY KEY,
network TEXT NOT NULL,
CONSTRAINT ck_ksp_store_identity_singleton CHECK (singleton = 1),
CONSTRAINT ck_ksp_store_identity_network CHECK (
octet_length(network) BETWEEN 1 AND 128
AND network ~ '^[A-Za-z0-9_.:-]+$'
)
);
CREATE TABLE ksp_raw_transactions (
signature BYTEA PRIMARY KEY,
slot NUMERIC(20, 0) NOT NULL,
block_time_unix_millis BIGINT NULL,
format_id TEXT NOT NULL,
format_version BIGINT NOT NULL,
content_hash BYTEA NOT NULL,
payload BYTEA NULL,
retention_state TEXT NOT NULL,
CONSTRAINT ck_ksp_raw_transactions_signature CHECK (octet_length(signature) = 64),
CONSTRAINT ck_ksp_raw_transactions_slot CHECK (slot BETWEEN 0 AND 18446744073709551615),
CONSTRAINT ck_ksp_raw_transactions_block_time CHECK (
block_time_unix_millis IS NULL
OR block_time_unix_millis BETWEEN 0 AND 253402300799999
),
CONSTRAINT ck_ksp_raw_transactions_format_id CHECK (
octet_length(format_id) BETWEEN 1 AND 128
AND format_id ~ '^[A-Za-z0-9_.:-]+$'
),
CONSTRAINT ck_ksp_raw_transactions_format_version CHECK (format_version BETWEEN 1 AND 4294967295),
CONSTRAINT ck_ksp_raw_transactions_content_hash CHECK (octet_length(content_hash) = 32),
CONSTRAINT ck_ksp_raw_transactions_payload CHECK (
payload IS NULL
OR octet_length(payload) BETWEEN 1 AND 16777216
),
CONSTRAINT ck_ksp_raw_transactions_retention_state CHECK (retention_state IN ('full', 'archived', 'purged')),
CONSTRAINT ck_ksp_raw_transactions_payload_state CHECK (
(retention_state = 'full' AND payload IS NOT NULL)
OR (retention_state IN ('archived', 'purged') AND payload IS NULL)
),
CONSTRAINT ck_ksp_raw_transactions_purged_block_time CHECK (
retention_state <> 'purged'
OR block_time_unix_millis IS NULL
)
);
CREATE TABLE ksp_raw_transaction_observations (
observation_key BYTEA PRIMARY KEY,
transaction_signature BYTEA NOT NULL REFERENCES ksp_raw_transactions(signature) ON DELETE RESTRICT,
provider TEXT NOT NULL,
protocol TEXT NOT NULL,
acquisition_method TEXT NOT NULL,
origin TEXT NOT NULL,
received_at_unix_millis BIGINT NOT NULL,
capture_session_id TEXT NULL,
commitment TEXT NULL,
endpoint_id TEXT NULL,
filter_id TEXT NULL,
observed_at_unix_millis BIGINT NULL,
source_payload_hash BYTEA NULL,
source_payload_size_bytes BIGINT NULL,
CONSTRAINT ck_ksp_raw_transaction_observations_key CHECK (octet_length(observation_key) = 32),
CONSTRAINT ck_ksp_raw_transaction_observations_signature CHECK (octet_length(transaction_signature) = 64),
CONSTRAINT ck_ksp_raw_transaction_observations_provider CHECK (
octet_length(provider) BETWEEN 1 AND 128
AND provider ~ '^[A-Za-z0-9_.:-]+$'
),
CONSTRAINT ck_ksp_raw_transaction_observations_protocol CHECK (
octet_length(protocol) BETWEEN 1 AND 128
AND protocol ~ '^[A-Za-z0-9_.:-]+$'
),
CONSTRAINT ck_ksp_raw_transaction_observations_method CHECK (
octet_length(acquisition_method) BETWEEN 1 AND 128
AND acquisition_method ~ '^[A-Za-z0-9_.:-]+$'
),
CONSTRAINT ck_ksp_raw_transaction_observations_origin CHECK (origin IN ('backfill', 'import', 'live', 'repair', 'replay')),
CONSTRAINT ck_ksp_raw_transaction_observations_received_at CHECK (received_at_unix_millis BETWEEN 0 AND 253402300799999),
CONSTRAINT ck_ksp_raw_transaction_observations_capture_session CHECK (
capture_session_id IS NULL
OR (
octet_length(capture_session_id) BETWEEN 1 AND 128
AND capture_session_id ~ '^[A-Za-z0-9_.:-]+$'
)
),
CONSTRAINT ck_ksp_raw_transaction_observations_commitment CHECK (
commitment IS NULL
OR (
octet_length(commitment) BETWEEN 1 AND 128
AND commitment ~ '^[A-Za-z0-9_.:-]+$'
)
),
CONSTRAINT ck_ksp_raw_transaction_observations_endpoint CHECK (
endpoint_id IS NULL
OR (
octet_length(endpoint_id) BETWEEN 1 AND 128
AND endpoint_id ~ '^[A-Za-z0-9_.:-]+$'
)
),
CONSTRAINT ck_ksp_raw_transaction_observations_filter CHECK (
filter_id IS NULL
OR (
octet_length(filter_id) BETWEEN 1 AND 128
AND filter_id ~ '^[A-Za-z0-9_.:-]+$'
)
),
CONSTRAINT ck_ksp_raw_transaction_observations_observed_at CHECK (
observed_at_unix_millis IS NULL
OR observed_at_unix_millis BETWEEN 0 AND 253402300799999
),
CONSTRAINT ck_ksp_raw_transaction_observations_time_order CHECK (
observed_at_unix_millis IS NULL
OR observed_at_unix_millis <= received_at_unix_millis
),
CONSTRAINT ck_ksp_raw_transaction_observations_source_hash CHECK (
source_payload_hash IS NULL
OR octet_length(source_payload_hash) = 32
),
CONSTRAINT ck_ksp_raw_transaction_observations_source_size CHECK (
source_payload_size_bytes IS NULL
OR source_payload_size_bytes BETWEEN 0 AND 67108864
)
);
CREATE TABLE ksp_raw_transaction_archive_payloads (
signature BYTEA PRIMARY KEY REFERENCES ksp_raw_transactions(signature) ON DELETE RESTRICT,
payload BYTEA NOT NULL,
CONSTRAINT ck_ksp_raw_transaction_archive_payloads_signature CHECK (octet_length(signature) = 64),
CONSTRAINT ck_ksp_raw_transaction_archive_payloads_payload CHECK (octet_length(payload) BETWEEN 1 AND 16777216)
);
CREATE INDEX ix_ksp_raw_transactions_slot_signature
ON ksp_raw_transactions (slot, signature)
WHERE retention_state <> 'purged';

View File

@@ -1,5 +1,9 @@
// file: crates/ksp-store-postgres-lib/src/error.rs
// version: 3
// version: 4
/// Stable KSP error code reserved for PostgreSQL retention transitions that require unsupported physical compaction.
pub const ERROR_CODE_POSTGRES_RETENTION_COMPACTION_UNSUPPORTED: ksp_store_api::ErrorCode =
ksp_store_api::ErrorCode::new("store", "postgres_retention_compaction_unsupported");
/// Safe backend-local classification used by the Store facade for stable error mapping.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/src/lib.rs
// version: 5
// version: 6
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -7,10 +7,11 @@
//! Official PostgreSQL backend implementation for KSP Store.
//!
//! `0.3.2-pre.007` owns the physical `tokio-postgres` connection, bounded
//! Deadpool pool, explicit Rustls TLS policy, private KSP migration/bootstrap
//! engine and safe lightweight health/readiness probe. Business persistence
//! remains absent from this foundation release.
//! The backend owns the physical `tokio-postgres` connection, bounded Deadpool
//! pool, explicit Rustls TLS policy, private KSP migration/bootstrap engine and
//! safe lightweight health/readiness probe. `0.3.3-pre.003` adds the immutable
//! V001 RawTransaction physical schema and mono-network database binding; the
//! business capability implementations remain deferred to later prereleases.
//!
//! This crate depends on `ksp-store-api` and never on `ksp-store-lib`. The
//! common facade consumes only this crate's narrow backend bridge and never
@@ -22,6 +23,8 @@ mod health;
mod migration;
mod runtime;
/// Stable KSP error code for unsupported PostgreSQL retention compaction.
pub use self::error::ERROR_CODE_POSTGRES_RETENTION_COMPACTION_UNSUPPORTED;
/// Safe backend-local error returned to the common Store facade.
pub use self::error::PostgresBackendError;
/// Safe backend-local error classification used by the common Store facade.

View File

@@ -1,18 +1,23 @@
// file: crates/ksp-store-postgres-lib/src/migration.rs
// version: 3
// version: 4
use sha2::Digest; // rust-rules: trait-import
const ADVISORY_LOCK_KEY: i64 = 0x4b53_5053_544f_5245;
const EMBEDDED_MIGRATIONS: &[EmbeddedMigration] = &[EmbeddedMigration {
hook: MigrationHook::None,
name: "bootstrap",
sql: include_str!("../migrations/V000__bootstrap.sql"),
version: 0,
}];
const EMBEDDED_MIGRATIONS: &[EmbeddedMigration] = &[
EmbeddedMigration { hook: MigrationHook::None, name: "bootstrap", sql: include_str!("../migrations/V000__bootstrap.sql"), version: 0 },
EmbeddedMigration {
hook: MigrationHook::StoreIdentity,
name: "raw_transaction",
sql: include_str!("../migrations/V001__raw_transaction.sql"),
version: 1,
},
];
const HEX_LOWER: &[u8; 16] = b"0123456789abcdef";
const HISTORY_INSERT_SQL: &str = "INSERT INTO ksp_store_schema_migrations (version, name, checksum, applied_at) VALUES ($1, $2, $3, CURRENT_TIMESTAMP)";
const HISTORY_LOAD_SQL: &str = "SELECT version, name, checksum FROM ksp_store_schema_migrations ORDER BY version";
const IDENTITY_INSERT_SQL: &str = "INSERT INTO ksp_store_identity (singleton, network) VALUES (1, $1)";
const IDENTITY_LOAD_SQL: &str = "SELECT singleton, network FROM ksp_store_identity ORDER BY singleton LIMIT 2";
const LOCK_POLL_INTERVAL_MS: u64 = 25;
const METADATA_EXISTS_SQL: &str = r#"SELECT EXISTS (
SELECT 1 FROM information_schema.tables
@@ -54,6 +59,7 @@ struct EmbeddedMigration {
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum MigrationHook {
None,
StoreIdentity,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -80,7 +86,11 @@ pub(crate) async fn bootstrap(
if let std::result::Result::Err(error) = registry_result {
return std::result::Result::Err(error);
}
let bounded = tokio::time::timeout(migration_timeout, bootstrap_inner(client, network, auto_migrate, migration_timeout, migration_lock_timeout)).await;
let bounded = tokio::time::timeout(
migration_timeout,
bootstrap_inner(client, network, auto_migrate, migration_timeout, migration_lock_timeout),
)
.await;
return match bounded {
std::result::Result::Ok(result) => result,
std::result::Result::Err(_) => {
@@ -310,16 +320,86 @@ async fn run_applied_migration_hooks(
}
async fn run_migration_hook(
_transaction: &deadpool_postgres::Transaction<'_>,
_network: &ksp_store_api::RawNetworkId,
transaction: &deadpool_postgres::Transaction<'_>,
network: &ksp_store_api::RawNetworkId,
hook: MigrationHook,
_context: MigrationHookContext,
context: MigrationHookContext,
) -> std::result::Result<(), crate::PostgresBackendError> {
return match hook {
MigrationHook::None => std::result::Result::Ok(()),
MigrationHook::StoreIdentity => bind_store_identity(transaction, network, context).await,
};
}
async fn bind_store_identity(
transaction: &deadpool_postgres::Transaction<'_>,
network: &ksp_store_api::RawNetworkId,
context: MigrationHookContext,
) -> std::result::Result<(), crate::PostgresBackendError> {
if context == MigrationHookContext::AppliedNow {
let insert_result = transaction.execute(IDENTITY_INSERT_SQL, &[&network.as_str()]).await;
match insert_result {
std::result::Result::Ok(1) => {},
std::result::Result::Ok(_) | std::result::Result::Err(_) => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationFailed,
"store_identity_insert",
));
},
}
}
let rows_result = transaction.query(IDENTITY_LOAD_SQL, &[]).await;
let rows = match rows_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationMismatch,
"store_identity_read",
));
},
};
if rows.len() != 1 {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationMismatch,
"store_identity_count",
));
}
let row = &rows[0];
let singleton_result = row.try_get::<usize, i16>(0);
let network_result = row.try_get::<usize, std::string::String>(1);
let (singleton, stored_network) = match (singleton_result, network_result) {
(std::result::Result::Ok(singleton), std::result::Result::Ok(stored_network)) => (singleton, stored_network),
_ => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationMismatch,
"store_identity_decode",
));
},
};
if singleton != 1 {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationMismatch,
"store_identity_singleton",
));
}
let stored_network = match ksp_store_api::RawNetworkId::new(stored_network) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationMismatch,
"store_identity_network",
));
},
};
if stored_network.as_str() != network.as_str() {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::MigrationMismatch,
"store_identity_network",
));
}
return std::result::Result::Ok(());
}
async fn set_statement_timeout(
transaction: &deadpool_postgres::Transaction<'_>,
timeout: std::time::Duration,
@@ -384,8 +464,12 @@ async fn verify_metadata_shape(transaction: &deadpool_postgres::Transaction<'_>)
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::MigrationFailed, "metadata_shape"));
},
};
const REQUIRED: [(&str, &str, &str); 4] =
[("version", "bigint", "NO"), ("name", "text", "NO"), ("checksum", "text", "NO"), ("applied_at", "timestamp with time zone", "NO")];
const REQUIRED: [(&str, &str, &str); 4] = [
("version", "bigint", "NO"),
("name", "text", "NO"),
("checksum", "text", "NO"),
("applied_at", "timestamp with time zone", "NO"),
];
let mut found = [false; REQUIRED.len()];
for row in rows {
let name_result = row.try_get::<usize, std::string::String>(0);

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/tests/dependency_boundary.rs
// version: 6
// version: 7
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -21,10 +21,10 @@ fn pre_005_backend_owns_exact_physical_runtime_dependencies_without_reverse_faca
let migration = include_str!("../src/migration.rs");
let bootstrap_sql = include_str!("../migrations/V000__bootstrap.sql");
assert!(migration.contains("include_str!(\"../migrations/V000__bootstrap.sql\")"));
assert!(migration.contains("include_str!(\"../migrations/V001__raw_transaction.sql\")"));
assert!(bootstrap_sql.contains("ksp_store_schema_migrations"));
for forbidden in ["RawTransaction", "RawAccountState", "raw_transaction", "raw_account", "CORE", "DECODE", "SPECIALIZED"] {
assert!(!migration.contains(forbidden), "business migration implementation leaked into foundation: {forbidden}");
assert!(!bootstrap_sql.contains(forbidden), "business schema leaked into foundation SQL: {forbidden}");
assert!(!bootstrap_sql.contains(forbidden), "business schema leaked into immutable V000 SQL: {forbidden}");
}
return;
}
@@ -84,17 +84,28 @@ fn pre_007_health_probe_remains_foundation_only_and_private_sql() {
}
#[test]
fn pre_002_migration_engine_is_registry_driven_and_network_hook_ready_without_v001_schema() {
fn pre_003_migration_engine_embeds_v001_and_binds_network_without_repository_scope() {
let migration = include_str!("../src/migration.rs");
let v001 = include_str!("../migrations/V001__raw_transaction.sql");
assert!(migration.contains("const EMBEDDED_MIGRATIONS: &[EmbeddedMigration]"));
assert!(migration.contains("MigrationHook::None"));
assert!(migration.contains("run_migration_hook(transaction, network, migration.hook, MigrationHookContext::AppliedNow).await"));
assert!(migration.contains("network: &ksp_store_api::RawNetworkId"));
assert!(migration.contains("MigrationHook::StoreIdentity"));
assert!(migration.contains("MigrationHookContext::AppliedNow"));
assert!(migration.contains("MigrationHookContext::Existing"));
assert!(migration.contains("run_applied_migration_hooks(&transaction, network, next_index).await"));
assert!(migration.contains("validate_history(history.as_slice(), EMBEDDED_MIGRATIONS)"));
assert!(migration.contains("apply_pending_migrations(&transaction, network, next_index).await"));
assert!(!migration.contains("V001__raw_transaction.sql"));
assert!(!migration.contains("ksp_store_identity"));
assert!(migration.contains("INSERT INTO ksp_store_identity (singleton, network) VALUES (1, $1)"));
assert!(migration.contains("SELECT singleton, network FROM ksp_store_identity ORDER BY singleton LIMIT 2"));
assert!(migration.contains("ksp_store_api::RawNetworkId::new(stored_network)"));
for required in [
"CREATE TABLE ksp_store_identity",
"CREATE TABLE ksp_raw_transactions",
"CREATE TABLE ksp_raw_transaction_observations",
"CREATE TABLE ksp_raw_transaction_archive_payloads",
"CREATE INDEX ix_ksp_raw_transactions_slot_signature",
] {
assert!(v001.contains(required), "missing V001 physical object: {required}");
}
for forbidden in ["impl ksp_store_api::RawTransaction", "repository", "sqlx", "RawAccountState"] {
assert!(!migration.contains(forbidden), "repository/cross-scope implementation leaked into migration engine: {forbidden}");
assert!(!v001.contains(forbidden), "forbidden V001 scope content detected: {forbidden}");
}
return;
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/tests/hardening_completeness.rs
// version: 1
// version: 2
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -115,6 +115,7 @@ fn pre_009_backend_modules_exports_and_manifest_dependencies_are_exact() {
assert!(!crate_root.contains("pub mod "));
let actual_exports = public_reexport_names(crate_root);
let mut expected_exports = [
"ERROR_CODE_POSTGRES_RETENTION_COMPACTION_UNSUPPORTED",
"PostgresBackend",
"PostgresBackendError",
"PostgresBackendErrorKind",
@@ -125,7 +126,7 @@ fn pre_009_backend_modules_exports_and_manifest_dependencies_are_exact() {
];
expected_exports.sort_unstable();
assert_eq!(actual_exports.as_slice(), expected_exports.as_slice());
assert_eq!(actual_exports.len(), 7);
assert_eq!(actual_exports.len(), 8);
let manifest = include_str!("../Cargo.toml");
let actual_dependencies = manifest_dependency_names(manifest);
let expected_dependencies = [

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs
// version: 1
// version: 2
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -8,12 +8,30 @@
//! Opt-in real PostgreSQL proof for the Store foundation runtime.
//!
//! The test reads one dedicated PostgreSQL URI from stdin, refuses to start
//! when the KSP migration metadata table already exists, never prints the URI,
//! creates no business table and cleans up only metadata it proved it created.
//! when any KSP Store table managed by V000/V001 already exists, never prints
//! the URI, and cleans up only the isolated schema it proved absent first. It
//! validates migration/bootstrap behavior, not RawTransaction capabilities.
const LIVE_BOOTSTRAP_SQL: &str = include_str!("../migrations/V000__bootstrap.sql");
const LIVE_BROKEN_CHECKSUM_A: &str = "0000000000000000000000000000000000000000000000000000000000000000";
const LIVE_BROKEN_CHECKSUM_B: &str = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_transaction_observations;
DROP TABLE IF EXISTS ksp_raw_transaction_archive_payloads;
DROP TABLE IF EXISTS ksp_raw_transactions;
DROP TABLE IF EXISTS ksp_store_identity;
DROP TABLE IF EXISTS ksp_store_schema_migrations;"#;
const LIVE_MANAGED_SCHEMA_EXISTS_SQL: &str = r#"SELECT EXISTS (
SELECT 1 FROM information_schema.tables
WHERE table_schema = current_schema()
AND table_name IN (
'ksp_store_schema_migrations',
'ksp_store_identity',
'ksp_raw_transactions',
'ksp_raw_transaction_observations',
'ksp_raw_transaction_archive_payloads'
)
AND table_type = 'BASE TABLE'
)"#;
const LIVE_MAX_URI_BYTES: usize = 4_096;
const LIVE_METADATA_EXISTS_SQL: &str = r#"SELECT EXISTS (
SELECT 1 FROM information_schema.tables
@@ -21,7 +39,6 @@ const LIVE_METADATA_EXISTS_SQL: &str = r#"SELECT EXISTS (
AND table_name = 'ksp_store_schema_migrations'
AND table_type = 'BASE TABLE'
)"#;
const LIVE_METADATA_DROP_SQL: &str = "DROP TABLE IF EXISTS ksp_store_schema_migrations";
const LIVE_SENTINEL_CHECKSUM_SQL: &str = "SELECT checksum FROM ksp_store_schema_migrations WHERE version = 0";
const LIVE_SENTINEL_INSERT_SQL: &str =
"INSERT INTO ksp_store_schema_migrations (version, name, checksum, applied_at) VALUES (0, 'bootstrap', 'pre008_rollback_injected', CURRENT_TIMESTAMP)";
@@ -83,13 +100,13 @@ async fn run_live_test(uri: &str) -> std::result::Result<(), LiveFailure> {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let preexisting_result = metadata_exists(&admin).await;
let preexisting_result = managed_schema_exists(&admin).await;
let preexisting = match preexisting_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if preexisting {
return std::result::Result::Err(LiveFailure::new("metadata_preexisting_refusal"));
return std::result::Result::Err(LiveFailure::new("managed_schema_preexisting_refusal"));
}
let major_result = postgres_major(&admin).await;
let major = match major_result {
@@ -100,16 +117,16 @@ async fn run_live_test(uri: &str) -> std::result::Result<(), LiveFailure> {
return std::result::Result::Err(LiveFailure::new("postgres_major_unsupported"));
}
eprintln!("KSP Store PostgreSQL live proof: server major {major}");
let mut owns_metadata = false;
let scenario = run_foundation_scenario(&mut admin, uri, &mut owns_metadata).await;
let cleanup = if owns_metadata { drop_metadata(&admin).await } else { std::result::Result::Ok(()) };
let mut owns_schema = false;
let scenario = run_foundation_scenario(&mut admin, uri, &mut owns_schema).await;
let cleanup = if owns_schema { drop_managed_schema(&admin).await } else { std::result::Result::Ok(()) };
if let std::result::Result::Err(error) = cleanup {
return std::result::Result::Err(error);
}
if let std::result::Result::Err(error) = scenario {
return std::result::Result::Err(error);
}
let remains_result = metadata_exists(&admin).await;
let remains_result = managed_schema_exists(&admin).await;
let remains = match remains_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
@@ -120,7 +137,7 @@ async fn run_live_test(uri: &str) -> std::result::Result<(), LiveFailure> {
return std::result::Result::Ok(());
}
async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str, owns_metadata: &mut bool) -> std::result::Result<(), LiveFailure> {
async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str, owns_schema: &mut bool) -> std::result::Result<(), LiveFailure> {
let initial_result = open_backend(uri).await;
let initial = match initial_result {
std::result::Result::Ok(value) => value,
@@ -134,9 +151,9 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
if !created {
return std::result::Result::Err(LiveFailure::new("initial_bootstrap_metadata"));
}
*owns_metadata = true;
*owns_schema = true;
let initial_health = initial.health().await;
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(0) || initial_health.pending_migration_count() != 0 {
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(1) || initial_health.pending_migration_count() != 0 {
return std::result::Result::Err(LiveFailure::new("initial_health"));
}
let initial_close = close_backend(initial).await;
@@ -150,7 +167,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
};
let idempotent_health = idempotent.health().await;
if !idempotent_health.is_ready()
|| idempotent_health.migration_version() != std::option::Option::Some(0)
|| idempotent_health.migration_version() != std::option::Option::Some(1)
|| idempotent_health.pending_migration_count() != 0
{
return std::result::Result::Err(LiveFailure::new("idempotent_health"));
@@ -159,7 +176,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
if let std::result::Result::Err(error) = idempotent_close {
return std::result::Result::Err(error);
}
let reset_result = drop_metadata(admin).await;
let reset_result = drop_managed_schema(admin).await;
if let std::result::Result::Err(error) = reset_result {
return std::result::Result::Err(error);
}
@@ -213,7 +230,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
if let std::result::Result::Err(error) = recovered_close {
return std::result::Result::Err(error);
}
let rollback_reset = drop_metadata(admin).await;
let rollback_reset = drop_managed_schema(admin).await;
if let std::result::Result::Err(error) = rollback_reset {
return std::result::Result::Err(error);
}
@@ -233,7 +250,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let final_health = final_backend.health().await;
if !final_health.is_ready() || final_health.migration_version() != std::option::Option::Some(0) || final_health.pending_migration_count() != 0 {
if !final_health.is_ready() || final_health.migration_version() != std::option::Option::Some(1) || final_health.pending_migration_count() != 0 {
return std::result::Result::Err(LiveFailure::new("final_health"));
}
return close_backend(final_backend).await;
@@ -354,6 +371,18 @@ async fn postgres_major(client: &tokio_postgres::Client) -> std::result::Result<
return std::result::Result::Ok(version_num / 10_000);
}
async fn managed_schema_exists(client: &tokio_postgres::Client) -> std::result::Result<bool, LiveFailure> {
let row_result = client.query_one(LIVE_MANAGED_SCHEMA_EXISTS_SQL, &[]).await;
let row = match row_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("managed_schema_probe")),
};
return match row.try_get::<usize, bool>(0) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("managed_schema_probe_decode")),
};
}
async fn metadata_exists(client: &tokio_postgres::Client) -> std::result::Result<bool, LiveFailure> {
let row_result = client.query_one(LIVE_METADATA_EXISTS_SQL, &[]).await;
let row = match row_result {
@@ -366,10 +395,10 @@ async fn metadata_exists(client: &tokio_postgres::Client) -> std::result::Result
};
}
async fn drop_metadata(client: &tokio_postgres::Client) -> std::result::Result<(), LiveFailure> {
return match client.batch_execute(LIVE_METADATA_DROP_SQL).await {
async fn drop_managed_schema(client: &tokio_postgres::Client) -> std::result::Result<(), LiveFailure> {
return match client.batch_execute(LIVE_MANAGED_SCHEMA_DROP_SQL).await {
std::result::Result::Ok(()) => std::result::Result::Ok(()),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("metadata_cleanup")),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("managed_schema_cleanup")),
};
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/tests/public_api.rs
// version: 3
// version: 4
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -58,3 +58,13 @@ fn pre_007_backend_health_bridge_exposes_only_safe_snapshot_types() {
let _health_probe = ksp_store_postgres_lib::PostgresBackend::health;
return;
}
#[test]
fn pre_003_retention_compaction_error_code_matches_store_contract_value() {
assert_eq!(ksp_store_postgres_lib::ERROR_CODE_POSTGRES_RETENTION_COMPACTION_UNSUPPORTED.domain(), "store");
assert_eq!(
ksp_store_postgres_lib::ERROR_CODE_POSTGRES_RETENTION_COMPACTION_UNSUPPORTED.code(),
"postgres_retention_compaction_unsupported",
);
return;
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/unit_tests/migration.rs
// version: 2
// version: 3
fn applied(version: i64, name: &str, checksum: &str) -> super::AppliedMigration {
return super::AppliedMigration { checksum: checksum.to_owned(), name: name.to_owned(), version };
@@ -10,36 +10,69 @@ fn embedded(version: i64, name: &'static str, sql: &'static str) -> super::Embed
}
#[test]
fn pre_002_embedded_registry_keeps_v000_immutable_and_current_version_registry_driven() {
assert_eq!(super::EMBEDDED_MIGRATIONS.len(), 1);
let migration = &super::EMBEDDED_MIGRATIONS[0];
assert_eq!(migration.version, 0);
assert_eq!(migration.name, "bootstrap");
assert_eq!(migration.hook, super::MigrationHook::None);
assert!(migration.sql.contains("CREATE TABLE ksp_store_schema_migrations"));
fn pre_003_embedded_registry_keeps_v000_immutable_and_adds_exact_v001() {
assert_eq!(super::EMBEDDED_MIGRATIONS.len(), 2);
let v000 = &super::EMBEDDED_MIGRATIONS[0];
assert_eq!(v000.version, 0);
assert_eq!(v000.name, "bootstrap");
assert_eq!(v000.hook, super::MigrationHook::None);
assert!(v000.sql.contains("CREATE TABLE ksp_store_schema_migrations"));
for forbidden in ["RawTransaction", "RawAccountState", "raw_transaction", "raw_account", "CORE", "DECODE", "SPECIALIZED"] {
assert!(!migration.sql.contains(forbidden), "business schema leaked into bootstrap SQL: {forbidden}");
assert!(!v000.sql.contains(forbidden), "business schema leaked into immutable V000 SQL: {forbidden}");
}
assert_eq!(super::migration_checksum(v000.sql), "d29068b8c13b9dc0cc9ef6aaadd0fa12d41e0fe4c56541a1118c4bfc846a1450");
let v001 = &super::EMBEDDED_MIGRATIONS[1];
assert_eq!(v001.version, 1);
assert_eq!(v001.name, "raw_transaction");
assert_eq!(v001.hook, super::MigrationHook::StoreIdentity);
assert_eq!(super::migration_checksum(v001.sql), "6fe57ed0313d2ed295280dd6e49f6d86695d4e4effee2724a25e36db1ea17761");
assert!(super::validate_embedded_registry(super::EMBEDDED_MIGRATIONS).is_ok());
assert_eq!(crate::current_migration_version(), 0);
let checksum = super::migration_checksum(migration.sql);
assert_eq!(checksum.len(), 64);
assert_eq!(checksum, "d29068b8c13b9dc0cc9ef6aaadd0fa12d41e0fe4c56541a1118c4bfc846a1450");
assert_eq!(crate::current_migration_version(), 1);
return;
}
#[test]
fn pre_002_ordered_registry_accepts_exact_history_prefix_and_full_history() {
fn pre_003_v001_inventory_indexes_and_api_bounds_are_exact() {
let sql = super::EMBEDDED_MIGRATIONS[1].sql;
for required in [
"CREATE TABLE ksp_store_identity",
"CREATE TABLE ksp_raw_transactions",
"CREATE TABLE ksp_raw_transaction_observations",
"CREATE TABLE ksp_raw_transaction_archive_payloads",
"CREATE INDEX ix_ksp_raw_transactions_slot_signature",
"WHERE retention_state <> 'purged'",
"octet_length(signature) = 64",
"slot BETWEEN 0 AND 18446744073709551615",
"block_time_unix_millis BETWEEN 0 AND 253402300799999",
"octet_length(content_hash) = 32",
"octet_length(payload) BETWEEN 1 AND 16777216",
"format_version BETWEEN 1 AND 4294967295",
"received_at_unix_millis BETWEEN 0 AND 253402300799999",
"source_payload_size_bytes BETWEEN 0 AND 67108864",
"origin IN ('backfill', 'import', 'live', 'repair', 'replay')",
"retention_state IN ('full', 'archived', 'purged')",
] {
assert!(sql.contains(required), "V001 physical contract is missing: {required}");
}
assert_eq!(sql.matches("PRIMARY KEY").count(), 4);
assert_eq!(sql.matches("REFERENCES ksp_raw_transactions(signature) ON DELETE RESTRICT").count(), 2);
assert!(!sql.contains("compacted"));
assert!(!sql.contains("BIGSERIAL"));
assert!(!sql.contains("slot BIGINT"));
assert!(sql.contains("network TEXT NOT NULL"));
return;
}
#[test]
fn pre_003_ordered_registry_accepts_v000_prefix_and_full_v001_history() {
let v000 = super::EMBEDDED_MIGRATIONS[0];
let v001 = embedded(1, "synthetic", "SELECT 1;");
let registry = [v000, v001];
assert!(super::validate_embedded_registry(&registry).is_ok());
let v001 = super::EMBEDDED_MIGRATIONS[1];
let v000_checksum = super::migration_checksum(v000.sql);
let prefix = [applied(0, v000.name, v000_checksum.as_str())];
assert_eq!(super::validate_history(&prefix, &registry).ok(), std::option::Option::Some(1));
assert_eq!(super::validate_history(&prefix, super::EMBEDDED_MIGRATIONS).ok(), std::option::Option::Some(1));
let v001_checksum = super::migration_checksum(v001.sql);
let full = [applied(0, v000.name, v000_checksum.as_str()), applied(1, v001.name, v001_checksum.as_str())];
assert_eq!(super::validate_history(&full, &registry).ok(), std::option::Option::Some(2));
assert_eq!(super::validate_history(&full, super::EMBEDDED_MIGRATIONS).ok(), std::option::Option::Some(2));
return;
}
@@ -58,10 +91,9 @@ fn pre_002_registry_rejects_empty_nonzero_gap_and_empty_metadata_entries() {
}
#[test]
fn pre_002_divergent_missing_or_gapped_history_is_terminal_mismatch() {
fn pre_003_divergent_missing_or_gapped_history_is_terminal_mismatch() {
let v000 = super::EMBEDDED_MIGRATIONS[0];
let v001 = embedded(1, "synthetic", "SELECT 1;");
let registry = [v000, v001];
let v001 = super::EMBEDDED_MIGRATIONS[1];
let v000_checksum = super::migration_checksum(v000.sql);
let v001_checksum = super::migration_checksum(v001.sql);
let wrong_name = [applied(0, "changed", v000_checksum.as_str())];
@@ -69,21 +101,24 @@ fn pre_002_divergent_missing_or_gapped_history_is_terminal_mismatch() {
let missing: [super::AppliedMigration; 0] = [];
let missing_v000 = [applied(1, v001.name, v001_checksum.as_str())];
for history in [&wrong_name[..], &wrong_checksum[..], &missing[..], &missing_v000[..]] {
let result = super::validate_history(history, &registry);
let result = super::validate_history(history, super::EMBEDDED_MIGRATIONS);
assert_eq!(result.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::MigrationMismatch));
}
return;
}
#[test]
fn pre_002_newer_history_is_rejected_without_down_migration() {
fn pre_003_newer_history_is_rejected_without_down_migration() {
let v000 = super::EMBEDDED_MIGRATIONS[0];
let v001 = embedded(1, "synthetic", "SELECT 1;");
let registry = [v000, v001];
let v001 = super::EMBEDDED_MIGRATIONS[1];
let v000_checksum = super::migration_checksum(v000.sql);
let v001_checksum = super::migration_checksum(v001.sql);
let history = [applied(0, v000.name, v000_checksum.as_str()), applied(1, v001.name, v001_checksum.as_str()), applied(2, "future", "future-checksum")];
let result = super::validate_history(&history, &registry);
let history = [
applied(0, v000.name, v000_checksum.as_str()),
applied(1, v001.name, v001_checksum.as_str()),
applied(2, "future", "future-checksum"),
];
let result = super::validate_history(&history, super::EMBEDDED_MIGRATIONS);
assert_eq!(result.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::SchemaNewer));
return;
}