v0.3.4-pre.004-fix.001

This commit is contained in:
2026-08-30 22:06:24 +02:00
parent fb8b585aa2
commit eb6dbc31e8
6 changed files with 283 additions and 43 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/src/raw_account.rs
// version: 1
// version: 2
const GET_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key, account_pubkey, account_slot::text AS account_slot_text, account_state_hash, 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, is_startup, transaction_signature, write_version::text AS write_version_text FROM ksp_raw_account_observations WHERE observation_key = $1";
const GET_ACCOUNT_STATE_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash, lamports::text AS lamports_text, owner, executable, rent_epoch::text AS rent_epoch_text, data FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3";
@@ -131,40 +131,139 @@ pub(crate) async fn get_raw_account_observation(
}
fn raw_account_observation_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawAccountObservationDbRow, crate::PostgresBackendError> {
let account_pubkey: std::vec::Vec<u8> = match row.try_get("account_pubkey") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let account_slot_text: std::string::String = match row.try_get("account_slot_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let account_state_hash: std::vec::Vec<u8> = match row.try_get("account_state_hash") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let acquisition_method: std::string::String = match row.try_get("acquisition_method") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let capture_session_id: std::option::Option<std::string::String> = match row.try_get("capture_session_id") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let commitment: std::option::Option<std::string::String> = match row.try_get("commitment") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let endpoint_id: std::option::Option<std::string::String> = match row.try_get("endpoint_id") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let filter_id: std::option::Option<std::string::String> = match row.try_get("filter_id") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let is_startup: std::option::Option<bool> = match row.try_get("is_startup") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let observation_key: std::vec::Vec<u8> = match row.try_get("observation_key") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let observed_at_unix_millis: std::option::Option<i64> = match row.try_get("observed_at_unix_millis") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let origin: std::string::String = match row.try_get("origin") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let protocol: std::string::String = match row.try_get("protocol") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let provider: std::string::String = match row.try_get("provider") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let received_at_unix_millis: i64 = match row.try_get("received_at_unix_millis") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let source_payload_hash: std::option::Option<std::vec::Vec<u8>> = match row.try_get("source_payload_hash") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let source_payload_size_bytes: std::option::Option<i64> = match row.try_get("source_payload_size_bytes") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let transaction_signature: std::option::Option<std::vec::Vec<u8>> = match row.try_get("transaction_signature") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
let write_version_text: std::option::Option<std::string::String> = match row.try_get("write_version_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")),
};
return std::result::Result::Ok(RawAccountObservationDbRow {
account_pubkey: row.try_get("account_pubkey").map_err(|_| data_invalid("raw_account_observation_decode"))?,
account_slot_text: row.try_get("account_slot_text").map_err(|_| data_invalid("raw_account_observation_decode"))?,
account_state_hash: row.try_get("account_state_hash").map_err(|_| data_invalid("raw_account_observation_decode"))?,
acquisition_method: row.try_get("acquisition_method").map_err(|_| data_invalid("raw_account_observation_decode"))?,
capture_session_id: row.try_get("capture_session_id").map_err(|_| data_invalid("raw_account_observation_decode"))?,
commitment: row.try_get("commitment").map_err(|_| data_invalid("raw_account_observation_decode"))?,
endpoint_id: row.try_get("endpoint_id").map_err(|_| data_invalid("raw_account_observation_decode"))?,
filter_id: row.try_get("filter_id").map_err(|_| data_invalid("raw_account_observation_decode"))?,
is_startup: row.try_get("is_startup").map_err(|_| data_invalid("raw_account_observation_decode"))?,
observation_key: row.try_get("observation_key").map_err(|_| data_invalid("raw_account_observation_decode"))?,
observed_at_unix_millis: row.try_get("observed_at_unix_millis").map_err(|_| data_invalid("raw_account_observation_decode"))?,
origin: row.try_get("origin").map_err(|_| data_invalid("raw_account_observation_decode"))?,
protocol: row.try_get("protocol").map_err(|_| data_invalid("raw_account_observation_decode"))?,
provider: row.try_get("provider").map_err(|_| data_invalid("raw_account_observation_decode"))?,
received_at_unix_millis: row.try_get("received_at_unix_millis").map_err(|_| data_invalid("raw_account_observation_decode"))?,
source_payload_hash: row.try_get("source_payload_hash").map_err(|_| data_invalid("raw_account_observation_decode"))?,
source_payload_size_bytes: row.try_get("source_payload_size_bytes").map_err(|_| data_invalid("raw_account_observation_decode"))?,
transaction_signature: row.try_get("transaction_signature").map_err(|_| data_invalid("raw_account_observation_decode"))?,
write_version_text: row.try_get("write_version_text").map_err(|_| data_invalid("raw_account_observation_decode"))?,
account_pubkey,
account_slot_text,
account_state_hash,
acquisition_method,
capture_session_id,
commitment,
endpoint_id,
filter_id,
is_startup,
observation_key,
observed_at_unix_millis,
origin,
protocol,
provider,
received_at_unix_millis,
source_payload_hash,
source_payload_size_bytes,
transaction_signature,
write_version_text,
});
}
fn raw_account_state_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawAccountStateDbRow, crate::PostgresBackendError> {
return std::result::Result::Ok(RawAccountStateDbRow {
data: row.try_get("data").map_err(|_| data_invalid("raw_account_state_decode"))?,
executable: row.try_get("executable").map_err(|_| data_invalid("raw_account_state_decode"))?,
lamports_text: row.try_get("lamports_text").map_err(|_| data_invalid("raw_account_state_decode"))?,
owner: row.try_get("owner").map_err(|_| data_invalid("raw_account_state_decode"))?,
pubkey: row.try_get("pubkey").map_err(|_| data_invalid("raw_account_state_decode"))?,
rent_epoch_text: row.try_get("rent_epoch_text").map_err(|_| data_invalid("raw_account_state_decode"))?,
slot_text: row.try_get("slot_text").map_err(|_| data_invalid("raw_account_state_decode"))?,
state_hash: row.try_get("state_hash").map_err(|_| data_invalid("raw_account_state_decode"))?,
});
let data: std::vec::Vec<u8> = match row.try_get("data") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let executable: bool = match row.try_get("executable") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let lamports_text: std::string::String = match row.try_get("lamports_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let owner: std::vec::Vec<u8> = match row.try_get("owner") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let pubkey: std::vec::Vec<u8> = match row.try_get("pubkey") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let rent_epoch_text: std::string::String = match row.try_get("rent_epoch_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let slot_text: std::string::String = match row.try_get("slot_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
let state_hash: std::vec::Vec<u8> = match row.try_get("state_hash") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")),
};
return std::result::Result::Ok(RawAccountStateDbRow { data, executable, lamports_text, owner, pubkey, rent_epoch_text, slot_text, state_hash });
}
fn decode_raw_account_observation_row(

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/tests/dependency_boundary.rs
// version: 17
// version: 18
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -127,7 +127,6 @@ fn pre_003_fix_001_migration_engine_uses_split_schema_contract_and_binds_network
#[test]
fn pre_003_v002_schema_is_complete_without_account_trait_or_write_dispatch() {
let crate_root = include_str!("../src/lib.rs");
let migration = include_str!("../src/migration.rs");
let schema = include_str!("../src/schema.rs");
let runtime = include_str!("../src/runtime.rs");