lib

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

sync_state.rs (11658B)


      1 #[cfg(not(feature = "std"))]
      2 use alloc::{collections::BTreeMap, string::String, string::ToString, vec::Vec};
      3 #[cfg(feature = "std")]
      4 use std::collections::BTreeMap;
      5 
      6 use radroots_replica_schema::farm::IFarmFindMany;
      7 use radroots_replica_schema::nostr_event_head::INostrEventHeadFindMany;
      8 use radroots_sql_core::SqlExecutor;
      9 
     10 use crate::error::RadrootsReplicaEventsError;
     11 use crate::event_head::{event_content_hash, event_head_key, tag_value};
     12 use crate::types::{RadrootsReplicaEventDraft, RadrootsReplicaFarmSelector};
     13 
     14 #[derive(Clone, Debug)]
     15 pub struct RadrootsReplicaSyncStatus {
     16     pub expected_count: usize,
     17     pub pending_count: usize,
     18 }
     19 
     20 #[derive(Clone, Debug)]
     21 pub struct RadrootsReplicaPendingPublishEvent {
     22     pub key: String,
     23     pub kind: u32,
     24     pub author: String,
     25     pub d_tag: String,
     26     pub content_hash: String,
     27     pub draft: RadrootsReplicaEventDraft,
     28 }
     29 
     30 #[derive(Clone, Debug)]
     31 pub struct RadrootsReplicaPendingPublishBatch {
     32     pub expected_count: usize,
     33     pub pending_count: usize,
     34     pub pending_events: Vec<RadrootsReplicaPendingPublishEvent>,
     35 }
     36 
     37 pub fn radroots_replica_sync_status<E: SqlExecutor>(
     38     exec: &E,
     39 ) -> Result<RadrootsReplicaSyncStatus, RadrootsReplicaEventsError> {
     40     let batch = radroots_replica_pending_publish_batch(exec)?;
     41     Ok(RadrootsReplicaSyncStatus {
     42         expected_count: batch.expected_count,
     43         pending_count: batch.pending_count,
     44     })
     45 }
     46 
     47 /// Computes the replica drafts that have not reached their expected heads.
     48 ///
     49 /// Lossy stored Profile projections are excluded from pending publication;
     50 /// Profile authoring must use `profile.build_authored_draft`.
     51 pub fn radroots_replica_pending_publish_batch<E: SqlExecutor>(
     52     exec: &E,
     53 ) -> Result<RadrootsReplicaPendingPublishBatch, RadrootsReplicaEventsError> {
     54     let farms =
     55         radroots_replica_store::farm::find_many(exec, &IFarmFindMany { filter: None })?.results;
     56     let mut expected: BTreeMap<String, RadrootsReplicaPendingPublishEvent> = BTreeMap::new();
     57 
     58     for farm in farms {
     59         let selector = RadrootsReplicaFarmSelector {
     60             id: Some(farm.id),
     61             d_tag: None,
     62             pubkey: None,
     63         };
     64         let bundle = crate::emit::radroots_replica_sync_all_with_options(exec, &selector, None)?;
     65         for event in bundle.events {
     66             let d_tag = tag_value(&event.tags, "d").unwrap_or("");
     67             let key = event_head_key(event.kind, &event.author, d_tag);
     68             let content_hash = draft_content_hash(&event)?;
     69             expected
     70                 .entry(key.clone())
     71                 .or_insert(RadrootsReplicaPendingPublishEvent {
     72                     key,
     73                     kind: event.kind,
     74                     author: event.author.clone(),
     75                     d_tag: d_tag.to_string(),
     76                     content_hash,
     77                     draft: event,
     78                 });
     79         }
     80     }
     81 
     82     let states_query = radroots_replica_store::nostr_event_head::find_many(
     83         exec,
     84         &INostrEventHeadFindMany { filter: None },
     85     );
     86     let states_result = states_query?;
     87     let states = states_result.results;
     88 
     89     let mut state_map: BTreeMap<String, String> = BTreeMap::new();
     90     for state in states {
     91         state_map.insert(state.key, state.content_hash);
     92     }
     93 
     94     let mut pending_events = Vec::new();
     95     for (key, event) in expected.iter() {
     96         match state_map.get(key) {
     97             Some(existing) if existing == &event.content_hash => {}
     98             _ => pending_events.push(event.clone()),
     99         }
    100     }
    101 
    102     Ok(RadrootsReplicaPendingPublishBatch {
    103         expected_count: expected.len(),
    104         pending_count: pending_events.len(),
    105         pending_events,
    106     })
    107 }
    108 
    109 fn draft_content_hash(
    110     event: &RadrootsReplicaEventDraft,
    111 ) -> Result<String, RadrootsReplicaEventsError> {
    112     #[cfg(test)]
    113     {
    114         event_content_hash(&event.content, &event.tags)
    115     }
    116     #[cfg(not(test))]
    117     {
    118         Ok(event_content_hash(&event.content, &event.tags))
    119     }
    120 }
    121 
    122 #[cfg(test)]
    123 mod tests {
    124     use super::{radroots_replica_pending_publish_batch, radroots_replica_sync_status};
    125     use crate::emit::radroots_replica_sync_all_with_options;
    126     use crate::event_head::{
    127         event_content_hash, event_content_hash_fail_next, event_head_key, tag_value,
    128     };
    129     use crate::types::RadrootsReplicaFarmSelector;
    130     use radroots_replica_schema::farm::IFarmFields;
    131     use radroots_replica_schema::nostr_event_head::INostrEventHeadFields;
    132     use radroots_replica_store::{farm, migrations, nostr_event_head};
    133     use radroots_sql_core::{SqlExecutor, SqlxSqliteExecutor};
    134 
    135     #[test]
    136     fn sync_status_empty_db_is_zero() {
    137         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    138         migrations::run_all_up(&exec).expect("migrations");
    139         let status = radroots_replica_sync_status(&exec).expect("status");
    140         assert_eq!(status.expected_count, 0);
    141         assert_eq!(status.pending_count, 0);
    142     }
    143 
    144     #[test]
    145     fn sync_status_tracks_expected_and_pending() {
    146         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    147         migrations::run_all_up(&exec).expect("migrations");
    148 
    149         let farm_row = farm::create(
    150             &exec,
    151             &IFarmFields {
    152                 d_tag: "AAAAAAAAAAAAAAAAAAAAAA".to_string(),
    153                 pubkey: "f".repeat(64),
    154                 name: "farm".to_string(),
    155                 about: None,
    156                 website: None,
    157                 picture: None,
    158                 banner: None,
    159                 location_primary: None,
    160                 location_city: None,
    161                 location_region: None,
    162                 location_country: None,
    163             },
    164         )
    165         .expect("farm")
    166         .result;
    167 
    168         let selector = RadrootsReplicaFarmSelector {
    169             id: Some(farm_row.id.clone()),
    170             d_tag: None,
    171             pubkey: None,
    172         };
    173         let bundle =
    174             radroots_replica_sync_all_with_options(&exec, &selector, None).expect("bundle");
    175         let expected_count = bundle.events.len();
    176         let first = bundle.events.first().expect("event");
    177         let d_tag = tag_value(&first.tags, "d").unwrap_or("");
    178         let key = event_head_key(first.kind, &first.author, d_tag);
    179         let content_hash = event_content_hash(&first.content, &first.tags).expect("hash");
    180         let fields = INostrEventHeadFields {
    181             key,
    182             kind: first.kind,
    183             pubkey: first.author.clone(),
    184             d_tag: d_tag.to_string(),
    185             last_event_id: format!("{:064x}", 1u64),
    186             last_created_at: 1,
    187             content_hash,
    188         };
    189         let _ = nostr_event_head::create(&exec, &fields).expect("state");
    190 
    191         let status = radroots_replica_sync_status(&exec).expect("status");
    192         assert_eq!(status.expected_count, expected_count);
    193         assert_eq!(status.pending_count, expected_count.saturating_sub(1));
    194     }
    195 
    196     #[test]
    197     fn pending_publish_batch_lists_only_missing_or_changed_expected_events() {
    198         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    199         migrations::run_all_up(&exec).expect("migrations");
    200 
    201         let farm_row = farm::create(
    202             &exec,
    203             &IFarmFields {
    204                 d_tag: "AAAAAAAAAAAAAAAAAAAAAA".to_string(),
    205                 pubkey: "a".repeat(64),
    206                 name: "farm".to_string(),
    207                 about: None,
    208                 website: None,
    209                 picture: None,
    210                 banner: None,
    211                 location_primary: None,
    212                 location_city: None,
    213                 location_region: None,
    214                 location_country: None,
    215             },
    216         )
    217         .expect("farm")
    218         .result;
    219 
    220         let selector = RadrootsReplicaFarmSelector {
    221             id: Some(farm_row.id.clone()),
    222             d_tag: None,
    223             pubkey: None,
    224         };
    225         let bundle =
    226             radroots_replica_sync_all_with_options(&exec, &selector, None).expect("bundle");
    227         let first = bundle.events.first().expect("event");
    228         let d_tag = tag_value(&first.tags, "d").unwrap_or("");
    229         let key = event_head_key(first.kind, &first.author, d_tag);
    230         let content_hash = event_content_hash(&first.content, &first.tags).expect("hash");
    231         let fields = INostrEventHeadFields {
    232             key: key.clone(),
    233             kind: first.kind,
    234             pubkey: first.author.clone(),
    235             d_tag: d_tag.to_string(),
    236             last_event_id: format!("{:064x}", 1u64),
    237             last_created_at: 1,
    238             content_hash,
    239         };
    240         let _ = nostr_event_head::create(&exec, &fields).expect("state");
    241 
    242         let batch = radroots_replica_pending_publish_batch(&exec).expect("batch");
    243 
    244         assert_eq!(batch.expected_count, bundle.events.len());
    245         assert_eq!(batch.pending_count, bundle.events.len().saturating_sub(1));
    246         for event in &batch.pending_events {
    247             assert_ne!(event.key, key);
    248             assert_eq!(event.content_hash.len(), 64);
    249         }
    250     }
    251 
    252     #[test]
    253     fn sync_status_reports_farm_query_errors() {
    254         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    255         let err = radroots_replica_sync_status(&exec).expect_err("farm query error");
    256         assert!(err.to_string().contains("invalid query"));
    257     }
    258 
    259     #[test]
    260     fn sync_status_reports_emit_errors() {
    261         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    262         migrations::run_all_up(&exec).expect("migrations");
    263         let _ = farm::create(
    264             &exec,
    265             &IFarmFields {
    266                 d_tag: "AAAAAAAAAAAAAAAAAAAAAA".to_string(),
    267                 pubkey: "b".repeat(64),
    268                 name: "farm".to_string(),
    269                 about: None,
    270                 website: None,
    271                 picture: None,
    272                 banner: None,
    273                 location_primary: None,
    274                 location_city: None,
    275                 location_region: None,
    276                 location_country: None,
    277             },
    278         )
    279         .expect("farm");
    280         let _ = exec
    281             .exec("DROP TABLE farm_tag;", "[]")
    282             .expect("drop farm_tag");
    283         let err = radroots_replica_sync_status(&exec).expect_err("emit error");
    284         assert!(err.to_string().contains("invalid query"));
    285     }
    286 
    287     #[test]
    288     fn sync_status_reports_content_hash_errors() {
    289         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    290         migrations::run_all_up(&exec).expect("migrations");
    291         let _ = farm::create(
    292             &exec,
    293             &IFarmFields {
    294                 d_tag: "AAAAAAAAAAAAAAAAAAAAAA".to_string(),
    295                 pubkey: "c".repeat(64),
    296                 name: "farm".to_string(),
    297                 about: None,
    298                 website: None,
    299                 picture: None,
    300                 banner: None,
    301                 location_primary: None,
    302                 location_city: None,
    303                 location_region: None,
    304                 location_country: None,
    305             },
    306         )
    307         .expect("farm");
    308         event_content_hash_fail_next();
    309         let err = radroots_replica_sync_status(&exec).expect_err("content hash error");
    310         assert!(err.to_string().contains("content_hash"));
    311     }
    312 
    313     #[test]
    314     fn sync_status_reports_state_query_errors() {
    315         let exec = SqlxSqliteExecutor::open_memory().expect("db");
    316         migrations::run_all_up(&exec).expect("migrations");
    317         let _ = exec
    318             .exec("DROP TABLE nostr_event_head;", "[]")
    319             .expect("drop nostr_event_head");
    320         let err = radroots_replica_sync_status(&exec).expect_err("state query error");
    321         assert!(err.to_string().contains("invalid query"));
    322     }
    323 }