nostrdb.rs (43726B)
1 use crate::config::RadrootsNostrdbConfig; 2 use crate::error::RadrootsNostrdbError; 3 use crate::filter::parse_hex_32; 4 use crate::ingest::RadrootsNostrdbIngestSource; 5 use crate::query::{RadrootsNostrdbNote, RadrootsNostrdbProfile, RadrootsNostrdbQuerySpec}; 6 use crate::subscription::{ 7 RadrootsNostrdbNoteKey, RadrootsNostrdbSubscriptionHandle, RadrootsNostrdbSubscriptionSpec, 8 RadrootsNostrdbSubscriptionStream, 9 }; 10 use radroots_nostr::event::Event as RadrootsNostrEvent; 11 use std::path::Path; 12 13 #[derive(Debug, Clone)] 14 pub struct RadrootsNostrdb { 15 db_dir: std::path::PathBuf, 16 pub(crate) inner: nostrdb::Ndb, 17 } 18 19 #[cfg(test)] 20 mod test_hooks { 21 use std::sync::atomic::{AtomicBool, Ordering}; 22 23 pub static FORCE_EVENT_JSON_ERROR: AtomicBool = AtomicBool::new(false); 24 pub static FORCE_PROCESS_EVENT_ERROR: AtomicBool = AtomicBool::new(false); 25 pub static FORCE_SUBSCRIBE_ERROR: AtomicBool = AtomicBool::new(false); 26 pub static FORCE_UNSUBSCRIBE_ERROR: AtomicBool = AtomicBool::new(false); 27 pub static FORCE_WAIT_ERROR: AtomicBool = AtomicBool::new(false); 28 pub static FORCE_TRANSACTION_ERROR: AtomicBool = AtomicBool::new(false); 29 pub static FORCE_QUERY_ERROR: AtomicBool = AtomicBool::new(false); 30 pub static FORCE_NOTE_JSON_ERROR: AtomicBool = AtomicBool::new(false); 31 pub static FORCE_PROFILE_QUERY_ERROR: AtomicBool = AtomicBool::new(false); 32 33 pub fn take(flag: &AtomicBool) -> bool { 34 flag.swap(false, Ordering::SeqCst) 35 } 36 } 37 38 fn map_profile_lookup_result<T>( 39 result: Result<T, nostrdb::Error>, 40 ) -> Result<Option<T>, RadrootsNostrdbError> { 41 match result { 42 Ok(value) => Ok(Some(value)), 43 Err(nostrdb::Error::NotFound) => Ok(None), 44 Err(source) => Err(source.into()), 45 } 46 } 47 48 impl RadrootsNostrdb { 49 fn serialize_event(event: &RadrootsNostrEvent) -> Result<String, RadrootsNostrdbError> { 50 #[cfg(test)] 51 if test_hooks::take(&test_hooks::FORCE_EVENT_JSON_ERROR) { 52 return Err(RadrootsNostrdbError::EventJsonEncode( 53 "forced event json error".into(), 54 )); 55 } 56 serde_json::to_string(event).map_err(Into::into) 57 } 58 59 fn process_event_with_inner( 60 &self, 61 json: &str, 62 metadata: nostrdb::IngestMetadata, 63 ) -> Result<(), RadrootsNostrdbError> { 64 #[cfg(test)] 65 if test_hooks::take(&test_hooks::FORCE_PROCESS_EVENT_ERROR) { 66 return Err(RadrootsNostrdbError::Nostrdb( 67 "forced process event error".into(), 68 )); 69 } 70 self.inner 71 .process_event_with(json, metadata) 72 .map_err(Into::into) 73 } 74 75 fn subscribe_inner( 76 &self, 77 filters: &[nostrdb::Filter], 78 ) -> Result<nostrdb::Subscription, RadrootsNostrdbError> { 79 #[cfg(test)] 80 if test_hooks::take(&test_hooks::FORCE_SUBSCRIBE_ERROR) { 81 return Err(RadrootsNostrdbError::Nostrdb( 82 "forced subscribe error".into(), 83 )); 84 } 85 self.inner.subscribe(filters).map_err(Into::into) 86 } 87 88 fn unsubscribe_inner( 89 &self, 90 subscription: nostrdb::Subscription, 91 ) -> Result<(), RadrootsNostrdbError> { 92 #[cfg(test)] 93 if test_hooks::take(&test_hooks::FORCE_UNSUBSCRIBE_ERROR) { 94 return Err(RadrootsNostrdbError::Nostrdb( 95 "forced unsubscribe error".into(), 96 )); 97 } 98 let mut inner = self.inner.clone(); 99 inner.unsubscribe(subscription).map_err(Into::into) 100 } 101 102 #[cfg(feature = "rt")] 103 async fn wait_for_notes_inner( 104 &self, 105 subscription: nostrdb::Subscription, 106 max_notes: u32, 107 ) -> Result<Vec<nostrdb::NoteKey>, RadrootsNostrdbError> { 108 #[cfg(test)] 109 if test_hooks::take(&test_hooks::FORCE_WAIT_ERROR) { 110 return Err(RadrootsNostrdbError::Nostrdb("forced wait error".into())); 111 } 112 self.inner 113 .wait_for_notes(subscription, max_notes) 114 .await 115 .map_err(Into::into) 116 } 117 118 fn open_txn(&self) -> Result<nostrdb::Transaction, RadrootsNostrdbError> { 119 #[cfg(test)] 120 if test_hooks::take(&test_hooks::FORCE_TRANSACTION_ERROR) { 121 return Err(RadrootsNostrdbError::Nostrdb( 122 "forced transaction error".into(), 123 )); 124 } 125 nostrdb::Transaction::new(&self.inner).map_err(Into::into) 126 } 127 128 fn query_inner<'a>( 129 &self, 130 txn: &'a nostrdb::Transaction, 131 filters: &[nostrdb::Filter], 132 max_results: i32, 133 ) -> Result<Vec<nostrdb::QueryResult<'a>>, RadrootsNostrdbError> { 134 #[cfg(test)] 135 if test_hooks::take(&test_hooks::FORCE_QUERY_ERROR) { 136 return Err(RadrootsNostrdbError::Nostrdb("forced query error".into())); 137 } 138 self.inner 139 .query(txn, filters, max_results) 140 .map_err(Into::into) 141 } 142 143 fn note_json_value(note: &nostrdb::Note) -> Result<String, RadrootsNostrdbError> { 144 #[cfg(test)] 145 if test_hooks::take(&test_hooks::FORCE_NOTE_JSON_ERROR) { 146 return Err(RadrootsNostrdbError::Nostrdb( 147 "forced note json error".into(), 148 )); 149 } 150 note.json().map_err(Into::into) 151 } 152 153 fn get_profile_record<'a>( 154 &self, 155 txn: &'a nostrdb::Transaction, 156 pubkey: &[u8; 32], 157 ) -> Result<Option<nostrdb::ProfileRecord<'a>>, RadrootsNostrdbError> { 158 #[cfg(test)] 159 if test_hooks::take(&test_hooks::FORCE_PROFILE_QUERY_ERROR) { 160 return map_profile_lookup_result(Err(nostrdb::Error::QueryError)); 161 } 162 map_profile_lookup_result(self.inner.get_profile_by_pubkey(txn, pubkey)) 163 } 164 165 pub fn open(config: RadrootsNostrdbConfig) -> Result<Self, RadrootsNostrdbError> { 166 let mut inner_config = nostrdb::Config::new().skip_validation(config.skip_validation()); 167 if let Some(mapsize_bytes) = config.mapsize_bytes() { 168 inner_config = inner_config.set_mapsize(mapsize_bytes); 169 } 170 if let Some(ingester_threads) = config.ingester_threads() { 171 inner_config = inner_config.set_ingester_threads(ingester_threads); 172 } 173 174 let db_dir = config.db_dir().to_path_buf(); 175 let db_dir_str = db_dir.to_str().ok_or(RadrootsNostrdbError::NonUtf8Path)?; 176 let inner = nostrdb::Ndb::new(db_dir_str, &inner_config)?; 177 178 Ok(Self { db_dir, inner }) 179 } 180 181 pub fn db_dir(&self) -> &Path { 182 &self.db_dir 183 } 184 185 pub fn ingest_event_json_with_source( 186 &self, 187 json: &str, 188 source: RadrootsNostrdbIngestSource, 189 ) -> Result<(), RadrootsNostrdbError> { 190 let metadata = source.to_nostrdb_metadata(); 191 self.process_event_with_inner(json, metadata)?; 192 Ok(()) 193 } 194 195 pub fn ingest_event_json(&self, json: &str) -> Result<(), RadrootsNostrdbError> { 196 self.ingest_event_json_with_source(json, RadrootsNostrdbIngestSource::default()) 197 } 198 199 pub fn ingest_event( 200 &self, 201 event: &RadrootsNostrEvent, 202 source: RadrootsNostrdbIngestSource, 203 ) -> Result<(), RadrootsNostrdbError> { 204 let json = Self::serialize_event(event)?; 205 self.ingest_event_json_with_source(json.as_str(), source) 206 } 207 208 #[cfg(feature = "giftwrap")] 209 pub fn add_giftwrap_secret_key(&self, secret_key: [u8; 32]) -> bool { 210 self.inner.add_key(&secret_key) 211 } 212 213 #[cfg(feature = "giftwrap")] 214 pub fn add_giftwrap_secret_key_hex( 215 &self, 216 secret_key_hex: &str, 217 ) -> Result<bool, RadrootsNostrdbError> { 218 let secret_key = parse_hex_32(secret_key_hex, "secret_key")?; 219 Ok(self.add_giftwrap_secret_key(secret_key)) 220 } 221 222 #[cfg(feature = "giftwrap")] 223 pub fn process_giftwraps(&self) -> Result<(), RadrootsNostrdbError> { 224 let txn = nostrdb::Transaction::new(&self.inner)?; 225 self.inner.process_giftwraps(&txn); 226 Ok(()) 227 } 228 229 pub fn subscribe( 230 &self, 231 spec: &RadrootsNostrdbSubscriptionSpec, 232 ) -> Result<RadrootsNostrdbSubscriptionHandle, RadrootsNostrdbError> { 233 let filters = spec 234 .filters() 235 .iter() 236 .map(|filter_spec| filter_spec.to_nostrdb_filter()) 237 .collect::<Result<Vec<_>, _>>()?; 238 let subscription = self.subscribe_inner(filters.as_slice())?; 239 Ok(RadrootsNostrdbSubscriptionHandle::new(subscription.id())) 240 } 241 242 pub fn unsubscribe( 243 &self, 244 handle: RadrootsNostrdbSubscriptionHandle, 245 ) -> Result<(), RadrootsNostrdbError> { 246 let subscription = nostrdb::Subscription::new(handle.id()); 247 self.unsubscribe_inner(subscription)?; 248 Ok(()) 249 } 250 251 pub fn poll_for_note_keys( 252 &self, 253 handle: RadrootsNostrdbSubscriptionHandle, 254 max_notes: u32, 255 ) -> Vec<RadrootsNostrdbNoteKey> { 256 self.inner 257 .poll_for_notes(nostrdb::Subscription::new(handle.id()), max_notes) 258 .into_iter() 259 .map(|note_key| RadrootsNostrdbNoteKey::new(note_key.as_u64())) 260 .collect() 261 } 262 263 #[cfg(feature = "rt")] 264 pub async fn wait_for_note_keys( 265 &self, 266 handle: RadrootsNostrdbSubscriptionHandle, 267 max_notes: u32, 268 ) -> Result<Vec<RadrootsNostrdbNoteKey>, RadrootsNostrdbError> { 269 let note_keys = self 270 .wait_for_notes_inner(nostrdb::Subscription::new(handle.id()), max_notes) 271 .await?; 272 Ok(note_keys 273 .into_iter() 274 .map(|note_key| RadrootsNostrdbNoteKey::new(note_key.as_u64())) 275 .collect()) 276 } 277 278 #[cfg(feature = "rt")] 279 pub fn subscription_stream( 280 &self, 281 handle: RadrootsNostrdbSubscriptionHandle, 282 notes_per_await: u32, 283 ) -> RadrootsNostrdbSubscriptionStream { 284 let stream = nostrdb::Subscription::new(handle.id()) 285 .stream(&self.inner) 286 .notes_per_await(notes_per_await.max(1)); 287 RadrootsNostrdbSubscriptionStream { inner: stream } 288 } 289 290 pub fn query_notes( 291 &self, 292 spec: &RadrootsNostrdbQuerySpec, 293 ) -> Result<Vec<RadrootsNostrdbNote>, RadrootsNostrdbError> { 294 if spec.filters().is_empty() { 295 return Ok(Vec::new()); 296 } 297 298 let filters = spec 299 .filters() 300 .iter() 301 .map(|filter_spec| filter_spec.to_nostrdb_filter()) 302 .collect::<Result<Vec<_>, _>>()?; 303 let txn = self.open_txn()?; 304 let query_results = 305 self.query_inner(&txn, filters.as_slice(), spec.max_results() as i32)?; 306 307 query_results 308 .into_iter() 309 .map(|query_result| { 310 let note = query_result.note; 311 let json = Self::note_json_value(¬e)?; 312 Ok(RadrootsNostrdbNote { 313 note_key: query_result.note_key.as_u64(), 314 id_hex: hex::encode(note.id()), 315 author_hex: hex::encode(note.pubkey()), 316 kind: note.kind(), 317 created_at_unix: note.created_at(), 318 content: note.content().to_owned(), 319 json, 320 }) 321 }) 322 .collect::<Result<Vec<_>, RadrootsNostrdbError>>() 323 } 324 325 pub fn get_profile_by_pubkey_hex( 326 &self, 327 pubkey_hex: &str, 328 ) -> Result<Option<RadrootsNostrdbProfile>, RadrootsNostrdbError> { 329 let pubkey = parse_hex_32(pubkey_hex, "pubkey")?; 330 let txn = self.open_txn()?; 331 let Some(profile_record) = self.get_profile_record(&txn, &pubkey)? else { 332 return Ok(None); 333 }; 334 335 let profile = profile_record.record().profile(); 336 let profile_key = profile_record.key().map(|key| key.as_u64()); 337 Ok(profile.map(|profile| RadrootsNostrdbProfile { 338 profile_key, 339 pubkey_hex: pubkey_hex.to_owned(), 340 name: profile.name().map(ToOwned::to_owned), 341 display_name: profile.display_name().map(ToOwned::to_owned), 342 about: profile.about().map(ToOwned::to_owned), 343 picture: profile.picture().map(ToOwned::to_owned), 344 banner: profile.banner().map(ToOwned::to_owned), 345 website: profile.website().map(ToOwned::to_owned), 346 nip05: profile.nip05().map(ToOwned::to_owned), 347 lud16: profile.lud16().map(ToOwned::to_owned), 348 })) 349 } 350 } 351 352 #[cfg(test)] 353 mod tests { 354 use super::*; 355 use crate::filter::RadrootsNostrdbFilterSpec; 356 use crate::ingest::RadrootsNostrdbIngestSource; 357 use crate::query::RadrootsNostrdbQuerySpec; 358 use crate::test_fixtures::{FIXTURE_ALICE_EMAIL, FIXTURE_ALICE_USERNAME}; 359 use futures::StreamExt; 360 use nostr::EventBuilder; 361 use nostr::Keys as RadrootsNostrKeys; 362 use radroots_nostr::event::Metadata as RadrootsNostrMetadata; 363 use std::sync::atomic::Ordering; 364 use std::sync::{Mutex, OnceLock}; 365 use std::time::Duration; 366 use tempfile::TempDir; 367 368 fn test_hooks_lock() -> &'static Mutex<()> { 369 static TEST_HOOKS_LOCK: OnceLock<Mutex<()>> = OnceLock::new(); 370 TEST_HOOKS_LOCK.get_or_init(|| Mutex::new(())) 371 } 372 373 fn test_hooks_guard() -> std::sync::MutexGuard<'static, ()> { 374 test_hooks_lock().lock().expect("test hooks lock") 375 } 376 377 fn reset_test_flags() { 378 test_hooks::FORCE_EVENT_JSON_ERROR.store(false, Ordering::SeqCst); 379 test_hooks::FORCE_PROCESS_EVENT_ERROR.store(false, Ordering::SeqCst); 380 test_hooks::FORCE_SUBSCRIBE_ERROR.store(false, Ordering::SeqCst); 381 test_hooks::FORCE_UNSUBSCRIBE_ERROR.store(false, Ordering::SeqCst); 382 test_hooks::FORCE_WAIT_ERROR.store(false, Ordering::SeqCst); 383 test_hooks::FORCE_TRANSACTION_ERROR.store(false, Ordering::SeqCst); 384 test_hooks::FORCE_QUERY_ERROR.store(false, Ordering::SeqCst); 385 test_hooks::FORCE_NOTE_JSON_ERROR.store(false, Ordering::SeqCst); 386 test_hooks::FORCE_PROFILE_QUERY_ERROR.store(false, Ordering::SeqCst); 387 } 388 389 #[test] 390 fn config_builder_tracks_values() { 391 let config = RadrootsNostrdbConfig::new("target/testdbs/nostrdb_config") 392 .with_mapsize_bytes(1024 * 1024) 393 .with_ingester_threads(2) 394 .with_skip_validation(true); 395 396 assert_eq!(config.mapsize_bytes(), Some(1024 * 1024)); 397 assert_eq!(config.ingester_threads(), Some(2)); 398 assert!(config.skip_validation()); 399 } 400 401 #[test] 402 fn map_profile_lookup_result_handles_all_error_kinds() { 403 let success = map_profile_lookup_result::<u64>(Ok(7)).expect("ok"); 404 assert_eq!(success, Some(7)); 405 406 let not_found = 407 map_profile_lookup_result::<u64>(Err(nostrdb::Error::NotFound)).expect("none"); 408 assert!(not_found.is_none()); 409 410 let query_error = map_profile_lookup_result::<u64>(Err(nostrdb::Error::QueryError)) 411 .expect_err("query error"); 412 assert!(query_error.to_string().starts_with("nostrdb error:")); 413 } 414 415 #[test] 416 fn open_creates_database() { 417 let tmp_dir = TempDir::new().expect("tempdir should open"); 418 let db_dir = tmp_dir.path().join("nostrdb"); 419 let config = RadrootsNostrdbConfig::new(&db_dir) 420 .with_mapsize_bytes(64 * 1024 * 1024) 421 .with_ingester_threads(1); 422 423 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 424 assert_eq!(nostrdb.db_dir(), db_dir.as_path()); 425 assert!(db_dir.exists()); 426 } 427 428 #[cfg(unix)] 429 #[test] 430 fn open_rejects_non_utf8_path() { 431 use std::os::unix::ffi::OsStrExt; 432 433 let path = std::path::PathBuf::from(std::ffi::OsStr::from_bytes(b"nostrdb-\xFF")); 434 let config = RadrootsNostrdbConfig::new(&path); 435 let err = RadrootsNostrdb::open(config).expect_err("non utf8 path"); 436 assert!(err.to_string().contains("utf-8")); 437 } 438 439 #[test] 440 fn open_reports_nostrdb_error_for_file_path() { 441 let tmp_dir = TempDir::new().expect("tempdir should open"); 442 let db_dir = tmp_dir.path().join("nostrdb"); 443 std::fs::write(&db_dir, "not a directory").expect("write db file"); 444 let config = RadrootsNostrdbConfig::new(&db_dir); 445 let err = RadrootsNostrdb::open(config).expect_err("file path should fail"); 446 assert!(err.to_string().starts_with("nostrdb error:")); 447 } 448 449 #[test] 450 fn ingest_source_builders_track_origin() { 451 assert_eq!( 452 RadrootsNostrdbIngestSource::default(), 453 RadrootsNostrdbIngestSource::client() 454 ); 455 assert_eq!( 456 RadrootsNostrdbIngestSource::relay("wss://radroots.org"), 457 RadrootsNostrdbIngestSource::Relay { 458 relay_url: Some("wss://radroots.org".into()) 459 } 460 ); 461 assert_eq!( 462 RadrootsNostrdbIngestSource::relay_unknown(), 463 RadrootsNostrdbIngestSource::Relay { relay_url: None } 464 ); 465 } 466 467 #[test] 468 fn ingest_event_accepts_signed_note() { 469 let tmp_dir = TempDir::new().expect("tempdir should open"); 470 let db_dir = tmp_dir.path().join("nostrdb"); 471 let config = RadrootsNostrdbConfig::new(&db_dir); 472 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 473 474 let keys = RadrootsNostrKeys::generate(); 475 let event = EventBuilder::text_note("hello from nostrdb") 476 .sign_with_keys(&keys) 477 .expect("event should sign"); 478 479 nostrdb 480 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 481 .expect("ingest should succeed"); 482 } 483 484 #[test] 485 fn ingest_event_json_accepts_signed_note() { 486 let tmp_dir = TempDir::new().expect("tempdir should open"); 487 let db_dir = tmp_dir.path().join("nostrdb"); 488 let config = RadrootsNostrdbConfig::new(&db_dir); 489 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 490 491 let keys = RadrootsNostrKeys::generate(); 492 let event = EventBuilder::text_note("hello from nostrdb json") 493 .sign_with_keys(&keys) 494 .expect("event should sign"); 495 let json = serde_json::to_string(&event).expect("event json"); 496 497 nostrdb 498 .ingest_event_json(&json) 499 .expect("json ingest should succeed"); 500 } 501 502 #[test] 503 fn ingest_event_json_rejects_invalid_json() { 504 let _guard = test_hooks_guard(); 505 reset_test_flags(); 506 let tmp_dir = TempDir::new().expect("tempdir should open"); 507 let db_dir = tmp_dir.path().join("nostrdb"); 508 let config = RadrootsNostrdbConfig::new(&db_dir); 509 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 510 511 test_hooks::FORCE_PROCESS_EVENT_ERROR.store(true, Ordering::SeqCst); 512 let err = nostrdb 513 .ingest_event_json_with_source("not json", RadrootsNostrdbIngestSource::client()) 514 .expect_err("process event error"); 515 assert!(err.to_string().starts_with("nostrdb error:")); 516 } 517 518 #[test] 519 fn ingest_event_reports_event_json_error() { 520 let _guard = test_hooks_guard(); 521 reset_test_flags(); 522 let tmp_dir = TempDir::new().expect("tempdir should open"); 523 let db_dir = tmp_dir.path().join("nostrdb"); 524 let config = RadrootsNostrdbConfig::new(&db_dir); 525 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 526 527 let keys = RadrootsNostrKeys::generate(); 528 let event = EventBuilder::text_note("forced json error") 529 .sign_with_keys(&keys) 530 .expect("event should sign"); 531 test_hooks::FORCE_EVENT_JSON_ERROR.store(true, Ordering::SeqCst); 532 533 let err = nostrdb 534 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 535 .expect_err("forced json error"); 536 assert!(err.to_string().starts_with("event json encode failed:")); 537 } 538 539 #[test] 540 fn subscribe_poll_and_unsubscribe_round_trip() { 541 let _guard = test_hooks_guard(); 542 reset_test_flags(); 543 let tmp_dir = TempDir::new().expect("tempdir should open"); 544 let db_dir = tmp_dir.path().join("nostrdb"); 545 let config = RadrootsNostrdbConfig::new(&db_dir); 546 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 547 let spec = RadrootsNostrdbSubscriptionSpec::single( 548 RadrootsNostrdbFilterSpec::new().with_kind(1).with_limit(10), 549 ); 550 let handle = nostrdb.subscribe(&spec).expect("subscribe should succeed"); 551 552 let keys = RadrootsNostrKeys::generate(); 553 let event = EventBuilder::text_note("subscription test") 554 .sign_with_keys(&keys) 555 .expect("event should sign"); 556 nostrdb 557 .ingest_event(&event, RadrootsNostrdbIngestSource::relay_unknown()) 558 .expect("ingest should succeed"); 559 560 let mut notes = Vec::new(); 561 for _ in 0..40 { 562 notes = nostrdb.poll_for_note_keys(handle, 32); 563 if !notes.is_empty() { 564 break; 565 } 566 std::thread::sleep(Duration::from_millis(25)); 567 } 568 569 assert!(!notes.is_empty()); 570 nostrdb 571 .unsubscribe(handle) 572 .expect("unsubscribe should succeed"); 573 } 574 575 #[test] 576 fn subscribe_reports_nostrdb_error() { 577 let _guard = test_hooks_guard(); 578 reset_test_flags(); 579 let tmp_dir = TempDir::new().expect("tempdir should open"); 580 let db_dir = tmp_dir.path().join("nostrdb"); 581 let config = RadrootsNostrdbConfig::new(&db_dir); 582 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 583 let spec = RadrootsNostrdbSubscriptionSpec::text_notes(Some(10), None); 584 test_hooks::FORCE_SUBSCRIBE_ERROR.store(true, Ordering::SeqCst); 585 586 let err = nostrdb 587 .subscribe(&spec) 588 .expect_err("forced subscribe error"); 589 assert!(err.to_string().starts_with("nostrdb error:")); 590 } 591 592 #[test] 593 fn unsubscribe_reports_nostrdb_error() { 594 let _guard = test_hooks_guard(); 595 reset_test_flags(); 596 let tmp_dir = TempDir::new().expect("tempdir should open"); 597 let db_dir = tmp_dir.path().join("nostrdb"); 598 let config = RadrootsNostrdbConfig::new(&db_dir); 599 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 600 let spec = RadrootsNostrdbSubscriptionSpec::text_notes(Some(10), None); 601 let handle = nostrdb.subscribe(&spec).expect("subscribe should succeed"); 602 test_hooks::FORCE_UNSUBSCRIBE_ERROR.store(true, Ordering::SeqCst); 603 604 let err = nostrdb 605 .unsubscribe(handle) 606 .expect_err("forced unsubscribe error"); 607 assert!(err.to_string().starts_with("nostrdb error:")); 608 } 609 610 #[tokio::test] 611 #[allow(clippy::await_holding_lock)] 612 async fn wait_for_note_keys_yields_results() { 613 let _guard = test_hooks_guard(); 614 reset_test_flags(); 615 let tmp_dir = TempDir::new().expect("tempdir should open"); 616 let db_dir = tmp_dir.path().join("nostrdb"); 617 let config = RadrootsNostrdbConfig::new(&db_dir); 618 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 619 let spec = RadrootsNostrdbSubscriptionSpec::text_notes(Some(10), None); 620 let handle = nostrdb.subscribe(&spec).expect("subscribe should succeed"); 621 622 let keys = RadrootsNostrKeys::generate(); 623 let event = EventBuilder::text_note("wait test") 624 .sign_with_keys(&keys) 625 .expect("event should sign"); 626 nostrdb 627 .ingest_event(&event, RadrootsNostrdbIngestSource::relay_unknown()) 628 .expect("ingest should succeed"); 629 630 let notes = nostrdb 631 .wait_for_note_keys(handle, 32) 632 .await 633 .expect("wait should succeed"); 634 assert!(!notes.is_empty()); 635 } 636 637 #[tokio::test] 638 #[allow(clippy::await_holding_lock)] 639 async fn wait_for_note_keys_reports_nostrdb_error() { 640 let _guard = test_hooks_guard(); 641 reset_test_flags(); 642 let tmp_dir = TempDir::new().expect("tempdir should open"); 643 let db_dir = tmp_dir.path().join("nostrdb"); 644 let config = RadrootsNostrdbConfig::new(&db_dir); 645 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 646 let spec = RadrootsNostrdbSubscriptionSpec::text_notes(Some(10), None); 647 let handle = nostrdb.subscribe(&spec).expect("subscribe should succeed"); 648 test_hooks::FORCE_WAIT_ERROR.store(true, Ordering::SeqCst); 649 650 let err = nostrdb 651 .wait_for_note_keys(handle, 1) 652 .await 653 .expect_err("forced wait error"); 654 assert!(err.to_string().starts_with("nostrdb error:")); 655 } 656 657 #[test] 658 fn query_notes_returns_ingested_results() { 659 let _guard = test_hooks_guard(); 660 reset_test_flags(); 661 let tmp_dir = TempDir::new().expect("tempdir should open"); 662 let db_dir = tmp_dir.path().join("nostrdb"); 663 let config = RadrootsNostrdbConfig::new(&db_dir); 664 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 665 666 let keys = RadrootsNostrKeys::generate(); 667 let event = EventBuilder::text_note("query note") 668 .sign_with_keys(&keys) 669 .expect("event should sign"); 670 nostrdb 671 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 672 .expect("ingest should succeed"); 673 674 let query_spec = RadrootsNostrdbQuerySpec::text_notes(Some(50), None, 50); 675 let mut notes = Vec::new(); 676 for _ in 0..40 { 677 notes = nostrdb 678 .query_notes(&query_spec) 679 .expect("query should succeed"); 680 if !notes.is_empty() { 681 break; 682 } 683 std::thread::sleep(Duration::from_millis(25)); 684 } 685 assert!(!notes.is_empty()); 686 let note_pairs = notes 687 .iter() 688 .map(|note| (note.id_hex.clone(), note.content.clone())) 689 .collect::<Vec<_>>(); 690 assert!(note_pairs.contains(&(event.id.to_hex(), "query note".to_string()))); 691 } 692 693 #[test] 694 fn empty_filter_preserves_match_all_query_and_subscription_semantics() { 695 let _guard = test_hooks_guard(); 696 reset_test_flags(); 697 let tmp_dir = TempDir::new().expect("tempdir should open"); 698 let db_dir = tmp_dir.path().join("nostrdb"); 699 let config = RadrootsNostrdbConfig::new(&db_dir); 700 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 701 702 let subscription = 703 RadrootsNostrdbSubscriptionSpec::single(RadrootsNostrdbFilterSpec::new()); 704 let handle = nostrdb 705 .subscribe(&subscription) 706 .expect("empty-filter subscription should succeed"); 707 708 let keys = RadrootsNostrKeys::generate(); 709 let event = EventBuilder::text_note("empty filter match-all") 710 .sign_with_keys(&keys) 711 .expect("event should sign"); 712 nostrdb 713 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 714 .expect("ingest should succeed"); 715 716 let query = RadrootsNostrdbQuerySpec::single(RadrootsNostrdbFilterSpec::new(), 50); 717 let mut queried = Vec::new(); 718 let mut notified = Vec::new(); 719 for _ in 0..40 { 720 queried = nostrdb.query_notes(&query).expect("query should succeed"); 721 notified = nostrdb.poll_for_note_keys(handle, 32); 722 if !queried.is_empty() && !notified.is_empty() { 723 break; 724 } 725 std::thread::sleep(Duration::from_millis(25)); 726 } 727 728 assert!(queried.iter().any(|note| note.id_hex == event.id.to_hex())); 729 assert!(!notified.is_empty()); 730 nostrdb 731 .unsubscribe(handle) 732 .expect("unsubscribe should succeed"); 733 } 734 735 #[test] 736 fn query_notes_empty_filters_returns_empty() { 737 let _guard = test_hooks_guard(); 738 reset_test_flags(); 739 let tmp_dir = TempDir::new().expect("tempdir should open"); 740 let db_dir = tmp_dir.path().join("nostrdb"); 741 let config = RadrootsNostrdbConfig::new(&db_dir); 742 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 743 744 let query_spec = RadrootsNostrdbQuerySpec::new(Vec::new(), 10); 745 let notes = nostrdb 746 .query_notes(&query_spec) 747 .expect("query should succeed"); 748 assert!(notes.is_empty()); 749 } 750 751 #[test] 752 fn query_notes_rejects_invalid_filters() { 753 let _guard = test_hooks_guard(); 754 reset_test_flags(); 755 let tmp_dir = TempDir::new().expect("tempdir should open"); 756 let db_dir = tmp_dir.path().join("nostrdb"); 757 let config = RadrootsNostrdbConfig::new(&db_dir); 758 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 759 760 let spec = RadrootsNostrdbQuerySpec::single( 761 RadrootsNostrdbFilterSpec::new().with_author_hex("not-hex"), 762 10, 763 ); 764 let err = nostrdb.query_notes(&spec).expect_err("invalid filter"); 765 assert!(err.to_string().contains("invalid hex")); 766 } 767 768 #[test] 769 fn query_notes_reports_transaction_error() { 770 let _guard = test_hooks_guard(); 771 reset_test_flags(); 772 let tmp_dir = TempDir::new().expect("tempdir should open"); 773 let db_dir = tmp_dir.path().join("nostrdb"); 774 let config = RadrootsNostrdbConfig::new(&db_dir); 775 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 776 let spec = RadrootsNostrdbQuerySpec::text_notes(Some(10), None, 10); 777 test_hooks::FORCE_TRANSACTION_ERROR.store(true, Ordering::SeqCst); 778 779 let err = nostrdb 780 .query_notes(&spec) 781 .expect_err("forced transaction error"); 782 assert!(err.to_string().starts_with("nostrdb error:")); 783 } 784 785 #[test] 786 fn query_notes_reports_query_error() { 787 let _guard = test_hooks_guard(); 788 reset_test_flags(); 789 let tmp_dir = TempDir::new().expect("tempdir should open"); 790 let db_dir = tmp_dir.path().join("nostrdb"); 791 let config = RadrootsNostrdbConfig::new(&db_dir); 792 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 793 let spec = RadrootsNostrdbQuerySpec::text_notes(Some(10), None, 10); 794 test_hooks::FORCE_QUERY_ERROR.store(true, Ordering::SeqCst); 795 796 let err = nostrdb.query_notes(&spec).expect_err("forced query error"); 797 assert!(err.to_string().starts_with("nostrdb error:")); 798 } 799 800 #[test] 801 fn query_notes_reports_note_json_error() { 802 let _guard = test_hooks_guard(); 803 reset_test_flags(); 804 let tmp_dir = TempDir::new().expect("tempdir should open"); 805 let db_dir = tmp_dir.path().join("nostrdb"); 806 let config = RadrootsNostrdbConfig::new(&db_dir); 807 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 808 809 let keys = RadrootsNostrKeys::generate(); 810 let event = EventBuilder::text_note("note json error") 811 .sign_with_keys(&keys) 812 .expect("event should sign"); 813 nostrdb 814 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 815 .expect("ingest should succeed"); 816 817 let query_spec = RadrootsNostrdbQuerySpec::text_notes(Some(50), None, 50); 818 for _ in 0..40 { 819 let notes = nostrdb 820 .query_notes(&query_spec) 821 .expect("query should succeed"); 822 if !notes.is_empty() { 823 break; 824 } 825 std::thread::sleep(Duration::from_millis(25)); 826 } 827 test_hooks::FORCE_NOTE_JSON_ERROR.store(true, Ordering::SeqCst); 828 829 let err = nostrdb 830 .query_notes(&query_spec) 831 .expect_err("forced note json error"); 832 assert!(err.to_string().starts_with("nostrdb error:")); 833 } 834 835 #[test] 836 fn profile_lookup_returns_metadata_fields() { 837 let _guard = test_hooks_guard(); 838 reset_test_flags(); 839 let tmp_dir = TempDir::new().expect("tempdir should open"); 840 let db_dir = tmp_dir.path().join("nostrdb"); 841 let config = RadrootsNostrdbConfig::new(&db_dir); 842 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 843 844 let keys = RadrootsNostrKeys::generate(); 845 let pubkey_hex = keys.public_key().to_hex(); 846 let metadata = RadrootsNostrMetadata::new() 847 .name(FIXTURE_ALICE_USERNAME) 848 .display_name(FIXTURE_ALICE_USERNAME) 849 .about("coffee operator") 850 .lud16(FIXTURE_ALICE_EMAIL); 851 let metadata_event = EventBuilder::metadata(&metadata) 852 .sign_with_keys(&keys) 853 .expect("metadata event should sign"); 854 nostrdb 855 .ingest_event(&metadata_event, RadrootsNostrdbIngestSource::client()) 856 .expect("ingest should succeed"); 857 858 let mut profile = None; 859 for _ in 0..40 { 860 profile = nostrdb 861 .get_profile_by_pubkey_hex(pubkey_hex.as_str()) 862 .expect("profile lookup should succeed"); 863 if profile.is_some() { 864 break; 865 } 866 std::thread::sleep(Duration::from_millis(25)); 867 } 868 let profile = profile.expect("profile should exist"); 869 assert_eq!(profile.pubkey_hex, pubkey_hex); 870 assert_eq!(profile.name.as_deref(), Some(FIXTURE_ALICE_USERNAME)); 871 assert_eq!( 872 profile.display_name.as_deref(), 873 Some(FIXTURE_ALICE_USERNAME) 874 ); 875 assert_eq!(profile.lud16.as_deref(), Some(FIXTURE_ALICE_EMAIL)); 876 } 877 878 #[test] 879 fn profile_lookup_returns_none_when_missing() { 880 let _guard = test_hooks_guard(); 881 reset_test_flags(); 882 let tmp_dir = TempDir::new().expect("tempdir should open"); 883 let db_dir = tmp_dir.path().join("nostrdb"); 884 let config = RadrootsNostrdbConfig::new(&db_dir); 885 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 886 887 let pubkey_hex = RadrootsNostrKeys::generate().public_key().to_hex(); 888 let profile = nostrdb 889 .get_profile_by_pubkey_hex(pubkey_hex.as_str()) 890 .expect("profile lookup"); 891 assert!(profile.is_none()); 892 } 893 894 #[test] 895 fn profile_lookup_reports_query_error() { 896 let _guard = test_hooks_guard(); 897 reset_test_flags(); 898 let tmp_dir = TempDir::new().expect("tempdir should open"); 899 let db_dir = tmp_dir.path().join("nostrdb"); 900 let config = RadrootsNostrdbConfig::new(&db_dir); 901 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 902 let pubkey_hex = RadrootsNostrKeys::generate().public_key().to_hex(); 903 test_hooks::FORCE_PROFILE_QUERY_ERROR.store(true, Ordering::SeqCst); 904 905 let err = nostrdb 906 .get_profile_by_pubkey_hex(pubkey_hex.as_str()) 907 .expect_err("forced profile query error"); 908 assert!(err.to_string().starts_with("nostrdb error:")); 909 } 910 911 #[test] 912 fn profile_lookup_reports_transaction_error() { 913 let _guard = test_hooks_guard(); 914 reset_test_flags(); 915 let tmp_dir = TempDir::new().expect("tempdir should open"); 916 let db_dir = tmp_dir.path().join("nostrdb"); 917 let config = RadrootsNostrdbConfig::new(&db_dir); 918 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 919 let pubkey_hex = RadrootsNostrKeys::generate().public_key().to_hex(); 920 test_hooks::FORCE_TRANSACTION_ERROR.store(true, Ordering::SeqCst); 921 922 let err = nostrdb 923 .get_profile_by_pubkey_hex(pubkey_hex.as_str()) 924 .expect_err("forced transaction error"); 925 assert!(err.to_string().starts_with("nostrdb error:")); 926 } 927 928 #[test] 929 fn profile_lookup_returns_none_without_metadata_record() { 930 let _guard = test_hooks_guard(); 931 reset_test_flags(); 932 let tmp_dir = TempDir::new().expect("tempdir should open"); 933 let db_dir = tmp_dir.path().join("nostrdb"); 934 let config = RadrootsNostrdbConfig::new(&db_dir); 935 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 936 937 let keys = RadrootsNostrKeys::generate(); 938 let pubkey_hex = keys.public_key().to_hex(); 939 let event = EventBuilder::text_note("non profile event") 940 .sign_with_keys(&keys) 941 .expect("event should sign"); 942 nostrdb 943 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 944 .expect("ingest should succeed"); 945 946 let profile = nostrdb 947 .get_profile_by_pubkey_hex(pubkey_hex.as_str()) 948 .expect("profile lookup"); 949 assert!(profile.is_none()); 950 } 951 952 #[test] 953 fn profile_lookup_invalid_metadata_content_returns_none() { 954 let _guard = test_hooks_guard(); 955 reset_test_flags(); 956 let tmp_dir = TempDir::new().expect("tempdir should open"); 957 let db_dir = tmp_dir.path().join("nostrdb"); 958 let config = RadrootsNostrdbConfig::new(&db_dir); 959 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 960 961 let keys = RadrootsNostrKeys::generate(); 962 let pubkey_hex = keys.public_key().to_hex(); 963 let event = EventBuilder::new( 964 radroots_nostr::event::Kind::Metadata, 965 "not valid metadata json", 966 ) 967 .sign_with_keys(&keys) 968 .expect("event should sign"); 969 nostrdb 970 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 971 .expect("ingest should succeed"); 972 973 let result = nostrdb.get_profile_by_pubkey_hex(pubkey_hex.as_str()); 974 assert!(result.expect("profile lookup").is_none()); 975 } 976 977 #[test] 978 fn subscribe_rejects_invalid_author_hex() { 979 let tmp_dir = TempDir::new().expect("tempdir should open"); 980 let db_dir = tmp_dir.path().join("nostrdb"); 981 let config = RadrootsNostrdbConfig::new(&db_dir); 982 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 983 984 let spec = RadrootsNostrdbSubscriptionSpec::single( 985 RadrootsNostrdbFilterSpec::new().with_author_hex("not-hex"), 986 ); 987 let err = nostrdb.subscribe(&spec).expect_err("subscribe should fail"); 988 assert!(err.to_string().contains("invalid hex for author")); 989 } 990 991 #[test] 992 fn profile_lookup_rejects_invalid_pubkey_length() { 993 let _guard = test_hooks_guard(); 994 reset_test_flags(); 995 let tmp_dir = TempDir::new().expect("tempdir should open"); 996 let db_dir = tmp_dir.path().join("nostrdb"); 997 let config = RadrootsNostrdbConfig::new(&db_dir); 998 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 999 1000 let err = nostrdb 1001 .get_profile_by_pubkey_hex("abcd") 1002 .expect_err("lookup should fail"); 1003 assert!(err.to_string().contains("invalid hex length for pubkey")); 1004 } 1005 1006 #[tokio::test] 1007 async fn subscription_stream_yields_events() { 1008 let tmp_dir = TempDir::new().expect("tempdir should open"); 1009 let db_dir = tmp_dir.path().join("nostrdb"); 1010 let config = RadrootsNostrdbConfig::new(&db_dir); 1011 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 1012 let spec = RadrootsNostrdbSubscriptionSpec::text_notes(Some(10), None); 1013 let handle = nostrdb.subscribe(&spec).expect("subscribe should succeed"); 1014 let mut stream = nostrdb.subscription_stream(handle, 0); 1015 1016 let pending = tokio::time::timeout(Duration::from_millis(20), stream.next()).await; 1017 assert!(pending.is_err()); 1018 1019 let keys = RadrootsNostrKeys::generate(); 1020 let event = EventBuilder::text_note("stream note") 1021 .sign_with_keys(&keys) 1022 .expect("event should sign"); 1023 nostrdb 1024 .ingest_event(&event, RadrootsNostrdbIngestSource::client()) 1025 .expect("ingest should succeed"); 1026 1027 let note_keys = tokio::time::timeout(Duration::from_secs(2), stream.next()) 1028 .await 1029 .expect("stream should wake") 1030 .expect("stream should yield note keys"); 1031 assert!(!note_keys.is_empty()); 1032 assert!(note_keys.iter().all(|key| key.as_u64() > 0)); 1033 } 1034 1035 #[test] 1036 fn concurrent_ingest_handles_parallel_writers() { 1037 let _guard = test_hooks_guard(); 1038 reset_test_flags(); 1039 let tmp_dir = TempDir::new().expect("tempdir should open"); 1040 let db_dir = tmp_dir.path().join("nostrdb"); 1041 let config = RadrootsNostrdbConfig::new(&db_dir); 1042 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 1043 1044 let worker_count = 4usize; 1045 let notes_per_worker = 20usize; 1046 let mut handles = Vec::new(); 1047 1048 for worker in 0..worker_count { 1049 let db = nostrdb.clone(); 1050 handles.push(std::thread::spawn(move || { 1051 let keys = RadrootsNostrKeys::generate(); 1052 for idx in 0..notes_per_worker { 1053 let content = format!("parallel-{worker}-{idx}"); 1054 let event = EventBuilder::text_note(content.as_str()) 1055 .sign_with_keys(&keys) 1056 .expect("event should sign"); 1057 db.ingest_event(&event, RadrootsNostrdbIngestSource::client()) 1058 .expect("ingest should succeed"); 1059 } 1060 })); 1061 } 1062 1063 for handle in handles { 1064 handle.join().expect("worker should complete"); 1065 } 1066 1067 let query_spec = RadrootsNostrdbQuerySpec::text_notes(Some(512), None, 512); 1068 let expected = worker_count * notes_per_worker; 1069 let mut observed = 0usize; 1070 let mut break_threshold = expected + 1; 1071 1072 for _ in 0..80 { 1073 let notes = nostrdb 1074 .query_notes(&query_spec) 1075 .expect("query should succeed"); 1076 observed = notes 1077 .iter() 1078 .filter(|note| note.content.starts_with("parallel-")) 1079 .count(); 1080 if observed >= break_threshold { 1081 break; 1082 } 1083 break_threshold = expected; 1084 std::thread::sleep(Duration::from_millis(25)); 1085 } 1086 1087 assert!( 1088 observed >= expected, 1089 "expected at least {expected} parallel notes, got {observed}" 1090 ); 1091 } 1092 1093 #[cfg(feature = "giftwrap")] 1094 #[test] 1095 fn giftwrap_secret_key_hex_validates_length() { 1096 let tmp_dir = TempDir::new().expect("tempdir should open"); 1097 let db_dir = tmp_dir.path().join("nostrdb"); 1098 let config = RadrootsNostrdbConfig::new(&db_dir); 1099 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 1100 1101 let result = nostrdb.add_giftwrap_secret_key_hex("abcd"); 1102 let err = result.expect_err("invalid giftwrap key"); 1103 assert!( 1104 err.to_string() 1105 .contains("invalid hex length for secret_key") 1106 ); 1107 } 1108 1109 #[cfg(feature = "giftwrap")] 1110 #[test] 1111 fn giftwrap_process_flow_executes() { 1112 let tmp_dir = TempDir::new().expect("tempdir should open"); 1113 let db_dir = tmp_dir.path().join("nostrdb"); 1114 let config = RadrootsNostrdbConfig::new(&db_dir); 1115 let nostrdb = RadrootsNostrdb::open(config).expect("database should open"); 1116 1117 let secret_key = [7u8; 32]; 1118 let _ = nostrdb.add_giftwrap_secret_key(secret_key); 1119 nostrdb 1120 .process_giftwraps() 1121 .expect("giftwrap processing should run"); 1122 } 1123 }