Files
khadhroony-bot3/kb-store/src/postgres/query/decode_pipeline_queries.rs
2026-07-24 14:23:58 +02:00

1341 lines
69 KiB
Rust

// file: kb-store/src/postgres/query/decode_pipeline_queries.rs
// version: 2
//! PostgreSQL queries for contextual decode, coverage and materialization persistence.
use sqlx::Row; // rust-rules: trait-import
#[derive(sqlx::FromRow)]
struct MaterializedEventDatabaseRow {
processor_name: std::string::String,
processor_version: std::string::String,
input_key: std::string::String,
output_key: std::string::String,
source_event_key: std::string::String,
source_decoder_name: std::string::String,
source_decoder_version: std::string::String,
signature: std::string::String,
slot: i64,
materialized_family: std::string::String,
payload_json: serde_json::Value,
created_at: std::string::String,
updated_at: std::string::String,
}
pub(in crate::postgres) async fn apply_decode_store_schema(
pool: &sqlx::PgPool,
) -> kb_core::Result<()> {
tracing::debug!(target: crate::TRACING_TARGET, action = "apply_decode_store_schema", statement_count = crate::decode_store_schema_statements().len(), "apply PostgreSQL decode store schema");
let validation_result = crate::validate_decode_store_table_names();
if let std::result::Result::Err(error) = validation_result {
return std::result::Result::Err(error);
}
let transaction_result = pool.begin().await;
let mut transaction = match transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode store schema transaction failed: {error}"
)));
},
};
let lock_result = sqlx::query("SELECT pg_advisory_xact_lock($1)")
.bind(crate::STORE_SCHEMA_ADVISORY_LOCK_ID)
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = lock_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode store schema advisory lock failed: {error}"
)));
}
for statement in crate::decode_store_schema_statements() {
let execution_result = sqlx::query(statement).execute(&mut *transaction).await;
if let std::result::Result::Err(error) = execution_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode store schema initialization failed: {error}"
)));
}
}
let commit_result = transaction.commit().await;
if let std::result::Result::Err(error) = commit_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode store schema commit failed: {error}"
)));
}
tracing::debug!(target: crate::TRACING_TARGET, action = "apply_decode_store_schema", committed = true, "PostgreSQL decode store schema applied");
return std::result::Result::Ok(());
}
pub(in crate::postgres) async fn list_decode_inputs(
pool: &sqlx::PgPool,
filter: &crate::DecodeSelectionFilter,
) -> kb_core::Result<std::vec::Vec<crate::MdCoreInstructionReplayInput>> {
tracing::debug!(target: crate::TRACING_TARGET, action = "list_decode_inputs", signature_count = filter.signatures.len(), signature_sample = ?filter.signatures.iter().take(5).map(std::string::String::as_str).collect::<std::vec::Vec<_>>(), processing_states = ?filter.processing_states, min_slot = ?filter.min_slot, max_slot = ?filter.max_slot, program_ids = ?filter.program_ids, instruction_paths = ?filter.instruction_paths, limit = filter.limit, "query PostgreSQL contextual decode inputs");
let result =
crate::postgres::query::core_queries::list_decode_replay_inputs(pool, filter).await;
return match result {
std::result::Result::Ok(inputs) => {
let selected_input_keys = inputs
.iter()
.take(10)
.map(|input| return input.replay_input_key.as_str())
.collect::<std::vec::Vec<_>>();
tracing::debug!(target: crate::TRACING_TARGET, action = "list_decode_inputs", selected_count = inputs.len(), input_key_sample = ?selected_input_keys, "PostgreSQL contextual decode inputs selected");
std::result::Result::Ok(inputs)
},
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "list_decode_inputs", error = %error, "PostgreSQL contextual decode input query failed");
std::result::Result::Err(error)
},
};
}
pub(in crate::postgres) async fn list_materialized_events(
pool: &sqlx::PgPool,
filter: &crate::MaterializedEventFilter,
) -> kb_core::Result<std::vec::Vec<crate::MaterializedEventQueryRow>> {
tracing::debug!(target: crate::TRACING_TARGET, action = "list_materialized_events", processor_name = ?filter.processor_name, materialized_family = ?filter.materialized_family, signature_contains = ?filter.signature_contains, limit = filter.limit, "query bounded PostgreSQL materialized events");
if filter.limit == 0 || filter.limit > crate::MAX_MATERIALIZED_EVENT_QUERY_ROWS {
return std::result::Result::Err(kb_core::Error::db(format!(
"materialized event query limit must be between 1 and {}",
crate::MAX_MATERIALIZED_EVENT_QUERY_ROWS
)));
}
let query_result = sqlx::query_as::<sqlx::Postgres, MaterializedEventDatabaseRow>(
"SELECT processor_name, processor_version, input_key, output_key, source_event_key, source_decoder_name, source_decoder_version, signature, slot, materialized_family, payload_jsonb AS payload_json, created_at::text AS created_at, updated_at::text AS updated_at FROM kb_sol_mat_events WHERE ($1::text IS NULL OR processor_name = $1) AND ($2::text IS NULL OR materialized_family = $2) AND ($3::text IS NULL OR POSITION(LOWER($3) IN LOWER(signature)) > 0) ORDER BY slot DESC, id DESC LIMIT $4",
)
.bind(filter.processor_name.as_deref())
.bind(filter.materialized_family.as_deref())
.bind(filter.signature_contains.as_deref())
.bind(i64::from(filter.limit))
.fetch_all(pool)
.await;
let rows = match query_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialized event query failed: {error}"
)));
},
};
let mut output = std::vec::Vec::with_capacity(rows.len());
for row in rows {
let slot_result = u64::try_from(row.slot);
let slot = match slot_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialized event slot conversion failed: {error}"
)));
},
};
output.push(crate::MaterializedEventQueryRow {
processor_name: row.processor_name,
processor_version: row.processor_version,
input_key: row.input_key,
output_key: row.output_key,
source_event_key: row.source_event_key,
source_decoder_name: row.source_decoder_name,
source_decoder_version: row.source_decoder_version,
signature: row.signature,
slot,
materialized_family: row.materialized_family,
payload_json: row.payload_json,
created_at: row.created_at,
updated_at: row.updated_at,
});
}
tracing::debug!(target: crate::TRACING_TARGET, action = "list_materialized_events", row_count = output.len(), "bounded PostgreSQL materialized events loaded");
return std::result::Result::Ok(output);
}
pub(in crate::postgres) async fn is_decode_current(
pool: &sqlx::PgPool,
identity: &crate::ProcessingLedgerIdentity,
) -> kb_core::Result<bool> {
tracing::debug!(target: crate::TRACING_TARGET, action = "is_decode_current", stage = %identity.stage, processor_name = %identity.processor_name, processor_version = %identity.processor_version, input_key = %identity.input_key, input_hash = %identity.input_hash, "query PostgreSQL processing ledger current state");
let query_result = sqlx::query_scalar::<sqlx::Postgres, bool>(
"SELECT EXISTS(SELECT 1 FROM kb_sol_ops_processing_ledger WHERE stage = $1 AND processor_name = $2 AND processor_version = $3 AND input_key = $4 AND input_hash = $5 AND status = 'succeeded')",
)
.bind(identity.stage.as_str())
.bind(identity.processor_name.as_str())
.bind(identity.processor_version.as_str())
.bind(identity.input_key.as_str())
.bind(identity.input_hash.as_str())
.fetch_one(pool)
.await;
return match query_result {
std::result::Result::Ok(value) => {
tracing::debug!(target: crate::TRACING_TARGET, action = "is_decode_current", stage = %identity.stage, processor_name = %identity.processor_name, processor_version = %identity.processor_version, input_key = %identity.input_key, input_hash = %identity.input_hash, current = value, "PostgreSQL processing ledger current state loaded");
std::result::Result::Ok(value)
},
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "is_decode_current", stage = %identity.stage, processor_name = %identity.processor_name, processor_version = %identity.processor_version, input_key = %identity.input_key, error = %error, "PostgreSQL processing ledger current check failed");
std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode ledger current check failed: {error}"
)))
},
};
}
pub(in crate::postgres) async fn persist_decode_coverage_declarations(
pool: &sqlx::PgPool,
declarations: &[crate::DecodeCoverageDeclarationInsert],
) -> kb_core::Result<crate::InsertOutcome> {
if declarations.is_empty() {
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_coverage_declarations", declaration_count = 0_usize, "skip empty PostgreSQL decode coverage declarations");
return std::result::Result::Ok(crate::InsertOutcome::new(0, 0, 0));
}
let processor_name = declarations[0].processor_name.as_str();
let processor_version = declarations[0].processor_version.as_str();
let program_ids = declarations
.iter()
.map(|entry| return entry.program_id.as_str())
.collect::<std::vec::Vec<_>>();
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_coverage_declarations", processor_name = %processor_name, processor_version = %processor_version, declaration_count = declarations.len(), program_ids = ?program_ids, "persist PostgreSQL decode coverage declarations");
if declarations.iter().any(|entry| {
return entry.processor_name != processor_name
|| entry.processor_version != processor_version
|| entry.program_id.trim().is_empty()
|| entry.entry_kind.trim().is_empty()
|| entry.entry_code.trim().is_empty();
}) {
return std::result::Result::Err(kb_core::Error::db(
"decode coverage declarations must share one valid processor identity",
));
}
let mut desired_keys = std::collections::BTreeSet::new();
for declaration in declarations {
let key = (
declaration.program_id.clone(),
declaration.surface_code.clone(),
declaration.entry_kind.clone(),
declaration.entry_code.clone(),
declaration.discriminator_hex.clone(),
);
if !desired_keys.insert(key) {
return std::result::Result::Err(kb_core::Error::db(
"decode coverage declarations contain a duplicate identity",
));
}
}
let transaction_result = pool.begin().await;
let mut transaction = match transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage transaction failed: {error}"
)));
},
};
let advisory_key = format!("decode_coverage:{processor_name}:{processor_version}");
let lock_result = sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended($1, 0))")
.bind(advisory_key.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = lock_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage declaration lock failed: {error}"
)));
}
let existing_result = sqlx::query(
"SELECT program_id, surface_code, entry_kind, entry_code, discriminator_hex, historical FROM kb_sol_decode_coverage_declarations WHERE processor_name = $1 AND processor_version = $2 FOR UPDATE",
)
.bind(processor_name)
.bind(processor_version)
.fetch_all(&mut *transaction)
.await;
let existing_rows = match existing_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage declaration lookup failed: {error}"
)));
},
};
let mut existing = std::collections::BTreeMap::new();
for row in existing_rows {
let key = (
row.get::<std::string::String, _>("program_id"),
row.get::<std::option::Option<std::string::String>, _>("surface_code"),
row.get::<std::string::String, _>("entry_kind"),
row.get::<std::string::String, _>("entry_code"),
row.get::<std::option::Option<std::string::String>, _>("discriminator_hex"),
);
existing.insert(key, row.get::<bool, _>("historical"));
}
let mut inserted_count = 0_u64;
let mut updated_count = 0_u64;
let mut skipped_count = 0_u64;
for declaration in declarations {
let key = (
declaration.program_id.clone(),
declaration.surface_code.clone(),
declaration.entry_kind.clone(),
declaration.entry_code.clone(),
declaration.discriminator_hex.clone(),
);
let existing_historical = existing.remove(&key);
match existing_historical {
std::option::Option::None => {
let insert_result = sqlx::query(
"INSERT INTO kb_sol_decode_coverage_declarations (processor_name, processor_version, program_id, surface_code, entry_kind, entry_code, discriminator_hex, historical) VALUES ($1, $2, $3, $4, $5, $6, $7, $8)",
)
.bind(declaration.processor_name.as_str())
.bind(declaration.processor_version.as_str())
.bind(declaration.program_id.as_str())
.bind(declaration.surface_code.as_deref())
.bind(declaration.entry_kind.as_str())
.bind(declaration.entry_code.as_str())
.bind(declaration.discriminator_hex.as_deref())
.bind(declaration.historical)
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = insert_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage declaration insert failed: {error}"
)));
}
inserted_count += 1;
},
std::option::Option::Some(historical) if historical == declaration.historical => {
skipped_count += 1;
},
std::option::Option::Some(_historical) => {
let update_result = sqlx::query(
"UPDATE kb_sol_decode_coverage_declarations SET historical = $1, updated_at = NOW() WHERE processor_name = $2 AND processor_version = $3 AND program_id = $4 AND COALESCE(surface_code, '') = COALESCE($5::text, '') AND entry_kind = $6 AND entry_code = $7 AND COALESCE(discriminator_hex, '') = COALESCE($8::text, '')",
)
.bind(declaration.historical)
.bind(declaration.processor_name.as_str())
.bind(declaration.processor_version.as_str())
.bind(declaration.program_id.as_str())
.bind(declaration.surface_code.as_deref())
.bind(declaration.entry_kind.as_str())
.bind(declaration.entry_code.as_str())
.bind(declaration.discriminator_hex.as_deref())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = update_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage declaration update failed: {error}"
)));
}
updated_count += 1;
},
}
}
for (key, _historical) in existing {
let delete_result = sqlx::query(
"DELETE FROM kb_sol_decode_coverage_declarations WHERE processor_name = $1 AND processor_version = $2 AND program_id = $3 AND COALESCE(surface_code, '') = COALESCE($4::text, '') AND entry_kind = $5 AND entry_code = $6 AND COALESCE(discriminator_hex, '') = COALESCE($7::text, '')",
)
.bind(processor_name)
.bind(processor_version)
.bind(key.0.as_str())
.bind(key.1.as_deref())
.bind(key.2.as_str())
.bind(key.3.as_str())
.bind(key.4.as_deref())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = delete_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres stale decode coverage declaration delete failed: {error}"
)));
}
updated_count += 1;
}
let commit_result = transaction.commit().await;
if let std::result::Result::Err(error) = commit_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage commit failed: {error}"
)));
}
let outcome = crate::InsertOutcome::new(inserted_count, updated_count, skipped_count);
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_coverage_declarations", processor_name = %processor_name, processor_version = %processor_version, outcome = ?outcome, "PostgreSQL decode coverage declarations persisted");
return std::result::Result::Ok(outcome);
}
pub(in crate::postgres) async fn persist_decode_result(
pool: &sqlx::PgPool,
bundle: &crate::DecodePersistenceBundle,
force_replay: bool,
) -> kb_core::Result<crate::InsertOutcome> {
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_decode_result", signature = %bundle.signature, instruction_path = %bundle.instruction_path, stage = %bundle.ledger_identity.stage, processor_name = %bundle.ledger_identity.processor_name, processor_version = %bundle.ledger_identity.processor_version, input_key = %bundle.ledger_identity.input_key, input_hash = %bundle.ledger_identity.input_hash, status = %bundle.status, observation_count = bundle.observations.len(), coverage_program_id = %bundle.coverage.program_id, force_replay, "persist PostgreSQL contextual decode result");
let validation_result = bundle.validate();
if let std::result::Result::Err(error) = validation_result {
return std::result::Result::Err(error);
}
let transaction_result = pool.begin().await;
let mut transaction = match transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode persistence transaction failed: {error}"
)));
},
};
let delete_events_result = sqlx::query(
"DELETE FROM kb_sol_decode_events WHERE processor_name = $1 AND processor_version = $2 AND input_key = $3",
)
.bind(bundle.ledger_identity.processor_name.as_str())
.bind(bundle.ledger_identity.processor_version.as_str())
.bind(bundle.ledger_identity.input_key.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = delete_events_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode output replacement failed: {error}"
)));
}
let delete_coverage_result = sqlx::query(
"DELETE FROM kb_sol_decode_coverage_observations WHERE processor_name = $1 AND processor_version = $2 AND input_key = $3",
)
.bind(bundle.ledger_identity.processor_name.as_str())
.bind(bundle.ledger_identity.processor_version.as_str())
.bind(bundle.ledger_identity.input_key.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = delete_coverage_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage replacement failed: {error}"
)));
}
for observation in &bundle.observations {
let slot_result = i64::try_from(observation.slot);
let slot = match slot_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"decoded observation slot conversion failed: {error}"
)));
},
};
let insert_result = sqlx::query(
"INSERT INTO kb_sol_decode_events (processor_name, processor_version, input_key, input_hash, event_key, signature, slot, instruction_path, program_id, protocol_code, surface_code, event_code, event_name, event_family, source_kind, confidence, proof_kind, proof_jsonb, payload_jsonb, transaction_failed, transaction_error_jsonb, observation_committed) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22) ON CONFLICT (processor_name, processor_version, input_key, event_key) DO UPDATE SET input_hash = EXCLUDED.input_hash, payload_jsonb = EXCLUDED.payload_jsonb, proof_jsonb = EXCLUDED.proof_jsonb, transaction_failed = EXCLUDED.transaction_failed, transaction_error_jsonb = EXCLUDED.transaction_error_jsonb, observation_committed = EXCLUDED.observation_committed, updated_at = NOW()",
)
.bind(observation.processor_name.as_str())
.bind(observation.processor_version.as_str())
.bind(observation.input_key.as_str())
.bind(observation.input_hash.as_str())
.bind(observation.event_key.as_str())
.bind(observation.signature.as_str())
.bind(slot)
.bind(observation.instruction_path.as_str())
.bind(observation.program_id.as_str())
.bind(observation.protocol_code.as_str())
.bind(observation.surface_code.as_str())
.bind(observation.event_code.as_str())
.bind(observation.event_name.as_str())
.bind(observation.event_family.as_str())
.bind(observation.source_kind.as_str())
.bind(observation.confidence.as_str())
.bind(observation.proof_kind.as_str())
.bind(&observation.proof_json)
.bind(&observation.payload_json)
.bind(observation.transaction_failed)
.bind(observation.transaction_error.clone())
.bind(observation.observation_committed)
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = insert_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decoded observation insert failed: {error}"
)));
}
}
let coverage_result =
crate::postgres::query::decode_pipeline_queries::upsert_coverage_observation(
&mut transaction,
&bundle.coverage,
)
.await;
if let std::result::Result::Err(error) = coverage_result {
return std::result::Result::Err(error);
}
let materialized_count_result =
crate::postgres::query::decode_pipeline_queries::refresh_materialized_coverage_count(
&mut transaction,
bundle.ledger_identity.processor_name.as_str(),
bundle.ledger_identity.processor_version.as_str(),
bundle.ledger_identity.input_key.as_str(),
)
.await;
if let std::result::Result::Err(error) = materialized_count_result {
return std::result::Result::Err(error);
}
let lifecycle_state = match bundle.status.as_str() {
"decoded" => "decoded",
"ignored" | "unsupported" => "ignored",
"failed" => "failed",
_unknown => "decoded",
};
let lifecycle_result = sqlx::query(
"UPDATE kb_sol_core_instructions SET processing_state = CASE WHEN $1 = 'decoded' AND processing_state <> 'materialized' THEN 'decoded' WHEN $1 IN ('ignored', 'failed') AND processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $1 ELSE processing_state END, processor_name = CASE WHEN $1 = 'decoded' AND processing_state <> 'materialized' THEN $2 WHEN $1 IN ('ignored', 'failed') AND processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $2 ELSE processor_name END, processor_version = CASE WHEN $1 = 'decoded' AND processing_state <> 'materialized' THEN $3 WHEN $1 IN ('ignored', 'failed') AND processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $3 ELSE processor_version END, lifecycle_reason = CASE WHEN $1 = 'decoded' AND processing_state <> 'materialized' THEN $4 WHEN $1 IN ('ignored', 'failed') AND processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $4 ELSE lifecycle_reason END, updated_at = NOW() WHERE signature = $5 AND instruction_path = $6",
)
.bind(lifecycle_state)
.bind(bundle.ledger_identity.processor_name.as_str())
.bind(bundle.ledger_identity.processor_version.as_str())
.bind(bundle.status.as_str())
.bind(bundle.signature.as_str())
.bind(bundle.instruction_path.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = lifecycle_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode instruction lifecycle update failed: {error}"
)));
}
let ledger_status = if bundle.status == "failed" { "failed" } else { "succeeded" };
let ledger_result = crate::postgres::query::decode_pipeline_queries::upsert_ledger_terminal(
&mut transaction,
&bundle.ledger_identity,
ledger_status,
bundle.error_code.as_deref(),
bundle.error_message.as_deref(),
)
.await;
if let std::result::Result::Err(error) = ledger_result {
return std::result::Result::Err(error);
}
let commit_result = transaction.commit().await;
if let std::result::Result::Err(error) = commit_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode persistence commit failed: {error}"
)));
}
let count_result = u64::try_from(bundle.observations.len());
let count = match count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"decoded observation count conversion failed: {error}"
)));
},
};
let outcome = crate::InsertOutcome::new(count, 1, 0);
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_decode_result", signature = %bundle.signature, instruction_path = %bundle.instruction_path, processor_name = %bundle.ledger_identity.processor_name, processor_version = %bundle.ledger_identity.processor_version, status = %bundle.status, outcome = ?outcome, committed = true, "PostgreSQL contextual decode result persisted");
return std::result::Result::Ok(outcome);
}
pub(in crate::postgres) async fn mark_decode_failed(
pool: &sqlx::PgPool,
failure: &crate::DecodeFailure,
) -> kb_core::Result<crate::InsertOutcome> {
tracing::error!(target: crate::TRACING_TARGET, action = "mark_decode_failed", signature = %failure.signature, instruction_path = %failure.instruction_path, processor_name = %failure.ledger_identity.processor_name, processor_version = %failure.ledger_identity.processor_version, input_key = %failure.ledger_identity.input_key, input_hash = %failure.ledger_identity.input_hash, error_code = %failure.error_code, error_message = %failure.error_message, "persist PostgreSQL contextual decode failure");
if failure.ledger_identity.stage != "instruction_decode"
|| failure.signature.trim().is_empty()
|| failure.instruction_path.trim().is_empty()
|| failure.error_code.trim().is_empty()
|| failure.error_message.trim().is_empty()
{
return std::result::Result::Err(kb_core::Error::db("decode failure identity is invalid"));
}
let transaction_result = pool.begin().await;
let mut transaction = match transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode failure transaction failed: {error}"
)));
},
};
let lifecycle_result = sqlx::query(
"UPDATE kb_sol_core_instructions SET processing_state = CASE WHEN processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN 'failed' ELSE processing_state END, processor_name = CASE WHEN processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $1 ELSE processor_name END, processor_version = CASE WHEN processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $2 ELSE processor_version END, lifecycle_reason = CASE WHEN processing_state IN ('pending', 'failed', 'replay_requested', 'ignored') THEN $3 ELSE lifecycle_reason END, updated_at = NOW() WHERE signature = $4 AND instruction_path = $5",
)
.bind(failure.ledger_identity.processor_name.as_str())
.bind(failure.ledger_identity.processor_version.as_str())
.bind(failure.error_code.as_str())
.bind(failure.signature.as_str())
.bind(failure.instruction_path.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = lifecycle_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode failure lifecycle update failed: {error}"
)));
}
let ledger_result = crate::postgres::query::decode_pipeline_queries::upsert_ledger_terminal(
&mut transaction,
&failure.ledger_identity,
"failed",
std::option::Option::Some(failure.error_code.as_str()),
std::option::Option::Some(failure.error_message.as_str()),
)
.await;
if let std::result::Result::Err(error) = ledger_result {
return std::result::Result::Err(error);
}
let commit_result = transaction.commit().await;
if let std::result::Result::Err(error) = commit_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode failure commit failed: {error}"
)));
}
let outcome = crate::InsertOutcome::new(0, 1, 0);
tracing::debug!(target: crate::TRACING_TARGET, action = "mark_decode_failed", signature = %failure.signature, instruction_path = %failure.instruction_path, processor_name = %failure.ledger_identity.processor_name, processor_version = %failure.ledger_identity.processor_version, outcome = ?outcome, committed = true, "PostgreSQL contextual decode failure persisted");
return std::result::Result::Ok(outcome);
}
pub(in crate::postgres) async fn persist_materialization_result(
pool: &sqlx::PgPool,
bundle: &crate::MaterializationPersistenceBundle,
force_replay: bool,
) -> kb_core::Result<crate::InsertOutcome> {
if bundle.status == "failed" {
tracing::error!(
target: crate::TRACING_TARGET,
action = "persist_materialization_failure",
signature = %bundle.signature,
instruction_path = %bundle.instruction_path,
processor_name = %bundle.ledger_identity.processor_name,
processor_version = %bundle.ledger_identity.processor_version,
input_key = %bundle.ledger_identity.input_key,
input_hash = %bundle.ledger_identity.input_hash,
source_decoder_name = %bundle.source_decoder_name,
source_decoder_version = %bundle.source_decoder_version,
error_code = ?bundle.error_code,
error_message = ?bundle.error_message,
"persist PostgreSQL materialization failure"
);
}
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_materialization_result", signature = %bundle.signature, instruction_path = %bundle.instruction_path, stage = %bundle.ledger_identity.stage, processor_name = %bundle.ledger_identity.processor_name, processor_version = %bundle.ledger_identity.processor_version, input_key = %bundle.ledger_identity.input_key, input_hash = %bundle.ledger_identity.input_hash, source_decoder_name = %bundle.source_decoder_name, source_decoder_version = %bundle.source_decoder_version, source_decode_input_key = %bundle.source_decode_input_key, status = %bundle.status, output_count = bundle.outputs.len(), force_replay, "persist PostgreSQL materialization result");
let validation_result = bundle.validate();
if let std::result::Result::Err(error) = validation_result {
return std::result::Result::Err(error);
}
let transaction_result = pool.begin().await;
let mut transaction = match transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialization transaction failed: {error}"
)));
},
};
let delete_result = sqlx::query(
"DELETE FROM kb_sol_mat_events WHERE processor_name = $1 AND processor_version = $2 AND input_key = $3",
)
.bind(bundle.ledger_identity.processor_name.as_str())
.bind(bundle.ledger_identity.processor_version.as_str())
.bind(bundle.ledger_identity.input_key.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = delete_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialization output replacement failed: {error}"
)));
}
for output in &bundle.outputs {
let slot_result = i64::try_from(output.slot);
let slot = match slot_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"materialized output slot conversion failed: {error}"
)));
},
};
let insert_result = sqlx::query(
"INSERT INTO kb_sol_mat_events (processor_name, processor_version, input_key, input_hash, output_key, source_event_key, source_decoder_name, source_decoder_version, source_decode_input_key, signature, slot, materialized_family, payload_jsonb) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13) ON CONFLICT (processor_name, processor_version, input_key, output_key) DO UPDATE SET input_hash = EXCLUDED.input_hash, source_decoder_name = EXCLUDED.source_decoder_name, source_decoder_version = EXCLUDED.source_decoder_version, source_decode_input_key = EXCLUDED.source_decode_input_key, payload_jsonb = EXCLUDED.payload_jsonb, updated_at = NOW()",
)
.bind(output.processor_name.as_str())
.bind(output.processor_version.as_str())
.bind(output.input_key.as_str())
.bind(output.input_hash.as_str())
.bind(output.output_key.as_str())
.bind(output.source_event_key.as_str())
.bind(bundle.source_decoder_name.as_str())
.bind(bundle.source_decoder_version.as_str())
.bind(bundle.source_decode_input_key.as_str())
.bind(output.signature.as_str())
.bind(slot)
.bind(output.materialized_family.as_str())
.bind(&output.payload_json)
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = insert_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialized output insert failed: {error}"
)));
}
}
let ledger_status = if bundle.status == "failed" { "failed" } else { "succeeded" };
let ledger_result = crate::postgres::query::decode_pipeline_queries::upsert_ledger_terminal(
&mut transaction,
&bundle.ledger_identity,
ledger_status,
bundle.error_code.as_deref(),
bundle.error_message.as_deref(),
)
.await;
if let std::result::Result::Err(error) = ledger_result {
return std::result::Result::Err(error);
}
if !bundle.outputs.is_empty() {
let lifecycle_result = sqlx::query(
"UPDATE kb_sol_core_instructions SET processing_state = 'materialized', processor_name = $1, processor_version = $2, lifecycle_reason = $3, updated_at = NOW() WHERE signature = $4 AND instruction_path = $5",
)
.bind(bundle.ledger_identity.processor_name.as_str())
.bind(bundle.ledger_identity.processor_version.as_str())
.bind(bundle.status.as_str())
.bind(bundle.signature.as_str())
.bind(bundle.instruction_path.as_str())
.execute(&mut *transaction)
.await;
if let std::result::Result::Err(error) = lifecycle_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialization lifecycle update failed: {error}"
)));
}
}
let coverage_result =
crate::postgres::query::decode_pipeline_queries::refresh_materialized_coverage_count(
&mut transaction,
bundle.source_decoder_name.as_str(),
bundle.source_decoder_version.as_str(),
bundle.source_decode_input_key.as_str(),
)
.await;
if let std::result::Result::Err(error) = coverage_result {
return std::result::Result::Err(error);
}
let commit_result = transaction.commit().await;
if let std::result::Result::Err(error) = commit_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialization commit failed: {error}"
)));
}
let count_result = u64::try_from(bundle.outputs.len());
let count = match count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"materialized output count conversion failed: {error}"
)));
},
};
let outcome = crate::InsertOutcome::new(count, 1, 0);
tracing::debug!(target: crate::TRACING_TARGET, action = "persist_materialization_result", signature = %bundle.signature, instruction_path = %bundle.instruction_path, processor_name = %bundle.ledger_identity.processor_name, processor_version = %bundle.ledger_identity.processor_version, status = %bundle.status, outcome = ?outcome, committed = true, "PostgreSQL materialization result persisted");
return std::result::Result::Ok(outcome);
}
pub(in crate::postgres) async fn list_decode_coverage_summary(
pool: &sqlx::PgPool,
processor_name: std::option::Option<&str>,
processor_version: std::option::Option<&str>,
limit: u32,
) -> kb_core::Result<std::vec::Vec<crate::DecodeCoverageSummaryRow>> {
tracing::debug!(target: crate::TRACING_TARGET, action = "list_decode_coverage_summary", processor_name = ?processor_name, processor_version = ?processor_version, limit, "query PostgreSQL decode coverage summary");
if limit == 0 {
return std::result::Result::Err(kb_core::Error::db(
"decode coverage summary limit must be greater than zero",
));
}
let query_result = sqlx::query(
"WITH declared AS (SELECT processor_name, processor_version, program_id, surface_code, entry_code, COUNT(*)::BIGINT AS declared_count FROM kb_sol_decode_coverage_declarations WHERE ($1::text IS NULL OR processor_name = $1) AND ($2::text IS NULL OR processor_version = $2) GROUP BY processor_name, processor_version, program_id, surface_code, entry_code), observed AS (SELECT processor_name, processor_version, program_id, surface_code, COALESCE(entry_code, 'unknown') AS entry_code, COUNT(*)::BIGINT AS observed_count, COUNT(*) FILTER (WHERE recognized)::BIGINT AS recognized_count, COALESCE(SUM(decoded_count), 0)::BIGINT AS decoded_count, COALESCE(SUM(materialized_count), 0)::BIGINT AS materialized_count, COALESCE(SUM(error_count), 0)::BIGINT AS error_count, COUNT(*) FILTER (WHERE status = 'unsupported' OR entry_code IS NULL)::BIGINT AS unknown_count, COUNT(*) FILTER (WHERE NOT transaction_failed)::BIGINT AS successful_transaction_count, COUNT(*) FILTER (WHERE transaction_failed)::BIGINT AS failed_transaction_count FROM kb_sol_decode_coverage_observations WHERE ($1::text IS NULL OR processor_name = $1) AND ($2::text IS NULL OR processor_version = $2) GROUP BY processor_name, processor_version, program_id, surface_code, COALESCE(entry_code, 'unknown')) SELECT COALESCE(d.processor_name, o.processor_name) AS processor_name, COALESCE(d.processor_version, o.processor_version) AS processor_version, COALESCE(d.program_id, o.program_id) AS program_id, COALESCE(d.surface_code, o.surface_code) AS surface_code, COALESCE(d.entry_code, o.entry_code) AS entry_code, COALESCE(d.declared_count, 0)::BIGINT AS declared_count, COALESCE(o.observed_count, 0)::BIGINT AS observed_count, COALESCE(o.recognized_count, 0)::BIGINT AS recognized_count, COALESCE(o.decoded_count, 0)::BIGINT AS decoded_count, COALESCE(o.materialized_count, 0)::BIGINT AS materialized_count, COALESCE(o.error_count, 0)::BIGINT AS error_count, COALESCE(o.unknown_count, 0)::BIGINT AS unknown_count, COALESCE(o.successful_transaction_count, 0)::BIGINT AS successful_transaction_count, COALESCE(o.failed_transaction_count, 0)::BIGINT AS failed_transaction_count FROM declared d FULL OUTER JOIN observed o ON d.processor_name = o.processor_name AND d.processor_version = o.processor_version AND d.program_id = o.program_id AND COALESCE(d.surface_code, '') = COALESCE(o.surface_code, '') AND d.entry_code = o.entry_code ORDER BY observed_count DESC, processor_name, program_id, entry_code LIMIT $3",
)
.bind(processor_name)
.bind(processor_version)
.bind(i64::from(limit))
.fetch_all(pool)
.await;
let rows = match query_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage summary query failed: {error}"
)));
},
};
let mut output = std::vec::Vec::with_capacity(rows.len());
for row in rows {
let mapped_result =
crate::postgres::query::decode_pipeline_queries::map_coverage_summary_row(&row);
match mapped_result {
std::result::Result::Ok(value) => output.push(value),
std::result::Result::Err(error) => return std::result::Result::Err(error),
}
}
tracing::debug!(target: crate::TRACING_TARGET, action = "list_decode_coverage_summary", processor_name = ?processor_name, processor_version = ?processor_version, row_count = output.len(), "PostgreSQL decode coverage summary loaded");
return std::result::Result::Ok(output);
}
async fn upsert_coverage_observation(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
coverage: &crate::DecodeCoverageObservationInsert,
) -> kb_core::Result<()> {
let slot_result = i64::try_from(coverage.slot);
let slot = match slot_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"decode coverage slot conversion failed: {error}"
)));
},
};
let decoded_count_result = i32::try_from(coverage.decoded_count);
let decoded_count = match decoded_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"decode coverage decoded count conversion failed: {error}"
)));
},
};
let materialized_count_result = i32::try_from(coverage.materialized_count);
let materialized_count = match materialized_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"decode coverage materialized count conversion failed: {error}"
)));
},
};
let error_count_result = i32::try_from(coverage.error_count);
let error_count = match error_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
return std::result::Result::Err(kb_core::Error::db(format!(
"decode coverage error count conversion failed: {error}"
)));
},
};
let query_result = sqlx::query(
"INSERT INTO kb_sol_decode_coverage_observations (processor_name, processor_version, input_key, input_hash, signature, slot, instruction_path, program_id, surface_code, entry_code, discriminator_hex, status, recognized, decoded_count, materialized_count, error_count, transaction_failed) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17) ON CONFLICT (processor_name, processor_version, input_key) DO UPDATE SET input_hash = EXCLUDED.input_hash, surface_code = EXCLUDED.surface_code, entry_code = EXCLUDED.entry_code, discriminator_hex = EXCLUDED.discriminator_hex, status = EXCLUDED.status, recognized = EXCLUDED.recognized, decoded_count = EXCLUDED.decoded_count, materialized_count = EXCLUDED.materialized_count, error_count = EXCLUDED.error_count, transaction_failed = EXCLUDED.transaction_failed, updated_at = NOW()",
)
.bind(coverage.processor_name.as_str())
.bind(coverage.processor_version.as_str())
.bind(coverage.input_key.as_str())
.bind(coverage.input_hash.as_str())
.bind(coverage.signature.as_str())
.bind(slot)
.bind(coverage.instruction_path.as_str())
.bind(coverage.program_id.as_str())
.bind(coverage.surface_code.as_deref())
.bind(coverage.entry_code.as_deref())
.bind(coverage.discriminator_hex.as_deref())
.bind(coverage.status.as_str())
.bind(coverage.recognized)
.bind(decoded_count)
.bind(materialized_count)
.bind(error_count)
.bind(coverage.transaction_failed)
.execute(&mut **transaction)
.await;
if let std::result::Result::Err(error) = query_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage observation upsert failed: {error}"
)));
}
return std::result::Result::Ok(());
}
async fn refresh_materialized_coverage_count(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
decoder_name: &str,
decoder_version: &str,
decode_input_key: &str,
) -> kb_core::Result<()> {
let query_result = sqlx::query(
"UPDATE kb_sol_decode_coverage_observations SET materialized_count = (SELECT COUNT(*)::INTEGER FROM kb_sol_mat_events WHERE source_decoder_name = $1 AND source_decoder_version = $2 AND source_decode_input_key = $3), updated_at = NOW() WHERE processor_name = $1 AND processor_version = $2 AND input_key = $3",
)
.bind(decoder_name)
.bind(decoder_version)
.bind(decode_input_key)
.execute(&mut **transaction)
.await;
if let std::result::Result::Err(error) = query_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres materialization coverage refresh failed: {error}"
)));
}
return std::result::Result::Ok(());
}
async fn upsert_ledger_terminal(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
identity: &crate::ProcessingLedgerIdentity,
status: &str,
error_code: std::option::Option<&str>,
error_message: std::option::Option<&str>,
) -> kb_core::Result<()> {
let query_result = sqlx::query(
"INSERT INTO kb_sol_ops_processing_ledger (stage, processor_name, processor_version, input_key, input_hash, status, attempt_count, started_at, finished_at, error_code, error_message) VALUES ($1, $2, $3, $4, $5, $6, 1, NOW(), NOW(), $7, $8) ON CONFLICT (stage, processor_name, processor_version, input_key) DO UPDATE SET input_hash = EXCLUDED.input_hash, status = EXCLUDED.status, attempt_count = kb_sol_ops_processing_ledger.attempt_count + 1, started_at = NOW(), finished_at = NOW(), error_code = EXCLUDED.error_code, error_message = EXCLUDED.error_message, updated_at = NOW()",
)
.bind(identity.stage.as_str())
.bind(identity.processor_name.as_str())
.bind(identity.processor_version.as_str())
.bind(identity.input_key.as_str())
.bind(identity.input_hash.as_str())
.bind(status)
.bind(error_code)
.bind(error_message)
.execute(&mut **transaction)
.await;
if let std::result::Result::Err(error) = query_result {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode processing ledger upsert failed: {error}"
)));
}
return std::result::Result::Ok(());
}
fn map_coverage_summary_row(
row: &sqlx::postgres::PgRow,
) -> kb_core::Result<crate::DecodeCoverageSummaryRow> {
let processor_name_result = row.try_get::<std::string::String, _>("processor_name");
let processor_name = match processor_name_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return map_error("processor_name", error),
};
let processor_version_result = row.try_get::<std::string::String, _>("processor_version");
let processor_version = match processor_version_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return map_error("processor_version", error),
};
let program_id_result = row.try_get::<std::string::String, _>("program_id");
let program_id = match program_id_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return map_error("program_id", error),
};
let surface_code_result =
row.try_get::<std::option::Option<std::string::String>, _>("surface_code");
let surface_code = match surface_code_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return map_error("surface_code", error),
};
let entry_code_result = row.try_get::<std::string::String, _>("entry_code");
let entry_code = match entry_code_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return map_error("entry_code", error),
};
let declared_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "declared_count");
let declared_count = match declared_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observed_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "observed_count");
let observed_count = match observed_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let recognized_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "recognized_count");
let recognized_count = match recognized_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let decoded_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "decoded_count");
let decoded_count = match decoded_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let materialized_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "materialized_count");
let materialized_count = match materialized_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let error_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "error_count");
let error_count = match error_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let unknown_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "unknown_count");
let unknown_count = match unknown_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let successful_transaction_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(
row,
"successful_transaction_count",
);
let successful_transaction_count = match successful_transaction_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let failed_transaction_count_result =
crate::postgres::query::decode_pipeline_queries::read_i64(row, "failed_transaction_count");
let failed_transaction_count = match failed_transaction_count_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
return std::result::Result::Ok(crate::DecodeCoverageSummaryRow {
processor_name,
processor_version,
program_id,
surface_code,
entry_code,
declared_count,
observed_count,
recognized_count,
decoded_count,
materialized_count,
error_count,
unknown_count,
successful_transaction_count,
failed_transaction_count,
});
}
fn read_i64(row: &sqlx::postgres::PgRow, column: &str) -> kb_core::Result<i64> {
let result = row.try_get::<i64, _>(column);
return match result {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(error) => {
crate::postgres::query::decode_pipeline_queries::map_error(column, error)
},
};
}
fn map_error<T>(column: &str, error: sqlx::Error) -> kb_core::Result<T> {
return std::result::Result::Err(kb_core::Error::db(format!(
"postgres decode coverage column {column} mapping failed: {error}"
)));
}
#[cfg(test)]
mod tests {
fn result_or_panic<T, E>(result: std::result::Result<T, E>) -> T
where
E: std::fmt::Display,
{
return match result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
panic!("postgres decode test operation failed: {error}")
},
};
}
async fn execute_sql(pool: &sqlx::PgPool, statement: &'static str) {
let result = sqlx::query(statement).execute(pool).await;
result_or_panic(result);
}
fn unique_processor_name(prefix: &str) -> std::string::String {
return format!(
"{}_{}_{}",
prefix,
std::process::id(),
chrono::Utc::now().timestamp_micros()
);
}
async fn test_pool_from_env() -> std::option::Option<sqlx::PgPool> {
let url = match std::env::var("KB_POSTGRES_TEST_URL") {
std::result::Result::Ok(value) if !value.trim().is_empty() => value,
_ => return std::option::Option::None,
};
let pool_result = sqlx::PgPool::connect(url.as_str()).await;
let pool = result_or_panic(pool_result);
result_or_panic(crate::postgres::query::raw_queries::apply_raw_store_schema(&pool).await);
result_or_panic(crate::postgres::query::core_queries::apply_core_store_schema(&pool).await);
result_or_panic(
crate::postgres::query::decode_pipeline_queries::apply_decode_store_schema(&pool).await,
);
return std::option::Option::Some(pool);
}
#[tokio::test]
async fn optional_postgres_coverage_declarations_report_insert_skip_and_update_from_env() {
let _postgres_guard = crate::postgres::test_serial::postgres_test_guard().await;
let pool = match test_pool_from_env().await {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let processor_name = unique_processor_name("decode_coverage_outcome_test");
let mut declaration = crate::DecodeCoverageDeclarationInsert {
processor_name: processor_name.clone(),
processor_version: "1".to_string(),
program_id: "11111111111111111111111111111111".to_string(),
surface_code: std::option::Option::Some("solana_native_system".to_string()),
entry_kind: "instruction".to_string(),
entry_code: "test".to_string(),
discriminator_hex: std::option::Option::None,
historical: false,
};
let first = result_or_panic(
crate::postgres::query::decode_pipeline_queries::persist_decode_coverage_declarations(
&pool,
&[declaration.clone()],
)
.await,
);
assert_eq!(first, crate::InsertOutcome::new(1, 0, 0));
let second = result_or_panic(
crate::postgres::query::decode_pipeline_queries::persist_decode_coverage_declarations(
&pool,
&[declaration.clone()],
)
.await,
);
assert_eq!(second, crate::InsertOutcome::new(0, 0, 1));
declaration.historical = true;
let third = result_or_panic(
crate::postgres::query::decode_pipeline_queries::persist_decode_coverage_declarations(
&pool,
&[declaration],
)
.await,
);
assert_eq!(third, crate::InsertOutcome::new(0, 1, 0));
let cleanup_result = sqlx::query(
"DELETE FROM kb_sol_decode_coverage_declarations WHERE processor_name = $1",
)
.bind(processor_name.as_str())
.execute(&pool)
.await;
result_or_panic(cleanup_result);
}
#[tokio::test]
async fn optional_postgres_same_version_and_hash_is_current_from_env() {
let _postgres_guard = crate::postgres::test_serial::postgres_test_guard().await;
let pool = match test_pool_from_env().await {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let processor_name = unique_processor_name("decode_current_test");
let input_key = format!("{}:0", processor_name);
let identity = result_or_panic(crate::ProcessingLedgerIdentity::new(
"instruction_decode",
processor_name.clone(),
"1",
input_key.clone(),
"stable-input-hash",
));
let coverage = crate::DecodeCoverageObservationInsert {
processor_name: processor_name.clone(),
processor_version: "1".to_string(),
input_key: input_key.clone(),
input_hash: "stable-input-hash".to_string(),
signature: processor_name.clone(),
slot: 1,
instruction_path: "0".to_string(),
program_id: "11111111111111111111111111111111".to_string(),
surface_code: std::option::Option::Some("solana_native_system".to_string()),
entry_code: std::option::Option::Some("unknown".to_string()),
discriminator_hex: std::option::Option::None,
status: "unsupported".to_string(),
recognized: true,
decoded_count: 0,
materialized_count: 0,
error_count: 0,
transaction_failed: false,
};
let bundle = crate::DecodePersistenceBundle {
ledger_identity: identity.clone(),
signature: processor_name.clone(),
instruction_path: "0".to_string(),
status: "unsupported".to_string(),
error_code: std::option::Option::None,
error_message: std::option::Option::None,
observations: std::vec::Vec::new(),
coverage,
};
result_or_panic(
crate::postgres::query::decode_pipeline_queries::persist_decode_result(
&pool, &bundle, true,
)
.await,
);
let current = result_or_panic(
crate::postgres::query::decode_pipeline_queries::is_decode_current(&pool, &identity)
.await,
);
assert!(current);
let cleanup_coverage = sqlx::query(
"DELETE FROM kb_sol_decode_coverage_observations WHERE processor_name = $1",
)
.bind(processor_name.as_str())
.execute(&pool)
.await;
result_or_panic(cleanup_coverage);
let cleanup_ledger =
sqlx::query("DELETE FROM kb_sol_ops_processing_ledger WHERE processor_name = $1")
.bind(processor_name.as_str())
.execute(&pool)
.await;
result_or_panic(cleanup_ledger);
}
#[tokio::test]
async fn optional_postgres_materialized_event_query_is_bounded_and_typed_from_env() {
let _postgres_guard = crate::postgres::test_serial::postgres_test_guard().await;
let pool = match test_pool_from_env().await {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let input_key = unique_processor_name("annotation_query_test");
let insert_result = sqlx::query(
"INSERT INTO kb_sol_mat_events (processor_name, processor_version, input_key, input_hash, output_key, source_event_key, source_decoder_name, source_decoder_version, source_decode_input_key, signature, slot, materialized_family, payload_jsonb) VALUES ('transaction_annotations', '0.4.3', $1, 'hash', 'annotation', 'memo:0', 'spl_memo', '0.4.3', 'decode-input', $2, 42, 'transaction_annotation', $3)",
)
.bind(input_key.as_str())
.bind(input_key.as_str())
.bind(serde_json::json!({"text":"postgres annotation"}))
.execute(&pool)
.await;
result_or_panic(insert_result);
let filter = result_or_panic(crate::MaterializedEventFilter::new(
std::option::Option::Some("transaction_annotations".to_string()),
std::option::Option::Some("transaction_annotation".to_string()),
std::option::Option::Some(input_key.clone()),
1,
));
let rows = result_or_panic(
crate::postgres::query::decode_pipeline_queries::list_materialized_events(
&pool, &filter,
)
.await,
);
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].slot, 42);
assert_eq!(rows[0].payload_json["text"], "postgres annotation");
let cleanup_result = sqlx::query("DELETE FROM kb_sol_mat_events WHERE input_key = $1")
.bind(input_key.as_str())
.execute(&pool)
.await;
result_or_panic(cleanup_result);
}
#[tokio::test]
async fn decoded_events_and_ledger_roll_back_together() {
let url = match std::env::var("KB_POSTGRES_TEST_URL") {
std::result::Result::Ok(value) if !value.trim().is_empty() => value,
_ => return,
};
let _postgres_guard = crate::postgres::test_serial::postgres_test_guard().await;
let pool_result = sqlx::PgPool::connect(url.as_str()).await;
let pool = result_or_panic(pool_result);
result_or_panic(crate::postgres::query::raw_queries::apply_raw_store_schema(&pool).await);
result_or_panic(crate::postgres::query::core_queries::apply_core_store_schema(&pool).await);
result_or_panic(
crate::postgres::query::decode_pipeline_queries::apply_decode_store_schema(&pool).await,
);
execute_sql(
&pool,
"DELETE FROM kb_sol_decode_events WHERE processor_name = 'decode_atomic_rollback_test'",
)
.await;
execute_sql(
&pool,
"DELETE FROM kb_sol_decode_coverage_observations WHERE processor_name = 'decode_atomic_rollback_test'",
)
.await;
execute_sql(
&pool,
"DELETE FROM kb_sol_ops_processing_ledger WHERE processor_name = 'decode_atomic_rollback_test'",
)
.await;
execute_sql(
&pool,
"CREATE OR REPLACE FUNCTION kb_test_reject_decode_ledger() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN IF NEW.processor_name = 'decode_atomic_rollback_test' THEN RAISE EXCEPTION 'forced decode ledger failure'; END IF; RETURN NEW; END; $$",
)
.await;
execute_sql(
&pool,
"DROP TRIGGER IF EXISTS kb_test_reject_decode_ledger_trigger ON kb_sol_ops_processing_ledger",
)
.await;
execute_sql(
&pool,
"CREATE TRIGGER kb_test_reject_decode_ledger_trigger BEFORE INSERT OR UPDATE ON kb_sol_ops_processing_ledger FOR EACH ROW EXECUTE FUNCTION kb_test_reject_decode_ledger()",
)
.await;
let identity_result = crate::ProcessingLedgerIdentity::new(
"instruction_decode",
"decode_atomic_rollback_test",
"1",
"rollback-signature:0",
"rollback-input-hash",
);
let identity = result_or_panic(identity_result);
let coverage = crate::DecodeCoverageObservationInsert {
processor_name: identity.processor_name.clone(),
processor_version: identity.processor_version.clone(),
input_key: identity.input_key.clone(),
input_hash: identity.input_hash.clone(),
signature: "rollback-signature".to_string(),
slot: 1,
instruction_path: "0".to_string(),
program_id: "11111111111111111111111111111111".to_string(),
surface_code: std::option::Option::Some("system_program".to_string()),
entry_code: std::option::Option::Some("test".to_string()),
discriminator_hex: std::option::Option::None,
status: "decoded".to_string(),
recognized: true,
decoded_count: 1,
materialized_count: 0,
error_count: 0,
transaction_failed: false,
};
let observation = crate::DecodeObservationInsert {
processor_name: identity.processor_name.clone(),
processor_version: identity.processor_version.clone(),
input_key: identity.input_key.clone(),
input_hash: identity.input_hash.clone(),
event_key: "test".to_string(),
signature: "rollback-signature".to_string(),
slot: 1,
instruction_path: "0".to_string(),
program_id: "11111111111111111111111111111111".to_string(),
protocol_code: "solana_core".to_string(),
surface_code: "system_program".to_string(),
event_code: "test".to_string(),
event_name: "test".to_string(),
event_family: "audit".to_string(),
source_kind: "instruction".to_string(),
confidence: "exact".to_string(),
proof_kind: "exact_layout".to_string(),
proof_json: serde_json::json!({"source": "test"}),
payload_json: serde_json::json!({"value": 1}),
transaction_failed: false,
transaction_error: std::option::Option::None,
observation_committed: true,
};
let bundle = crate::DecodePersistenceBundle {
ledger_identity: identity,
signature: "rollback-signature".to_string(),
instruction_path: "0".to_string(),
status: "decoded".to_string(),
error_code: std::option::Option::None,
error_message: std::option::Option::None,
observations: std::vec![observation],
coverage,
};
let persistence_result =
crate::postgres::query::decode_pipeline_queries::persist_decode_result(
&pool, &bundle, true,
)
.await;
assert!(persistence_result.is_err());
let event_count_result = sqlx::query_scalar::<sqlx::Postgres, i64>(
"SELECT COUNT(*) FROM kb_sol_decode_events WHERE processor_name = 'decode_atomic_rollback_test'",
)
.fetch_one(&pool)
.await;
let event_count = result_or_panic(event_count_result);
let ledger_count_result = sqlx::query_scalar::<sqlx::Postgres, i64>(
"SELECT COUNT(*) FROM kb_sol_ops_processing_ledger WHERE processor_name = 'decode_atomic_rollback_test'",
)
.fetch_one(&pool)
.await;
let ledger_count = result_or_panic(ledger_count_result);
assert_eq!(event_count, 0);
assert_eq!(ledger_count, 0);
execute_sql(
&pool,
"DROP TRIGGER IF EXISTS kb_test_reject_decode_ledger_trigger ON kb_sol_ops_processing_ledger",
)
.await;
execute_sql(&pool, "DROP FUNCTION IF EXISTS kb_test_reject_decode_ledger()").await;
}
}