v0.3.3-pre.007

This commit is contained in:
2026-08-30 13:37:22 +02:00
parent 10996f11f7
commit aa56a12846
13 changed files with 683 additions and 41 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/src/error.rs
// version: 7
// version: 8
/// 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 =
@@ -33,6 +33,8 @@ pub enum PostgresBackendErrorKind {
ReadFailed,
/// A RAW write requires an existing canonical reference that is not durable.
ReferenceNotFound,
/// The requested RAW retention transition requires a compacted representation unsupported by PostgreSQL.
RetentionCompactionUnsupported,
/// The database schema history contains a migration newer than this runtime understands.
SchemaNewer,
/// Explicit backend shutdown did not drain inside the supplied deadline.

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/src/lib.rs
// version: 12
// version: 13
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -17,6 +17,8 @@
//! checks and safe conflict classification without exposing PostgreSQL rows or
//! SQL through the public bridge. `0.3.3-pre.006` adds deterministic keyset
//! pagination with a fixed opaque cursor bound to network, range and direction.
//! `0.3.3-pre.007` adds atomic `Full -> Archived -> Purged` retention transitions
//! with compare-and-transition outcomes and explicit rejection of `Compacted`.
//!
//! 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
@@ -75,6 +77,8 @@ pub(crate) use self::raw_transaction::list_raw_transactions;
pub(crate) use self::raw_transaction::persist_raw_transaction_acquisition;
/// Private additional RAW transaction observation writer consumed by the physical backend runtime.
pub(crate) use self::raw_transaction::record_raw_transaction_observation;
/// Private RAW transaction retention transition writer consumed by the physical backend runtime.
pub(crate) use self::raw_transaction::transition_raw_transaction_retention;
/// Private Deadpool error mapper shared with the health probe.
pub(crate) use self::runtime::map_pool_error;
/// Private Deadpool status projector shared with the health probe.

View File

@@ -1,22 +1,29 @@
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
// version: 3
// version: 4
pub(crate) mod cursor;
const DELETE_ARCHIVE_PAYLOAD_SQL: &str = "DELETE FROM ksp_raw_transaction_archive_payloads WHERE signature = $1";
const GET_ARCHIVE_PAYLOAD_SQL: &str = "SELECT payload FROM ksp_raw_transaction_archive_payloads WHERE signature = $1";
const GET_OBSERVATION_SQL: &str = "SELECT observation_key, transaction_signature, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes FROM ksp_raw_transaction_observations WHERE observation_key = $1";
const GET_RETENTION_SQL: &str = "SELECT retention_state FROM ksp_raw_transactions WHERE signature = $1";
const GET_TOMBSTONE_SQL: &str = "SELECT signature, slot::text AS slot_text, block_time_unix_millis, format_id, format_version, content_hash, retention_state FROM ksp_raw_transactions WHERE signature = $1";
const GET_TRANSACTION_SQL: &str = "SELECT transaction_row.signature, transaction_row.slot::text AS slot_text, transaction_row.block_time_unix_millis, transaction_row.format_id, transaction_row.format_version, transaction_row.content_hash, transaction_row.payload, transaction_row.retention_state, archive_row.payload AS archive_payload FROM ksp_raw_transactions AS transaction_row LEFT JOIN ksp_raw_transaction_archive_payloads AS archive_row ON archive_row.signature = transaction_row.signature WHERE transaction_row.signature = $1";
const INSERT_ARCHIVE_PAYLOAD_SQL: &str = "INSERT INTO ksp_raw_transaction_archive_payloads (signature, payload) VALUES ($1, $2)";
const INSERT_OBSERVATION_SQL: &str = "INSERT INTO ksp_raw_transaction_observations (observation_key, transaction_signature, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14) ON CONFLICT (observation_key) DO NOTHING RETURNING observation_key";
const INSERT_TRANSACTION_SQL: &str = "INSERT INTO ksp_raw_transactions (signature, slot, block_time_unix_millis, format_id, format_version, content_hash, payload, retention_state) VALUES ($1, $2::TEXT::NUMERIC, $3, $4, $5, $6, $7, 'full') ON CONFLICT (signature) DO NOTHING RETURNING signature";
const LIST_TRANSACTIONS_ASC_SQL: &str = "SELECT signature, slot::text AS slot_text FROM ksp_raw_transactions WHERE retention_state <> 'purged' AND ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR (slot, signature) > ($3::TEXT::NUMERIC, $4::BYTEA)) ORDER BY slot ASC, signature ASC LIMIT $5";
const LIST_TRANSACTIONS_DESC_SQL: &str = "SELECT signature, slot::text AS slot_text FROM ksp_raw_transactions WHERE retention_state <> 'purged' AND ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR (slot, signature) < ($3::TEXT::NUMERIC, $4::BYTEA)) ORDER BY slot DESC, signature DESC LIMIT $5";
const LOCK_OBSERVATION_SQL: &str = "SELECT observation_key, transaction_signature, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes FROM ksp_raw_transaction_observations WHERE observation_key = $1 FOR UPDATE";
const LOCK_RETENTION_TRANSACTION_SQL: &str =
"SELECT block_time_unix_millis, payload, retention_state FROM ksp_raw_transactions WHERE signature = $1 FOR UPDATE";
const LOCK_TRANSACTION_SQL: &str = "SELECT signature, slot::text AS slot_text, block_time_unix_millis, format_id, format_version, content_hash, payload, retention_state, NULL::BYTEA AS archive_payload FROM ksp_raw_transactions WHERE signature = $1 FOR UPDATE";
const LOCK_TRANSACTION_STATE_SQL: &str = "SELECT retention_state FROM ksp_raw_transactions WHERE signature = $1 FOR UPDATE";
const REHYDRATE_TRANSACTION_SQL: &str =
"UPDATE ksp_raw_transactions SET block_time_unix_millis = $2, payload = $3, retention_state = 'full' WHERE signature = $1";
const UPDATE_ARCHIVED_TRANSACTION_SQL: &str =
"UPDATE ksp_raw_transactions SET payload = NULL, retention_state = 'archived' WHERE signature = $1 AND retention_state = 'full'";
const UPDATE_PURGED_TRANSACTION_SQL: &str = "UPDATE ksp_raw_transactions SET block_time_unix_millis = NULL, payload = NULL, retention_state = 'purged' WHERE signature = $1 AND retention_state = 'archived'";
struct RawListDbRow {
signature: std::vec::Vec<u8>,
@@ -40,6 +47,12 @@ struct RawObservationDbRow {
transaction_signature: std::vec::Vec<u8>,
}
struct RawRetentionDbRow {
block_time_unix_millis: std::option::Option<i64>,
payload: std::option::Option<std::vec::Vec<u8>>,
retention_state: std::string::String,
}
struct RawTransactionDbRow {
archive_payload: std::option::Option<std::vec::Vec<u8>>,
block_time_unix_millis: std::option::Option<i64>,
@@ -450,12 +463,250 @@ pub(crate) async fn record_raw_transaction_observation(
return std::result::Result::Ok(outcome);
}
/// Applies one atomic compare-and-transition RAW transaction retention mutation.
pub(crate) async fn transition_raw_transaction_retention(
pool: &deadpool_postgres::Pool,
network: &ksp_store_api::RawNetworkId,
transition: ksp_store_api::RawTransactionRetentionTransition,
) -> std::result::Result<ksp_store_api::RawRetentionWriteOutcome, crate::PostgresBackendError> {
let input_result = ensure_retention_transition_inputs(network, &transition);
if let std::result::Result::Err(error) = input_result {
return std::result::Result::Err(error);
}
let client_result = pool.get().await;
let mut client = match client_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)),
};
let sql_transaction_result = client.transaction().await;
let sql_transaction = match sql_transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_begin")),
};
let signature = transition.reference().signature();
let signature_bytes: &[u8] = signature.as_bytes();
let row_result = sql_transaction.query_opt(LOCK_RETENTION_TRANSACTION_SQL, &[&signature_bytes]).await;
let row = match row_result {
std::result::Result::Ok(std::option::Option::Some(value)) => value,
std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(reference_not_found("raw_retention_transition_reference")),
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_lock_transaction")),
};
let physical = match raw_retention_db_row(&row) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let current = match decode_retention_state(physical.retention_state.as_str()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let shape_result = validate_retention_shape(&sql_transaction, signature_bytes, &physical, current).await;
if let std::result::Result::Err(error) = shape_result {
return std::result::Result::Err(error);
}
let decision = match retention_transition_decision(current, transition.expected(), transition.target()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match decision {
RetentionTransitionDecision::AlreadyAtTarget => {
return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::AlreadyAtTarget).await;
},
RetentionTransitionDecision::ExpectedStateMismatch => {
return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::ExpectedStateMismatch).await;
},
RetentionTransitionDecision::Archive => {
let payload = match physical.payload.as_ref() {
std::option::Option::Some(value) => value.as_slice(),
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_full_payload")),
};
let archive_result = archive_full_transaction(&sql_transaction, signature_bytes, payload).await;
if let std::result::Result::Err(error) = archive_result {
return std::result::Result::Err(error);
}
},
RetentionTransitionDecision::Purge => {
let purge_result = purge_archived_transaction(&sql_transaction, signature_bytes).await;
if let std::result::Result::Err(error) = purge_result {
return std::result::Result::Err(error);
}
},
}
return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::Applied).await;
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ExistingTransactionMatch {
Active,
Purged,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum RetentionTransitionDecision {
AlreadyAtTarget,
Archive,
ExpectedStateMismatch,
Purge,
}
async fn archive_full_transaction(
sql_transaction: &deadpool_postgres::Transaction<'_>,
signature_bytes: &[u8],
payload: &[u8],
) -> std::result::Result<(), crate::PostgresBackendError> {
let insert_result = sql_transaction.execute(INSERT_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes, &payload]).await;
match insert_result {
std::result::Result::Ok(1) => {},
std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_retention_archive_insert_cardinality")),
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_archive_insert")),
}
let update_result = sql_transaction.execute(UPDATE_ARCHIVED_TRANSACTION_SQL, &[&signature_bytes]).await;
return match update_result {
std::result::Result::Ok(1) => std::result::Result::Ok(()),
std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_retention_archive_update_cardinality")),
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_retention_archive_update")),
};
}
async fn commit_retention_outcome(
sql_transaction: deadpool_postgres::Transaction<'_>,
outcome: ksp_store_api::RawRetentionWriteOutcome,
) -> std::result::Result<ksp_store_api::RawRetentionWriteOutcome, crate::PostgresBackendError> {
let commit_result = sql_transaction.commit().await;
if commit_result.is_err() {
return std::result::Result::Err(write_failed("raw_retention_commit"));
}
return std::result::Result::Ok(outcome);
}
async fn load_retention_archive_payload(
sql_transaction: &deadpool_postgres::Transaction<'_>,
signature_bytes: &[u8],
) -> std::result::Result<std::option::Option<std::vec::Vec<u8>>, crate::PostgresBackendError> {
let row_result = sql_transaction.query_opt(GET_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await;
let row = match row_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_archive_query")),
};
let archive_row = match row {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Ok(std::option::Option::None),
};
let payload_result = archive_row.try_get::<_, std::vec::Vec<u8>>("payload");
return match payload_result {
std::result::Result::Ok(value) => std::result::Result::Ok(std::option::Option::Some(value)),
std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_retention_archive_decode")),
};
}
async fn purge_archived_transaction(
sql_transaction: &deadpool_postgres::Transaction<'_>,
signature_bytes: &[u8],
) -> std::result::Result<(), crate::PostgresBackendError> {
let delete_result = sql_transaction.execute(DELETE_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await;
match delete_result {
std::result::Result::Ok(1) => {},
std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_retention_purge_archive_cardinality")),
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_purge_archive")),
}
let update_result = sql_transaction.execute(UPDATE_PURGED_TRANSACTION_SQL, &[&signature_bytes]).await;
return match update_result {
std::result::Result::Ok(1) => std::result::Result::Ok(()),
std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_retention_purge_update_cardinality")),
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_retention_purge_update")),
};
}
fn raw_retention_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawRetentionDbRow, crate::PostgresBackendError> {
let block_time_unix_millis = match row.try_get::<_, std::option::Option<i64>>("block_time_unix_millis") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_transition_decode")),
};
let payload = match row.try_get::<_, std::option::Option<std::vec::Vec<u8>>>("payload") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_transition_decode")),
};
let retention_state = match row.try_get::<_, std::string::String>("retention_state") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_transition_decode")),
};
return std::result::Result::Ok(RawRetentionDbRow { block_time_unix_millis, payload, retention_state });
}
fn retention_transition_decision(
current: ksp_store_api::RawRetentionState,
expected: ksp_store_api::RawRetentionState,
target: ksp_store_api::RawRetentionState,
) -> std::result::Result<RetentionTransitionDecision, crate::PostgresBackendError> {
if current == target {
return std::result::Result::Ok(RetentionTransitionDecision::AlreadyAtTarget);
}
if current != expected {
return std::result::Result::Ok(RetentionTransitionDecision::ExpectedStateMismatch);
}
return match (expected, target) {
(ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Archived) => std::result::Result::Ok(RetentionTransitionDecision::Archive),
(ksp_store_api::RawRetentionState::Archived, ksp_store_api::RawRetentionState::Purged) => std::result::Result::Ok(RetentionTransitionDecision::Purge),
(ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Compacted)
| (ksp_store_api::RawRetentionState::Compacted, ksp_store_api::RawRetentionState::Archived) => {
std::result::Result::Err(retention_compaction_unsupported("raw_retention_compaction"))
},
_ => std::result::Result::Err(data_invalid("raw_retention_transition_unreachable")),
};
}
async fn validate_retention_shape(
sql_transaction: &deadpool_postgres::Transaction<'_>,
signature_bytes: &[u8],
row: &RawRetentionDbRow,
state: ksp_store_api::RawRetentionState,
) -> std::result::Result<(), crate::PostgresBackendError> {
let archive_payload = match load_retention_archive_payload(sql_transaction, signature_bytes).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
return match state {
ksp_store_api::RawRetentionState::Full => {
let payload = match row.payload.as_ref() {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_full_shape")),
};
let payload_result = validate_retained_payload_bytes(payload.as_slice(), "raw_retention_full_payload");
if let std::result::Result::Err(error) = payload_result {
return std::result::Result::Err(error);
}
if archive_payload.is_some() {
return std::result::Result::Err(data_invalid("raw_retention_full_archive_residue"));
}
std::result::Result::Ok(())
},
ksp_store_api::RawRetentionState::Archived => {
if row.payload.is_some() {
return std::result::Result::Err(data_invalid("raw_retention_archived_hot_payload"));
}
let payload = match archive_payload.as_ref() {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_archived_payload_missing")),
};
validate_retained_payload_bytes(payload.as_slice(), "raw_retention_archived_payload")
},
ksp_store_api::RawRetentionState::Purged => {
if row.payload.is_some() || archive_payload.is_some() || row.block_time_unix_millis.is_some() {
return std::result::Result::Err(data_invalid("raw_retention_purged_shape"));
}
std::result::Result::Ok(())
},
ksp_store_api::RawRetentionState::Compacted => std::result::Result::Err(retention_compaction_unsupported("raw_retention_compaction")),
_ => std::result::Result::Err(data_invalid("raw_retention_state")),
};
}
fn validate_retained_payload_bytes(payload: &[u8], phase: &'static str) -> std::result::Result<(), crate::PostgresBackendError> {
if payload.is_empty() || payload.len() > ksp_store_api::MAX_RAW_PAYLOAD_BYTES {
return std::result::Result::Err(data_invalid(phase));
}
return std::result::Result::Ok(());
}
async fn insert_canonical_transaction(
sql_transaction: &deadpool_postgres::Transaction<'_>,
raw_transaction: &ksp_store_api::RawTransaction,
@@ -717,6 +968,20 @@ fn ensure_acquisition_inputs(
return std::result::Result::Ok(());
}
fn ensure_retention_transition_inputs(
network: &ksp_store_api::RawNetworkId,
transition: &ksp_store_api::RawTransactionRetentionTransition,
) -> std::result::Result<(), crate::PostgresBackendError> {
let network_result = ensure_network(network, transition.reference(), "raw_retention_transition_network");
if let std::result::Result::Err(error) = network_result {
return std::result::Result::Err(error);
}
if transition.expected() == ksp_store_api::RawRetentionState::Compacted || transition.target() == ksp_store_api::RawRetentionState::Compacted {
return std::result::Result::Err(retention_compaction_unsupported("raw_retention_compaction"));
}
return std::result::Result::Ok(());
}
fn encode_origin(origin: ksp_store_api::RawAcquisitionOrigin) -> std::result::Result<&'static str, crate::PostgresBackendError> {
return match origin {
ksp_store_api::RawAcquisitionOrigin::Backfill => std::result::Result::Ok("backfill"),
@@ -732,6 +997,10 @@ fn conflict(phase: &'static str) -> crate::PostgresBackendError {
return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::Conflict, phase);
}
fn retention_compaction_unsupported(phase: &'static str) -> crate::PostgresBackendError {
return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::RetentionCompactionUnsupported, phase);
}
fn reference_not_found(phase: &'static str) -> crate::PostgresBackendError {
return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReferenceNotFound, phase);
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/src/runtime.rs
// version: 8
// version: 9
const APPLICATION_NAME: &str = "ksp-store";
const MAX_CONNECTION_URI_BYTES: usize = 4_096;
@@ -390,6 +390,14 @@ impl PostgresBackend {
return crate::record_raw_transaction_observation(&self.pool, &self.network, observation).await;
}
/// Applies one policy-authorized atomic RAW transaction retention transition.
pub async fn transition_raw_transaction_retention(
&self,
transition: ksp_store_api::RawTransactionRetentionTransition,
) -> std::result::Result<ksp_store_api::RawRetentionWriteOutcome, crate::PostgresBackendError> {
return crate::transition_raw_transaction_retention(&self.pool, &self.network, transition).await;
}
/// Explicitly closes the pool and waits for all owned pooled objects to drain inside the supplied bound.
pub async fn close(self, timeout: std::time::Duration) -> std::result::Result<(), crate::PostgresBackendError> {
self.pool.close();