use super::*; pub async fn mark_selected_transactions_processed( batch: &SelectedMempoolBatch, block_number: u32, ) -> Result<()> { // Mark each selected mempool row as processed under the saved block // number so it can be cleaned up or restored later if needed. let client_handle = db_client().await?; let client = client_handle.as_ref(); let bn = block_number as i32; // Selected batches are grouped by table, then marked with one UPDATE per // table instead of touching rows one at a time. mark_rows_by_ids(client, "transfer", &ids_for_table(batch, "transfer"), bn).await?; mark_rows_by_ids(client, "token", &ids_for_table(batch, "token"), bn).await?; mark_rows_by_ids( &client, "issue_token", &ids_for_table(batch, "issue_token"), bn, ) .await?; mark_rows_by_ids(client, "burn", &ids_for_table(batch, "burn"), bn).await?; mark_rows_by_ids(client, "nft", &ids_for_table(batch, "nft"), bn).await?; mark_rows_by_ids(client, "marketing", &ids_for_table(batch, "marketing"), bn).await?; mark_rows_by_ids( &client, "vanity_address", &ids_for_table(batch, "vanity_address"), bn, ) .await?; mark_rows_by_ids(client, "swap", &ids_for_table(batch, "swap"), bn).await?; mark_rows_by_ids( &client, "loan_contract", &ids_for_table(batch, "loan_contract"), bn, ) .await?; mark_rows_by_ids( &client, "loan_payment", &ids_for_table(batch, "loan_payment"), bn, ) .await?; mark_rows_by_ids( &client, "collateral_claim", &ids_for_table(batch, "collateral_claim"), bn, ) .await?; Ok(()) } pub async fn restore_selected_transactions_processed(batch: &SelectedMempoolBatch) -> Result<()> { // If block commit fails after selected rows were marked processed, // restore them before the chain height can acknowledge the block. let client_handle = db_client().await?; let client = client_handle.as_ref(); unmark_rows_by_ids(client, "transfer", &ids_for_table(batch, "transfer")).await?; unmark_rows_by_ids(client, "token", &ids_for_table(batch, "token")).await?; unmark_rows_by_ids(client, "issue_token", &ids_for_table(batch, "issue_token")).await?; unmark_rows_by_ids(client, "burn", &ids_for_table(batch, "burn")).await?; unmark_rows_by_ids(client, "nft", &ids_for_table(batch, "nft")).await?; unmark_rows_by_ids(client, "marketing", &ids_for_table(batch, "marketing")).await?; unmark_rows_by_ids( &client, "vanity_address", &ids_for_table(batch, "vanity_address"), ) .await?; unmark_rows_by_ids(client, "swap", &ids_for_table(batch, "swap")).await?; unmark_rows_by_ids( &client, "loan_contract", &ids_for_table(batch, "loan_contract"), ) .await?; unmark_rows_by_ids( &client, "loan_payment", &ids_for_table(batch, "loan_payment"), ) .await?; unmark_rows_by_ids( &client, "collateral_claim", &ids_for_table(batch, "collateral_claim"), ) .await?; Ok(()) } pub async fn restore_processed_by_signatures(signatures: &[String]) -> Result { // Orphan correction can revive recently processed mempool rows by // signature when a saved block is rolled back out of the chain. if signatures.is_empty() { return Ok(false); } let client_handle = db_client().await?; let client = client_handle.as_ref(); let mut restored = 0_u64; // Each table keeps its own signature columns, so rollback unmarks every // column that could contain one of the rolled-back signatures. restored += unmark_by_signatures(client, "transfer", "signature", signatures).await?; restored += unmark_by_signatures(client, "token", "signature", signatures).await?; restored += unmark_by_signatures(client, "issue_token", "signature", signatures).await?; restored += unmark_by_signatures(client, "burn", "signature", signatures).await?; restored += unmark_by_signatures(client, "nft", "signature", signatures).await?; restored += unmark_by_signatures(client, "marketing", "signature", signatures).await?; restored += unmark_by_signatures(client, "vanity_address", "signature", signatures).await?; restored += unmark_by_signatures(client, "swap", "signature1", signatures).await?; restored += unmark_by_signatures(client, "swap", "signature2", signatures).await?; restored += unmark_by_signatures(client, "loan_contract", "signature1", signatures).await?; restored += unmark_by_signatures(client, "loan_contract", "signature2", signatures).await?; restored += unmark_by_signatures(client, "loan_payment", "signature", signatures).await?; restored += unmark_by_signatures(client, "collateral_claim", "signature", signatures).await?; Ok(restored > 0) } pub fn spawn_processed_cleanup(saved_block_number: u32) { // Cleanup trails the chain tip by a small depth so recent processed // mempool rows can still be restored during short orphan events. if saved_block_number <= CLEANUP_DEPTH { return; } if CLEANUP_RUNNING .compare_exchange( false, true, crate::AtomicOrdering::SeqCst, crate::AtomicOrdering::SeqCst, ) .is_err() { return; } task::spawn(async move { let safe_block = saved_block_number.saturating_sub(CLEANUP_DEPTH); // Cleanup is deliberately delayed behind the tip so short reorgs can // still restore recently processed rows. if let Err(err) = delete_processed_before_or_at(safe_block, CLEANUP_BATCH_LIMIT).await { eprintln!( "[mempool_cleanup] failed: saved_block={saved_block_number} safe_block={safe_block} err={err}" ); } CLEANUP_RUNNING.store(false, crate::AtomicOrdering::SeqCst); }); } pub async fn mark_processed_by_signatures(signatures: &[String], block_number: u32) -> Result<()> { // Synced blocks arrive with signatures instead of selected-row IDs, // so processed marking on the updating path works by signature. if signatures.is_empty() { return Ok(()); } let client_handle = db_client().await?; let client = client_handle.as_ref(); let bn = block_number as i32; // Remote/synced blocks do not know local row IDs, so they mark by // transaction signatures instead. client .execute( "UPDATE transfer SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE token SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE issue_token SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE burn SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE nft SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE marketing SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE vanity_address SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE swap SET processed=true, processed_block_number=$1 WHERE signature1 = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE swap SET processed=true, processed_block_number=$1 WHERE signature2 = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE loan_contract SET processed=true, processed_block_number=$1 WHERE signature1 = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE loan_contract SET processed=true, processed_block_number=$1 WHERE signature2 = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE loan_payment SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; client .execute( "UPDATE collateral_claim SET processed=true, processed_block_number=$1 WHERE signature = ANY($2)", &[&bn, &signatures], ) .await?; Ok(()) } pub async fn delete_by_signatures(signatures: &[String]) -> Result<()> { // Some validation failures need to remove mempool rows directly by // signature regardless of which table the transaction lives in. if signatures.is_empty() { return Ok(()); } let client_handle = db_client().await?; let client = client_handle.as_ref(); // Failed validation removes every matching pending row no matter which // transaction table currently owns the signature. client .execute( "DELETE FROM transfer WHERE signature = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM token WHERE signature = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM issue_token WHERE signature = ANY($1)", &[&signatures], ) .await?; client .execute("DELETE FROM burn WHERE signature = ANY($1)", &[&signatures]) .await?; client .execute("DELETE FROM nft WHERE signature = ANY($1)", &[&signatures]) .await?; client .execute( "DELETE FROM marketing WHERE signature = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM vanity_address WHERE signature = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM swap WHERE signature1 = ANY($1) OR signature2 = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM loan_contract WHERE signature1 = ANY($1) OR signature2 = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM loan_payment WHERE signature = ANY($1)", &[&signatures], ) .await?; client .execute( "DELETE FROM collateral_claim WHERE signature = ANY($1)", &[&signatures], ) .await?; Ok(()) }