authored_delivery_fact_tests.rs (10126B)
1 use super::signed_fact_fixture::{RAW, claim, ids, record}; 2 use super::*; 3 use crate::OpenMode; 4 use radroots_storage::{ 5 authored::WorkClaim, 6 authored_atomic::{CancelAuthoredWork, ClaimAuthoredWork, RecordDeliveryFact}, 7 }; 8 use radroots_transport::{DeliveryReceipt, outcome::DeliveryOutcome, sink::DeliveryTargetReceipt}; 9 use tempfile::TempDir; 10 11 async fn prepared(temp: &TempDir) -> (SqliteStorage, WorkClaim, AuthoredAtomicReceipt) { 12 let (store, event, signing) = super::signed_fact_tests::prepared(temp).await; 13 store 14 .execute_authored(record(event, signing, 12)) 15 .await 16 .unwrap(); 17 let plan = store 18 .authored_delivery_plan(ids().2) 19 .await 20 .unwrap() 21 .unwrap(); 22 let active = claim(plan.revision(), 5, 13); 23 let original = store 24 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 25 ClaimAuthoredTarget::DeliveryPlan(ids().2), 26 active.clone(), 27 ))) 28 .await 29 .unwrap(); 30 (store, active, original) 31 } 32 33 pub(super) fn fact(plan: &AuthoredDeliveryPlan, active: WorkClaim) -> AuthoredAtomicCommand { 34 let request = plan.request().unwrap(); 35 let receipt = DeliveryReceipt::for_request( 36 request, 37 request 38 .target_set() 39 .targets() 40 .iter() 41 .cloned() 42 .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted())) 43 .collect(), 44 ) 45 .unwrap(); 46 AuthoredAtomicCommand::RecordDelivery( 47 RecordDeliveryFact::new( 48 ids().2, 49 ids().1, 50 active, 51 DeliveryAttemptOutcome::Receipt(receipt), 52 50, 53 ) 54 .unwrap(), 55 ) 56 } 57 58 #[tokio::test] 59 async fn stopped_late_delivery_reopens_and_preserves_original_receipt_and_exact_raw() { 60 let temp = TempDir::new().unwrap(); 61 let (store, active, original) = prepared(&temp).await; 62 let plan = store 63 .authored_delivery_plan(ids().2) 64 .await 65 .unwrap() 66 .unwrap(); 67 let command = fact(&plan, active); 68 store 69 .execute_authored(AuthoredAtomicCommand::Cancel( 70 CancelAuthoredWork::new( 71 CancelAuthoredTarget::DeliveryPlan(ids().2), 72 plan.revision(), 73 20, 74 ) 75 .unwrap(), 76 )) 77 .await 78 .unwrap(); 79 let receipt = store.execute_authored(command.clone()).await.unwrap(); 80 let retained = store 81 .authored_delivery_plan(ids().2) 82 .await 83 .unwrap() 84 .unwrap(); 85 assert_eq!(retained.state(), AuthoredDeliveryState::Cancelled); 86 assert_eq!(retained.stop_requested_at_unix_ms(), Some(20)); 87 assert_eq!( 88 retained.delivery_satisfaction().unwrap(), 89 SatisfactionState::Satisfied 90 ); 91 assert_eq!( 92 retained.request().unwrap().payload().event().raw_json(), 93 RAW 94 ); 95 assert_eq!(retained.delivery_facts().len(), 1); 96 for sql in [ 97 "UPDATE radroots_runtime_authored_delivery_facts SET observed_at_unix_ms = 99", 98 "DELETE FROM radroots_runtime_authored_delivery_facts", 99 "UPDATE radroots_runtime_authored_delivery_plans SET stop_requested_at_unix_ms = NULL", 100 ] { 101 assert!(sqlx::query(sql).execute(store.pool()).await.is_err()); 102 } 103 store.close().await.unwrap(); 104 let store = super::signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await; 105 assert_eq!( 106 store 107 .authored_delivery_plan(ids().2) 108 .await 109 .unwrap() 110 .unwrap(), 111 retained 112 ); 113 assert_eq!( 114 store 115 .authored_receipt(original.commit_id()) 116 .await 117 .unwrap() 118 .unwrap(), 119 original 120 ); 121 let replay = store.execute_authored(command).await.unwrap(); 122 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 123 assert_eq!(replay.outcome(), receipt.outcome()); 124 assert_eq!( 125 store 126 .authored_delivery_plan(ids().2) 127 .await 128 .unwrap() 129 .unwrap(), 130 retained 131 ); 132 store.close().await.unwrap(); 133 } 134 135 #[tokio::test] 136 async fn real_delivery_commit_failure_rolls_back_facts_snapshot_and_receipt() { 137 let temp = TempDir::new().unwrap(); 138 let (store, active, _) = prepared(&temp).await; 139 let before = store 140 .authored_delivery_plan(ids().2) 141 .await 142 .unwrap() 143 .unwrap(); 144 let command = fact(&before, active); 145 sqlx::query("CREATE TABLE delivery_fact_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)").execute(store.pool()).await.unwrap(); 146 sqlx::query("CREATE TRIGGER delivery_fact_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_atomic_commits WHEN NEW.phase IN ('delivery', 'cancel') BEGIN INSERT INTO delivery_fact_commit_fault VALUES (x'99999999999999999999999999999999'); END").execute(store.pool()).await.unwrap(); 147 let stop = AuthoredAtomicCommand::Cancel( 148 CancelAuthoredWork::new( 149 CancelAuthoredTarget::DeliveryPlan(ids().2), 150 before.revision(), 151 20, 152 ) 153 .unwrap(), 154 ); 155 assert!(store.execute_authored(stop.clone()).await.is_err()); 156 assert_eq!( 157 store 158 .authored_delivery_plan(ids().2) 159 .await 160 .unwrap() 161 .unwrap(), 162 before 163 ); 164 assert!( 165 store 166 .authored_receipt(stop.commit_id()) 167 .await 168 .unwrap() 169 .is_none() 170 ); 171 assert!(store.execute_authored(command.clone()).await.is_err()); 172 assert_eq!( 173 store 174 .authored_delivery_plan(ids().2) 175 .await 176 .unwrap() 177 .unwrap(), 178 before 179 ); 180 assert!( 181 store 182 .authored_receipt(command.commit_id()) 183 .await 184 .unwrap() 185 .is_none() 186 ); 187 assert_eq!( 188 sqlx::query_scalar::<_, i64>( 189 "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_facts" 190 ) 191 .fetch_one(store.pool()) 192 .await 193 .unwrap(), 194 0 195 ); 196 assert_eq!( 197 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_fact_commit_fault") 198 .fetch_one(store.pool()) 199 .await 200 .unwrap(), 201 0 202 ); 203 sqlx::query("DROP TRIGGER delivery_fact_commit_fault_trigger") 204 .execute(store.pool()) 205 .await 206 .unwrap(); 207 sqlx::query("DROP TABLE delivery_fact_commit_fault") 208 .execute(store.pool()) 209 .await 210 .unwrap(); 211 store.close().await.unwrap(); 212 let store = super::signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await; 213 assert_eq!( 214 store 215 .authored_delivery_plan(ids().2) 216 .await 217 .unwrap() 218 .unwrap(), 219 before 220 ); 221 store.execute_authored(command.clone()).await.unwrap(); 222 assert_eq!( 223 store 224 .authored_delivery_plan(ids().2) 225 .await 226 .unwrap() 227 .unwrap() 228 .delivery_satisfaction() 229 .unwrap(), 230 SatisfactionState::Satisfied 231 ); 232 store.close().await.unwrap(); 233 } 234 235 #[tokio::test] 236 async fn forged_claim_and_normalized_fact_corruption_fail_closed() { 237 let temp = TempDir::new().unwrap(); 238 let (store, active, _) = prepared(&temp).await; 239 let before = store 240 .authored_delivery_plan(ids().2) 241 .await 242 .unwrap() 243 .unwrap(); 244 let forged = WorkClaim::new( 245 *active.token(), 246 "forged-owner", 247 active.generation(), 248 active.acquired_at_unix_ms(), 249 active.expires_at_unix_ms(), 250 active.row_revision(), 251 ) 252 .unwrap(); 253 assert_eq!( 254 store.execute_authored(fact(&before, forged)).await, 255 Err(Error::AtomicWorkflowMismatch) 256 ); 257 assert_eq!( 258 store 259 .authored_delivery_plan(ids().2) 260 .await 261 .unwrap() 262 .unwrap(), 263 before 264 ); 265 store.execute_authored(fact(&before, active)).await.unwrap(); 266 sqlx::query("DROP TRIGGER radroots_runtime_authored_delivery_facts_update_guard") 267 .execute(store.pool()) 268 .await 269 .unwrap(); 270 sqlx::query("UPDATE radroots_runtime_authored_delivery_facts SET observed_at_unix_ms = 99") 271 .execute(store.pool()) 272 .await 273 .unwrap(); 274 assert_eq!( 275 store.authored_delivery_plan(ids().2).await, 276 Err(Error::InvalidAuthoredDeliveryPlan) 277 ); 278 store.close().await.unwrap(); 279 } 280 281 #[tokio::test] 282 async fn delivery_read_snapshot_remains_consistent_across_a_late_commit() { 283 let temp = TempDir::new().unwrap(); 284 let (store, active, _) = prepared(&temp).await; 285 let before = store 286 .authored_delivery_plan(ids().2) 287 .await 288 .unwrap() 289 .unwrap(); 290 let mut read = store.pool().begin().await.unwrap(); 291 let frozen = load_optional_plan_tx(&mut read, ids().2) 292 .await 293 .unwrap() 294 .unwrap(); 295 assert_eq!(frozen, before); 296 store.execute_authored(fact(&before, active)).await.unwrap(); 297 validate_plan_children_tx(&mut read, &frozen).await.unwrap(); 298 assert_eq!( 299 load_optional_plan_tx(&mut read, ids().2) 300 .await 301 .unwrap() 302 .unwrap(), 303 frozen 304 ); 305 read.rollback().await.unwrap(); 306 let current = store 307 .authored_delivery_plan(ids().2) 308 .await 309 .unwrap() 310 .unwrap(); 311 assert_eq!(current.delivery_facts().len(), 1); 312 assert_eq!(current.revision(), before.revision()); 313 let mut stale = store.pool().begin().await.unwrap(); 314 assert_eq!( 315 delivery_facts::persist(&mut stale, &before).await, 316 Err(Error::InvalidAuthoredDeliveryPlan) 317 ); 318 stale.rollback().await.unwrap(); 319 assert_eq!( 320 store 321 .authored_delivery_plan(ids().2) 322 .await 323 .unwrap() 324 .unwrap(), 325 current 326 ); 327 store.close().await.unwrap(); 328 }