authored_durability_tests.rs (11473B)
1 use super::*; 2 use crate::{OpenMode, OpenOptions, Paths}; 3 use radroots_storage::event::SourceGeneration; 4 use std::path::Path; 5 use tempfile::TempDir; 6 7 #[path = "authored_durability_crash_tests.rs"] 8 mod crash; 9 10 async fn open(directory: &Path, mode: OpenMode) -> SqliteStorage { 11 let mut options = OpenOptions::new(Paths::from_directory(directory).unwrap(), mode); 12 if matches!(mode, OpenMode::Create) { 13 options = options 14 .with_source_generation(SourceGeneration::new([91; 32]).unwrap(), 9) 15 .unwrap(); 16 } 17 SqliteStorage::open(options).await.unwrap() 18 } 19 20 fn first() -> AuthoredDraft { 21 AuthoredDraft::initial( 22 AuthoredDraftId::new([41; 16]).unwrap(), 23 [7; 32], 24 "radroots.durability-fixture.v1", 25 b"acknowledged baseline".to_vec(), 26 AuthoredDraftStage::Draft, 27 None, 28 10, 29 ) 30 .unwrap() 31 } 32 33 fn next(previous: &AuthoredDraft, payload: Vec<u8>) -> AuthoredDraft { 34 previous 35 .successor(payload, AuthoredDraftStage::Draft, None, 11) 36 .unwrap() 37 } 38 39 async fn assert_head(store: &SqliteStorage, expected: &AuthoredDraft) { 40 assert_eq!( 41 store 42 .authored_draft_head(expected.draft_id()) 43 .await 44 .unwrap(), 45 Some(expected.clone()) 46 ); 47 } 48 49 #[tokio::test] 50 async fn authored_durability_commit_fault_never_acknowledges_and_exact_retry_recovers() { 51 let temp = TempDir::new().unwrap(); 52 let store = open(temp.path(), OpenMode::Create).await; 53 let baseline = first(); 54 let pending = next(&baseline, b"pending complete revision".to_vec()); 55 store 56 .append_authored_draft(baseline.clone(), None) 57 .await 58 .unwrap(); 59 sqlx::query("CREATE TABLE authored_commit_parent (id INTEGER PRIMARY KEY)") 60 .execute(store.pool()) 61 .await 62 .unwrap(); 63 sqlx::query("CREATE TABLE authored_commit_fault (id INTEGER REFERENCES authored_commit_parent(id) DEFERRABLE INITIALLY DEFERRED)") 64 .execute(store.pool()).await.unwrap(); 65 sqlx::query("CREATE TRIGGER authored_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_draft_revisions BEGIN INSERT INTO authored_commit_fault VALUES (99); END") 66 .execute(store.pool()).await.unwrap(); 67 68 // Establish that the insertion succeeds and the actual COMMIT is the fault. 69 let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); 70 insert_draft_tx(&mut transaction, &pending).await.unwrap(); 71 let failure = transaction.commit().await.unwrap_err(); 72 assert_eq!( 73 failure.as_database_error().unwrap().kind(), 74 sqlx::error::ErrorKind::ForeignKeyViolation 75 ); 76 assert_eq!( 77 store 78 .append_authored_draft(pending.clone(), Some(baseline.revision())) 79 .await, 80 Err(Error::BackendUnavailable) 81 ); 82 assert_head(&store, &baseline).await; 83 assert!( 84 store 85 .authored_draft_revision(pending.draft_id(), pending.revision()) 86 .await 87 .unwrap() 88 .is_none() 89 ); 90 91 for statement in [ 92 "DROP TRIGGER authored_commit_fault_trigger", 93 "DROP TABLE authored_commit_fault", 94 "DROP TABLE authored_commit_parent", 95 ] { 96 sqlx::query(statement).execute(store.pool()).await.unwrap(); 97 } 98 let receipt = store 99 .append_authored_draft(pending.clone(), Some(baseline.revision())) 100 .await 101 .unwrap(); 102 assert_eq!(receipt.disposition(), DraftAppendDisposition::Inserted); 103 store.close().await.unwrap(); 104 let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; 105 assert_head(&reopened, &pending).await; 106 assert_eq!( 107 reopened 108 .append_authored_draft(pending, Some(baseline.revision())) 109 .await 110 .unwrap() 111 .disposition(), 112 DraftAppendDisposition::Replay 113 ); 114 reopened.close().await.unwrap(); 115 } 116 117 #[tokio::test] 118 async fn authored_durability_sqlite_capacity_failure_preserves_the_acknowledged_head() { 119 let temp = TempDir::new().unwrap(); 120 let store = open(temp.path(), OpenMode::Create).await; 121 let baseline = first(); 122 store 123 .append_authored_draft(baseline.clone(), None) 124 .await 125 .unwrap(); 126 let mut connections = Vec::new(); 127 for _ in 0..4 { 128 connections.push(store.pool().acquire().await.unwrap()); 129 } 130 let mut original_limits = Vec::new(); 131 for connection in &mut connections { 132 original_limits.push( 133 sqlx::query_scalar::<_, i64>("PRAGMA max_page_count") 134 .fetch_one(&mut **connection) 135 .await 136 .unwrap(), 137 ); 138 let pages = sqlx::query_scalar::<_, i64>("PRAGMA page_count") 139 .fetch_one(&mut **connection) 140 .await 141 .unwrap(); 142 // PRAGMA assignments do not accept bind parameters; only an i64 is interpolated. 143 let limit = sqlx::query_scalar::<_, i64>(sqlx::AssertSqlSafe(format!( 144 "PRAGMA max_page_count = {}", 145 pages + 2 146 ))) 147 .fetch_one(&mut **connection) 148 .await 149 .unwrap(); 150 assert_eq!(limit, pages + 2); 151 } 152 drop(connections); 153 // This bounded allocation proves SQLITE_FULL without filling the host disk. 154 let failure = 155 sqlx::query("CREATE TABLE authored_capacity_probe AS SELECT zeroblob(1048576) AS data") 156 .execute(store.pool()) 157 .await 158 .unwrap_err(); 159 assert_eq!( 160 failure.as_database_error().unwrap().code().as_deref(), 161 Some("13") 162 ); 163 let pending = next(&baseline, vec![42; 256 * 1024]); 164 assert_eq!( 165 store 166 .append_authored_draft(pending.clone(), Some(baseline.revision())) 167 .await, 168 Err(Error::SpaceInsufficient) 169 ); 170 assert_head(&store, &baseline).await; 171 172 let mut connections = Vec::new(); 173 for _ in 0..4 { 174 connections.push(store.pool().acquire().await.unwrap()); 175 } 176 for (connection, limit) in connections.iter_mut().zip(original_limits) { 177 sqlx::query(sqlx::AssertSqlSafe(format!( 178 "PRAGMA max_page_count = {limit}" 179 ))) 180 .execute(&mut **connection) 181 .await 182 .unwrap(); 183 } 184 drop(connections); 185 store.close().await.unwrap(); 186 let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; 187 assert_head(&reopened, &baseline).await; 188 assert_eq!( 189 reopened 190 .append_authored_draft(pending.clone(), Some(baseline.revision())) 191 .await 192 .unwrap() 193 .disposition(), 194 DraftAppendDisposition::Inserted 195 ); 196 assert_head(&reopened, &pending).await; 197 reopened.close().await.unwrap(); 198 } 199 200 #[tokio::test] 201 async fn authored_durability_denied_writes_have_no_receipt_or_head_advance() { 202 let temp = TempDir::new().unwrap(); 203 let store = open(temp.path(), OpenMode::Create).await; 204 let baseline = first(); 205 let pending = next(&baseline, b"denied revision".to_vec()); 206 store 207 .append_authored_draft(baseline.clone(), None) 208 .await 209 .unwrap(); 210 let mut connections = Vec::new(); 211 for _ in 0..4 { 212 connections.push(store.pool().acquire().await.unwrap()); 213 } 214 for connection in &mut connections { 215 sqlx::query("PRAGMA query_only = ON") 216 .execute(&mut **connection) 217 .await 218 .unwrap(); 219 } 220 drop(connections); 221 assert_eq!( 222 store 223 .append_authored_draft(pending, Some(baseline.revision())) 224 .await, 225 Err(Error::BackendUnavailable) 226 ); 227 assert_head(&store, &baseline).await; 228 store.close().await.unwrap(); 229 let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; 230 assert_head(&reopened, &baseline).await; 231 reopened.close().await.unwrap(); 232 } 233 234 #[tokio::test] 235 async fn authored_durability_busy_writer_never_returns_a_success_receipt() { 236 let temp = TempDir::new().unwrap(); 237 let store = SqliteStorage::open( 238 OpenOptions::new( 239 Paths::from_directory(temp.path()).unwrap(), 240 OpenMode::Create, 241 ) 242 .with_source_generation(SourceGeneration::new([91; 32]).unwrap(), 9) 243 .unwrap() 244 .with_busy_timeout(std::time::Duration::from_millis(10)) 245 .unwrap(), 246 ) 247 .await 248 .unwrap(); 249 let baseline = first(); 250 let pending = next(&baseline, b"blocked complete revision".to_vec()); 251 store 252 .append_authored_draft(baseline.clone(), None) 253 .await 254 .unwrap(); 255 let transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); 256 assert_eq!( 257 store 258 .append_authored_draft(pending.clone(), Some(baseline.revision())) 259 .await, 260 Err(Error::BackendUnavailable) 261 ); 262 assert_head(&store, &baseline).await; 263 transaction.rollback().await.unwrap(); 264 assert_eq!( 265 store 266 .append_authored_draft(pending.clone(), Some(baseline.revision())) 267 .await 268 .unwrap() 269 .disposition(), 270 DraftAppendDisposition::Inserted 271 ); 272 assert_head(&store, &pending).await; 273 store.close().await.unwrap(); 274 } 275 276 #[tokio::test] 277 async fn authored_durability_read_only_and_closed_stores_cannot_acknowledge() { 278 let temp = TempDir::new().unwrap(); 279 let store = open(temp.path(), OpenMode::Create).await; 280 let baseline = first(); 281 let pending = next(&baseline, b"not writable".to_vec()); 282 store 283 .append_authored_draft(baseline.clone(), None) 284 .await 285 .unwrap(); 286 store.close().await.unwrap(); 287 assert_eq!( 288 store 289 .append_authored_draft(pending.clone(), Some(baseline.revision())) 290 .await, 291 Err(Error::BackendUnavailable) 292 ); 293 let read_only = open(temp.path(), OpenMode::ReadOnly).await; 294 assert_head(&read_only, &baseline).await; 295 assert_eq!( 296 read_only 297 .append_authored_draft(pending, Some(baseline.revision())) 298 .await, 299 Err(Error::BackendUnavailable) 300 ); 301 read_only.close().await.unwrap(); 302 } 303 304 #[test] 305 fn authored_durability_contract_distinguishes_crash_and_power_loss() { 306 let policy: toml::Value = toml::from_str(include_str!( 307 "../../../contracts/storage/failure_injection_policy_v1.toml" 308 )) 309 .unwrap(); 310 let authored = &policy["authored_write"]; 311 assert_eq!( 312 authored["acknowledgment"].as_str(), 313 Some("only_after_successful_commit_or_exact_committed_replay") 314 ); 315 assert_eq!( 316 authored["commit_fault"].as_str(), 317 Some("deferred_foreign_key_at_actual_commit") 318 ); 319 assert_eq!( 320 authored["capacity_fault"].as_str(), 321 Some("bounded_sqlite_max_page_count") 322 ); 323 assert_eq!( 324 authored["write_denied_fault"].as_str(), 325 Some("owned_connection_query_only") 326 ); 327 assert_eq!( 328 authored["busy_fault"].as_str(), 329 Some("owned_begin_immediate_with_bounded_timeout") 330 ); 331 assert_eq!( 332 authored["process_termination_points"] 333 .as_array() 334 .unwrap() 335 .iter() 336 .map(|value| value.as_str().unwrap()) 337 .collect::<Vec<_>>(), 338 ["after_acknowledgment", "during_uncommitted_insert"] 339 ); 340 assert_eq!(authored["power_loss_qualified"].as_bool(), Some(false)); 341 assert_eq!( 342 authored["protected_data_policy_owner"].as_str(), 343 Some("native_host") 344 ); 345 }