authored_draft_query.rs (15691B)
1 use futures_executor::block_on; 2 use radroots_storage::{ 3 Error, 4 authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, 5 authored_draft_query::{ 6 AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES, AuthoredDraftCursor, AuthoredDraftPage, 7 AuthoredDraftQuery, AuthoredDraftQueryRecord, AuthoredDraftScope, 8 }, 9 memory::MemoryStorage, 10 }; 11 12 fn draft( 13 id: u8, 14 schema: &str, 15 scope: Option<AuthoredDraftScope>, 16 payload: Vec<u8>, 17 ) -> AuthoredDraft { 18 let draft = AuthoredDraft::initial( 19 AuthoredDraftId::new([id; 16]).unwrap(), 20 [9; 32], 21 schema, 22 payload, 23 AuthoredDraftStage::Draft, 24 None, 25 10, 26 ) 27 .unwrap(); 28 scope.map_or_else( 29 || draft.clone(), 30 |scope| draft.clone().with_scope(scope).unwrap(), 31 ) 32 } 33 fn query(scope: Option<AuthoredDraftScope>, limit: u16) -> AuthoredDraftQuery { 34 AuthoredDraftQuery::new([9; 32], "fixture.composer.v1", scope, limit).unwrap() 35 } 36 37 #[test] 38 fn author_wide_cursor_requires_explicit_selection_and_independent_authority() { 39 let q = AuthoredDraftQuery::for_author([9; 32], "fixture.composer.v1", 256).unwrap(); 40 assert!(q.is_author_wide()); 41 assert_eq!(q.scope(), None); 42 assert!(!query(None, 1).is_author_wide()); 43 assert!(AuthoredDraftQuery::for_author([0; 32], "fixture.composer.v1", 1).is_err()); 44 let cursor = q.cursor_after([3; 16]); 45 let bytes = serde_json::to_vec(&cursor).unwrap(); 46 let decoded: AuthoredDraftCursor = serde_json::from_slice(&bytes).unwrap(); 47 assert_eq!(decoded, cursor); 48 assert_eq!( 49 q.clone().with_cursor(&decoded).unwrap().after(), 50 Some([3; 16]) 51 ); 52 assert_eq!(serde_json::to_vec(&decoded).unwrap(), bytes); 53 let wire = serde_json::to_value(&cursor).unwrap(); 54 assert_eq!(wire["schema_version"], 2); 55 assert_eq!(wire["selection"], "author_schema"); 56 assert!(wire["scope"].is_null()); 57 for changed in [ 58 query(None, 1), 59 query(Some(AuthoredDraftScope::new([7; 32]).unwrap()), 1), 60 AuthoredDraftQuery::for_author([8; 32], q.payload_schema().unwrap(), 1).unwrap(), 61 AuthoredDraftQuery::for_author([9; 32], "fixture.other.v1", 1).unwrap(), 62 ] { 63 assert!(changed.with_cursor(&decoded).is_err()); 64 } 65 assert!( 66 q.with_cursor(&query(None, 1).cursor_after([3; 16])) 67 .is_err() 68 ); 69 for (field, value) in [ 70 ("schema_version", serde_json::json!(1)), 71 ("schema_version", serde_json::json!(3)), 72 ("selection", serde_json::Value::Null), 73 ("selection", serde_json::json!("other")), 74 ("scope", serde_json::json!([7; 32].to_vec())), 75 ("author", serde_json::json!([0; 32].to_vec())), 76 ("payload_schema", serde_json::json!("")), 77 ("unknown", serde_json::json!(true)), 78 ] { 79 let mut forged = wire.clone(); 80 forged[field] = value; 81 assert!( 82 serde_json::from_value::<AuthoredDraftCursor>(forged).is_err(), 83 "{field}" 84 ); 85 } 86 let mut absent = wire; 87 let mut absent_scope = absent.clone(); 88 absent_scope.as_object_mut().unwrap().remove("scope"); 89 assert!(serde_json::from_value::<AuthoredDraftCursor>(absent_scope).is_err()); 90 absent.as_object_mut().unwrap().remove("selection"); 91 assert!(serde_json::from_value::<AuthoredDraftCursor>(absent).is_err()); 92 let old = query(None, 1).cursor_after([3; 16]); 93 let mut forbidden = serde_json::to_value(&old).unwrap(); 94 forbidden["selection"] = serde_json::Value::Null; 95 assert!(serde_json::from_value::<AuthoredDraftCursor>(forbidden).is_err()); 96 let expected = format!( 97 "{{\"schema_version\":1,\"author\":{},\"payload_schema\":\"fixture.composer.v1\",\"scope\":null,\"after_id\":{}}}", 98 serde_json::to_string(&[9; 32]).unwrap(), 99 serde_json::to_string(&[3; 16]).unwrap() 100 ); 101 assert_eq!(serde_json::to_string(&old).unwrap(), expected); 102 } 103 104 #[test] 105 fn author_wide_memory_pages_cover_a_thousand_scopes_and_preserve_bounds() { 106 let store = MemoryStorage::default(); 107 let scope = AuthoredDraftScope::new([7; 32]).unwrap(); 108 for id in 1_u128..=1000 { 109 let value = AuthoredDraft::initial( 110 AuthoredDraftId::new(id.to_be_bytes()).unwrap(), 111 [9; 32], 112 "fixture.composer.v1", 113 vec![1], 114 AuthoredDraftStage::Draft, 115 None, 116 10, 117 ) 118 .unwrap(); 119 let value = if id % 2 == 0 { 120 value.with_scope(scope).unwrap() 121 } else { 122 value 123 }; 124 block_on(store.append_authored_draft(value, None)).unwrap(); 125 } 126 for (id, author, schema) in [ 127 (1001_u128, [8; 32], "fixture.composer.v1"), 128 (1002, [9; 32], "fixture.other.v1"), 129 ] { 130 let value = AuthoredDraft::initial( 131 AuthoredDraftId::new(id.to_be_bytes()).unwrap(), 132 author, 133 schema, 134 vec![1], 135 AuthoredDraftStage::Draft, 136 None, 137 10, 138 ) 139 .unwrap(); 140 block_on(store.append_authored_draft(value, None)).unwrap(); 141 } 142 let q = AuthoredDraftQuery::for_author([9; 32], "fixture.composer.v1", 37).unwrap(); 143 let mut cursor = None; 144 let mut expected = 1_u128; 145 loop { 146 let query = cursor.as_ref().map_or_else( 147 || q.clone(), 148 |cursor| q.clone().with_cursor(cursor).unwrap(), 149 ); 150 let page = block_on(store.query_authored_drafts(query)).unwrap(); 151 assert!(page.records().len() <= 37); 152 for record in page.records() { 153 assert_eq!(record.draft_key(), expected.to_be_bytes()); 154 expected += 1; 155 } 156 cursor = page.next_cursor().cloned(); 157 if cursor.is_none() { 158 break; 159 } 160 } 161 assert_eq!(expected, 1001); 162 for id in 1003_u128..=1005 { 163 let value = AuthoredDraft::initial( 164 AuthoredDraftId::new(id.to_be_bytes()).unwrap(), 165 [9; 32], 166 "fixture.large.v1", 167 vec![1; 2 * 1024 * 1024], 168 AuthoredDraftStage::Draft, 169 None, 170 10, 171 ) 172 .unwrap(); 173 let value = if id % 2 == 0 { 174 value.with_scope(scope).unwrap() 175 } else { 176 value 177 }; 178 block_on(store.append_authored_draft(value, None)).unwrap(); 179 } 180 let q = AuthoredDraftQuery::for_author([9; 32], "fixture.large.v1", 256).unwrap(); 181 let first = block_on(store.query_authored_drafts(q.clone())).unwrap(); 182 assert_eq!(first.records().len(), 2); 183 let second = 184 block_on(store.query_authored_drafts(q.with_cursor(first.next_cursor().unwrap()).unwrap())) 185 .unwrap(); 186 assert_eq!(second.records().len(), 1); 187 assert!(second.next_cursor().is_none()); 188 } 189 190 #[test] 191 fn scoped_pages_preserve_revisions_and_do_not_mix_other_schemas_or_scopes() { 192 let store = MemoryStorage::default(); 193 let scope = AuthoredDraftScope::new([7; 32]).unwrap(); 194 let first = draft( 195 2, 196 "fixture.composer.v1", 197 Some(scope), 198 b"unfinished 1.".to_vec(), 199 ); 200 let second = draft( 201 4, 202 "fixture.composer.v1", 203 Some(scope), 204 b"partial date 2026-".to_vec(), 205 ); 206 for value in [ 207 first.clone(), 208 second.clone(), 209 draft(1, "fixture.profile.v1", Some(scope), b"profile".to_vec()), 210 draft(3, "fixture.composer.v1", None, b"other context".to_vec()), 211 ] { 212 block_on(store.append_authored_draft(value, None)).unwrap(); 213 } 214 let page = block_on(store.query_authored_drafts(query(Some(scope), 1))).unwrap(); 215 assert_eq!( 216 page.records(), 217 [AuthoredDraftQueryRecord::Draft(first.clone())] 218 ); 219 assert_eq!(page.records()[0].draft_id().unwrap(), first.draft_id()); 220 assert_eq!(page.records()[0].revision(), first.revision()); 221 let cursor = page.next_cursor().unwrap(); 222 let changed = first 223 .successor(b"later edit".to_vec(), AuthoredDraftStage::Draft, None, 11) 224 .unwrap(); 225 assert_eq!(changed.scope(), Some(scope)); 226 block_on(store.append_authored_draft(changed.clone(), Some(first.revision()))).unwrap(); 227 let next = 228 block_on(store.query_authored_drafts(query(Some(scope), 2).with_cursor(cursor).unwrap())) 229 .unwrap(); 230 assert_eq!(next.records(), [AuthoredDraftQueryRecord::Draft(second)]); 231 assert!(next.next_cursor().is_none()); 232 let fresh = block_on(store.query_authored_drafts(query(Some(scope), 1))).unwrap(); 233 assert_eq!( 234 fresh.records(), 235 [AuthoredDraftQueryRecord::Draft(changed.clone())] 236 ); 237 assert!(changed.with_scope(scope).is_err()); 238 assert!(first.clone().with_scope(scope).is_err()); 239 let mut forged = serde_json::to_value( 240 first 241 .successor(b"next".to_vec(), AuthoredDraftStage::Draft, None, 11) 242 .unwrap(), 243 ) 244 .unwrap(); 245 forged["scope"] = serde_json::json!([8; 32].to_vec()); 246 let forged: AuthoredDraft = serde_json::from_value(forged).unwrap(); 247 assert_eq!( 248 forged.validate_successor_of(&first), 249 Err(Error::DraftRevisionConflict) 250 ); 251 } 252 253 #[test] 254 fn query_and_cursor_reject_invalid_or_changed_authority() { 255 assert!(AuthoredDraftScope::new([0; 32]).is_err()); 256 for schema in ["", " x", "x\n", &"x".repeat(129)] { 257 assert!(AuthoredDraftQuery::new([9; 32], schema, None, 1).is_err()); 258 } 259 for (author, limit) in [([0; 32], 1), ([9; 32], 0), ([9; 32], 257)] { 260 assert!(AuthoredDraftQuery::new(author, "fixture.composer.v1", None, limit).is_err()); 261 } 262 let q = query(None, 256); 263 let cursor = q.cursor_after([0; 16]); 264 for changed in [ 265 AuthoredDraftQuery::new([8; 32], q.payload_schema().unwrap(), None, 1).unwrap(), 266 AuthoredDraftQuery::new([9; 32], "fixture.profile.v1", None, 1).unwrap(), 267 query(Some(AuthoredDraftScope::new([7; 32]).unwrap()), 1), 268 ] { 269 assert!(changed.with_cursor(&cursor).is_err()); 270 } 271 let bytes = serde_json::to_vec(&cursor).unwrap(); 272 let decoded: AuthoredDraftCursor = serde_json::from_slice(&bytes).unwrap(); 273 assert_eq!(decoded, cursor); 274 assert_eq!( 275 q.clone().with_cursor(&decoded).unwrap().after(), 276 Some([0; 16]) 277 ); 278 for (field, value) in [ 279 ("schema_version", serde_json::json!(2)), 280 ("payload_schema", serde_json::json!("")), 281 ("scope", serde_json::json!([0; 32].to_vec())), 282 ("unknown", serde_json::json!(true)), 283 ] { 284 let mut wire = serde_json::to_value(&cursor).unwrap(); 285 wire[field] = value; 286 assert!(serde_json::from_value::<AuthoredDraftCursor>(wire).is_err()); 287 } 288 } 289 290 #[test] 291 fn legacy_unscoped_wire_and_scoped_serde_keep_exact_payload_bytes() { 292 let legacy = draft(1, "fixture.composer.v1", None, b"0.\n2026-".to_vec()); 293 let bytes = serde_json::to_vec(&legacy).unwrap(); 294 assert!( 295 !String::from_utf8(bytes.clone()) 296 .unwrap() 297 .contains("\"scope\"") 298 ); 299 let restored: AuthoredDraft = serde_json::from_slice(&bytes).unwrap(); 300 assert_eq!(restored, legacy); 301 assert_eq!(serde_json::to_vec(&restored).unwrap(), bytes); 302 let scope = AuthoredDraftScope::new([7; 32]).unwrap(); 303 assert_eq!( 304 AuthoredDraftScope::try_from(<[u8; 32]>::from(scope)).unwrap(), 305 scope 306 ); 307 let scoped = legacy.with_scope(scope).unwrap(); 308 let scoped_bytes = serde_json::to_vec(&scoped).unwrap(); 309 assert_eq!( 310 serde_json::from_slice::<AuthoredDraft>(&scoped_bytes).unwrap(), 311 scoped 312 ); 313 let mut malformed = serde_json::to_value(&scoped).unwrap(); 314 malformed["scope"] = serde_json::json!([0; 32].to_vec()); 315 assert!(serde_json::from_value::<AuthoredDraft>(malformed).is_err()); 316 } 317 318 #[test] 319 fn page_payload_budget_advances_without_losing_the_next_large_record() { 320 let store = MemoryStorage::default(); 321 for id in [1, 2] { 322 block_on(store.append_authored_draft( 323 draft(id, "fixture.composer.v1", None, vec![id; 3 * 1024 * 1024]), 324 None, 325 )) 326 .unwrap(); 327 } 328 let q = query(None, 256); 329 let first = block_on(store.query_authored_drafts(q.clone())).unwrap(); 330 assert_eq!(first.records().len(), 1); 331 let next = block_on( 332 store.query_authored_drafts(q.clone().with_cursor(first.next_cursor().unwrap()).unwrap()), 333 ) 334 .unwrap(); 335 assert_eq!(next.records()[0].draft_key(), [2; 16]); 336 assert!(next.next_cursor().is_none()); 337 let mut records = first.into_records(); 338 records.extend(next.into_records()); 339 let total_payload: usize = records 340 .iter() 341 .map(|record| match record { 342 AuthoredDraftQueryRecord::Draft(draft) => draft.payload().len(), 343 AuthoredDraftQueryRecord::Corrupt { .. } => 0, 344 }) 345 .sum(); 346 assert!(total_payload > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES); 347 assert!(AuthoredDraftPage::new(&q, records, false).is_err()); 348 assert!(AuthoredDraftPage::new(&q, vec![], true).is_err()); 349 let record = AuthoredDraftQueryRecord::Draft(draft(1, "fixture.profile.v1", None, vec![1])); 350 assert!(AuthoredDraftPage::new(&q, vec![record], false).is_err()); 351 let record = AuthoredDraftQueryRecord::Draft(draft(1, "fixture.composer.v1", None, vec![1])); 352 assert!(AuthoredDraftPage::new(&q, vec![record.clone(), record.clone()], false).is_err()); 353 assert!(AuthoredDraftPage::new(&query(None, 1), vec![record.clone(), record], false).is_err()); 354 } 355 356 #[test] 357 fn bounded_reference_pages_scan_a_thousand_reversed_heads_and_preserve_newer_revisions() { 358 let store = MemoryStorage::default(); 359 for number in (1u16..=1000).rev() { 360 let mut id = [0; 16]; 361 id[14..].copy_from_slice(&number.to_be_bytes()); 362 let value = AuthoredDraft::initial( 363 radroots_storage::authored_draft::AuthoredDraftId::new(id).unwrap(), 364 [9; 32], 365 "fixture.composer.v1", 366 number.to_be_bytes().to_vec(), 367 AuthoredDraftStage::Draft, 368 None, 369 10, 370 ) 371 .unwrap(); 372 block_on(store.append_authored_draft(value, None)).unwrap(); 373 } 374 let q = query(None, 127); 375 let mut cursor = None; 376 let mut found = Vec::new(); 377 loop { 378 let current = cursor.as_ref().map_or_else( 379 || q.clone(), 380 |cursor| q.clone().with_cursor(cursor).unwrap(), 381 ); 382 let page = block_on(store.query_authored_drafts(current)).unwrap(); 383 assert!(page.records().len() <= usize::from(q.limit())); 384 for record in page.records() { 385 let key = record.draft_key(); 386 found.push(u16::from_be_bytes([key[14], key[15]])); 387 } 388 cursor = page.next_cursor().cloned(); 389 if found.len() == 127 { 390 let AuthoredDraftQueryRecord::Draft(first) = &page.records()[0] else { 391 panic!("valid fixture") 392 }; 393 let edited = first 394 .successor( 395 b"newer source".to_vec(), 396 AuthoredDraftStage::Draft, 397 None, 398 11, 399 ) 400 .unwrap(); 401 block_on(store.append_authored_draft(edited, Some(first.revision()))).unwrap(); 402 } 403 if cursor.is_none() { 404 break; 405 } 406 } 407 assert_eq!(found, (1u16..=1000).collect::<Vec<_>>()); 408 let fresh = block_on(store.query_authored_drafts(query(None, 1))).unwrap(); 409 let AuthoredDraftQueryRecord::Draft(first) = &fresh.records()[0] else { 410 panic!("valid fixture") 411 }; 412 assert_eq!(first.payload(), b"newer source"); 413 }