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 }