authored_delivery_reconciliation.rs (7202B)
1 //! Indexed issued claims and immutable fact-to-attempt provenance. 2 3 use super::{ 4 column, decode_receipt_row, load_artifact_tx, load_optional_plan_tx, map_backend, persist_plan, 5 }; 6 use radroots_storage::{ 7 Error, 8 atomic::AtomicCommitId, 9 authored_atomic::{ 10 AuthoredAtomicCommand, AuthoredAtomicReceipt, ClaimAuthoredTarget, ClaimAuthoredWork, 11 ReconcileDeliveryFacts, 12 }, 13 authored_delivery::{ 14 AuthoredDeliveryHistory, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, 15 DELIVERY_PLAN_ATTEMPTS_MAX, 16 }, 17 }; 18 use sqlx::{Sqlite, sqlite::SqliteRow}; 19 20 const MARKERS: &str = "SELECT attempt, claim_id FROM radroots_runtime_authored_delivery_reconciliations WHERE plan_id = ? ORDER BY attempt LIMIT 1025"; 21 22 async fn receipt( 23 transaction: &mut sqlx::Transaction<'_, Sqlite>, 24 id: &[u8], 25 ) -> Result<AuthoredAtomicReceipt, Error> { 26 let row = sqlx::query("SELECT commit_id, commit_digest, requested_at_unix_ms, committed_at_unix_ms, receipt FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?") 27 .bind(id).fetch_one(&mut **transaction).await.map_err(map_backend)?; 28 decode_receipt_row(&row) 29 } 30 31 pub(super) async fn history( 32 transaction: &mut sqlx::Transaction<'_, Sqlite>, 33 plan_id: AuthoredDeliveryPlanId, 34 ) -> Result<Option<AuthoredDeliveryHistory>, Error> { 35 let Some(plan) = load_optional_plan_tx(transaction, plan_id).await? else { 36 return Ok(None); 37 }; 38 let operation_id = load_artifact_tx(transaction, plan.artifact_id()) 39 .await? 40 .operation_id(); 41 let preparations: Vec<Vec<u8>> = sqlx::query_scalar("SELECT commit_id FROM radroots_runtime_authored_atomic_commits WHERE target_id = ? AND phase = 'prepare' AND json_type(CAST(receipt AS TEXT), '$.outcome.prepared') = 'object' ORDER BY commit_id LIMIT 2") 42 .bind(operation_id.as_bytes().as_slice()).fetch_all(&mut **transaction).await.map_err(map_backend)?; 43 if preparations.len() > 1 { 44 return Err(Error::AtomicWorkflowMismatch); 45 } 46 let original = match preparations.first() { 47 Some(id) => Some(receipt(transaction, id).await?), 48 None => None, 49 }; 50 let mut history = AuthoredDeliveryHistory::new(plan, original.as_ref())?; 51 drop(original); 52 // Fetch only bounded IDs, then decode one bounded original receipt at a time. 53 let ids: Vec<Vec<u8>> = sqlx::query_scalar("SELECT claim_id FROM radroots_runtime_authored_delivery_claims WHERE plan_id = ? ORDER BY claim_id LIMIT 1025") 54 .bind(plan_id.as_bytes().as_slice()).fetch_all(&mut **transaction).await.map_err(map_backend)?; 55 if ids.len() > DELIVERY_PLAN_ATTEMPTS_MAX as usize { 56 history.mark_truncated(); 57 } 58 for id in ids.iter().take(DELIVERY_PLAN_ATTEMPTS_MAX as usize) { 59 history.push_claim(&receipt(transaction, id).await?)?; 60 } 61 history.validate()?; 62 Ok(Some(history)) 63 } 64 65 pub(super) async fn record_claim( 66 transaction: &mut sqlx::Transaction<'_, Sqlite>, 67 command: &AuthoredAtomicCommand, 68 ) -> Result<(), Error> { 69 let AuthoredAtomicCommand::Claim(value) = command else { 70 return Ok(()); 71 }; 72 let ClaimAuthoredTarget::DeliveryPlan(plan_id) = value.target() else { 73 return Ok(()); 74 }; 75 let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM (SELECT 1 FROM radroots_runtime_authored_delivery_claims WHERE plan_id = ? LIMIT 1024)") 76 .bind(plan_id.as_bytes().as_slice()).fetch_one(&mut **transaction).await.map_err(map_backend)?; 77 if count >= i64::from(DELIVERY_PLAN_ATTEMPTS_MAX) { 78 return Err(Error::DeliveryAttemptOverflow); 79 } 80 sqlx::query( 81 "INSERT INTO radroots_runtime_authored_delivery_claims (plan_id, claim_id) VALUES (?, ?)", 82 ) 83 .bind(plan_id.as_bytes().as_slice()) 84 .bind(command.commit_id().as_bytes().as_slice()) 85 .execute(&mut **transaction) 86 .await 87 .map_err(map_backend)?; 88 Ok(()) 89 } 90 91 pub(super) async fn reconcile( 92 transaction: &mut sqlx::Transaction<'_, Sqlite>, 93 command: &ReconcileDeliveryFacts, 94 ) -> Result<AuthoredDeliveryPlan, Error> { 95 let history = history(transaction, command.plan_id()) 96 .await? 97 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 98 let plan = command.apply_to(&history)?; 99 persist_plan(transaction, &plan).await?; 100 Ok(plan) 101 } 102 103 fn marker_id( 104 plan_id: AuthoredDeliveryPlanId, 105 claim: &radroots_storage::authored::WorkClaim, 106 ) -> AtomicCommitId { 107 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 108 ClaimAuthoredTarget::DeliveryPlan(plan_id), 109 claim.clone(), 110 )) 111 .commit_id() 112 } 113 114 fn validate_existing( 115 plan: &AuthoredDeliveryPlan, 116 rows: &[SqliteRow], 117 ) -> Result<std::collections::BTreeSet<u32>, Error> { 118 let mut existing = std::collections::BTreeSet::new(); 119 for row in rows { 120 let ordinal = u32::try_from(column::<i64>(row, "attempt")?) 121 .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?; 122 let index = ordinal 123 .checked_sub(1) 124 .and_then(|index| usize::try_from(index).ok()) 125 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 126 let claim = plan 127 .attempts() 128 .get(index) 129 .and_then(|attempt| attempt.claim_evidence()) 130 .ok_or(Error::InvalidAuthoredDeliveryPlan)?; 131 if !existing.insert(ordinal) 132 || column::<Vec<u8>>(row, "claim_id")?.as_slice() 133 != marker_id(plan.plan_id(), claim).as_bytes() 134 { 135 return Err(Error::InvalidAuthoredDeliveryPlan); 136 } 137 } 138 Ok(existing) 139 } 140 141 pub(super) async fn persist( 142 transaction: &mut sqlx::Transaction<'_, Sqlite>, 143 plan: &AuthoredDeliveryPlan, 144 ) -> Result<(), Error> { 145 let existing = sqlx::query(MARKERS) 146 .bind(plan.plan_id().as_bytes().as_slice()) 147 .fetch_all(&mut **transaction) 148 .await 149 .map_err(map_backend)?; 150 let existing = validate_existing(plan, &existing)?; 151 for (attempt, claim) in plan 152 .attempts() 153 .iter() 154 .filter_map(|attempt| { 155 attempt 156 .claim_evidence() 157 .map(|claim| (attempt.attempt(), claim)) 158 }) 159 .filter(|(attempt, _)| !existing.contains(&attempt.get())) 160 { 161 sqlx::query("INSERT INTO radroots_runtime_authored_delivery_reconciliations (plan_id, attempt, claim_id) VALUES (?, ?, ?)") 162 .bind(plan.plan_id().as_bytes().as_slice()).bind(i64::from(attempt.get())) 163 .bind(marker_id(plan.plan_id(), claim).as_bytes().as_slice()) 164 .execute(&mut **transaction).await.map_err(map_backend)?; 165 } 166 Ok(()) 167 } 168 169 pub(super) async fn validate( 170 transaction: &mut sqlx::Transaction<'_, Sqlite>, 171 plan: &AuthoredDeliveryPlan, 172 ) -> Result<(), Error> { 173 let rows = sqlx::query(MARKERS) 174 .bind(plan.plan_id().as_bytes().as_slice()) 175 .fetch_all(&mut **transaction) 176 .await 177 .map_err(map_backend)?; 178 let count = plan 179 .attempts() 180 .iter() 181 .filter(|attempt| attempt.claim_evidence().is_some()) 182 .count(); 183 if rows.len() != count { 184 return Err(Error::InvalidAuthoredDeliveryPlan); 185 } 186 validate_existing(plan, &rows).map(|_| ()) 187 }