diff --git a/Cargo.toml b/Cargo.toml index 5e64328..5bb71eb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 351 +# version: 352 [workspace] resolver = "3" members = ["crates/ksp-app-config-desk", "crates/ksp-app-solprices-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-interface-lib", "crates/ksp-logging-lib", "crates/ksp-offchain-transport-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-program-api", "crates/ksp-store-api", "crates/ksp-store-lib", "crates/ksp-store-postgres-lib", "crates/ksp-wallet-lib"] [workspace.package] -version = "0.3.3-pre.4" +version = "0.3.3-pre.5" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-store-postgres-lib/README.md b/crates/ksp-store-postgres-lib/README.md index a505b06..afd9d00 100644 --- a/crates/ksp-store-postgres-lib/README.md +++ b/crates/ksp-store-postgres-lib/README.md @@ -1,5 +1,5 @@ - + # ksp-store-postgres-lib @@ -83,7 +83,7 @@ ksp_store_schema_migrations Le moteur vérifie version, nom et checksum SHA-256, sérialise les runners par advisory transaction lock et refuse une history divergente ou plus récente que le runtime. -V001 possède désormais le schéma physique `RawTransaction` et son contrat de compatibilité ; les opérations métier restent introduites par tranches afin de préserver des gates courts et vérifiables. +V001 possède désormais le schéma physique `RawTransaction` et son contrat de compatibilité. Les lectures exactes sont acquises depuis `pre.004`; `pre.005` ajoute les écritures atomiques transaction + observation, l'idempotence réelle et la classification de conflit. ## Health et erreurs @@ -119,13 +119,28 @@ Le SQL et les rows restent privés au backend. Le mapping PostgreSQL est fallibl `Full` lit le payload chaud, `Archived` le reconstruit depuis la relation archive et `Purged` retourne `None`; le tombstone reste accessible séparément pour `Purged`. +## Écritures RAW `0.3.3-pre.005` + +Le backend expose deux écritures étroites : + +```text +persist_raw_transaction_acquisition +record_raw_transaction_observation +``` + +L'acquisition canonique et son observation sont commises dans une seule transaction PostgreSQL. L'insertion utilise les clés uniques physiques sans prélecture `has_*`; après un conflit unique, le backend verrouille la ligne gagnante et compare le contenu réel avant de conclure `AlreadyPresent` ou `Conflict`. Les octets du payload sont comparés lorsqu'ils existent encore : le hash seul ne constitue jamais une preuve d'idempotence. + +Un tombstone `Purged` compatible produit `SkippedPurged/NotRecorded` en mode normal. `ForceRehydrate` restaure `Full` et l'observation dans la même transaction. Une observation supplémentaire ne crée jamais implicitement son canonique ; une référence absente est classée `ReferenceNotFound` et un canonical purgé retourne `NotRecorded`. + +Les erreurs physiques d'écriture sont réduites à `WriteFailed`; aucune erreur serveur, SQLSTATE, query ou valeur de bind n'est conservée. + ## Hors périmètre actuel La crate ne contient encore : -- aucune écriture PostgreSQL `RawTransaction*` ; - aucune pagination/listing `RawTransaction` ; - aucune implémentation complète des traits `RawTransaction*` de `ksp-store-api` tant que `list_raw_transactions` manque ; +- aucune transition mutante de rétention ; - aucune implémentation PostgreSQL des capabilities `RawAccount*` ; - aucune orchestration worker/job ; - aucun transport d'acquisition ou decoder Program. diff --git a/crates/ksp-store-postgres-lib/USAGE.md b/crates/ksp-store-postgres-lib/USAGE.md index 3598b47..4cbed4e 100644 --- a/crates/ksp-store-postgres-lib/USAGE.md +++ b/crates/ksp-store-postgres-lib/USAGE.md @@ -1,5 +1,5 @@ - + # Utilisation de ksp-store-postgres-lib @@ -123,15 +123,18 @@ Le runner est transactionnel et sérialisé par advisory transaction lock. Une d match error.kind() { ksp_store_postgres_lib::PostgresBackendErrorKind::ConfigInvalid => {} ksp_store_postgres_lib::PostgresBackendErrorKind::ConnectFailed => {} + ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {} ksp_store_postgres_lib::PostgresBackendErrorKind::DataInvalid => {} ksp_store_postgres_lib::PostgresBackendErrorKind::PoolTimeout => {} ksp_store_postgres_lib::PostgresBackendErrorKind::HealthFailed => {} ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationFailed => {} ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationMismatch => {} ksp_store_postgres_lib::PostgresBackendErrorKind::ReadFailed => {} + ksp_store_postgres_lib::PostgresBackendErrorKind::ReferenceNotFound => {} ksp_store_postgres_lib::PostgresBackendErrorKind::SchemaNewer => {} ksp_store_postgres_lib::PostgresBackendErrorKind::ShutdownTimeout => {} ksp_store_postgres_lib::PostgresBackendErrorKind::TlsFailed => {} + ksp_store_postgres_lib::PostgresBackendErrorKind::WriteFailed => {} ksp_store_postgres_lib::PostgresBackendErrorKind::WrongNetwork => {} _ => {} } @@ -165,13 +168,32 @@ absent -> None Le tombstone `Purged` reste lisible séparément. -## 8. Ce que cette crate ne permet pas encore +## 8. Écritures RAW transaction + +Depuis `0.3.3-pre.005`, le backend physique expose également : + +```rust +let acquisition = backend + .persist_raw_transaction_acquisition(transaction, observation, mode) + .await; + +let observation = backend + .record_raw_transaction_observation(additional_observation) + .await; +``` + +La première opération est atomique : transaction canonique et observation sont toutes deux durables ou aucune ne l'est. Les doublons ne sont pas détectés par une prélecture `has_*` : l'insert unique est tenté directement, puis un conflit relit/verrouille la ligne gagnante et compare son contenu. Une divergence sous la même signature ou la même `observation_key` produit `PostgresBackendErrorKind::Conflict`. + +Pour une transaction purgée, le mode normal retourne `SkippedPurged/NotRecorded` lorsque le tombstone est compatible. `ForceRehydrate` restaure le payload `Full` puis enregistre l'observation dans la même transaction. `record_raw_transaction_observation` ne crée jamais de canonique : référence absente -> `ReferenceNotFound`, canonique `Purged` -> `NotRecorded`. + +Les références réseau-scopées sont toujours validées avant `pool.get()`. + +## 9. Ce que cette crate ne permet pas encore La tranche ne fournit pas encore : ```text list_raw_transactions / cursor -écritures canonique + observation transitions de rétention implémentations complètes des six traits RawTransaction* capabilities RawAccount* diff --git a/crates/ksp-store-postgres-lib/src/error.rs b/crates/ksp-store-postgres-lib/src/error.rs index fa557b5..a558df2 100644 --- a/crates/ksp-store-postgres-lib/src/error.rs +++ b/crates/ksp-store-postgres-lib/src/error.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/src/error.rs -// version: 5 +// version: 6 /// 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 = @@ -17,6 +17,8 @@ pub enum PostgresBackendErrorKind { PoolTimeout, /// A lightweight PostgreSQL health/readiness probe failed without exposing server text or SQL. HealthFailed, + /// A canonical RAW identity or observation key already exists with divergent durable content. + Conflict, /// PostgreSQL returned stored RAW data that cannot be represented by the stable Store API contract. DataInvalid, /// PostgreSQL migration/bootstrap execution failed without exposing server text or SQL. @@ -25,12 +27,16 @@ pub enum PostgresBackendErrorKind { MigrationMismatch, /// A PostgreSQL RAW read statement failed without exposing server text, SQL or bind values. ReadFailed, + /// A RAW write requires an existing canonical reference that is not durable. + ReferenceNotFound, /// The database schema history contains a migration newer than this runtime understands. SchemaNewer, /// Explicit backend shutdown did not drain inside the supplied deadline. ShutdownTimeout, /// Verified TLS configuration or negotiation could not be established. TlsFailed, + /// A PostgreSQL RAW write statement or transaction failed without exposing server text, SQL or bind values. + WriteFailed, /// A network-scoped RAW operation targeted a network different from the backend binding. WrongNetwork, } diff --git a/crates/ksp-store-postgres-lib/src/lib.rs b/crates/ksp-store-postgres-lib/src/lib.rs index c940f04..117573f 100644 --- a/crates/ksp-store-postgres-lib/src/lib.rs +++ b/crates/ksp-store-postgres-lib/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/src/lib.rs -// version: 9 +// version: 10 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -12,8 +12,10 @@ //! safe lightweight health/readiness probe. `0.3.3-pre.003-fix.001` splits //! migrations into versioned physical resources and verifies the effective //! PostgreSQL schema contract before readiness. `0.3.3-pre.004` adds exact -//! backend-private RAW transaction/observation/retention read mapping without -//! exposing PostgreSQL rows or SQL through the public bridge. +//! backend-private RAW transaction/observation/retention read mapping. +//! `0.3.3-pre.005` adds atomic canonical/observation writes, real idempotence +//! checks and safe conflict classification without exposing PostgreSQL rows or +//! SQL through the public bridge. //! //! 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 @@ -60,6 +62,10 @@ pub(crate) use self::raw_transaction::get_raw_transaction_observation; pub(crate) use self::raw_transaction::get_raw_transaction_retention_state; /// Private RAW transaction tombstone reader consumed by the physical backend runtime. pub(crate) use self::raw_transaction::get_raw_transaction_tombstone; +/// Private atomic RAW transaction acquisition writer consumed by the physical backend runtime. +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 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. diff --git a/crates/ksp-store-postgres-lib/src/raw_transaction.rs b/crates/ksp-store-postgres-lib/src/raw_transaction.rs index ec4862b..0813a6a 100644 --- a/crates/ksp-store-postgres-lib/src/raw_transaction.rs +++ b/crates/ksp-store-postgres-lib/src/raw_transaction.rs @@ -1,10 +1,18 @@ // file: crates/ksp-store-postgres-lib/src/raw_transaction.rs -// version: 1 +// version: 2 +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_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 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_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"; struct RawObservationDbRow { acquisition_method: std::string::String, @@ -227,6 +235,426 @@ pub(crate) async fn get_raw_transaction_tombstone( return std::result::Result::Ok(std::option::Option::Some(tombstone)); } +/// Persists one canonical RAW transaction and one observation atomically. +pub(crate) async fn persist_raw_transaction_acquisition( + pool: &deadpool_postgres::Pool, + network: &ksp_store_api::RawNetworkId, + raw_transaction: ksp_store_api::RawTransaction, + observation: ksp_store_api::RawTransactionObservation, + mode: ksp_store_api::RawTransactionAcquisitionMode, +) -> std::result::Result { + let input_result = ensure_acquisition_inputs(network, &raw_transaction, &observation); + 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_acquisition_begin")), + }; + let insert_result = insert_canonical_transaction(&sql_transaction, &raw_transaction).await; + let inserted = match insert_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let entity_outcome = if inserted { + ksp_store_api::RawEntityWriteOutcome::Inserted + } else { + let locked_result = load_locked_transaction_row(&sql_transaction, raw_transaction.reference()).await; + let locked = match locked_result { + std::result::Result::Ok(std::option::Option::Some(value)) => value, + std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(data_invalid("raw_acquisition_conflict_missing")), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let comparison_result = compare_existing_transaction(network, locked, &raw_transaction); + let comparison = match comparison_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + match comparison { + ExistingTransactionMatch::Active => ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, + ExistingTransactionMatch::Purged => { + if mode == ksp_store_api::RawTransactionAcquisitionMode::ForceRehydrate { + let rehydrate_result = rehydrate_transaction(&sql_transaction, &raw_transaction).await; + if let std::result::Result::Err(error) = rehydrate_result { + return std::result::Result::Err(error); + } + ksp_store_api::RawEntityWriteOutcome::Rehydrated + } else { + let commit_result = sql_transaction.commit().await; + if commit_result.is_err() { + return std::result::Result::Err(write_failed("raw_acquisition_commit")); + } + return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new( + ksp_store_api::RawEntityWriteOutcome::SkippedPurged, + ksp_store_api::RawObservationWriteOutcome::NotRecorded, + )); + } + }, + } + }; + let observation_result = persist_observation_row(&sql_transaction, network, &observation).await; + let observation_outcome = match observation_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let commit_result = sql_transaction.commit().await; + if commit_result.is_err() { + return std::result::Result::Err(write_failed("raw_acquisition_commit")); + } + return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new(entity_outcome, observation_outcome)); +} + +/// Persists one additional observation for an already known RAW transaction. +pub(crate) async fn record_raw_transaction_observation( + pool: &deadpool_postgres::Pool, + network: &ksp_store_api::RawNetworkId, + observation: ksp_store_api::RawTransactionObservation, +) -> std::result::Result { + let network_result = ensure_network(network, observation.transaction(), "raw_observation_write_network"); + if let std::result::Result::Err(error) = network_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_observation_begin")), + }; + let signature = observation.transaction().signature(); + let signature_bytes: &[u8] = signature.as_bytes(); + let state_row_result = sql_transaction.query_opt(LOCK_TRANSACTION_STATE_SQL, &[&signature_bytes]).await; + let state_row = match state_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_observation_transaction")), + std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_observation_lock_transaction")), + }; + let state_result = state_row.try_get::<_, std::string::String>("retention_state"); + let state = match state_result { + std::result::Result::Ok(value) => match decode_retention_state(value.as_str()) { + std::result::Result::Ok(decoded) => decoded, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_observation_transaction_state")), + }; + if state == ksp_store_api::RawRetentionState::Purged { + let commit_result = sql_transaction.commit().await; + if commit_result.is_err() { + return std::result::Result::Err(write_failed("raw_observation_commit")); + } + return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::NotRecorded); + } + let observation_result = persist_observation_row(&sql_transaction, network, &observation).await; + let outcome = match observation_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let commit_result = sql_transaction.commit().await; + if commit_result.is_err() { + return std::result::Result::Err(write_failed("raw_observation_commit")); + } + return std::result::Result::Ok(outcome); +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum ExistingTransactionMatch { + Active, + Purged, +} + +async fn insert_canonical_transaction( + sql_transaction: &deadpool_postgres::Transaction<'_>, + raw_transaction: &ksp_store_api::RawTransaction, +) -> std::result::Result { + let signature = raw_transaction.reference().signature(); + let signature_bytes: &[u8] = signature.as_bytes(); + let slot_text = raw_transaction.slot().to_string(); + let block_time = match raw_transaction.block_time() { + std::option::Option::Some(value) => match i64::try_from(value.unix_millis()) { + std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_acquisition_block_time")), + }, + std::option::Option::None => std::option::Option::None, + }; + let format_version = i64::from(raw_transaction.payload().format_version()); + let content_hash = raw_transaction.payload().content_hash(); + let content_hash_bytes: &[u8] = content_hash.as_bytes(); + let payload_bytes = raw_transaction.payload().bytes(); + let row_result = sql_transaction + .query_opt( + INSERT_TRANSACTION_SQL, + &[ + &signature_bytes, + &slot_text.as_str(), + &block_time, + &raw_transaction.payload().format_id().as_str(), + &format_version, + &content_hash_bytes, + &payload_bytes, + ], + ) + .await; + return match row_result { + std::result::Result::Ok(std::option::Option::Some(_)) => std::result::Result::Ok(true), + std::result::Result::Ok(std::option::Option::None) => std::result::Result::Ok(false), + std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_acquisition_insert_transaction")), + }; +} + +async fn load_locked_transaction_row( + sql_transaction: &deadpool_postgres::Transaction<'_>, + reference: &ksp_store_api::RawTransactionReference, +) -> std::result::Result, crate::PostgresBackendError> { + let signature = reference.signature(); + let signature_bytes: &[u8] = signature.as_bytes(); + let row_result = sql_transaction.query_opt(LOCK_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::Ok(std::option::Option::None), + std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_acquisition_lock_transaction")), + }; + let mut physical = match raw_transaction_db_row(&row) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let archive_result = sql_transaction.query_opt(GET_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await; + physical.archive_payload = match archive_result { + std::result::Result::Ok(std::option::Option::Some(archive_row)) => match archive_row.try_get::<_, std::vec::Vec>("payload") { + std::result::Result::Ok(value) => std::option::Option::Some(value), + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_acquisition_archive_decode")), + }, + std::result::Result::Ok(std::option::Option::None) => std::option::Option::None, + std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_acquisition_archive_query")), + }; + return std::result::Result::Ok(std::option::Option::Some(physical)); +} + +fn compare_existing_transaction( + network: &ksp_store_api::RawNetworkId, + row: RawTransactionDbRow, + incoming: &ksp_store_api::RawTransaction, +) -> std::result::Result { + let state = match decode_retention_state(row.retention_state.as_str()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + if state == ksp_store_api::RawRetentionState::Purged { + if row.payload.is_some() || row.archive_payload.is_some() || row.block_time_unix_millis.is_some() { + return std::result::Result::Err(data_invalid("raw_acquisition_purged_shape")); + } + let signature = match fixed_bytes::<64>(row.signature) { + std::result::Result::Ok(value) => ksp_store_api::RawTransactionSignature::new(value), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let slot = match decode_u64_decimal(row.slot_text.as_str(), "raw_acquisition_purged_slot") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let format_id = match decode_format_id(row.format_id) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let format_version = match decode_u32_i64(row.format_version, "raw_acquisition_purged_format_version") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let content_hash = match fixed_bytes::<32>(row.content_hash) { + std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let reference = ksp_store_api::RawTransactionReference::new(network.clone(), signature); + let matches = reference.eq(incoming.reference()) + && slot == incoming.slot() + && format_id.as_str() == incoming.payload().format_id().as_str() + && format_version == incoming.payload().format_version() + && content_hash == incoming.payload().content_hash(); + if matches { + return std::result::Result::Ok(ExistingTransactionMatch::Purged); + } + return std::result::Result::Err(conflict("raw_acquisition_purged_conflict")); + } + let stored_result = decode_raw_transaction_row(network, row); + let stored = match stored_result { + std::result::Result::Ok(std::option::Option::Some(value)) => value, + std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(data_invalid("raw_acquisition_active_shape")), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + if raw_transactions_equal(&stored, incoming) { + return std::result::Result::Ok(ExistingTransactionMatch::Active); + } + return std::result::Result::Err(conflict("raw_acquisition_content_conflict")); +} + +fn raw_transactions_equal(left: &ksp_store_api::RawTransaction, right: &ksp_store_api::RawTransaction) -> bool { + return left.reference() == right.reference() + && left.slot() == right.slot() + && left.block_time() == right.block_time() + && left.payload().format_id() == right.payload().format_id() + && left.payload().format_version() == right.payload().format_version() + && left.payload().content_hash() == right.payload().content_hash() + && left.payload().bytes() == right.payload().bytes(); +} + +async fn rehydrate_transaction( + sql_transaction: &deadpool_postgres::Transaction<'_>, + raw_transaction: &ksp_store_api::RawTransaction, +) -> std::result::Result<(), crate::PostgresBackendError> { + let signature = raw_transaction.reference().signature(); + let signature_bytes: &[u8] = signature.as_bytes(); + let block_time = match raw_transaction.block_time() { + std::option::Option::Some(value) => match i64::try_from(value.unix_millis()) { + std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_rehydrate_block_time")), + }, + std::option::Option::None => std::option::Option::None, + }; + let payload_bytes = raw_transaction.payload().bytes(); + let update_result = sql_transaction.execute(REHYDRATE_TRANSACTION_SQL, &[&signature_bytes, &block_time, &payload_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_rehydrate_cardinality")), + std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_rehydrate_update")), + }; +} + +async fn persist_observation_row( + sql_transaction: &deadpool_postgres::Transaction<'_>, + network: &ksp_store_api::RawNetworkId, + observation: &ksp_store_api::RawTransactionObservation, +) -> std::result::Result { + let provenance = observation.provenance(); + let observation_key = observation.observation_key(); + let observation_key_bytes: &[u8] = observation_key.as_bytes(); + let signature = observation.transaction().signature(); + let signature_bytes: &[u8] = signature.as_bytes(); + let origin = match encode_origin(provenance.origin()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let received_at = match i64::try_from(provenance.received_at().unix_millis()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_observation_received_at_encode")), + }; + let observed_at = match provenance.observed_at() { + std::option::Option::Some(value) => match i64::try_from(value.unix_millis()) { + std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_observation_observed_at_encode")), + }, + std::option::Option::None => std::option::Option::None, + }; + let source_payload_size = match provenance.source_payload_size_bytes() { + std::option::Option::Some(value) => match i64::try_from(value) { + std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_observation_source_size_encode")), + }, + std::option::Option::None => std::option::Option::None, + }; + let capture_session_id = provenance.capture_session_id().map(|value| return value.as_str()); + let commitment = provenance.commitment().map(|value| return value.as_str()); + let endpoint_id = provenance.endpoint_id().map(|value| return value.as_str()); + let filter_id = provenance.filter_id().map(|value| return value.as_str()); + let source_payload_hash = provenance.source_payload_hash(); + let source_payload_hash_bytes: std::option::Option<&[u8]> = source_payload_hash.as_ref().map(|value| return &value.as_bytes()[..]); + let insert_result = sql_transaction + .query_opt( + INSERT_OBSERVATION_SQL, + &[ + &observation_key_bytes, + &signature_bytes, + &provenance.provider().as_str(), + &provenance.protocol().as_str(), + &provenance.acquisition_method().as_str(), + &origin, + &received_at, + &capture_session_id, + &commitment, + &endpoint_id, + &filter_id, + &observed_at, + &source_payload_hash_bytes, + &source_payload_size, + ], + ) + .await; + let inserted = match insert_result { + std::result::Result::Ok(std::option::Option::Some(_)) => true, + std::result::Result::Ok(std::option::Option::None) => false, + std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_observation_insert")), + }; + if inserted { + return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::Inserted); + } + let existing_result = sql_transaction.query_opt(LOCK_OBSERVATION_SQL, &[&observation_key_bytes]).await; + let existing_row = match existing_result { + std::result::Result::Ok(std::option::Option::Some(value)) => value, + std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(data_invalid("raw_observation_conflict_missing")), + std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_observation_conflict_query")), + }; + let physical = match raw_observation_db_row(&existing_row) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let stored = match decode_raw_observation_row(network, physical) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + if stored.eq(observation) { + return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::AlreadyPresent); + } + return std::result::Result::Err(conflict("raw_observation_content_conflict")); +} + +fn ensure_acquisition_inputs( + network: &ksp_store_api::RawNetworkId, + raw_transaction: &ksp_store_api::RawTransaction, + observation: &ksp_store_api::RawTransactionObservation, +) -> std::result::Result<(), crate::PostgresBackendError> { + let transaction_network_result = ensure_network(network, raw_transaction.reference(), "raw_acquisition_transaction_network"); + if let std::result::Result::Err(error) = transaction_network_result { + return std::result::Result::Err(error); + } + let observation_network_result = ensure_network(network, observation.transaction(), "raw_acquisition_observation_network"); + if let std::result::Result::Err(error) = observation_network_result { + return std::result::Result::Err(error); + } + if observation.transaction() != raw_transaction.reference() { + return std::result::Result::Err(conflict("raw_acquisition_reference_mismatch")); + } + 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"), + ksp_store_api::RawAcquisitionOrigin::Import => std::result::Result::Ok("import"), + ksp_store_api::RawAcquisitionOrigin::Live => std::result::Result::Ok("live"), + ksp_store_api::RawAcquisitionOrigin::Repair => std::result::Result::Ok("repair"), + ksp_store_api::RawAcquisitionOrigin::Replay => std::result::Result::Ok("replay"), + _ => std::result::Result::Err(data_invalid("raw_observation_origin_encode")), + }; +} + +fn conflict(phase: &'static str) -> crate::PostgresBackendError { + return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::Conflict, phase); +} + +fn reference_not_found(phase: &'static str) -> crate::PostgresBackendError { + return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReferenceNotFound, phase); +} + +fn write_failed(phase: &'static str) -> crate::PostgresBackendError { + return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WriteFailed, phase); +} + fn raw_transaction_db_row(row: &tokio_postgres::Row) -> std::result::Result { let signature = match row.try_get::<_, std::vec::Vec>("signature") { std::result::Result::Ok(value) => value, diff --git a/crates/ksp-store-postgres-lib/src/runtime.rs b/crates/ksp-store-postgres-lib/src/runtime.rs index 549109b..dc141cd 100644 --- a/crates/ksp-store-postgres-lib/src/runtime.rs +++ b/crates/ksp-store-postgres-lib/src/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/src/runtime.rs -// version: 6 +// version: 7 const APPLICATION_NAME: &str = "ksp-store"; const MAX_CONNECTION_URI_BYTES: usize = 4_096; @@ -364,6 +364,24 @@ impl PostgresBackend { return crate::get_raw_transaction_tombstone(&self.pool, &self.network, reference).await; } + /// Persists one canonical RAW transaction and its acquisition observation atomically. + pub async fn persist_raw_transaction_acquisition( + &self, + raw_transaction: ksp_store_api::RawTransaction, + observation: ksp_store_api::RawTransactionObservation, + mode: ksp_store_api::RawTransactionAcquisitionMode, + ) -> std::result::Result { + return crate::persist_raw_transaction_acquisition(&self.pool, &self.network, raw_transaction, observation, mode).await; + } + + /// Persists one additional acquisition observation for an existing RAW transaction. + pub async fn record_raw_transaction_observation( + &self, + observation: ksp_store_api::RawTransactionObservation, + ) -> std::result::Result { + return crate::record_raw_transaction_observation(&self.pool, &self.network, observation).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(); diff --git a/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs b/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs index 352c5b3..cbe6c85 100644 --- a/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs +++ b/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/tests/dependency_boundary.rs -// version: 9 +// version: 10 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -125,7 +125,7 @@ fn pre_003_fix_001_migration_engine_uses_split_schema_contract_and_binds_network } #[test] -fn pre_004_raw_read_sql_and_mapping_remain_backend_private_and_read_only() { +fn pre_004_raw_read_sql_and_mapping_remain_backend_private() { let crate_root = include_str!("../src/lib.rs"); let raw = include_str!("../src/raw_transaction.rs"); assert!(crate_root.contains("mod raw_transaction;")); @@ -141,8 +141,41 @@ fn pre_004_raw_read_sql_and_mapping_remain_backend_private_and_read_only() { ] { assert!(raw.contains(required), "missing private RAW read mapping contract: {required}"); } - for forbidden in ["INSERT INTO", "UPDATE ", "DELETE FROM", "std::env", "dotenv", "ksp_store_lib", "ksp_config_lib"] { - assert!(!raw.contains(forbidden), "pre.004 RAW read module contains forbidden ownership/write material: {forbidden}"); + for forbidden in ["std::env", "dotenv", "ksp_store_lib", "ksp_config_lib"] { + assert!(!raw.contains(forbidden), "RAW module contains forbidden ownership material: {forbidden}"); + } + return; +} + +#[test] +fn pre_005_raw_write_sql_is_atomic_idempotent_and_keeps_later_scope_closed() { + let raw = include_str!("../src/raw_transaction.rs"); + for required in [ + "INSERT INTO ksp_raw_transactions", + "ON CONFLICT (signature) DO NOTHING RETURNING signature", + "FOR UPDATE", + "INSERT INTO ksp_raw_transaction_observations", + "ON CONFLICT (observation_key) DO NOTHING RETURNING observation_key", + "REHYDRATE_TRANSACTION_SQL", + "RawEntityWriteOutcome::SkippedPurged", + "RawEntityWriteOutcome::Rehydrated", + "RawObservationWriteOutcome::NotRecorded", + "PostgresBackendErrorKind::Conflict", + "PostgresBackendErrorKind::ReferenceNotFound", + "PostgresBackendErrorKind::WriteFailed", + ] { + assert!(raw.contains(required), "missing pre.005 RAW write contract: {required}"); + } + for forbidden in [ + "DELETE FROM", + "list_raw_transactions", + "RawPage", + "RawCursor", + "transition_raw_transaction_retention", + "impl ksp_store_api::RawTransactionWrite", + "impl ksp_store_api::RawTransactionObservationWrite", + ] { + assert!(!raw.contains(forbidden), "pre.005 opened later RAW scope prematurely: {forbidden}"); } return; } diff --git a/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs b/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs index 332a7d1..90be626 100644 --- a/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs +++ b/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/tests/hardening_completeness.rs -// version: 4 +// version: 5 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -177,7 +177,7 @@ fn pre_009_backend_error_bridge_cannot_retain_external_error_or_secret_text() { } #[test] -fn pre_009_backend_has_no_env_bypass_or_business_write_capability() { +fn pre_009_backend_has_no_env_bypass_or_direct_store_trait_implementation() { let production = std::format!( "{} {} @@ -214,7 +214,7 @@ fn pre_009_backend_has_no_env_bypass_or_business_write_capability() { "impl ksp_store_api::RawTransaction", "impl ksp_store_api::RawAccount", ] { - assert!(!production.contains(forbidden), "forbidden backend ownership/capability material detected: {forbidden}"); + assert!(!production.contains(forbidden), "forbidden backend ownership/direct-trait material detected: {forbidden}"); } let bootstrap_sql = include_str!("../migrations/v000_bootstrap/tables/001_ksp_store_schema_migrations.sql"); assert!(bootstrap_sql.contains("ksp_store_schema_migrations")); diff --git a/crates/ksp-store-postgres-lib/tests/public_api.rs b/crates/ksp-store-postgres-lib/tests/public_api.rs index 5f73206..2e8a6bf 100644 --- a/crates/ksp-store-postgres-lib/tests/public_api.rs +++ b/crates/ksp-store-postgres-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/tests/public_api.rs -// version: 5 +// version: 6 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -38,18 +38,21 @@ fn pre_005_backend_error_projection_is_safe_and_static() { let kinds = [ ksp_store_postgres_lib::PostgresBackendErrorKind::ConfigInvalid, ksp_store_postgres_lib::PostgresBackendErrorKind::ConnectFailed, + ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict, ksp_store_postgres_lib::PostgresBackendErrorKind::DataInvalid, ksp_store_postgres_lib::PostgresBackendErrorKind::PoolTimeout, ksp_store_postgres_lib::PostgresBackendErrorKind::HealthFailed, ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationFailed, ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationMismatch, ksp_store_postgres_lib::PostgresBackendErrorKind::ReadFailed, + ksp_store_postgres_lib::PostgresBackendErrorKind::ReferenceNotFound, ksp_store_postgres_lib::PostgresBackendErrorKind::SchemaNewer, ksp_store_postgres_lib::PostgresBackendErrorKind::ShutdownTimeout, ksp_store_postgres_lib::PostgresBackendErrorKind::TlsFailed, + ksp_store_postgres_lib::PostgresBackendErrorKind::WriteFailed, ksp_store_postgres_lib::PostgresBackendErrorKind::WrongNetwork, ]; - assert_eq!(kinds.len(), 12); + assert_eq!(kinds.len(), 15); return; } @@ -77,3 +80,10 @@ fn pre_004_raw_read_bridge_uses_only_backend_independent_models() { let _tombstone = ksp_store_postgres_lib::PostgresBackend::get_raw_transaction_tombstone; return; } + +#[test] +fn pre_005_raw_write_bridge_uses_only_backend_independent_models_and_outcomes() { + let _acquisition = ksp_store_postgres_lib::PostgresBackend::persist_raw_transaction_acquisition; + let _observation = ksp_store_postgres_lib::PostgresBackend::record_raw_transaction_observation; + return; +} diff --git a/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs b/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs index 1f07cad..b698ccc 100644 --- a/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs +++ b/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs -// version: 1 +// version: 2 fn network() -> ksp_store_api::RawNetworkId { return match ksp_store_api::RawNetworkId::new("devnet") { @@ -181,3 +181,175 @@ fn pre_004_wrong_network_is_rejected_by_the_private_pre_io_guard() { assert_eq!(rejected.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::WrongNetwork)); return; } + +fn raw_transaction(signature_byte: u8, payload_bytes: &[u8], content_hash_byte: u8) -> ksp_store_api::RawTransaction { + let reference = ksp_store_api::RawTransactionReference::new(network(), ksp_store_api::RawTransactionSignature::new([signature_byte; 64])); + let format_id = match ksp_store_api::RawFormatId::new("ksp.raw.transaction") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid test format rejected: {error:?}"), + }; + let payload = match ksp_store_api::RawPayload::try_new( + format_id, + 1, + payload_bytes.to_vec().into_boxed_slice(), + ksp_store_api::RawContentHash::new([content_hash_byte; 32]), + ) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid test payload rejected: {error:?}"), + }; + return ksp_store_api::RawTransaction::new(reference, 42, std::option::Option::None, payload); +} + +fn observation(reference: ksp_store_api::RawTransactionReference, key_byte: u8, provider: &str) -> ksp_store_api::RawTransactionObservation { + let provider = match ksp_store_api::RawProvenanceCode::new(provider) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid test provider rejected: {error:?}"), + }; + let protocol = match ksp_store_api::RawProvenanceCode::new("solana-ws") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid test protocol rejected: {error:?}"), + }; + let method = match ksp_store_api::RawProvenanceCode::new("transactionSubscribe") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid test method rejected: {error:?}"), + }; + let received_at = match ksp_store_api::RawTimestamp::from_unix_millis(1_700_000_000_000) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid test timestamp rejected: {error:?}"), + }; + let provenance = ksp_store_api::RawAcquisitionProvenance::new(provider, protocol, method, ksp_store_api::RawAcquisitionOrigin::Live, received_at); + return ksp_store_api::RawTransactionObservation::new(ksp_store_api::RawObservationKey::new([key_byte; 32]), reference, provenance); +} + +#[test] +fn pre_005_atomic_acquisition_pre_io_guard_requires_backend_network_and_exact_reference() { + let backend_network = network(); + let raw_transaction = raw_transaction(1, &[1, 2, 3], 7); + let matching = observation(raw_transaction.reference().clone(), 3, "publicnode"); + assert!(super::ensure_acquisition_inputs(&backend_network, &raw_transaction, &matching).is_ok()); + let other_reference = ksp_store_api::RawTransactionReference::new(backend_network.clone(), ksp_store_api::RawTransactionSignature::new([2; 64])); + let mismatched = observation(other_reference, 4, "publicnode"); + let mismatch = super::ensure_acquisition_inputs(&backend_network, &raw_transaction, &mismatched); + assert_eq!(mismatch.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::Conflict)); + let other_network = match ksp_store_api::RawNetworkId::new("mainnet-beta") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid alternate network rejected: {error:?}"), + }; + let foreign_reference = ksp_store_api::RawTransactionReference::new(other_network, ksp_store_api::RawTransactionSignature::new([1; 64])); + let foreign = observation(foreign_reference, 5, "publicnode"); + let wrong_network = super::ensure_acquisition_inputs(&backend_network, &raw_transaction, &foreign); + assert_eq!(wrong_network.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::WrongNetwork)); + return; +} + +#[test] +fn pre_005_existing_full_content_requires_exact_payload_equality_not_hash_only() { + let network = network(); + let incoming = raw_transaction(9, &[1, 2, 3, 4], 7); + let matching = super::compare_existing_transaction(&network, transaction_row("full"), &incoming); + assert!(matching.is_err(), "fixture intentionally differs in slot/content and must conflict"); + let mut exact_row = transaction_row("full"); + exact_row.signature = vec![9; 64]; + exact_row.slot_text = "42".to_owned(); + exact_row.block_time_unix_millis = std::option::Option::None; + exact_row.format_version = 1; + exact_row.payload = std::option::Option::Some(vec![1, 2, 3, 4]); + let exact = super::compare_existing_transaction(&network, exact_row, &incoming); + assert!(matches!(exact, std::result::Result::Ok(super::ExistingTransactionMatch::Active))); + let mut same_hash_different_bytes = transaction_row("full"); + same_hash_different_bytes.signature = vec![9; 64]; + same_hash_different_bytes.slot_text = "42".to_owned(); + same_hash_different_bytes.block_time_unix_millis = std::option::Option::None; + same_hash_different_bytes.format_version = 1; + same_hash_different_bytes.payload = std::option::Option::Some(vec![9, 9, 9, 9]); + let conflict = super::compare_existing_transaction(&network, same_hash_different_bytes, &incoming); + assert_eq!(conflict.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::Conflict)); + return; +} + +#[test] +fn pre_005_purged_tombstone_matches_only_retained_identity_metadata() { + let network = network(); + let incoming = raw_transaction(9, &[1, 2, 3, 4], 7); + let mut purged = transaction_row("purged"); + purged.signature = vec![9; 64]; + purged.slot_text = "42".to_owned(); + purged.block_time_unix_millis = std::option::Option::None; + purged.format_version = 1; + purged.payload = std::option::Option::None; + purged.archive_payload = std::option::Option::None; + let compatible = super::compare_existing_transaction(&network, purged, &incoming); + assert!(matches!(compatible, std::result::Result::Ok(super::ExistingTransactionMatch::Purged))); + let mut divergent = transaction_row("purged"); + divergent.signature = vec![9; 64]; + divergent.slot_text = "42".to_owned(); + divergent.block_time_unix_millis = std::option::Option::None; + divergent.format_version = 1; + divergent.content_hash = vec![8; 32]; + divergent.payload = std::option::Option::None; + divergent.archive_payload = std::option::Option::None; + let conflict = super::compare_existing_transaction(&network, divergent, &incoming); + assert_eq!(conflict.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::Conflict)); + return; +} + +#[test] +fn pre_005_observation_idempotence_compares_reference_and_complete_provenance() { + let incoming_transaction = raw_transaction(8, &[1, 2, 3], 4); + let incoming = observation(incoming_transaction.reference().clone(), 3, "publicnode"); + let matching_row = super::RawObservationDbRow { + acquisition_method: "transactionSubscribe".to_owned(), + capture_session_id: std::option::Option::None, + commitment: std::option::Option::None, + endpoint_id: std::option::Option::None, + filter_id: std::option::Option::None, + observation_key: vec![3; 32], + observed_at_unix_millis: std::option::Option::None, + origin: "live".to_owned(), + protocol: "solana-ws".to_owned(), + provider: "publicnode".to_owned(), + received_at_unix_millis: 1_700_000_000_000, + source_payload_hash: std::option::Option::None, + source_payload_size_bytes: std::option::Option::None, + transaction_signature: vec![8; 64], + }; + let matching = super::decode_raw_observation_row(&network(), matching_row); + let matching = match matching { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid matching observation rejected: {error:?}"), + }; + assert!(matching.eq(&incoming)); + let divergent_row = super::RawObservationDbRow { + acquisition_method: "transactionSubscribe".to_owned(), + capture_session_id: std::option::Option::None, + commitment: std::option::Option::None, + endpoint_id: std::option::Option::None, + filter_id: std::option::Option::None, + observation_key: vec![3; 32], + observed_at_unix_millis: std::option::Option::None, + origin: "live".to_owned(), + protocol: "solana-ws".to_owned(), + provider: "another-provider".to_owned(), + received_at_unix_millis: 1_700_000_000_000, + source_payload_hash: std::option::Option::None, + source_payload_size_bytes: std::option::Option::None, + transaction_signature: vec![8; 64], + }; + let divergent = super::decode_raw_observation_row(&network(), divergent_row); + let divergent = match divergent { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid divergent observation rejected: {error:?}"), + }; + assert!(!divergent.eq(&incoming)); + return; +} + +#[test] +fn pre_005_observation_origin_encoding_is_exact_and_static() { + assert_eq!(super::encode_origin(ksp_store_api::RawAcquisitionOrigin::Backfill), std::result::Result::Ok("backfill")); + assert_eq!(super::encode_origin(ksp_store_api::RawAcquisitionOrigin::Import), std::result::Result::Ok("import")); + assert_eq!(super::encode_origin(ksp_store_api::RawAcquisitionOrigin::Live), std::result::Result::Ok("live")); + assert_eq!(super::encode_origin(ksp_store_api::RawAcquisitionOrigin::Repair), std::result::Result::Ok("repair")); + assert_eq!(super::encode_origin(ksp_store_api::RawAcquisitionOrigin::Replay), std::result::Result::Ok("replay")); + return; +} diff --git a/deltas/0.3.3/pre.005.md b/deltas/0.3.3/pre.005.md new file mode 100644 index 0000000..b134f16 --- /dev/null +++ b/deltas/0.3.3/pre.005.md @@ -0,0 +1,324 @@ + + + +# Delta `0.3.3-pre.005` — écriture atomique RawTransaction + observation + +## 1. Base requise + +```text +0.3.3-pre.4 +``` + +Le gate opérateur fourni pour `pre.004` est entièrement vert : + +```text +cargo fmt --all PASS +audit Rust général / exports / workspace PASS +audit Markdown PASS — 214 tables / 137 files +cargo check --workspace PASS +cargo clippy --workspace --all-targets PASS +cargo test -p ksp-store-api PASS +cargo test -p ksp-store-lib PASS +cargo test -p ksp-store-postgres-lib PASS — 25 unit tests + canaris, live ignored +cargo test -p ksp-config-lib PASS — 128 unit tests + ownership/public API +cargo check -p ksp-store-lib --no-default-features PASS +``` + +Les lectures physiques `RawTransaction` de `pre.004` sont donc acquises. + +## 2. Objectif + +Implémenter la tranche écriture de la vertical slice PostgreSQL sans ouvrir encore pagination ni transitions mutantes de rétention : + +```text +persist_raw_transaction_acquisition +record_raw_transaction_observation +``` + +La première opération doit rendre le canonical et son observation atomiques. Les deux opérations doivent être réellement idempotentes et ne jamais confondre une collision de clé avec un contenu identique. + +## 3. Version + +Le workspace passe à : + +```text +0.3.3-pre.5 +``` + +## 4. Pré-I/O + +Avant `pool.get()` : + +```text +transaction.reference.network == backend.network +observation.transaction.network == backend.network +observation.transaction == transaction.reference +``` + +Un mauvais réseau retourne `WrongNetwork`. Une acquisition composée de deux références logiques différentes retourne `Conflict` sans toucher PostgreSQL. + +`record_raw_transaction_observation` valide également son réseau avant toute acquisition du pool. + +## 5. Acquisition atomique canonique + observation + +L'algorithme est : + +```text +BEGIN + INSERT canonical + ON CONFLICT (signature) DO NOTHING + RETURNING signature + + si inséré + -> candidate Inserted + + sinon + SELECT canonical FOR UPDATE + SELECT archive payload éventuel + comparer le contenu réel + + écrire/idempotenter observation + +COMMIT +``` + +Toute erreur après l'insert canonique, notamment une collision d'observation divergente, fait sortir sans commit : PostgreSQL rollbacke donc l'ensemble de l'acquisition. + +Aucun `has_*` ni SELECT préventif n'est introduit. + +## 6. Idempotence canonique + +Pour `Full` et `Archived`, l'égalité exige exactement : + +```text +reference +slot +block_time +format_id +format_version +content_hash +payload bytes +``` + +Le hash n'est jamais utilisé seul lorsque les octets restent disponibles. + +Résultats : + +```text +première insertion -> Inserted +même identité + contenu identique -> AlreadyPresent +même identité + contenu divergent -> Conflict +``` + +Le conflit backend sera projeté vers `ERROR_CODE_RAW_CONFLICT` par la façade dans `pre.008`. + +## 7. Tombstone et ForceRehydrate + +Lorsqu'une signature existe en `Purged`, le backend exige que la forme physique soit réellement minimale : + +```text +payload hot absent +payload archive absent +block_time absent +``` + +Il compare ensuite uniquement les métadonnées conservées par le tombstone : + +```text +reference +slot +format_id +format_version +content_hash +``` + +Une divergence produit `Conflict`. + +Pour un tombstone compatible : + +```text +Normal + -> SkippedPurged / NotRecorded + +ForceRehydrate + -> UPDATE payload + block_time + state=full + -> observation dans la même transaction + -> Rehydrated +``` + +Le `ForceRehydrate` ne masque donc jamais une divergence détectable. + +## 8. Observation idempotente + +L'observation utilise : + +```text +INSERT ... +ON CONFLICT (observation_key) DO NOTHING +RETURNING observation_key +``` + +Après collision : + +```text +SELECT observation ... FOR UPDATE +``` + +Puis le modèle reconstruit est comparé intégralement à l'observation entrante : référence transaction + provenance complète. + +Résultats : + +```text +key absente + insert réussi -> Inserted +key présente + contenu identique -> AlreadyPresent +key présente + contenu divergent -> Conflict +``` + +## 9. Observation write séparée + +`record_raw_transaction_observation` ne crée jamais un canonical implicite. + +Sous transaction : + +```text +SELECT retention_state FROM canonical FOR UPDATE +``` + +Puis : + +```text +canonical absent -> ReferenceNotFound +canonical Purged -> NotRecorded +Full / Archived -> insertion/idempotence observation +``` + +Le verrou canonical prépare aussi la sérialisation future avec purge/rehydrate de `pre.007`. + +## 10. Erreurs backend + +`PostgresBackendErrorKind` ajoute : + +```text +Conflict +ReferenceNotFound +WriteFailed +``` + +Ces erreurs restent composées exclusivement de : + +```text +kind +phase &'static str +``` + +Aucun texte PostgreSQL, SQLSTATE, query, bind, URI, signature, hash ou payload n'est retenu. + +## 11. Surface backend + +`PostgresBackend` expose désormais en plus : + +```text +persist_raw_transaction_acquisition +record_raw_transaction_observation +``` + +Les traits `RawTransactionWrite` / `RawTransactionObservationWrite` ne sont pas encore implémentés directement : la tranche finale des six traits attend `list_raw_transactions` en `pre.006`, puis la rétention mutante en `pre.007` et le dispatch façade en `pre.008`. + +## 12. Tests déterministes + +Le miroir `unit_tests/raw_transaction.rs` couvre en plus : + +```text +pré-I/O réseau + référence acquisition exacte +égalité canonique stricte +même hash mais payload différent -> Conflict +tombstone compatible -> Purged match +tombstone divergent -> Conflict +mapping exact des origins d'observation +``` + +Les canaris d'intégration figent : + +```text +INSERT unique direct sans has_* +FOR UPDATE après collision +observation unique idempotente +ForceRehydrate limité à UPDATE +aucun DELETE +aucune pagination/cursor +aucune transition retention mutante +aucune impl complète de trait prématurée +``` + +La concurrence PostgreSQL réelle et les preuves de rollback restent réservées au test live `pre.009`. + +## 13. Migrations + +Aucune ressource de migration n'est modifiée. + +Checksums inchangés : + +```text +V000 d29068b8c13b9dc0cc9ef6aaadd0fa12d41e0fe4c56541a1118c4bfc846a1450 +V001 31488cda2f08f3f46c4cdbdbb6c18c243662fada02eac4487040c8735d72cc51 +``` + +## 14. Fichiers modifiés + +```text +Cargo.toml +crates/ksp-store-postgres-lib/README.md +crates/ksp-store-postgres-lib/USAGE.md +crates/ksp-store-postgres-lib/src/error.rs +crates/ksp-store-postgres-lib/src/lib.rs +crates/ksp-store-postgres-lib/src/raw_transaction.rs +crates/ksp-store-postgres-lib/src/runtime.rs +crates/ksp-store-postgres-lib/tests/dependency_boundary.rs +crates/ksp-store-postgres-lib/tests/hardening_completeness.rs +crates/ksp-store-postgres-lib/tests/public_api.rs +crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs +docs/plans/024-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION_PLAN.md +docs/validation/020-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION.md +``` + +## 15. Fichiers ajoutés + +```text +deltas/0.3.3/pre.005.md +``` + +## 16. Fichiers supprimés + +Aucun. + +## 17. Hors scope confirmé + +Aucun changement n'est apporté à : + +```text +ksp-store-api contrats +ksp-store-lib dispatch métier +list_raw_transactions / cursor +transition Full -> Archived -> Purged +Compacted physique +RawAccountState +worker/job/app +migration SQL V000/V001 +``` + +## 18. Gate opérateur + +```bash +cargo fmt --all +python3 scripts/audit_rust_workspace_rules.py +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas/0.3.3 +cargo check --workspace +cargo clippy --workspace --all-targets +cargo test -p ksp-store-api +cargo test -p ksp-store-lib +cargo test -p ksp-store-postgres-lib +cargo test -p ksp-config-lib +cargo check -p ksp-store-lib --no-default-features +``` + +Le test PostgreSQL live de fondation reste `#[ignore]`. Aucune URI réelle n'est requise pour cette tranche. diff --git a/docs/plans/024-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION_PLAN.md b/docs/plans/024-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION_PLAN.md index 7d294c1..bef221f 100644 --- a/docs/plans/024-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION_PLAN.md +++ b/docs/plans/024-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION_PLAN.md @@ -1,5 +1,5 @@ - + # Plan `0.3.3` — Store/PostgreSQL RawTransaction vertical slice @@ -1090,11 +1090,21 @@ Tranche matérialisée : ### `0.3.3-pre.005` — écriture atomique et observations -- acquisition atomique canonique + observation ; -- idempotence réelle ; -- divergence -> conflict ; -- observation write séparée ; -- concurrence logique testable sans API `has_*`. +Tranche matérialisée : + +- `PostgresBackend::persist_raw_transaction_acquisition` et `record_raw_transaction_observation` utilisent uniquement les modèles/outcomes backend-agnostiques ; +- validation pré-I/O des réseaux et de l'égalité `observation.transaction == transaction.reference` ; +- acquisition canonique + observation dans une transaction PostgreSQL unique ; +- insert canonique direct `ON CONFLICT (signature) DO NOTHING RETURNING`, sans API/prélecture `has_*` ; +- après conflit unique, `SELECT ... FOR UPDATE` puis comparaison de `slot`, `block_time`, format, version, hash et payload réel lorsque disponible ; +- `Full`/`Archived` identiques -> `AlreadyPresent`, divergence -> `Conflict` ; +- tombstone `Purged` compatible + `Normal` -> `SkippedPurged/NotRecorded` ; +- tombstone compatible + `ForceRehydrate` -> restauration `Full` puis observation dans la même transaction ; +- observation insert direct `ON CONFLICT (observation_key) DO NOTHING RETURNING`, puis verrou/lecture et comparaison de toute la référence/provenance ; +- observation canonique absente -> `ReferenceNotFound`, canonique purgé -> `NotRecorded` ; +- nouveaux kinds backend sûrs `Conflict`, `ReferenceNotFound`, `WriteFailed` ; +- aucun `DELETE`, aucune pagination/cursor et aucune transition de rétention mutante dans cette tranche ; +- V000/V001 et leurs checksums restent byte-identiques. ### `0.3.3-pre.006` — pagination/cursor diff --git a/docs/validation/020-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION.md b/docs/validation/020-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION.md index b9c186d..6333779 100644 --- a/docs/validation/020-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION.md +++ b/docs/validation/020-V0_3_3_STORE_POSTGRES_RAW_TRANSACTION.md @@ -1,5 +1,5 @@ - + # Validation `0.3.3` — Store/PostgreSQL RawTransaction vertical slice @@ -562,18 +562,25 @@ cap 500/1000 dans Store pagination ### `pre.004` -- quatre lectures RAW backend-specific sans fuite de row/SQL : PASS statique ; -- conversion `NUMERIC(20,0) -> u64` jusqu'à `u64::MAX` et conversions entières/timestamps fallibles : PASS unit design ; -- `Full/Archived/Purged`, observation complète, rétention/tombstone : PASS unit design ; -- mauvais réseau avant pool I/O : PASS unit design ; -- malformed DB -> `DataInvalid`, SELECT -> `ReadFailed`, aucune valeur hostile retenue : PASS statique/unit design ; -- SQL strictement read-only dans cette tranche : PASS canari source ; -- gate Cargo opérateur : À EXÉCUTER. +- quatre lectures RAW backend-specific sans fuite de row/SQL : PASS ; +- conversion `NUMERIC(20,0) -> u64` jusqu'à `u64::MAX` et conversions entières/timestamps fallibles : PASS ; +- `Full/Archived/Purged`, observation complète, rétention/tombstone : PASS ; +- mauvais réseau avant pool I/O : PASS ; +- malformed DB -> `DataInvalid`, SELECT -> `ReadFailed`, aucune valeur hostile retenue : PASS ; +- SQL strictement read-only dans cette tranche : PASS ; +- gate opérateur complet du 2026-08-30 : PASS (`check`, Clippy, Store/API/PostgreSQL/Config, no-default-features ; 25 tests unit backend). ### `pre.005` -- atomic writes/idempotence/conflict ; -- observation writes. +- acquisition transaction + observation atomique sous transaction SQL : PASS statique/design ; +- insert unique sans `has_*`, puis verrouillage `FOR UPDATE` et comparaison réelle : PASS statique/unit design ; +- même contenu -> `AlreadyPresent`, divergence -> `Conflict` : PASS unit design ; +- tombstone compatible `Normal` -> `SkippedPurged/NotRecorded` : PASS statique/design ; +- tombstone compatible `ForceRehydrate` -> `Rehydrated` + observation atomique : PASS statique/design ; +- observation write séparée, key identique/divergente, absent/purged : PASS statique/unit design ; +- `Conflict` / `ReferenceNotFound` / `WriteFailed` sans texte serveur : PASS statique ; +- preuve PostgreSQL concurrente/rollback réelle : différée à `pre.009` ; +- gate Cargo opérateur : À EXÉCUTER. ### `pre.006` @@ -615,7 +622,7 @@ cap 500/1000 dans Store pagination - publication stable. -## 22. Gate courant `pre.004` +## 22. Gate courant `pre.005` ```bash cargo fmt --all @@ -667,4 +674,17 @@ cargo check -p ksp-store-lib --no-default-features - [PASS] wrong-network gardé avant pool I/O ; - [PASS] `ReadFailed` / `DataInvalid` / `WrongNetwork` sans texte externe ; - [PASS] tests déterministes de row mapping et hostile data ajoutés ; -- [À FAIRE] gate Cargo opérateur complet de `pre.004`. +- [PASS] gate Cargo opérateur complet de `pre.004` fourni le 2026-08-30. + +### `pre.005` — état de la tranche + +- [PASS] écritures backend-specific canonique + observation ajoutées sans implémenter encore les traits incomplets ; +- [PASS] validation réseau et cohérence transaction/observation avant `pool.get()` ; +- [PASS] idempotence par insert unique + comparaison réelle sous `FOR UPDATE`, sans `has_*` ; +- [PASS] payload comparé byte-for-byte lorsque disponible ; +- [PASS] `Purged` normal skip et `ForceRehydrate` atomique selon tombstone ; +- [PASS] observation write séparée : `Inserted` / `AlreadyPresent` / `NotRecorded` / conflit ; +- [PASS] référence canonique absente classée `ReferenceNotFound` ; +- [PASS] erreurs physiques d'écriture classées `WriteFailed` sans conserver le texte PostgreSQL ; +- [PASS] aucune migration SQL modifiée et aucun scope pagination/rétention mutante ouvert ; +- [À FAIRE] gate Cargo opérateur complet de `pre.005`.