lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }