lib

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

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 }