v0.3.4-pre.006

This commit is contained in:
2026-08-30 22:31:02 +02:00
parent 50d4142797
commit fccb7d876c
10 changed files with 364 additions and 35 deletions

View File

@@ -1,11 +1,13 @@
// file: crates/ksp-store-postgres-lib/src/raw_account.rs
// version: 3
// version: 4
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";
const INSERT_ACCOUNT_OBSERVATION_SQL: &str = "INSERT INTO ksp_raw_account_observations (observation_key, account_pubkey, account_slot, 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) VALUES ($1, $2, $3::TEXT::NUMERIC, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19::TEXT::NUMERIC) ON CONFLICT (observation_key) DO NOTHING RETURNING observation_key";
const INSERT_ACCOUNT_STATE_SQL: &str = "INSERT INTO ksp_raw_account_states (pubkey, slot, state_hash, lamports, owner, executable, rent_epoch, data) VALUES ($1, $2::TEXT::NUMERIC, $3, $4::TEXT::NUMERIC, $5, $6, $7::TEXT::NUMERIC, $8) ON CONFLICT (pubkey, slot, state_hash) DO NOTHING RETURNING pubkey";
const LOCK_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 FOR UPDATE";
const LOCK_ACCOUNT_REFERENCE_SQL: &str =
"SELECT 1 FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3 FOR KEY SHARE";
const LOCK_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 FOR UPDATE";
struct RawAccountObservationDbRow {
@@ -107,10 +109,7 @@ pub(crate) async fn get_raw_account_observation(
let rows = match rows_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::ReadFailed,
"raw_account_observation_query",
));
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReadFailed, "raw_account_observation_query"));
},
};
if rows.is_empty() {
@@ -189,6 +188,53 @@ pub(crate) async fn persist_raw_account_acquisition(
return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new(entity_outcome, observation_outcome));
}
/// Persists one additional acquisition observation for an already durable RAW account state.
pub(crate) async fn record_raw_account_observation(
pool: &deadpool_postgres::Pool,
network: &ksp_store_api::RawNetworkId,
observation: ksp_store_api::RawAccountObservation,
) -> std::result::Result<ksp_store_api::RawObservationWriteOutcome, crate::PostgresBackendError> {
let input_result = ensure_observation_write_input(network, &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_account_observation_begin")),
};
let account = observation.account();
let account_pubkey_bytes: &[u8] = account.pubkey().as_ref();
let account_slot_text = account.slot().to_string();
let account_state_hash = account.state_hash();
let account_state_hash_bytes: &[u8] = account_state_hash.as_bytes();
let reference_result = sql_transaction
.query_opt(LOCK_ACCOUNT_REFERENCE_SQL, &[&account_pubkey_bytes, &account_slot_text.as_str(), &account_state_hash_bytes])
.await;
match reference_result {
std::result::Result::Ok(std::option::Option::Some(_)) => {},
std::result::Result::Ok(std::option::Option::None) => {
return std::result::Result::Err(reference_not_found("raw_account_observation_reference"));
},
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_observation_lock_reference")),
}
let observation_result = persist_account_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_account_observation_commit"));
}
return std::result::Result::Ok(outcome);
}
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,
@@ -508,16 +554,7 @@ async fn insert_account_state(
let row_result = sql_transaction
.query_opt(
INSERT_ACCOUNT_STATE_SQL,
&[
&pubkey_bytes,
&slot_text.as_str(),
&state_hash_bytes,
&lamports_text.as_str(),
&owner_bytes,
&executable,
&rent_epoch_text.as_str(),
&data,
],
&[&pubkey_bytes, &slot_text.as_str(), &state_hash_bytes, &lamports_text.as_str(), &owner_bytes, &executable, &rent_epoch_text.as_str(), &data],
)
.await;
return match row_result {
@@ -670,6 +707,13 @@ fn ensure_acquisition_inputs(
return std::result::Result::Ok(());
}
fn ensure_observation_write_input(
network: &ksp_store_api::RawNetworkId,
observation: &ksp_store_api::RawAccountObservation,
) -> std::result::Result<(), crate::PostgresBackendError> {
return ensure_network(network, observation.account(), "raw_account_observation_write_network");
}
fn raw_account_states_equal(left: &ksp_store_api::RawAccountState, right: &ksp_store_api::RawAccountState) -> bool {
return left.reference() == right.reference()
&& left.lamports() == right.lamports()
@@ -759,6 +803,10 @@ fn data_invalid(phase: &'static str) -> crate::PostgresBackendError {
return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::DataInvalid, 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);
}