authored_delivery_facts.rs (3369B)
1 //! Normalized, append-only delivery facts in the authored transaction. 2 3 use crate::authored::{column, decode_snapshot, encode_snapshot, i64_from_u64, map_backend}; 4 use radroots_storage::{ 5 Error, 6 authored_atomic::{AuthoredAtomicCommand, ClaimAuthoredTarget, ClaimAuthoredWork}, 7 authored_delivery::{AuthoredDeliveryFact, AuthoredDeliveryPlan}, 8 }; 9 use sqlx::{Sqlite, sqlite::SqliteRow}; 10 11 pub(crate) const SELECT: &str = "SELECT ordinal, claim_id, observed_at_unix_ms, fact_snapshot FROM radroots_runtime_authored_delivery_facts WHERE plan_id = ? ORDER BY ordinal"; 12 13 fn claim_id( 14 plan: &AuthoredDeliveryPlan, 15 fact: &AuthoredDeliveryFact, 16 ) -> radroots_storage::atomic::AtomicCommitId { 17 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 18 ClaimAuthoredTarget::DeliveryPlan(plan.plan_id()), 19 fact.claim().clone(), 20 )) 21 .commit_id() 22 } 23 24 pub(crate) async fn persist( 25 transaction: &mut sqlx::Transaction<'_, Sqlite>, 26 plan: &AuthoredDeliveryPlan, 27 ) -> Result<(), Error> { 28 sqlx::query("UPDATE radroots_runtime_authored_delivery_plans SET stop_requested_at_unix_ms = ? WHERE plan_id = ?") 29 .bind(plan.stop_requested_at_unix_ms().map(i64_from_u64).transpose()?) 30 .bind(plan.plan_id().as_bytes().as_slice()) 31 .execute(&mut **transaction).await.map_err(map_backend)?; 32 let existing = sqlx::query(SELECT) 33 .bind(plan.plan_id().as_bytes().as_slice()) 34 .fetch_all(&mut **transaction) 35 .await 36 .map_err(map_backend)?; 37 if existing.len() > plan.delivery_facts().len() { 38 return Err(Error::InvalidAuthoredDeliveryPlan); 39 } 40 validate_prefix(plan, &existing)?; 41 for (ordinal, fact) in plan 42 .delivery_facts() 43 .iter() 44 .enumerate() 45 .skip(existing.len()) 46 { 47 sqlx::query("INSERT INTO radroots_runtime_authored_delivery_facts (plan_id, ordinal, claim_id, observed_at_unix_ms, fact_snapshot) VALUES (?, ?, ?, ?, ?)") 48 .bind(plan.plan_id().as_bytes().as_slice()) 49 .bind(i64::try_from(ordinal).map_err(|_| Error::InvalidAuthoredDeliveryPlan)?) 50 .bind(claim_id(plan, fact).as_bytes().as_slice()) 51 .bind(i64_from_u64(fact.observed_at_unix_ms())?) 52 .bind(encode_snapshot(fact)?) 53 .execute(&mut **transaction).await.map_err(map_backend)?; 54 } 55 Ok(()) 56 } 57 58 pub(crate) fn validate(plan: &AuthoredDeliveryPlan, rows: &[SqliteRow]) -> Result<(), Error> { 59 if rows.len() != plan.delivery_facts().len() { 60 return Err(Error::InvalidAuthoredDeliveryPlan); 61 } 62 validate_prefix(plan, rows) 63 } 64 65 fn validate_prefix(plan: &AuthoredDeliveryPlan, rows: &[SqliteRow]) -> Result<(), Error> { 66 for (ordinal, (row, expected)) in rows.iter().zip(plan.delivery_facts()).enumerate() { 67 let decoded = decode_snapshot::<AuthoredDeliveryFact>(column(row, "fact_snapshot")?)?; 68 if column::<i64>(row, "ordinal")? 69 != i64::try_from(ordinal).map_err(|_| Error::InvalidAuthoredDeliveryPlan)? 70 || column::<Vec<u8>>(row, "claim_id")?.as_slice() != claim_id(plan, expected).as_bytes() 71 || column::<i64>(row, "observed_at_unix_ms")? 72 != i64_from_u64(expected.observed_at_unix_ms())? 73 || decoded != *expected 74 { 75 return Err(Error::InvalidAuthoredDeliveryPlan); 76 } 77 } 78 Ok(()) 79 }