transport_publish.rs (228802B)
1 use std::collections::{BTreeMap, BTreeSet}; 2 use std::fmt; 3 use std::future::Future; 4 use std::net::IpAddr; 5 use std::path::{Path, PathBuf}; 6 use std::pin::Pin; 7 use std::str::FromStr; 8 use std::sync::{Arc, Mutex}; 9 use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; 10 11 use crate::host_nostr::{DaemonNostrClient, Filter, Kind, PublicKey}; 12 use crate::transport::relay_publish::{ 13 LiveRelayPublishAdapter as DaemonNostrClientPublishAdapter, 14 RelayOutcome as RadrootsRelayOutcome, RelayOutcomeKind as RadrootsRelayOutcomeKind, 15 RelayPublishAdapter as RadrootsRelayPublishAdapter, 16 RelayPublishReceipt as RadrootsRelayPublishRelayReceipt, 17 RelayPublishRequest as RadrootsRelayPublishRequest, RelayTargetSet as RadrootsRelayTargetSet, 18 RelayUrl as RadrootsRelayUrl, RelayUrlPolicy as RadrootsRelayUrlPolicy, 19 }; 20 use radroots_event::draft::{SignedEvent, SignedEventError}; 21 use radroots_event::wire::{EventWireError, Nip01EventWire}; 22 use radroots_protocol::radrootsd::transport_publish::v5::{ 23 DeliveryPolicy, EventRequest, EventResponse, Job, JobStatus, NostrTargetSourcePolicy, 24 OutcomeKind, RETICULUM_ENDPOINT_URI as RADROOTS_RETICULUM_ENDPOINT_URI, 25 RETICULUM_UNAVAILABLE_MESSAGE as RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, Target, 26 TargetFingerprint as ProtocolTargetFingerprint, TargetOutcome, TargetPolicy, TargetPolicyName, 27 TargetSource, 28 }; 29 use radroots_transport::{ 30 Error as TransportError, Target as TransportTarget, TransportId, 31 target::{TargetFingerprint, TargetLabel, TargetScope}, 32 }; 33 use serde::{Deserialize, Serialize}; 34 use sha2::{Digest, Sha256}; 35 use sqlx::sqlite::{SqliteConnectOptions, SqliteConnection, SqliteRow}; 36 use sqlx::{Connection as _, Row}; 37 use thiserror::Error; 38 use tokio::sync::{OwnedSemaphorePermit, Semaphore}; 39 use uuid::Uuid; 40 41 use crate::app::config::TransportPublishConfig; 42 43 const TOKEN_PREFIX: &str = "rrd_tp_"; 44 const TOKEN_HASH_PREFIX: &str = "sha256:"; 45 const SCHEMA_VERSION: i64 = 4; 46 const TRANSPORT_KIND_NOSTR: &str = "nostr"; 47 const TRANSPORT_KIND_RETICULUM: &str = "reticulum"; 48 const TRANSPORT_PUBLISH_SCHEMA_SQL: &str = r#" 49 CREATE TABLE IF NOT EXISTS transport_publish_principals ( 50 principal_id TEXT PRIMARY KEY NOT NULL, 51 label TEXT NOT NULL, 52 token_hash TEXT NOT NULL UNIQUE, 53 allowed_pubkeys_json TEXT NOT NULL, 54 allowed_kinds_json TEXT NOT NULL, 55 allowed_target_policies_json TEXT NOT NULL, 56 allowed_explicit_transport_kinds_json TEXT NOT NULL, 57 allowed_nostr_source_policies_json TEXT NOT NULL, 58 allow_request_targets INTEGER NOT NULL, 59 job_visibility TEXT NOT NULL, 60 expires_at_unix INTEGER, 61 revoked_at_unix INTEGER, 62 created_at_unix INTEGER NOT NULL 63 ); 64 CREATE TABLE IF NOT EXISTS transport_publish_jobs ( 65 job_id TEXT PRIMARY KEY NOT NULL, 66 principal_id TEXT NOT NULL, 67 idempotency_key TEXT, 68 request_fingerprint TEXT NOT NULL, 69 status TEXT NOT NULL, 70 event_id TEXT NOT NULL, 71 event_pubkey TEXT NOT NULL, 72 event_kind INTEGER NOT NULL, 73 target_policy_json TEXT NOT NULL, 74 delivery_policy_json TEXT NOT NULL, 75 requested_target_count INTEGER NOT NULL, 76 effective_target_count INTEGER NOT NULL, 77 request_json TEXT NOT NULL, 78 requested_at_ms INTEGER NOT NULL, 79 updated_at_ms INTEGER NOT NULL, 80 completed_at_ms INTEGER, 81 last_error TEXT, 82 FOREIGN KEY(principal_id) REFERENCES transport_publish_principals(principal_id) 83 ); 84 CREATE UNIQUE INDEX IF NOT EXISTS transport_publish_jobs_principal_idempotency_idx 85 ON transport_publish_jobs(principal_id, idempotency_key) 86 WHERE idempotency_key IS NOT NULL; 87 CREATE TABLE IF NOT EXISTS transport_publish_target_results ( 88 job_id TEXT NOT NULL, 89 transport_kind TEXT NOT NULL, 90 endpoint_uri TEXT NOT NULL, 91 target_scope TEXT NOT NULL, 92 target_label TEXT, 93 source TEXT NOT NULL, 94 attempted INTEGER NOT NULL, 95 outcome_kind TEXT NOT NULL, 96 message TEXT, 97 latency_ms INTEGER, 98 updated_at_ms INTEGER NOT NULL, 99 PRIMARY KEY(job_id, transport_kind, endpoint_uri, target_scope), 100 FOREIGN KEY(job_id) REFERENCES transport_publish_jobs(job_id) 101 ); 102 CREATE TABLE IF NOT EXISTS transport_publish_target_snapshots ( 103 job_id TEXT NOT NULL, 104 target_index INTEGER NOT NULL, 105 transport_kind TEXT NOT NULL, 106 endpoint_uri TEXT NOT NULL, 107 target_scope TEXT NOT NULL, 108 target_label TEXT, 109 source TEXT NOT NULL, 110 attempted INTEGER NOT NULL, 111 outcome_kind TEXT NOT NULL, 112 message TEXT, 113 latency_ms INTEGER, 114 created_at_ms INTEGER NOT NULL, 115 PRIMARY KEY(job_id, target_index), 116 FOREIGN KEY(job_id) REFERENCES transport_publish_jobs(job_id) 117 ); 118 CREATE TABLE IF NOT EXISTS transport_publish_nostr_author_cache ( 119 pubkey TEXT PRIMARY KEY NOT NULL, 120 relays_json TEXT NOT NULL, 121 updated_at_ms INTEGER NOT NULL 122 ); 123 "#; 124 const TRANSPORT_PUBLISH_TABLES: &[&str] = &[ 125 "transport_publish_principals", 126 "transport_publish_jobs", 127 "transport_publish_target_results", 128 "transport_publish_target_snapshots", 129 "transport_publish_nostr_author_cache", 130 ]; 131 132 #[derive(Debug, Error)] 133 pub enum TransportPublishError { 134 #[error("transport publish storage error: {0}")] 135 Sqlite(#[from] sqlx::Error), 136 #[error("transport publish json error: {0}")] 137 Json(#[from] serde_json::Error), 138 #[error("transport publish io error: {0}")] 139 Io(#[from] std::io::Error), 140 #[error("invalid transport publish scope: {0}")] 141 InvalidScope(String), 142 #[error("invalid signed Nostr event: {0}")] 143 InvalidSignedEvent(String), 144 #[error("signed event wire error: {0}")] 145 EventWire(#[from] EventWireError), 146 #[error("signed event conversion error: {0}")] 147 SignedEvent(#[from] SignedEventError), 148 #[error("transport publish relay error: {0}")] 149 Relay(String), 150 #[error("transport publish transport error: {0}")] 151 Transport(String), 152 #[error("transport publish schema incompatible for table `{table}`: {detail}")] 153 Schema { table: &'static str, detail: String }, 154 #[error("transport publish job state validation failed: {0}")] 155 InvalidPublishJobState(String), 156 #[error("transport publish concurrency limit reached")] 157 ConcurrencyLimit, 158 #[error("transport publish idempotency conflict for key `{0}`")] 159 IdempotencyConflict(String), 160 } 161 162 #[derive(Clone)] 163 pub struct TransportPublish { 164 pub config: TransportPublishConfig, 165 pub store: TransportPublishStore, 166 publisher: Option<Arc<dyn RadrootsRelayPublishAdapter>>, 167 resolver: Arc<dyn PublishRelayResolver>, 168 author_relay_discovery: Arc<dyn PublishAuthorRelayDiscovery>, 169 publish_jobs: Arc<Semaphore>, 170 } 171 172 impl TransportPublish { 173 #[cfg(not(test))] 174 pub fn open(config: TransportPublishConfig) -> Result<Self, TransportPublishError> { 175 let store = TransportPublishStore::open(config.database_path.clone())?; 176 let publish_jobs = Arc::new(Semaphore::new(config.max_concurrent_publish_jobs)); 177 Ok(Self { 178 config, 179 store, 180 publisher: None, 181 resolver: Arc::new(SystemPublishRelayResolver), 182 author_relay_discovery: Arc::new(NostrPublishAuthorRelayDiscovery), 183 publish_jobs, 184 }) 185 } 186 187 #[cfg(test)] 188 pub fn memory(config: TransportPublishConfig) -> Result<Self, TransportPublishError> { 189 let store = TransportPublishStore::memory()?; 190 let publish_jobs = Arc::new(Semaphore::new(config.max_concurrent_publish_jobs)); 191 Ok(Self { 192 config, 193 store, 194 publisher: None, 195 resolver: Arc::new(SystemPublishRelayResolver), 196 author_relay_discovery: Arc::new(NostrPublishAuthorRelayDiscovery), 197 publish_jobs, 198 }) 199 } 200 201 #[cfg(test)] 202 pub(crate) fn with_publisher( 203 mut self, 204 publisher: Arc<dyn RadrootsRelayPublishAdapter>, 205 ) -> Self { 206 self.publisher = Some(publisher); 207 self 208 } 209 210 #[cfg(test)] 211 pub(crate) fn with_relay_resolver(mut self, resolver: Arc<dyn PublishRelayResolver>) -> Self { 212 self.resolver = resolver; 213 self 214 } 215 216 #[cfg(test)] 217 fn with_author_relay_discovery( 218 mut self, 219 author_relay_discovery: Arc<dyn PublishAuthorRelayDiscovery>, 220 ) -> Self { 221 self.author_relay_discovery = author_relay_discovery; 222 self 223 } 224 225 fn acquire_publish_permit(&self) -> Result<OwnedSemaphorePermit, TransportPublishError> { 226 self.publish_jobs 227 .clone() 228 .try_acquire_owned() 229 .map_err(|_| TransportPublishError::ConcurrencyLimit) 230 } 231 232 pub async fn publish_event( 233 &self, 234 principal: &PublishPrincipal, 235 request: EventRequest, 236 ) -> Result<EventResponse, TransportPublishError> { 237 request 238 .validate(self.config.max_targets_per_request) 239 .map_err(|error| { 240 TransportPublishError::InvalidSignedEvent(format!( 241 "publish request validation failed: {error}" 242 )) 243 })?; 244 if request.raw_event_json.len() > self.config.max_event_bytes { 245 return Err(TransportPublishError::InvalidSignedEvent( 246 "signed event exceeds transport_publish max_event_bytes".to_owned(), 247 )); 248 } 249 let signed_event = signed_event_from_raw_json(request.raw_event_json.as_str())?; 250 principal.allows_event(&signed_event, &request)?; 251 let effective_timeout_ms = effective_publish_timeout_ms(&self.config, request.timeout_ms)?; 252 let _permit = self.acquire_publish_permit()?; 253 let request_fingerprint = request_intent_fingerprint( 254 principal.principal_id.as_str(), 255 signed_event.raw_json(), 256 &request, 257 effective_timeout_ms, 258 )?; 259 let signed_pubkey = signed_event.pubkey().to_hex(); 260 let resolution = self 261 .resolve_targets_for_request(signed_pubkey.as_str(), &request) 262 .await?; 263 validate_delivery_policy_for_resolution(&request.delivery_policy, &resolution)?; 264 let target_snapshots = target_snapshots_from_resolution(&resolution); 265 let response = self.store.record_publish_job(PublishJobInsert { 266 principal_id: principal.principal_id.clone(), 267 idempotency_key: request.idempotency_key.clone(), 268 event: PublishEventMetadata::from_signed_event(&signed_event), 269 request: request.clone(), 270 request_fingerprint, 271 effective_target_count: resolution.target_count(), 272 target_snapshots, 273 })?; 274 if response.deduplicated { 275 return Ok(response); 276 } 277 let completed = self 278 .complete_job_execution( 279 response.job.job_id.as_str(), 280 signed_event, 281 request.delivery_policy.clone(), 282 effective_timeout_ms, 283 resolution, 284 ) 285 .await?; 286 Ok(EventResponse { 287 deduplicated: false, 288 job: completed, 289 }) 290 } 291 292 pub async fn resolve_targets_for_request( 293 &self, 294 pubkey: &str, 295 request: &EventRequest, 296 ) -> Result<PublishRelayResolution, TransportPublishError> { 297 match &request.target_policy { 298 TargetPolicy::ExplicitTargets { targets } => { 299 self.resolve_explicit_targets(targets).await 300 } 301 TargetPolicy::Nostr { 302 source_policy, 303 relay_urls, 304 } => match source_policy { 305 NostrTargetSourcePolicy::ExplicitOnly => { 306 self.resolve_request_relays(relay_urls).await 307 } 308 NostrTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault => { 309 if !relay_urls.is_empty() { 310 self.resolve_request_relays(relay_urls).await 311 } else { 312 self.resolve_author_or_default_relays(pubkey).await 313 } 314 } 315 NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault => { 316 self.resolve_author_or_default_relays(pubkey).await 317 } 318 NostrTargetSourcePolicy::DaemonDefaultOnly => { 319 self.resolve_daemon_default_relays().await 320 } 321 }, 322 } 323 } 324 325 async fn resolve_explicit_targets( 326 &self, 327 targets: &[Target], 328 ) -> Result<PublishRelayResolution, TransportPublishError> { 329 let mut resolved = Vec::new(); 330 let mut outcomes = Vec::new(); 331 for (index, target) in targets.iter().enumerate() { 332 match TransportId::parse_canonical(target.transport_kind.as_str()).map_err(|error| { 333 TransportPublishError::InvalidSignedEvent(format!( 334 "transport target {index} kind is invalid: {error}" 335 )) 336 })? { 337 TransportId::NOSTR => { 338 self.resolve_request_target(&mut resolved, &mut outcomes, target) 339 .await; 340 } 341 TransportId::RETICULUM => { 342 outcomes.push(reticulum_unavailable_outcome(target)); 343 } 344 _ => outcomes.push(unsupported_transport_outcome(target)), 345 } 346 } 347 Ok(PublishRelayResolution { 348 targets: resolved, 349 outcomes, 350 }) 351 } 352 353 async fn resolve_request_target( 354 &self, 355 targets: &mut Vec<ResolvedPublishRelay>, 356 outcomes: &mut Vec<TargetOutcome>, 357 target: &Target, 358 ) { 359 match RadrootsRelayUrl::parse(target.endpoint_uri.as_str(), relay_url_policy(&self.config)) 360 { 361 Ok(url) => { 362 self.push_checked_relay_target( 363 targets, 364 outcomes, 365 url, 366 TargetSource::Request, 367 PublishTargetMetadata::from_target(target), 368 ) 369 .await; 370 } 371 Err(error) => outcomes.push(TargetOutcome { 372 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 373 endpoint_uri: target.endpoint_uri.clone(), 374 target_scope: target.target_scope.clone(), 375 target_label: target.target_label.clone(), 376 source: TargetSource::Request, 377 attempted: false, 378 outcome_kind: OutcomeKind::TargetRejected, 379 message: Some(error.to_string()), 380 latency_ms: None, 381 }), 382 } 383 } 384 async fn resolve_author_or_default_relays( 385 &self, 386 pubkey: &str, 387 ) -> Result<PublishRelayResolution, TransportPublishError> { 388 let mut author_relays = self.resolve_author_write_relays(pubkey).await?; 389 if author_relays.targets.is_empty() { 390 let mut daemon_defaults = self.resolve_daemon_default_relays().await?; 391 daemon_defaults.outcomes.append(&mut author_relays.outcomes); 392 Ok(daemon_defaults) 393 } else { 394 Ok(author_relays) 395 } 396 } 397 398 async fn resolve_request_relays( 399 &self, 400 relays: &[String], 401 ) -> Result<PublishRelayResolution, TransportPublishError> { 402 let mut targets = Vec::new(); 403 let mut outcomes = Vec::new(); 404 for relay in relays { 405 match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { 406 Ok(url) => { 407 self.push_checked_relay_target( 408 &mut targets, 409 &mut outcomes, 410 url, 411 TargetSource::Request, 412 PublishTargetMetadata::default(), 413 ) 414 .await; 415 } 416 Err(error) => outcomes.push(TargetOutcome { 417 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 418 endpoint_uri: relay.clone(), 419 target_scope: None, 420 target_label: None, 421 source: TargetSource::Request, 422 attempted: false, 423 outcome_kind: OutcomeKind::TargetRejected, 424 message: Some(error.to_string()), 425 latency_ms: None, 426 }), 427 } 428 } 429 Ok(PublishRelayResolution { targets, outcomes }) 430 } 431 432 async fn resolve_author_write_relays( 433 &self, 434 pubkey: &str, 435 ) -> Result<PublishRelayResolution, TransportPublishError> { 436 let cached = self.store.cached_author_write_relays(pubkey)?; 437 let mut cached_resolution = self.resolve_author_relay_inputs(&cached).await?; 438 if !cached_resolution.targets.is_empty() { 439 return Ok(cached_resolution); 440 } 441 if self.config.nostr.author_relay_discovery_relays.is_empty() { 442 return Ok(cached_resolution); 443 } 444 let mut discovery_targets = self 445 .resolve_config_relays( 446 &self.config.nostr.author_relay_discovery_relays, 447 TargetSource::DaemonDefault, 448 ) 449 .await?; 450 if discovery_targets.targets.is_empty() { 451 discovery_targets 452 .outcomes 453 .append(&mut cached_resolution.outcomes); 454 return Ok(discovery_targets); 455 } 456 let discovered = self 457 .author_relay_discovery 458 .fetch_author_write_relays( 459 pubkey, 460 std::mem::take(&mut discovery_targets.targets), 461 self.config.connect_timeout_secs, 462 ) 463 .await?; 464 self.store.cache_author_write_relays(pubkey, &discovered)?; 465 let mut discovered_resolution = self.resolve_author_relay_inputs(&discovered).await?; 466 discovered_resolution 467 .outcomes 468 .append(&mut cached_resolution.outcomes); 469 discovered_resolution 470 .outcomes 471 .append(&mut discovery_targets.outcomes); 472 Ok(discovered_resolution) 473 } 474 475 async fn resolve_author_relay_inputs( 476 &self, 477 relays: &[String], 478 ) -> Result<PublishRelayResolution, TransportPublishError> { 479 let mut targets = Vec::new(); 480 let mut outcomes = Vec::new(); 481 for relay in relays { 482 match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { 483 Ok(url) => { 484 self.push_checked_relay_target( 485 &mut targets, 486 &mut outcomes, 487 url, 488 TargetSource::NostrAuthorWrite, 489 PublishTargetMetadata::default(), 490 ) 491 .await; 492 } 493 Err(_error) => {} 494 } 495 } 496 Ok(PublishRelayResolution { targets, outcomes }) 497 } 498 499 async fn resolve_daemon_default_relays( 500 &self, 501 ) -> Result<PublishRelayResolution, TransportPublishError> { 502 self.resolve_config_relays( 503 &self.config.nostr.daemon_default_relays, 504 TargetSource::DaemonDefault, 505 ) 506 .await 507 } 508 509 async fn resolve_config_relays( 510 &self, 511 relays: &[String], 512 source: TargetSource, 513 ) -> Result<PublishRelayResolution, TransportPublishError> { 514 let mut targets = Vec::new(); 515 let mut outcomes = Vec::new(); 516 for relay in relays { 517 match RadrootsRelayUrl::parse(relay, relay_url_policy(&self.config)) { 518 Ok(url) => { 519 self.push_checked_relay_target( 520 &mut targets, 521 &mut outcomes, 522 url, 523 source, 524 PublishTargetMetadata::default(), 525 ) 526 .await; 527 } 528 Err(_error) => {} 529 } 530 } 531 Ok(PublishRelayResolution { targets, outcomes }) 532 } 533 534 async fn push_checked_relay_target( 535 &self, 536 targets: &mut Vec<ResolvedPublishRelay>, 537 outcomes: &mut Vec<TargetOutcome>, 538 url: RadrootsRelayUrl, 539 source: TargetSource, 540 metadata: PublishTargetMetadata, 541 ) { 542 if relay_url_policy(&self.config) == RadrootsRelayUrlPolicy::Localhost { 543 push_resolved_relay(targets, url, source, metadata); 544 return; 545 } 546 match self.resolver.resolve(&url).await { 547 Ok(addresses) if addresses.is_empty() => { 548 outcomes.push(relay_resolution_connection_failure( 549 url.as_str(), 550 source, 551 &metadata, 552 "dns lookup returned no addresses", 553 )); 554 } 555 Ok(addresses) => match url.validate_public_resolved_ip_addrs(addresses) { 556 Ok(()) => push_resolved_relay(targets, url, source, metadata), 557 Err(error) => outcomes.push(TargetOutcome { 558 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 559 endpoint_uri: url.as_str().to_owned(), 560 target_scope: metadata.target_scope, 561 target_label: metadata.target_label, 562 source, 563 attempted: false, 564 outcome_kind: OutcomeKind::TargetRejected, 565 message: Some(error.to_string()), 566 latency_ms: None, 567 }), 568 }, 569 Err(error) => outcomes.push(relay_resolution_connection_failure( 570 url.as_str(), 571 source, 572 &metadata, 573 format!("dns lookup failed: {error}"), 574 )), 575 } 576 } 577 578 async fn complete_job_execution( 579 &self, 580 job_id: &str, 581 signed_event: SignedEvent, 582 delivery_policy: DeliveryPolicy, 583 timeout_ms: u64, 584 resolution: PublishRelayResolution, 585 ) -> Result<Job, TransportPublishError> { 586 let target_count = resolution.target_count(); 587 if resolution.targets.is_empty() { 588 let status = if resolution.outcomes.is_empty() { 589 JobStatus::Rejected 590 } else { 591 delivery_status(&delivery_policy, target_count, &resolution.outcomes) 592 }; 593 let last_error = last_error_for_status(status); 594 self.store.complete_publish_job( 595 job_id, 596 status, 597 resolution.outcomes, 598 last_error.map(str::to_owned), 599 )?; 600 return self.store.job_by_id(job_id); 601 } 602 let required_target_count = delivery_policy.required_target_count(target_count); 603 if required_target_count > target_count { 604 self.store.complete_publish_job( 605 job_id, 606 JobStatus::Rejected, 607 resolution.outcomes, 608 Some("delivery_quorum_exceeds_target_count".to_owned()), 609 )?; 610 return self.store.job_by_id(job_id); 611 } 612 let target_set = RadrootsRelayTargetSet::from_urls( 613 resolution 614 .targets 615 .iter() 616 .map(|target| target.url.clone()) 617 .collect(), 618 ) 619 .map_err(|error| TransportPublishError::Relay(error.to_string()))?; 620 let publish_relays = target_set.relays().to_vec(); 621 let publish_request = 622 RadrootsRelayPublishRequest::new(signed_event, target_set, current_unix_millis()); 623 let started = Instant::now(); 624 let publish_timeout = Duration::from_millis(timeout_ms); 625 let receipts = 626 match tokio::time::timeout(publish_timeout, self.publish_with_adapter(publish_request)) 627 .await 628 { 629 Ok(Ok(receipts)) => receipts, 630 Ok(Err(error)) => transport_error_receipts(publish_relays.as_slice(), error), 631 Err(_) => timeout_receipts(publish_relays.as_slice()), 632 }; 633 let latency_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX); 634 let mut outcomes = resolution.outcomes; 635 outcomes.extend(receipts.into_iter().flat_map(|receipt| { 636 publish_outcomes_from_receipt(receipt, &resolution.targets, Some(latency_ms)) 637 })); 638 let status = delivery_status(&delivery_policy, target_count, &outcomes); 639 let last_error = last_error_for_status(status).map(str::to_owned); 640 self.store 641 .complete_publish_job(job_id, status, outcomes, last_error)?; 642 self.store.job_by_id(job_id) 643 } 644 645 async fn publish_with_adapter( 646 &self, 647 request: RadrootsRelayPublishRequest, 648 ) -> Result<Vec<RadrootsRelayPublishRelayReceipt>, TransportPublishError> { 649 if let Some(publisher) = &self.publisher { 650 return publisher 651 .publish(request) 652 .await 653 .map_err(|error| TransportPublishError::Relay(error.to_string())); 654 } 655 let adapter = DaemonNostrClientPublishAdapter::new(DaemonNostrClient::signerless()); 656 adapter 657 .publish(request) 658 .await 659 .map_err(|error| TransportPublishError::Relay(error.to_string())) 660 } 661 } 662 663 #[derive(Clone)] 664 pub struct TransportPublishStore { 665 inner: Arc<Mutex<SqliteConnection>>, 666 } 667 668 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] 669 #[serde(rename_all = "snake_case")] 670 pub enum PublishJobVisibility { 671 Own, 672 Admin, 673 } 674 675 impl FromStr for PublishJobVisibility { 676 type Err = TransportPublishError; 677 678 fn from_str(value: &str) -> Result<Self, Self::Err> { 679 match value { 680 "own" => Ok(Self::Own), 681 "admin" => Ok(Self::Admin), 682 other => Err(TransportPublishError::InvalidScope(format!( 683 "unknown job visibility `{other}`" 684 ))), 685 } 686 } 687 } 688 689 impl fmt::Display for PublishJobVisibility { 690 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { 691 match self { 692 Self::Own => f.write_str("own"), 693 Self::Admin => f.write_str("admin"), 694 } 695 } 696 } 697 698 #[derive(Debug, Clone, PartialEq, Eq)] 699 pub struct PublishPrincipalInit { 700 pub label: String, 701 pub token_hash: String, 702 pub allowed_pubkeys: Vec<String>, 703 pub allowed_kinds: Vec<u32>, 704 pub allowed_target_policies: Vec<TargetPolicyName>, 705 pub allowed_explicit_transport_kinds: Vec<String>, 706 pub allowed_nostr_source_policies: Vec<NostrTargetSourcePolicy>, 707 pub allow_request_targets: bool, 708 pub job_visibility: PublishJobVisibility, 709 pub expires_at_unix: Option<i64>, 710 } 711 712 #[derive(Debug, Clone, PartialEq, Eq)] 713 pub struct PublishPrincipal { 714 pub principal_id: String, 715 pub label: String, 716 pub allowed_pubkeys: Vec<String>, 717 pub allowed_kinds: Vec<u32>, 718 pub allowed_target_policies: Vec<TargetPolicyName>, 719 pub allowed_explicit_transport_kinds: Vec<String>, 720 pub allowed_nostr_source_policies: Vec<NostrTargetSourcePolicy>, 721 pub allow_request_targets: bool, 722 pub job_visibility: PublishJobVisibility, 723 pub expires_at_unix: Option<i64>, 724 } 725 726 impl PublishPrincipal { 727 pub fn allows_event( 728 &self, 729 signed_event: &SignedEvent, 730 request: &EventRequest, 731 ) -> Result<(), TransportPublishError> { 732 let pubkey = signed_event.pubkey().to_hex(); 733 ensure_lower_hex("pubkey", pubkey.as_str(), 64)?; 734 if !self 735 .allowed_pubkeys 736 .iter() 737 .any(|allowed_pubkey| allowed_pubkey == pubkey.as_str()) 738 { 739 return Err(TransportPublishError::InvalidScope( 740 "principal is not allowed to publish for event pubkey".to_owned(), 741 )); 742 } 743 if !self.allowed_kinds.contains(&signed_event.kind()) { 744 return Err(TransportPublishError::InvalidScope( 745 "principal is not allowed to publish event kind".to_owned(), 746 )); 747 } 748 match &request.target_policy { 749 TargetPolicy::ExplicitTargets { targets } => { 750 if !self 751 .allowed_target_policies 752 .contains(&TargetPolicyName::ExplicitTargets) 753 { 754 return Err(TransportPublishError::InvalidScope( 755 "principal is not allowed to use explicit transport targets".to_owned(), 756 )); 757 } 758 if !self.allow_request_targets && !targets.is_empty() { 759 return Err(TransportPublishError::InvalidScope( 760 "principal is not allowed to provide request targets".to_owned(), 761 )); 762 } 763 for target in targets { 764 let kind = TransportId::parse_canonical(target.transport_kind.as_str()) 765 .map_err(|error| { 766 TransportPublishError::InvalidScope(format!( 767 "principal explicit target kind check failed: {error}" 768 )) 769 })?; 770 let transport_kind = kind.canonical_label(); 771 if !self 772 .allowed_explicit_transport_kinds 773 .iter() 774 .any(|allowed| allowed == &transport_kind) 775 { 776 return Err(TransportPublishError::InvalidScope(format!( 777 "principal is not allowed to use explicit transport target kind `{transport_kind}`" 778 ))); 779 } 780 } 781 } 782 TargetPolicy::Nostr { 783 source_policy, 784 relay_urls, 785 } => { 786 if !self 787 .allowed_target_policies 788 .contains(&TargetPolicyName::Nostr) 789 { 790 return Err(TransportPublishError::InvalidScope( 791 "principal is not allowed to use Nostr target policy".to_owned(), 792 )); 793 } 794 if !self.allowed_nostr_source_policies.contains(source_policy) { 795 return Err(TransportPublishError::InvalidScope( 796 "principal is not allowed to use requested Nostr source policy".to_owned(), 797 )); 798 } 799 if !self.allow_request_targets && !relay_urls.is_empty() { 800 return Err(TransportPublishError::InvalidScope( 801 "principal is not allowed to provide request targets".to_owned(), 802 )); 803 } 804 } 805 } 806 Ok(()) 807 } 808 809 fn can_read_job(&self, principal_id: &str) -> bool { 810 self.job_visibility == PublishJobVisibility::Admin || self.principal_id == principal_id 811 } 812 } 813 814 #[derive(Debug, Clone)] 815 pub struct PublishJobInsert { 816 pub principal_id: String, 817 pub idempotency_key: Option<String>, 818 pub event: PublishEventMetadata, 819 pub request: EventRequest, 820 pub request_fingerprint: String, 821 pub effective_target_count: usize, 822 pub target_snapshots: Vec<TargetOutcome>, 823 } 824 825 #[derive(Debug, Clone, PartialEq, Eq)] 826 pub struct PublishEventMetadata { 827 pub event_id: String, 828 pub pubkey: String, 829 pub kind: u32, 830 } 831 832 impl PublishEventMetadata { 833 fn from_signed_event(signed_event: &SignedEvent) -> Self { 834 Self { 835 event_id: signed_event.id_str().to_owned(), 836 pubkey: signed_event.pubkey().to_hex(), 837 kind: signed_event.kind(), 838 } 839 } 840 } 841 842 #[derive(Debug, Clone, PartialEq, Eq)] 843 pub struct ResolvedPublishRelay { 844 pub(crate) url: RadrootsRelayUrl, 845 pub source: TargetSource, 846 target_scope: Option<String>, 847 target_label: Option<String>, 848 } 849 850 #[derive(Debug, Clone, PartialEq, Eq)] 851 pub struct PublishRelayResolution { 852 pub targets: Vec<ResolvedPublishRelay>, 853 pub outcomes: Vec<TargetOutcome>, 854 } 855 856 impl PublishRelayResolution { 857 fn target_count(&self) -> usize { 858 self.targets.len() + self.outcomes.len() 859 } 860 861 fn target_fingerprints(&self) -> Result<Vec<TargetFingerprint>, TransportPublishError> { 862 let mut fingerprints = Vec::with_capacity(self.target_count()); 863 for target in &self.targets { 864 fingerprints.push(target.fingerprint()?); 865 } 866 for (index, outcome) in self.outcomes.iter().enumerate() { 867 fingerprints.push(target_outcome_fingerprint(outcome, index)?); 868 } 869 Ok(fingerprints) 870 } 871 } 872 873 impl ResolvedPublishRelay { 874 fn fingerprint(&self) -> Result<TargetFingerprint, TransportPublishError> { 875 let scope = self 876 .target_scope 877 .as_deref() 878 .map(TargetScope::parse) 879 .transpose() 880 .map_err(|error| TransportPublishError::Transport(error.to_string()))?; 881 let label = self 882 .target_label 883 .as_deref() 884 .map(TargetLabel::parse) 885 .transpose() 886 .map_err(|error| TransportPublishError::Transport(error.to_string()))?; 887 let target = TransportTarget::nostr_relay_with_metadata(self.url.as_str(), scope, label) 888 .map_err(|error| TransportPublishError::Transport(error.to_string()))?; 889 Ok(target.fingerprint().clone()) 890 } 891 } 892 893 #[derive(Clone, Debug, Default, PartialEq, Eq)] 894 struct PublishTargetMetadata { 895 target_scope: Option<String>, 896 target_label: Option<String>, 897 } 898 899 impl PublishTargetMetadata { 900 fn from_target(target: &Target) -> Self { 901 Self { 902 target_scope: target.target_scope.clone(), 903 target_label: target.target_label.clone(), 904 } 905 } 906 } 907 908 pub(crate) type PublishRelayResolveFuture<'a> = 909 Pin<Box<dyn Future<Output = Result<Vec<IpAddr>, std::io::Error>> + Send + 'a>>; 910 911 pub(crate) trait PublishRelayResolver: Send + Sync { 912 fn resolve<'a>(&'a self, url: &'a RadrootsRelayUrl) -> PublishRelayResolveFuture<'a>; 913 } 914 915 type PublishAuthorRelayDiscoveryFuture<'a> = 916 Pin<Box<dyn Future<Output = Result<Vec<String>, TransportPublishError>> + Send + 'a>>; 917 918 trait PublishAuthorRelayDiscovery: Send + Sync { 919 fn fetch_author_write_relays<'a>( 920 &'a self, 921 pubkey: &'a str, 922 discovery_targets: Vec<ResolvedPublishRelay>, 923 connect_timeout_secs: u64, 924 ) -> PublishAuthorRelayDiscoveryFuture<'a>; 925 } 926 927 #[derive(Debug)] 928 struct SystemPublishRelayResolver; 929 930 impl PublishRelayResolver for SystemPublishRelayResolver { 931 fn resolve<'a>(&'a self, url: &'a RadrootsRelayUrl) -> PublishRelayResolveFuture<'a> { 932 Box::pin(async move { 933 let (host, port) = relay_socket_target(url)?; 934 let addrs = tokio::net::lookup_host((host.as_str(), port)).await?; 935 Ok(addrs.map(|addr| addr.ip()).collect()) 936 }) 937 } 938 } 939 940 #[derive(Debug)] 941 struct NostrPublishAuthorRelayDiscovery; 942 943 impl PublishAuthorRelayDiscovery for NostrPublishAuthorRelayDiscovery { 944 fn fetch_author_write_relays<'a>( 945 &'a self, 946 pubkey: &'a str, 947 discovery_targets: Vec<ResolvedPublishRelay>, 948 connect_timeout_secs: u64, 949 ) -> PublishAuthorRelayDiscoveryFuture<'a> { 950 Box::pin(async move { 951 let Ok(public_key) = PublicKey::from_hex(pubkey) else { 952 return Ok(Vec::new()); 953 }; 954 let client = DaemonNostrClient::signerless(); 955 for target in discovery_targets { 956 if client.add_read_relay(target.url.as_str()).await.is_err() { 957 return Ok(Vec::new()); 958 } 959 } 960 let filter = Filter::new() 961 .author(public_key) 962 .kind(Kind::Custom(10_002)) 963 .limit(10); 964 let timeout = Duration::from_secs(connect_timeout_secs); 965 let Ok(events) = client.fetch_events(filter, timeout).await else { 966 return Ok(Vec::new()); 967 }; 968 let Some(event) = events.into_iter().max_by(|left, right| { 969 left.created_at 970 .as_secs() 971 .cmp(&right.created_at.as_secs()) 972 .then_with(|| left.id.to_hex().cmp(&right.id.to_hex())) 973 }) else { 974 return Ok(Vec::new()); 975 }; 976 Ok(author_write_relays_from_nip65_event(&event)) 977 }) 978 } 979 } 980 981 impl TransportPublishStore { 982 pub fn open(path: PathBuf) -> Result<Self, TransportPublishError> { 983 if let Some(parent) = path 984 .parent() 985 .filter(|parent| !parent.as_os_str().is_empty()) 986 { 987 std::fs::create_dir_all(parent)?; 988 } 989 let connection = connect_sqlite( 990 SqliteConnectOptions::new() 991 .filename(path) 992 .create_if_missing(true), 993 )?; 994 Self::from_connection(connection) 995 } 996 997 #[cfg(test)] 998 pub fn memory() -> Result<Self, TransportPublishError> { 999 Self::from_connection(connect_sqlite(SqliteConnectOptions::new().in_memory(true))?) 1000 } 1001 1002 fn from_connection(mut connection: SqliteConnection) -> Result<Self, TransportPublishError> { 1003 execute_sql(&mut connection, "PRAGMA foreign_keys = ON")?; 1004 match transport_publish_schema_state(&mut connection)? { 1005 TransportPublishSchemaState::Fresh => { 1006 execute_raw_sql(&mut connection, TRANSPORT_PUBLISH_SCHEMA_SQL)?; 1007 execute_sql( 1008 &mut connection, 1009 format!("PRAGMA user_version = {SCHEMA_VERSION}").as_str(), 1010 )?; 1011 validate_transport_publish_schema(&mut connection)?; 1012 } 1013 TransportPublishSchemaState::Existing => { 1014 validate_transport_publish_schema_version(&mut connection)?; 1015 validate_transport_publish_schema(&mut connection)?; 1016 } 1017 } 1018 recover_interrupted_publish_jobs(&mut connection)?; 1019 Ok(Self { 1020 inner: Arc::new(Mutex::new(connection)), 1021 }) 1022 } 1023 1024 pub fn create_principal( 1025 &self, 1026 input: PublishPrincipalInit, 1027 ) -> Result<PublishPrincipal, TransportPublishError> { 1028 validate_principal_init(&input)?; 1029 let principal_id = Uuid::new_v4().to_string(); 1030 let now = current_unix_secs(); 1031 let connection = self 1032 .inner 1033 .lock() 1034 .unwrap_or_else(std::sync::PoisonError::into_inner); 1035 let mut connection = connection; 1036 block_on_sqlite( 1037 sqlx::query( 1038 r#" 1039 INSERT INTO transport_publish_principals ( 1040 principal_id, 1041 label, 1042 token_hash, 1043 allowed_pubkeys_json, 1044 allowed_kinds_json, 1045 allowed_target_policies_json, 1046 allowed_explicit_transport_kinds_json, 1047 allowed_nostr_source_policies_json, 1048 allow_request_targets, 1049 job_visibility, 1050 expires_at_unix, 1051 revoked_at_unix, 1052 created_at_unix 1053 ) 1054 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL, ?12) 1055 "#, 1056 ) 1057 .bind(principal_id.as_str()) 1058 .bind(input.label.trim()) 1059 .bind(input.token_hash.as_str()) 1060 .bind(serde_json::to_string(&input.allowed_pubkeys)?) 1061 .bind(serde_json::to_string(&input.allowed_kinds)?) 1062 .bind(serde_json::to_string(&input.allowed_target_policies)?) 1063 .bind(serde_json::to_string( 1064 &input.allowed_explicit_transport_kinds, 1065 )?) 1066 .bind(serde_json::to_string(&input.allowed_nostr_source_policies)?) 1067 .bind(input.allow_request_targets) 1068 .bind(input.job_visibility.to_string()) 1069 .bind(input.expires_at_unix) 1070 .bind(now) 1071 .execute(&mut *connection), 1072 )?; 1073 drop(connection); 1074 self.principal_by_id(principal_id.as_str())?.ok_or_else(|| { 1075 TransportPublishError::InvalidScope("created principal missing".to_owned()) 1076 }) 1077 } 1078 1079 pub fn principal_for_token_hash( 1080 &self, 1081 token_hash: &str, 1082 ) -> Result<Option<PublishPrincipal>, TransportPublishError> { 1083 let now = current_unix_secs(); 1084 let mut connection = self 1085 .inner 1086 .lock() 1087 .unwrap_or_else(std::sync::PoisonError::into_inner); 1088 let principal = block_on_sqlite( 1089 sqlx::query( 1090 r#" 1091 SELECT 1092 principal_id, 1093 label, 1094 allowed_pubkeys_json, 1095 allowed_kinds_json, 1096 allowed_target_policies_json, 1097 allowed_explicit_transport_kinds_json, 1098 allowed_nostr_source_policies_json, 1099 allow_request_targets, 1100 job_visibility, 1101 expires_at_unix 1102 FROM transport_publish_principals 1103 WHERE token_hash = ?1 1104 AND revoked_at_unix IS NULL 1105 AND (expires_at_unix IS NULL OR expires_at_unix > ?2) 1106 "#, 1107 ) 1108 .bind(token_hash) 1109 .bind(now) 1110 .fetch_optional(&mut *connection), 1111 )? 1112 .map(|row| principal_from_row(&row)) 1113 .transpose()?; 1114 Ok(principal) 1115 } 1116 1117 pub fn principal_by_id( 1118 &self, 1119 principal_id: &str, 1120 ) -> Result<Option<PublishPrincipal>, TransportPublishError> { 1121 let mut connection = self 1122 .inner 1123 .lock() 1124 .unwrap_or_else(std::sync::PoisonError::into_inner); 1125 let principal = block_on_sqlite( 1126 sqlx::query( 1127 r#" 1128 SELECT 1129 principal_id, 1130 label, 1131 allowed_pubkeys_json, 1132 allowed_kinds_json, 1133 allowed_target_policies_json, 1134 allowed_explicit_transport_kinds_json, 1135 allowed_nostr_source_policies_json, 1136 allow_request_targets, 1137 job_visibility, 1138 expires_at_unix 1139 FROM transport_publish_principals 1140 WHERE principal_id = ?1 1141 "#, 1142 ) 1143 .bind(principal_id) 1144 .fetch_optional(&mut *connection), 1145 )? 1146 .map(|row| principal_from_row(&row)) 1147 .transpose()?; 1148 Ok(principal) 1149 } 1150 1151 pub fn record_publish_job( 1152 &self, 1153 insert: PublishJobInsert, 1154 ) -> Result<EventResponse, TransportPublishError> { 1155 if insert.effective_target_count != insert.target_snapshots.len() { 1156 return Err(TransportPublishError::InvalidScope( 1157 "publish job target snapshot count must match effective target count".to_owned(), 1158 )); 1159 } 1160 if let Some(idempotency_key) = insert.idempotency_key.as_deref() 1161 && let Some(existing) = 1162 self.job_for_principal_id_and_key(insert.principal_id.as_str(), idempotency_key)? 1163 { 1164 if existing.request_fingerprint != insert.request_fingerprint { 1165 return Err(TransportPublishError::IdempotencyConflict( 1166 idempotency_key.to_owned(), 1167 )); 1168 } 1169 return Ok(EventResponse { 1170 deduplicated: true, 1171 job: existing.view, 1172 }); 1173 } 1174 1175 let job_id = Uuid::new_v4().to_string(); 1176 let now = current_unix_millis(); 1177 let request_json = serde_json::to_string(&insert.request)?; 1178 let requested_target_count = storage_count_i64( 1179 insert.request.target_policy.request_target_count(), 1180 "requested_target_count", 1181 )?; 1182 let effective_target_count = 1183 storage_count_i64(insert.effective_target_count, "effective_target_count")?; 1184 let mut connection = self 1185 .inner 1186 .lock() 1187 .unwrap_or_else(std::sync::PoisonError::into_inner); 1188 execute_sql(&mut connection, "BEGIN")?; 1189 let transaction_result = (|| -> Result<(), TransportPublishError> { 1190 let insert_result = block_on_sqlite( 1191 sqlx::query( 1192 r#" 1193 INSERT INTO transport_publish_jobs ( 1194 job_id, 1195 principal_id, 1196 idempotency_key, 1197 request_fingerprint, 1198 status, 1199 event_id, 1200 event_pubkey, 1201 event_kind, 1202 target_policy_json, 1203 delivery_policy_json, 1204 requested_target_count, 1205 effective_target_count, 1206 request_json, 1207 requested_at_ms, 1208 updated_at_ms 1209 ) 1210 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15) 1211 "#, 1212 ) 1213 .bind(job_id.as_str()) 1214 .bind(insert.principal_id.as_str()) 1215 .bind(insert.idempotency_key.as_deref()) 1216 .bind(insert.request_fingerprint.as_str()) 1217 .bind(serde_json::to_string(&JobStatus::Publishing)?) 1218 .bind(insert.event.event_id.as_str()) 1219 .bind(insert.event.pubkey.as_str()) 1220 .bind(i64::from(insert.event.kind)) 1221 .bind(serde_json::to_string(&insert.request.target_policy)?) 1222 .bind(serde_json::to_string(&insert.request.delivery_policy)?) 1223 .bind(requested_target_count) 1224 .bind(effective_target_count) 1225 .bind(request_json.as_str()) 1226 .bind(now) 1227 .bind(now) 1228 .execute(&mut *connection), 1229 ); 1230 match insert_result { 1231 Ok(_) => {} 1232 Err(error) if is_sqlite_constraint_error(&error) => { 1233 return Err(TransportPublishError::IdempotencyConflict( 1234 "idempotency key conflicts with an existing publish job".to_owned(), 1235 )); 1236 } 1237 Err(error) => return Err(error), 1238 } 1239 insert_target_snapshots( 1240 &mut connection, 1241 job_id.as_str(), 1242 &insert.target_snapshots, 1243 now, 1244 )?; 1245 Ok(()) 1246 })(); 1247 match transaction_result { 1248 Ok(()) => { 1249 execute_sql(&mut connection, "COMMIT")?; 1250 } 1251 Err(error) => { 1252 let _ = execute_sql(&mut connection, "ROLLBACK"); 1253 return Err(error); 1254 } 1255 } 1256 drop(connection); 1257 let job = self.job_by_id(job_id.as_str())?; 1258 Ok(EventResponse { 1259 deduplicated: false, 1260 job, 1261 }) 1262 } 1263 1264 pub fn job_by_id_for_principal( 1265 &self, 1266 job_id: &str, 1267 principal: &PublishPrincipal, 1268 ) -> Result<Option<Job>, TransportPublishError> { 1269 let mut connection = self 1270 .inner 1271 .lock() 1272 .unwrap_or_else(std::sync::PoisonError::into_inner); 1273 let sql = job_select_sql("WHERE job_id = ?1"); 1274 let row = block_on_sqlite( 1275 sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 1276 .bind(job_id) 1277 .fetch_optional(&mut *connection), 1278 )? 1279 .map(|row| job_from_row(&row)) 1280 .transpose()?; 1281 drop(connection); 1282 let Some(job) = row else { 1283 return Ok(None); 1284 }; 1285 if !principal.can_read_job(job.principal_id.as_str()) { 1286 return Ok(None); 1287 } 1288 let job = self.finalize_job_row_for_egress(job)?; 1289 Ok(Some(job.view)) 1290 } 1291 1292 pub fn list_jobs_for_principal( 1293 &self, 1294 principal: &PublishPrincipal, 1295 limit: usize, 1296 ) -> Result<Vec<Job>, TransportPublishError> { 1297 let limit = i64::try_from(limit.clamp(1, 200)).unwrap_or(200); 1298 let mut connection = self 1299 .inner 1300 .lock() 1301 .unwrap_or_else(std::sync::PoisonError::into_inner); 1302 let sql = if principal.job_visibility == PublishJobVisibility::Admin { 1303 job_select_sql("ORDER BY requested_at_ms DESC, job_id DESC LIMIT ?1") 1304 } else { 1305 job_select_sql( 1306 "WHERE principal_id = ?1 ORDER BY requested_at_ms DESC, job_id DESC LIMIT ?2", 1307 ) 1308 }; 1309 let rows = if principal.job_visibility == PublishJobVisibility::Admin { 1310 block_on_sqlite( 1311 sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 1312 .bind(limit) 1313 .fetch_all(&mut *connection), 1314 )? 1315 } else { 1316 block_on_sqlite( 1317 sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 1318 .bind(principal.principal_id.as_str()) 1319 .bind(limit) 1320 .fetch_all(&mut *connection), 1321 )? 1322 }; 1323 let rows = rows 1324 .iter() 1325 .map(job_from_row) 1326 .collect::<Result<Vec<_>, _>>()?; 1327 drop(connection); 1328 1329 rows.into_iter() 1330 .map(|row| { 1331 let row = self.finalize_job_row_for_egress(row)?; 1332 Ok(row.view) 1333 }) 1334 .collect() 1335 } 1336 1337 fn job_for_principal_id_and_key( 1338 &self, 1339 principal_id: &str, 1340 idempotency_key: &str, 1341 ) -> Result<Option<PublishJobRow>, TransportPublishError> { 1342 let mut connection = self 1343 .inner 1344 .lock() 1345 .unwrap_or_else(std::sync::PoisonError::into_inner); 1346 let sql = job_select_sql("WHERE principal_id = ?1 AND idempotency_key = ?2"); 1347 let row = block_on_sqlite( 1348 sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 1349 .bind(principal_id) 1350 .bind(idempotency_key) 1351 .fetch_optional(&mut *connection), 1352 )? 1353 .map(|row| job_from_row(&row)) 1354 .transpose()?; 1355 drop(connection); 1356 let Some(job) = row else { 1357 return Ok(None); 1358 }; 1359 let job = self.finalize_job_row_for_egress(job)?; 1360 Ok(Some(job)) 1361 } 1362 1363 pub fn job_by_id(&self, job_id: &str) -> Result<Job, TransportPublishError> { 1364 let mut connection = self 1365 .inner 1366 .lock() 1367 .unwrap_or_else(std::sync::PoisonError::into_inner); 1368 let sql = job_select_sql("WHERE job_id = ?1"); 1369 let row = block_on_sqlite( 1370 sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 1371 .bind(job_id) 1372 .fetch_optional(&mut *connection), 1373 )? 1374 .map(|row| job_from_row(&row)) 1375 .transpose()?; 1376 drop(connection); 1377 let Some(job) = row else { 1378 return Err(TransportPublishError::InvalidScope( 1379 "unknown publish job".to_owned(), 1380 )); 1381 }; 1382 let job = self.finalize_job_row_for_egress(job)?; 1383 Ok(job.view) 1384 } 1385 1386 pub fn complete_publish_job( 1387 &self, 1388 job_id: &str, 1389 status: JobStatus, 1390 outcomes: Vec<TargetOutcome>, 1391 last_error: Option<String>, 1392 ) -> Result<(), TransportPublishError> { 1393 let now = current_unix_millis(); 1394 let target_count = storage_count_i64(outcomes.len(), "effective_target_count")?; 1395 let mut connection = self 1396 .inner 1397 .lock() 1398 .unwrap_or_else(std::sync::PoisonError::into_inner); 1399 block_on_sqlite( 1400 sqlx::query( 1401 r#" 1402 UPDATE transport_publish_jobs 1403 SET status = ?2, 1404 updated_at_ms = ?3, 1405 completed_at_ms = ?4, 1406 last_error = ?5, 1407 effective_target_count = ?6 1408 WHERE job_id = ?1 1409 "#, 1410 ) 1411 .bind(job_id) 1412 .bind(serde_json::to_string(&status)?) 1413 .bind(now) 1414 .bind(now) 1415 .bind(last_error.as_deref()) 1416 .bind(target_count) 1417 .execute(&mut *connection), 1418 )?; 1419 replace_target_outcomes(&mut connection, job_id, &outcomes, now)?; 1420 Ok(()) 1421 } 1422 1423 pub fn cached_author_write_relays( 1424 &self, 1425 pubkey: &str, 1426 ) -> Result<Vec<String>, TransportPublishError> { 1427 let mut connection = self 1428 .inner 1429 .lock() 1430 .unwrap_or_else(std::sync::PoisonError::into_inner); 1431 let relays_json = block_on_sqlite( 1432 sqlx::query_scalar::<_, String>( 1433 "SELECT relays_json FROM transport_publish_nostr_author_cache WHERE pubkey = ?1", 1434 ) 1435 .bind(pubkey) 1436 .fetch_optional(&mut *connection), 1437 )?; 1438 relays_json 1439 .map(|value| serde_json::from_str(value.as_str()).map_err(TransportPublishError::from)) 1440 .unwrap_or_else(|| Ok(Vec::new())) 1441 } 1442 1443 pub fn cache_author_write_relays( 1444 &self, 1445 pubkey: &str, 1446 relays: &[String], 1447 ) -> Result<(), TransportPublishError> { 1448 let now = current_unix_millis(); 1449 let mut connection = self 1450 .inner 1451 .lock() 1452 .unwrap_or_else(std::sync::PoisonError::into_inner); 1453 block_on_sqlite( 1454 sqlx::query( 1455 r#" 1456 INSERT INTO transport_publish_nostr_author_cache (pubkey, relays_json, updated_at_ms) 1457 VALUES (?1, ?2, ?3) 1458 ON CONFLICT(pubkey) DO UPDATE SET 1459 relays_json = excluded.relays_json, 1460 updated_at_ms = excluded.updated_at_ms 1461 "#, 1462 ) 1463 .bind(pubkey) 1464 .bind(serde_json::to_string(relays)?) 1465 .bind(now) 1466 .execute(&mut *connection), 1467 )?; 1468 Ok(()) 1469 } 1470 1471 fn target_outcomes(&self, job_id: &str) -> Result<Vec<TargetOutcome>, TransportPublishError> { 1472 let mut connection = self 1473 .inner 1474 .lock() 1475 .unwrap_or_else(std::sync::PoisonError::into_inner); 1476 let rows = block_on_sqlite( 1477 sqlx::query( 1478 r#" 1479 SELECT transport_kind, endpoint_uri, target_scope, target_label, source, attempted, outcome_kind, message, latency_ms 1480 FROM transport_publish_target_results 1481 WHERE job_id = ?1 1482 ORDER BY transport_kind, endpoint_uri, target_scope 1483 "#, 1484 ) 1485 .bind(job_id) 1486 .fetch_all(&mut *connection), 1487 )?; 1488 let outcomes = rows 1489 .iter() 1490 .map(target_outcome_from_row) 1491 .collect::<Result<Vec<_>, _>>()?; 1492 Ok(outcomes) 1493 } 1494 1495 fn finalize_job_row_for_egress( 1496 &self, 1497 mut job: PublishJobRow, 1498 ) -> Result<PublishJobRow, TransportPublishError> { 1499 job.view.targets = self.target_outcomes(job.view.job_id.as_str())?; 1500 finalize_job_view(&mut job.view); 1501 job.view 1502 .validate() 1503 .map_err(|error| TransportPublishError::InvalidPublishJobState(error.to_string()))?; 1504 Ok(job) 1505 } 1506 } 1507 1508 struct PublishJobRow { 1509 principal_id: String, 1510 request_fingerprint: String, 1511 view: Job, 1512 } 1513 1514 enum TransportPublishSchemaState { 1515 Fresh, 1516 Existing, 1517 } 1518 1519 fn connect_sqlite( 1520 options: SqliteConnectOptions, 1521 ) -> Result<SqliteConnection, TransportPublishError> { 1522 block_on_sqlite(SqliteConnection::connect_with(&options)) 1523 } 1524 1525 fn block_on_sqlite<T>( 1526 future: impl Future<Output = Result<T, sqlx::Error>>, 1527 ) -> Result<T, TransportPublishError> { 1528 Ok(futures_executor::block_on(future)?) 1529 } 1530 1531 fn execute_sql(connection: &mut SqliteConnection, sql: &str) -> Result<u64, TransportPublishError> { 1532 Ok(block_on_sqlite(sqlx::query(sqlx::AssertSqlSafe(sql)).execute(connection))?.rows_affected()) 1533 } 1534 1535 fn execute_raw_sql( 1536 connection: &mut SqliteConnection, 1537 sql: &str, 1538 ) -> Result<(), TransportPublishError> { 1539 block_on_sqlite(sqlx::raw_sql(sqlx::AssertSqlSafe(sql)).execute(connection))?; 1540 Ok(()) 1541 } 1542 1543 fn fetch_all_sql( 1544 connection: &mut SqliteConnection, 1545 sql: &str, 1546 ) -> Result<Vec<SqliteRow>, TransportPublishError> { 1547 block_on_sqlite(sqlx::query(sqlx::AssertSqlSafe(sql)).fetch_all(connection)) 1548 } 1549 1550 fn is_sqlite_constraint_error(error: &TransportPublishError) -> bool { 1551 match error { 1552 TransportPublishError::Sqlite(sqlx::Error::Database(error)) => { 1553 error.is_unique_violation() 1554 || error 1555 .code() 1556 .as_deref() 1557 .is_some_and(|code| matches!(code, "1555" | "2067" | "19")) 1558 } 1559 _ => false, 1560 } 1561 } 1562 1563 fn transport_publish_schema_state( 1564 connection: &mut SqliteConnection, 1565 ) -> Result<TransportPublishSchemaState, TransportPublishError> { 1566 let rows = block_on_sqlite( 1567 sqlx::query( 1568 "SELECT name FROM sqlite_schema WHERE type = 'table' AND name LIKE 'transport_publish_%'", 1569 ) 1570 .fetch_all(&mut *connection), 1571 )?; 1572 let names = rows 1573 .iter() 1574 .map(|row| row.try_get::<String, _>(0)) 1575 .collect::<Result<BTreeSet<_>, _>>()?; 1576 if names.is_empty() { 1577 let version = transport_publish_schema_version(connection)?; 1578 if version == 0 { 1579 Ok(TransportPublishSchemaState::Fresh) 1580 } else { 1581 Err(TransportPublishError::Schema { 1582 table: "transport_publish_schema", 1583 detail: format!( 1584 "fresh schema initialization requires user_version 0, got {version}" 1585 ), 1586 }) 1587 } 1588 } else { 1589 Ok(TransportPublishSchemaState::Existing) 1590 } 1591 } 1592 1593 fn transport_publish_schema_version( 1594 connection: &mut SqliteConnection, 1595 ) -> Result<i64, TransportPublishError> { 1596 block_on_sqlite(sqlx::query_scalar::<_, i64>("PRAGMA user_version").fetch_one(connection)) 1597 } 1598 1599 fn validate_transport_publish_schema_version( 1600 connection: &mut SqliteConnection, 1601 ) -> Result<(), TransportPublishError> { 1602 let version = transport_publish_schema_version(connection)?; 1603 if version == SCHEMA_VERSION { 1604 Ok(()) 1605 } else { 1606 Err(TransportPublishError::Schema { 1607 table: "transport_publish_schema", 1608 detail: format!("user_version must be {SCHEMA_VERSION}, got {version}"), 1609 }) 1610 } 1611 } 1612 1613 fn validate_transport_publish_schema( 1614 connection: &mut SqliteConnection, 1615 ) -> Result<(), TransportPublishError> { 1616 validate_foreign_keys_enabled(connection)?; 1617 validate_table_columns( 1618 connection, 1619 "transport_publish_principals", 1620 &[ 1621 RequiredColumn::text_not_null("principal_id"), 1622 RequiredColumn::text_not_null("label"), 1623 RequiredColumn::text_not_null("token_hash"), 1624 RequiredColumn::text_not_null("allowed_pubkeys_json"), 1625 RequiredColumn::text_not_null("allowed_kinds_json"), 1626 RequiredColumn::text_not_null("allowed_target_policies_json"), 1627 RequiredColumn::text_not_null("allowed_explicit_transport_kinds_json"), 1628 RequiredColumn::text_not_null("allowed_nostr_source_policies_json"), 1629 RequiredColumn::integer_not_null("allow_request_targets"), 1630 RequiredColumn::text_not_null("job_visibility"), 1631 RequiredColumn::integer_nullable("expires_at_unix"), 1632 RequiredColumn::integer_nullable("revoked_at_unix"), 1633 RequiredColumn::integer_not_null("created_at_unix"), 1634 ], 1635 )?; 1636 validate_primary_key( 1637 connection, 1638 "transport_publish_principals", 1639 &["principal_id"], 1640 )?; 1641 validate_unique_index( 1642 connection, 1643 "transport_publish_principals", 1644 &["token_hash"], 1645 None, 1646 )?; 1647 validate_table_columns( 1648 connection, 1649 "transport_publish_jobs", 1650 &[ 1651 RequiredColumn::text_not_null("job_id"), 1652 RequiredColumn::text_not_null("principal_id"), 1653 RequiredColumn::text_nullable("idempotency_key"), 1654 RequiredColumn::text_not_null("request_fingerprint"), 1655 RequiredColumn::text_not_null("status"), 1656 RequiredColumn::text_not_null("event_id"), 1657 RequiredColumn::text_not_null("event_pubkey"), 1658 RequiredColumn::integer_not_null("event_kind"), 1659 RequiredColumn::text_not_null("target_policy_json"), 1660 RequiredColumn::text_not_null("delivery_policy_json"), 1661 RequiredColumn::integer_not_null("requested_target_count"), 1662 RequiredColumn::integer_not_null("effective_target_count"), 1663 RequiredColumn::text_not_null("request_json"), 1664 RequiredColumn::integer_not_null("requested_at_ms"), 1665 RequiredColumn::integer_not_null("updated_at_ms"), 1666 RequiredColumn::integer_nullable("completed_at_ms"), 1667 RequiredColumn::text_nullable("last_error"), 1668 ], 1669 )?; 1670 validate_primary_key(connection, "transport_publish_jobs", &["job_id"])?; 1671 validate_foreign_key( 1672 connection, 1673 "transport_publish_jobs", 1674 &["principal_id"], 1675 "transport_publish_principals", 1676 &["principal_id"], 1677 )?; 1678 validate_unique_index( 1679 connection, 1680 "transport_publish_jobs", 1681 &["principal_id", "idempotency_key"], 1682 Some("WHERE idempotency_key IS NOT NULL"), 1683 )?; 1684 validate_table_columns( 1685 connection, 1686 "transport_publish_target_results", 1687 &[ 1688 RequiredColumn::text_not_null("job_id"), 1689 RequiredColumn::text_not_null("transport_kind"), 1690 RequiredColumn::text_not_null("endpoint_uri"), 1691 RequiredColumn::text_not_null("target_scope"), 1692 RequiredColumn::text_nullable("target_label"), 1693 RequiredColumn::text_not_null("source"), 1694 RequiredColumn::integer_not_null("attempted"), 1695 RequiredColumn::text_not_null("outcome_kind"), 1696 RequiredColumn::text_nullable("message"), 1697 RequiredColumn::integer_nullable("latency_ms"), 1698 RequiredColumn::integer_not_null("updated_at_ms"), 1699 ], 1700 )?; 1701 validate_primary_key( 1702 connection, 1703 "transport_publish_target_results", 1704 &["job_id", "transport_kind", "endpoint_uri", "target_scope"], 1705 )?; 1706 validate_foreign_key( 1707 connection, 1708 "transport_publish_target_results", 1709 &["job_id"], 1710 "transport_publish_jobs", 1711 &["job_id"], 1712 )?; 1713 validate_table_columns( 1714 connection, 1715 "transport_publish_target_snapshots", 1716 &[ 1717 RequiredColumn::text_not_null("job_id"), 1718 RequiredColumn::integer_not_null("target_index"), 1719 RequiredColumn::text_not_null("transport_kind"), 1720 RequiredColumn::text_not_null("endpoint_uri"), 1721 RequiredColumn::text_not_null("target_scope"), 1722 RequiredColumn::text_nullable("target_label"), 1723 RequiredColumn::text_not_null("source"), 1724 RequiredColumn::integer_not_null("attempted"), 1725 RequiredColumn::text_not_null("outcome_kind"), 1726 RequiredColumn::text_nullable("message"), 1727 RequiredColumn::integer_nullable("latency_ms"), 1728 RequiredColumn::integer_not_null("created_at_ms"), 1729 ], 1730 )?; 1731 validate_primary_key( 1732 connection, 1733 "transport_publish_target_snapshots", 1734 &["job_id", "target_index"], 1735 )?; 1736 validate_foreign_key( 1737 connection, 1738 "transport_publish_target_snapshots", 1739 &["job_id"], 1740 "transport_publish_jobs", 1741 &["job_id"], 1742 )?; 1743 validate_table_columns( 1744 connection, 1745 "transport_publish_nostr_author_cache", 1746 &[ 1747 RequiredColumn::text_not_null("pubkey"), 1748 RequiredColumn::text_not_null("relays_json"), 1749 RequiredColumn::integer_not_null("updated_at_ms"), 1750 ], 1751 )?; 1752 validate_primary_key( 1753 connection, 1754 "transport_publish_nostr_author_cache", 1755 &["pubkey"], 1756 )?; 1757 for table in TRANSPORT_PUBLISH_TABLES { 1758 validate_table_present(connection, table)?; 1759 } 1760 Ok(()) 1761 } 1762 1763 #[derive(Clone, Copy)] 1764 struct RequiredColumn { 1765 name: &'static str, 1766 column_type: &'static str, 1767 not_null: bool, 1768 } 1769 1770 impl RequiredColumn { 1771 const fn text_not_null(name: &'static str) -> Self { 1772 Self { 1773 name, 1774 column_type: "TEXT", 1775 not_null: true, 1776 } 1777 } 1778 1779 const fn text_nullable(name: &'static str) -> Self { 1780 Self { 1781 name, 1782 column_type: "TEXT", 1783 not_null: false, 1784 } 1785 } 1786 1787 const fn integer_not_null(name: &'static str) -> Self { 1788 Self { 1789 name, 1790 column_type: "INTEGER", 1791 not_null: true, 1792 } 1793 } 1794 1795 const fn integer_nullable(name: &'static str) -> Self { 1796 Self { 1797 name, 1798 column_type: "INTEGER", 1799 not_null: false, 1800 } 1801 } 1802 } 1803 1804 struct TableColumnInfo { 1805 column_type: String, 1806 not_null: bool, 1807 primary_key_position: i64, 1808 } 1809 1810 struct ForeignKeyEntry { 1811 target_table: String, 1812 from_columns: Vec<String>, 1813 to_columns: Vec<String>, 1814 } 1815 1816 fn validate_foreign_keys_enabled( 1817 connection: &mut SqliteConnection, 1818 ) -> Result<(), TransportPublishError> { 1819 let enabled = 1820 block_on_sqlite(sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys").fetch_one(connection))?; 1821 if enabled == 1 { 1822 Ok(()) 1823 } else { 1824 Err(TransportPublishError::Schema { 1825 table: "transport_publish_schema", 1826 detail: "foreign key enforcement must be enabled".to_owned(), 1827 }) 1828 } 1829 } 1830 1831 fn validate_table_present( 1832 connection: &mut SqliteConnection, 1833 table: &'static str, 1834 ) -> Result<(), TransportPublishError> { 1835 let exists = block_on_sqlite( 1836 sqlx::query_scalar::<_, i64>( 1837 "SELECT 1 FROM sqlite_schema WHERE type = 'table' AND name = ?1", 1838 ) 1839 .bind(table) 1840 .fetch_optional(connection), 1841 )? 1842 .is_some(); 1843 if exists { 1844 Ok(()) 1845 } else { 1846 Err(TransportPublishError::Schema { 1847 table, 1848 detail: "table is missing".to_owned(), 1849 }) 1850 } 1851 } 1852 1853 fn table_columns( 1854 connection: &mut SqliteConnection, 1855 table: &'static str, 1856 ) -> Result<BTreeMap<String, TableColumnInfo>, TransportPublishError> { 1857 let sql = format!("PRAGMA table_info({table})"); 1858 let rows = fetch_all_sql(connection, sql.as_str())?; 1859 Ok(rows 1860 .iter() 1861 .map(|row| { 1862 Ok(( 1863 row.try_get::<String, _>(1)?, 1864 TableColumnInfo { 1865 column_type: row.try_get::<String, _>(2)?.to_ascii_uppercase(), 1866 not_null: row.try_get::<i64, _>(3)? != 0, 1867 primary_key_position: row.try_get::<i64, _>(5)?, 1868 }, 1869 )) 1870 }) 1871 .collect::<Result<BTreeMap<_, _>, sqlx::Error>>()?) 1872 } 1873 1874 fn validate_primary_key( 1875 connection: &mut SqliteConnection, 1876 table: &'static str, 1877 expected_columns: &[&'static str], 1878 ) -> Result<(), TransportPublishError> { 1879 let columns = table_columns(connection, table)?; 1880 let mut primary_key = columns 1881 .iter() 1882 .filter_map(|(name, column)| { 1883 (column.primary_key_position > 0) 1884 .then_some((column.primary_key_position, name.as_str())) 1885 }) 1886 .collect::<Vec<_>>(); 1887 primary_key.sort_by_key(|(position, _)| *position); 1888 let actual_columns = primary_key 1889 .into_iter() 1890 .map(|(_, name)| name) 1891 .collect::<Vec<_>>(); 1892 if actual_columns == expected_columns { 1893 Ok(()) 1894 } else { 1895 Err(TransportPublishError::Schema { 1896 table, 1897 detail: format!("primary key must be ({})", expected_columns.join(", ")), 1898 }) 1899 } 1900 } 1901 1902 fn validate_unique_index( 1903 connection: &mut SqliteConnection, 1904 table: &'static str, 1905 expected_columns: &[&'static str], 1906 partial_where: Option<&'static str>, 1907 ) -> Result<(), TransportPublishError> { 1908 let sql = format!("PRAGMA index_list({table})"); 1909 let rows = fetch_all_sql(connection, sql.as_str())?; 1910 let indexes = rows 1911 .iter() 1912 .map(|row| { 1913 Ok(( 1914 row.try_get::<String, _>(1)?, 1915 row.try_get::<i64, _>(2)? != 0, 1916 row.try_get::<i64, _>(4)? != 0, 1917 )) 1918 }) 1919 .collect::<Result<Vec<_>, sqlx::Error>>()?; 1920 for (index_name, unique, partial) in indexes { 1921 if !unique { 1922 continue; 1923 } 1924 if partial_where.is_some() != partial { 1925 continue; 1926 } 1927 let index_columns = index_columns(connection, index_name.as_str())?; 1928 if !columns_match(index_columns.as_slice(), expected_columns) { 1929 continue; 1930 } 1931 let where_matches = match partial_where { 1932 Some(required_where) => { 1933 index_sql_contains_where(connection, table, index_name.as_str(), required_where)? 1934 } 1935 None => true, 1936 }; 1937 if where_matches { 1938 return Ok(()); 1939 } 1940 } 1941 Err(TransportPublishError::Schema { 1942 table, 1943 detail: format!( 1944 "missing required unique index on ({})", 1945 expected_columns.join(", ") 1946 ), 1947 }) 1948 } 1949 1950 fn index_columns( 1951 connection: &mut SqliteConnection, 1952 index_name: &str, 1953 ) -> Result<Vec<String>, TransportPublishError> { 1954 let sql = format!("PRAGMA index_info({index_name})"); 1955 let rows = fetch_all_sql(connection, sql.as_str())?; 1956 let mut columns = rows 1957 .iter() 1958 .map(|row| Ok((row.try_get::<i64, _>(0)?, row.try_get::<String, _>(2)?))) 1959 .collect::<Result<Vec<_>, sqlx::Error>>()?; 1960 columns.sort_by_key(|(position, _)| *position); 1961 Ok(columns.into_iter().map(|(_, name)| name).collect()) 1962 } 1963 1964 fn index_sql_contains_where( 1965 connection: &mut SqliteConnection, 1966 table: &'static str, 1967 index_name: &str, 1968 required_where: &str, 1969 ) -> Result<bool, TransportPublishError> { 1970 let sql = block_on_sqlite( 1971 sqlx::query_scalar::<_, Option<String>>( 1972 "SELECT sql FROM sqlite_schema WHERE type = 'index' AND tbl_name = ?1 AND name = ?2", 1973 ) 1974 .bind(table) 1975 .bind(index_name) 1976 .fetch_optional(connection), 1977 )? 1978 .flatten() 1979 .unwrap_or_default(); 1980 Ok(normalized_sql(sql.as_str()).contains(normalized_sql(required_where).as_str())) 1981 } 1982 1983 fn normalized_sql(sql: &str) -> String { 1984 sql.split_whitespace() 1985 .collect::<Vec<_>>() 1986 .join(" ") 1987 .to_ascii_uppercase() 1988 } 1989 1990 fn validate_foreign_key( 1991 connection: &mut SqliteConnection, 1992 table: &'static str, 1993 expected_from_columns: &[&'static str], 1994 expected_target_table: &'static str, 1995 expected_to_columns: &[&'static str], 1996 ) -> Result<(), TransportPublishError> { 1997 for foreign_key in foreign_keys(connection, table)? { 1998 if foreign_key.target_table == expected_target_table 1999 && columns_match(foreign_key.from_columns.as_slice(), expected_from_columns) 2000 && columns_match(foreign_key.to_columns.as_slice(), expected_to_columns) 2001 { 2002 return Ok(()); 2003 } 2004 } 2005 Err(TransportPublishError::Schema { 2006 table, 2007 detail: format!( 2008 "missing required foreign key ({}) references {}({})", 2009 expected_from_columns.join(", "), 2010 expected_target_table, 2011 expected_to_columns.join(", ") 2012 ), 2013 }) 2014 } 2015 2016 fn foreign_keys( 2017 connection: &mut SqliteConnection, 2018 table: &'static str, 2019 ) -> Result<Vec<ForeignKeyEntry>, TransportPublishError> { 2020 let mut groups = BTreeMap::<i64, (String, Vec<(i64, String, String)>)>::new(); 2021 let sql = format!("PRAGMA foreign_key_list({table})"); 2022 let rows = fetch_all_sql(connection, sql.as_str())?; 2023 let rows = rows 2024 .iter() 2025 .map(|row| { 2026 Ok(( 2027 row.try_get::<i64, _>(0)?, 2028 row.try_get::<i64, _>(1)?, 2029 row.try_get::<String, _>(2)?, 2030 row.try_get::<String, _>(3)?, 2031 row.try_get::<String, _>(4)?, 2032 )) 2033 }) 2034 .collect::<Result<Vec<_>, sqlx::Error>>()?; 2035 for (id, seq, target_table, from_column, to_column) in rows { 2036 let entry = groups 2037 .entry(id) 2038 .or_insert_with(|| (target_table, Vec::new())); 2039 entry.1.push((seq, from_column, to_column)); 2040 } 2041 let mut foreign_keys = Vec::new(); 2042 for (_, (target_table, mut columns)) in groups { 2043 columns.sort_by_key(|(seq, _, _)| *seq); 2044 foreign_keys.push(ForeignKeyEntry { 2045 target_table, 2046 from_columns: columns 2047 .iter() 2048 .map(|(_, from_column, _)| from_column.clone()) 2049 .collect(), 2050 to_columns: columns 2051 .into_iter() 2052 .map(|(_, _, to_column)| to_column) 2053 .collect(), 2054 }); 2055 } 2056 Ok(foreign_keys) 2057 } 2058 2059 fn columns_match(actual: &[String], expected: &[&'static str]) -> bool { 2060 actual.len() == expected.len() 2061 && actual 2062 .iter() 2063 .map(String::as_str) 2064 .zip(expected.iter().copied()) 2065 .all(|(actual, expected)| actual == expected) 2066 } 2067 2068 fn validate_table_columns( 2069 connection: &mut SqliteConnection, 2070 table: &'static str, 2071 required_columns: &[RequiredColumn], 2072 ) -> Result<(), TransportPublishError> { 2073 let columns = table_columns(connection, table)?; 2074 if columns.is_empty() { 2075 return Err(TransportPublishError::Schema { 2076 table, 2077 detail: "table is missing".to_owned(), 2078 }); 2079 } 2080 for required in required_columns { 2081 let Some(column) = columns.get(required.name) else { 2082 return Err(TransportPublishError::Schema { 2083 table, 2084 detail: format!("missing required column `{}`", required.name), 2085 }); 2086 }; 2087 if column.column_type != required.column_type { 2088 return Err(TransportPublishError::Schema { 2089 table, 2090 detail: format!( 2091 "column `{}` must have type {}, got {}", 2092 required.name, required.column_type, column.column_type 2093 ), 2094 }); 2095 } 2096 if column.not_null != required.not_null { 2097 return Err(TransportPublishError::Schema { 2098 table, 2099 detail: format!( 2100 "column `{}` must be {}", 2101 required.name, 2102 if required.not_null { 2103 "NOT NULL" 2104 } else { 2105 "nullable" 2106 } 2107 ), 2108 }); 2109 } 2110 } 2111 if columns.len() != required_columns.len() { 2112 let required = required_columns 2113 .iter() 2114 .map(|column| column.name) 2115 .collect::<BTreeSet<_>>(); 2116 let extras = columns 2117 .keys() 2118 .filter(|column| !required.contains(column.as_str())) 2119 .cloned() 2120 .collect::<Vec<_>>(); 2121 return Err(TransportPublishError::Schema { 2122 table, 2123 detail: format!("unexpected columns: {}", extras.join(", ")), 2124 }); 2125 } 2126 Ok(()) 2127 } 2128 2129 fn recover_interrupted_publish_jobs( 2130 connection: &mut SqliteConnection, 2131 ) -> Result<(), TransportPublishError> { 2132 let now = current_unix_millis(); 2133 let publishing = serde_json::to_string(&JobStatus::Publishing)?; 2134 let sql = job_select_sql("WHERE status = ?1"); 2135 let rows = block_on_sqlite( 2136 sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 2137 .bind(publishing.as_str()) 2138 .fetch_all(&mut *connection), 2139 )?; 2140 let rows = rows 2141 .iter() 2142 .map(job_from_row) 2143 .collect::<Result<Vec<_>, _>>()?; 2144 for row in rows { 2145 let job_id = row.view.job_id.clone(); 2146 let snapshots = target_snapshot_outcomes(connection, job_id.as_str())?; 2147 if snapshots.is_empty() { 2148 block_on_sqlite( 2149 sqlx::query( 2150 r#" 2151 UPDATE transport_publish_jobs 2152 SET status = ?2, 2153 updated_at_ms = ?3, 2154 completed_at_ms = ?4, 2155 last_error = ?5, 2156 effective_target_count = 0 2157 WHERE job_id = ?1 2158 "#, 2159 ) 2160 .bind(job_id.as_str()) 2161 .bind(serde_json::to_string(&JobStatus::Rejected)?) 2162 .bind(now) 2163 .bind(now) 2164 .bind("publish_attempt_interrupted_missing_target_snapshot") 2165 .execute(&mut *connection), 2166 )?; 2167 replace_target_outcomes(connection, job_id.as_str(), &[], now)?; 2168 continue; 2169 } 2170 let status = delivery_status(&row.view.delivery_policy, snapshots.len(), &snapshots); 2171 let effective_target_count = storage_count_i64(snapshots.len(), "effective_target_count")?; 2172 let last_error = if status == JobStatus::DeliveryUnsatisfiedRetryable { 2173 Some("publish_attempt_interrupted".to_owned()) 2174 } else { 2175 last_error_for_status(status).map(str::to_owned) 2176 }; 2177 block_on_sqlite( 2178 sqlx::query( 2179 r#" 2180 UPDATE transport_publish_jobs 2181 SET status = ?2, 2182 updated_at_ms = ?3, 2183 completed_at_ms = ?4, 2184 last_error = ?5, 2185 effective_target_count = ?6 2186 WHERE job_id = ?1 2187 "#, 2188 ) 2189 .bind(job_id.as_str()) 2190 .bind(serde_json::to_string(&status)?) 2191 .bind(now) 2192 .bind(now) 2193 .bind(last_error.as_deref()) 2194 .bind(effective_target_count) 2195 .execute(&mut *connection), 2196 )?; 2197 replace_target_outcomes(connection, job_id.as_str(), &snapshots, now)?; 2198 } 2199 Ok(()) 2200 } 2201 2202 fn insert_target_snapshots( 2203 connection: &mut SqliteConnection, 2204 job_id: &str, 2205 outcomes: &[TargetOutcome], 2206 now: i64, 2207 ) -> Result<(), TransportPublishError> { 2208 for (target_index, outcome) in outcomes.iter().enumerate() { 2209 block_on_sqlite( 2210 sqlx::query( 2211 r#" 2212 INSERT INTO transport_publish_target_snapshots ( 2213 job_id, 2214 target_index, 2215 transport_kind, 2216 endpoint_uri, 2217 target_scope, 2218 target_label, 2219 source, 2220 attempted, 2221 outcome_kind, 2222 message, 2223 latency_ms, 2224 created_at_ms 2225 ) 2226 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12) 2227 "#, 2228 ) 2229 .bind(job_id) 2230 .bind(i64::try_from(target_index).unwrap_or(i64::MAX)) 2231 .bind(outcome.transport_kind.as_str()) 2232 .bind(outcome.endpoint_uri.as_str()) 2233 .bind(storage_target_scope(outcome.target_scope.as_deref())) 2234 .bind(outcome.target_label.as_deref()) 2235 .bind(serde_json::to_string(&outcome.source)?) 2236 .bind(outcome.attempted) 2237 .bind(serde_json::to_string(&outcome.outcome_kind)?) 2238 .bind(outcome.message.as_deref()) 2239 .bind( 2240 outcome 2241 .latency_ms 2242 .and_then(|value| i64::try_from(value).ok()), 2243 ) 2244 .bind(now) 2245 .execute(&mut *connection), 2246 )?; 2247 } 2248 Ok(()) 2249 } 2250 2251 fn replace_target_outcomes( 2252 connection: &mut SqliteConnection, 2253 job_id: &str, 2254 outcomes: &[TargetOutcome], 2255 now: i64, 2256 ) -> Result<(), TransportPublishError> { 2257 block_on_sqlite( 2258 sqlx::query("DELETE FROM transport_publish_target_results WHERE job_id = ?1") 2259 .bind(job_id) 2260 .execute(&mut *connection), 2261 )?; 2262 for outcome in outcomes { 2263 block_on_sqlite( 2264 sqlx::query( 2265 r#" 2266 INSERT OR REPLACE INTO transport_publish_target_results ( 2267 job_id, 2268 transport_kind, 2269 endpoint_uri, 2270 target_scope, 2271 target_label, 2272 source, 2273 attempted, 2274 outcome_kind, 2275 message, 2276 latency_ms, 2277 updated_at_ms 2278 ) 2279 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11) 2280 "#, 2281 ) 2282 .bind(job_id) 2283 .bind(outcome.transport_kind.as_str()) 2284 .bind(outcome.endpoint_uri.as_str()) 2285 .bind(storage_target_scope(outcome.target_scope.as_deref())) 2286 .bind(outcome.target_label.as_deref()) 2287 .bind(serde_json::to_string(&outcome.source)?) 2288 .bind(outcome.attempted) 2289 .bind(serde_json::to_string(&outcome.outcome_kind)?) 2290 .bind(outcome.message.as_deref()) 2291 .bind( 2292 outcome 2293 .latency_ms 2294 .and_then(|value| i64::try_from(value).ok()), 2295 ) 2296 .bind(now) 2297 .execute(&mut *connection), 2298 )?; 2299 } 2300 Ok(()) 2301 } 2302 2303 fn target_snapshot_outcomes( 2304 connection: &mut SqliteConnection, 2305 job_id: &str, 2306 ) -> Result<Vec<TargetOutcome>, TransportPublishError> { 2307 let rows = block_on_sqlite( 2308 sqlx::query( 2309 r#" 2310 SELECT transport_kind, endpoint_uri, target_scope, target_label, source, attempted, outcome_kind, message, latency_ms 2311 FROM transport_publish_target_snapshots 2312 WHERE job_id = ?1 2313 ORDER BY target_index 2314 "#, 2315 ) 2316 .bind(job_id) 2317 .fetch_all(connection), 2318 )?; 2319 let outcomes = rows 2320 .iter() 2321 .map(target_outcome_from_row) 2322 .collect::<Result<Vec<_>, _>>()?; 2323 Ok(outcomes) 2324 } 2325 2326 fn job_select_sql(tail: &str) -> String { 2327 format!( 2328 r#" 2329 SELECT 2330 job_id, 2331 principal_id, 2332 request_fingerprint, 2333 status, 2334 event_id, 2335 event_pubkey, 2336 event_kind, 2337 target_policy_json, 2338 delivery_policy_json, 2339 effective_target_count, 2340 requested_at_ms, 2341 completed_at_ms, 2342 last_error 2343 FROM transport_publish_jobs 2344 {tail} 2345 "# 2346 ) 2347 } 2348 2349 fn principal_from_row(row: &SqliteRow) -> Result<PublishPrincipal, TransportPublishError> { 2350 let visibility: String = row.try_get(8)?; 2351 Ok(PublishPrincipal { 2352 principal_id: row.try_get(0)?, 2353 label: row.try_get(1)?, 2354 allowed_pubkeys: json_column(row, 2, "allowed_pubkeys_json")?, 2355 allowed_kinds: json_column(row, 3, "allowed_kinds_json")?, 2356 allowed_target_policies: json_column(row, 4, "allowed_target_policies_json")?, 2357 allowed_explicit_transport_kinds: json_column( 2358 row, 2359 5, 2360 "allowed_explicit_transport_kinds_json", 2361 )?, 2362 allowed_nostr_source_policies: json_column(row, 6, "allowed_nostr_source_policies_json")?, 2363 allow_request_targets: row.try_get(7)?, 2364 job_visibility: PublishJobVisibility::from_str(visibility.as_str())?, 2365 expires_at_unix: row.try_get(9)?, 2366 }) 2367 } 2368 2369 fn job_from_row(row: &SqliteRow) -> Result<PublishJobRow, TransportPublishError> { 2370 let status: JobStatus = json_text(row, 3, "status")?; 2371 let target_policy: TargetPolicy = json_text(row, 7, "target_policy_json")?; 2372 let delivery_policy: DeliveryPolicy = json_text(row, 8, "delivery_policy_json")?; 2373 Ok(PublishJobRow { 2374 principal_id: row.try_get(1)?, 2375 request_fingerprint: row.try_get(2)?, 2376 view: Job { 2377 job_id: row.try_get(0)?, 2378 status, 2379 terminal: false, 2380 delivery_satisfied: false, 2381 event_id: row.try_get(4)?, 2382 pubkey: row.try_get(5)?, 2383 event_kind: checked_event_kind_column(row, 6)?, 2384 target_policy, 2385 delivery_policy, 2386 target_count: checked_usize_column(row, 9, "effective_target_count")?, 2387 acknowledged_count: 0, 2388 retryable_count: 0, 2389 terminal_count: 0, 2390 requested_at_ms: row.try_get(10)?, 2391 completed_at_ms: row.try_get(11)?, 2392 last_error: row.try_get(12)?, 2393 targets: Vec::new(), 2394 }, 2395 }) 2396 } 2397 2398 fn target_outcome_from_row(row: &SqliteRow) -> Result<TargetOutcome, TransportPublishError> { 2399 let source: TargetSource = json_text(row, 4, "source")?; 2400 let outcome_kind: OutcomeKind = json_text(row, 6, "outcome_kind")?; 2401 Ok(TargetOutcome { 2402 transport_kind: row.try_get(0)?, 2403 endpoint_uri: row.try_get(1)?, 2404 target_scope: storage_target_scope_to_protocol(row.try_get::<String, _>(2)?), 2405 target_label: row.try_get(3)?, 2406 source, 2407 attempted: row.try_get(5)?, 2408 outcome_kind, 2409 message: row.try_get(7)?, 2410 latency_ms: checked_optional_u64_column(row, 8, "latency_ms")?, 2411 }) 2412 } 2413 2414 fn finalize_job_view(view: &mut Job) { 2415 view.acknowledged_count = view 2416 .targets 2417 .iter() 2418 .filter(|relay| relay.outcome_kind.counts_toward_accepted_delivery()) 2419 .count(); 2420 view.retryable_count = view 2421 .targets 2422 .iter() 2423 .filter(|relay| relay.outcome_kind.is_retryable()) 2424 .count(); 2425 view.terminal_count = view 2426 .targets 2427 .iter() 2428 .filter(|relay| relay.outcome_kind.is_terminal_failure()) 2429 .count(); 2430 view.terminal = matches!( 2431 view.status, 2432 JobStatus::DeliverySatisfied 2433 | JobStatus::DeliveryUnsatisfiedTerminal 2434 | JobStatus::DeliveryDeferred 2435 | JobStatus::DeliveryDeferredUntilImplemented 2436 | JobStatus::Rejected 2437 ); 2438 view.delivery_satisfied = view.status == JobStatus::DeliverySatisfied; 2439 } 2440 2441 fn validate_principal_init(input: &PublishPrincipalInit) -> Result<(), TransportPublishError> { 2442 if input.label.trim().is_empty() { 2443 return Err(TransportPublishError::InvalidScope( 2444 "principal label must not be empty".to_owned(), 2445 )); 2446 } 2447 if !input.token_hash.starts_with(TOKEN_HASH_PREFIX) { 2448 return Err(TransportPublishError::InvalidScope( 2449 "principal token hash must use sha256 prefix".to_owned(), 2450 )); 2451 } 2452 if input.allowed_pubkeys.is_empty() { 2453 return Err(TransportPublishError::InvalidScope( 2454 "principal must include at least one allowed pubkey".to_owned(), 2455 )); 2456 } 2457 for pubkey in &input.allowed_pubkeys { 2458 ensure_lower_hex("allowed_pubkey", pubkey, 64)?; 2459 } 2460 if input.allowed_kinds.is_empty() { 2461 return Err(TransportPublishError::InvalidScope( 2462 "principal must include at least one allowed kind".to_owned(), 2463 )); 2464 } 2465 if input.allowed_target_policies.is_empty() { 2466 return Err(TransportPublishError::InvalidScope( 2467 "principal must include at least one allowed target policy".to_owned(), 2468 )); 2469 } 2470 let allows_explicit_targets = input 2471 .allowed_target_policies 2472 .contains(&TargetPolicyName::ExplicitTargets); 2473 if allows_explicit_targets && input.allowed_explicit_transport_kinds.is_empty() { 2474 return Err(TransportPublishError::InvalidScope( 2475 "principal must include at least one allowed explicit transport kind".to_owned(), 2476 )); 2477 } 2478 if !allows_explicit_targets && !input.allowed_explicit_transport_kinds.is_empty() { 2479 return Err(TransportPublishError::InvalidScope( 2480 "principal cannot include explicit transport kinds without explicit target policy" 2481 .to_owned(), 2482 )); 2483 } 2484 let mut explicit_transport_kinds = BTreeSet::new(); 2485 for transport_kind in &input.allowed_explicit_transport_kinds { 2486 let canonical = parse_explicit_transport_kind(transport_kind)?; 2487 if canonical != *transport_kind || !explicit_transport_kinds.insert(canonical) { 2488 return Err(TransportPublishError::InvalidScope( 2489 "allowed explicit transport kinds must be canonical and unique".to_owned(), 2490 )); 2491 } 2492 } 2493 if input 2494 .allowed_target_policies 2495 .contains(&TargetPolicyName::Nostr) 2496 && input.allowed_nostr_source_policies.is_empty() 2497 { 2498 return Err(TransportPublishError::InvalidScope( 2499 "principal must include at least one allowed Nostr source policy".to_owned(), 2500 )); 2501 } 2502 Ok(()) 2503 } 2504 2505 pub fn generate_bearer_token() -> String { 2506 let bytes: [u8; 32] = rand::random(); 2507 format!("{TOKEN_PREFIX}{}", hex_lower(&bytes)) 2508 } 2509 2510 pub fn hash_bearer_token(token: &str) -> String { 2511 let mut hasher = Sha256::new(); 2512 hasher.update(token.as_bytes()); 2513 format!("{TOKEN_HASH_PREFIX}{}", hex_lower(&hasher.finalize())) 2514 } 2515 2516 fn hex_lower(bytes: &[u8]) -> String { 2517 let mut output = String::with_capacity(bytes.len() * 2); 2518 for byte in bytes { 2519 use std::fmt::Write; 2520 let _ = write!(&mut output, "{byte:02x}"); 2521 } 2522 output 2523 } 2524 2525 pub fn parse_nostr_source_policy( 2526 value: &str, 2527 ) -> Result<NostrTargetSourcePolicy, TransportPublishError> { 2528 match value { 2529 "explicit_only" => Ok(NostrTargetSourcePolicy::ExplicitOnly), 2530 "request_then_author_write_then_daemon_default" => { 2531 Ok(NostrTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault) 2532 } 2533 "author_write_then_daemon_default" => { 2534 Ok(NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault) 2535 } 2536 "daemon_default_only" => Ok(NostrTargetSourcePolicy::DaemonDefaultOnly), 2537 other => Err(TransportPublishError::InvalidScope(format!( 2538 "unknown Nostr source policy `{other}`" 2539 ))), 2540 } 2541 } 2542 2543 pub fn parse_target_policy(value: &str) -> Result<TargetPolicyName, TransportPublishError> { 2544 match value { 2545 "explicit_targets" => Ok(TargetPolicyName::ExplicitTargets), 2546 "nostr" => Ok(TargetPolicyName::Nostr), 2547 other => Err(TransportPublishError::InvalidScope(format!( 2548 "unknown target policy `{other}`" 2549 ))), 2550 } 2551 } 2552 2553 pub fn parse_explicit_transport_kind(value: &str) -> Result<String, TransportPublishError> { 2554 let kind = TransportId::parse_canonical(value).map_err(|error| { 2555 TransportPublishError::InvalidScope(format!( 2556 "unknown explicit transport kind `{value}`: {error}" 2557 )) 2558 })?; 2559 let canonical = kind.canonical_label(); 2560 if !matches!(canonical.as_str(), "local" | "nostr" | "reticulum") { 2561 return Err(TransportPublishError::InvalidScope(format!( 2562 "unknown explicit transport kind `{value}`" 2563 ))); 2564 } 2565 Ok(canonical) 2566 } 2567 2568 fn signed_event_from_raw_json(raw_json: &str) -> Result<SignedEvent, TransportPublishError> { 2569 let wire = Nip01EventWire::parse_json(raw_json)?; 2570 let signed_event = SignedEvent::from_wire_verified_id(wire, raw_json.to_owned())?; 2571 if radroots_nostr::event::verify(signed_event.envelope()) 2572 != radroots_nostr::event::Verification::Verified 2573 { 2574 return Err(TransportPublishError::InvalidSignedEvent( 2575 "signature verification failed".to_owned(), 2576 )); 2577 } 2578 Ok(signed_event) 2579 } 2580 2581 fn request_intent_fingerprint( 2582 principal_id: &str, 2583 canonical_event_json: &str, 2584 request: &EventRequest, 2585 effective_timeout_ms: u64, 2586 ) -> Result<String, TransportPublishError> { 2587 #[derive(Serialize)] 2588 struct FingerprintInput<'a> { 2589 principal_id: &'a str, 2590 canonical_event_json: &'a str, 2591 target_policy: &'a TargetPolicy, 2592 delivery_policy: &'a DeliveryPolicy, 2593 effective_timeout_ms: u64, 2594 } 2595 2596 let input = FingerprintInput { 2597 principal_id, 2598 canonical_event_json, 2599 target_policy: &request.target_policy, 2600 delivery_policy: &request.delivery_policy, 2601 effective_timeout_ms, 2602 }; 2603 let bytes = serde_json::to_vec(&input)?; 2604 let mut hasher = Sha256::new(); 2605 hasher.update(bytes); 2606 Ok(hex_lower(&hasher.finalize())) 2607 } 2608 2609 fn effective_publish_timeout_ms( 2610 config: &TransportPublishConfig, 2611 timeout_ms: Option<u64>, 2612 ) -> Result<u64, TransportPublishError> { 2613 let max_timeout_ms = config.connect_timeout_secs.saturating_mul(1_000); 2614 match timeout_ms { 2615 Some(0) => Err(TransportPublishError::InvalidSignedEvent( 2616 "timeout_ms must be greater than zero".to_owned(), 2617 )), 2618 Some(timeout_ms) if timeout_ms > max_timeout_ms => { 2619 Err(TransportPublishError::InvalidSignedEvent(format!( 2620 "timeout_ms must be at most {max_timeout_ms}" 2621 ))) 2622 } 2623 Some(timeout_ms) => Ok(timeout_ms), 2624 None => Ok(max_timeout_ms), 2625 } 2626 } 2627 2628 fn push_resolved_relay( 2629 targets: &mut Vec<ResolvedPublishRelay>, 2630 url: RadrootsRelayUrl, 2631 source: TargetSource, 2632 metadata: PublishTargetMetadata, 2633 ) { 2634 if !targets 2635 .iter() 2636 .any(|target| target.url == url && target.target_scope == metadata.target_scope) 2637 { 2638 targets.push(ResolvedPublishRelay { 2639 url, 2640 source, 2641 target_scope: metadata.target_scope, 2642 target_label: metadata.target_label, 2643 }); 2644 } 2645 } 2646 2647 fn reticulum_unavailable_outcome(target: &Target) -> TargetOutcome { 2648 TargetOutcome { 2649 transport_kind: TRANSPORT_KIND_RETICULUM.to_owned(), 2650 endpoint_uri: target.endpoint_uri.trim().to_owned(), 2651 target_scope: target.target_scope.clone(), 2652 target_label: target.target_label.clone(), 2653 source: TargetSource::Reticulum, 2654 attempted: false, 2655 outcome_kind: OutcomeKind::DeferredUntilImplemented, 2656 message: Some(RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE.to_owned()), 2657 latency_ms: None, 2658 } 2659 } 2660 2661 fn unsupported_transport_outcome(target: &Target) -> TargetOutcome { 2662 TargetOutcome { 2663 transport_kind: target.transport_kind.trim().to_owned(), 2664 endpoint_uri: target.endpoint_uri.trim().to_owned(), 2665 target_scope: target.target_scope.clone(), 2666 target_label: target.target_label.clone(), 2667 source: TargetSource::Request, 2668 attempted: false, 2669 outcome_kind: OutcomeKind::Unsupported, 2670 message: Some("transport kind is not supported by radrootsd transport publish".to_owned()), 2671 latency_ms: None, 2672 } 2673 } 2674 2675 fn relay_resolution_connection_failure( 2676 relay_url: impl Into<String>, 2677 source: TargetSource, 2678 metadata: &PublishTargetMetadata, 2679 message: impl Into<String>, 2680 ) -> TargetOutcome { 2681 TargetOutcome { 2682 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 2683 endpoint_uri: relay_url.into(), 2684 target_scope: metadata.target_scope.clone(), 2685 target_label: metadata.target_label.clone(), 2686 source, 2687 attempted: false, 2688 outcome_kind: OutcomeKind::ConnectionFailed, 2689 message: Some(message.into()), 2690 latency_ms: None, 2691 } 2692 } 2693 2694 fn target_snapshots_from_resolution(resolution: &PublishRelayResolution) -> Vec<TargetOutcome> { 2695 let mut snapshots = resolution.outcomes.clone(); 2696 snapshots.extend(resolution.targets.iter().map(|target| TargetOutcome { 2697 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 2698 endpoint_uri: target.url.as_str().to_owned(), 2699 target_scope: target.target_scope.clone(), 2700 target_label: target.target_label.clone(), 2701 source: target.source, 2702 attempted: true, 2703 outcome_kind: OutcomeKind::ConnectionFailed, 2704 message: Some("publish_attempt_interrupted".to_owned()), 2705 latency_ms: None, 2706 })); 2707 snapshots 2708 } 2709 2710 fn relay_socket_target(url: &RadrootsRelayUrl) -> Result<(String, u16), std::io::Error> { 2711 let parsed = url::Url::parse(url.as_str()) 2712 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidInput, error))?; 2713 let host = parsed 2714 .host_str() 2715 .filter(|host| !host.is_empty()) 2716 .ok_or_else(|| { 2717 std::io::Error::new( 2718 std::io::ErrorKind::InvalidInput, 2719 "relay URL must include a DNS host", 2720 ) 2721 })? 2722 .to_owned(); 2723 let port = parsed.port_or_known_default().ok_or_else(|| { 2724 std::io::Error::new( 2725 std::io::ErrorKind::InvalidInput, 2726 "relay URL scheme must have a default port", 2727 ) 2728 })?; 2729 Ok((host, port)) 2730 } 2731 2732 fn relay_url_policy(config: &TransportPublishConfig) -> RadrootsRelayUrlPolicy { 2733 match config.nostr.relay_url_policy { 2734 crate::app::config::NostrRelayUrlPolicy::Public => RadrootsRelayUrlPolicy::Public, 2735 crate::app::config::NostrRelayUrlPolicy::Localhost => RadrootsRelayUrlPolicy::Localhost, 2736 } 2737 } 2738 2739 fn author_write_relays_from_nip65_event(event: &crate::host_nostr::Event) -> Vec<String> { 2740 event 2741 .tags 2742 .iter() 2743 .filter_map(|tag| { 2744 let values = tag.as_slice(); 2745 if values.first().map(String::as_str) != Some("r") { 2746 return None; 2747 } 2748 let relay = values.get(1)?.trim(); 2749 if relay.is_empty() { 2750 return None; 2751 } 2752 if values.get(2).map(String::as_str) == Some("read") { 2753 return None; 2754 } 2755 Some(relay.to_owned()) 2756 }) 2757 .collect() 2758 } 2759 2760 fn publish_outcomes_from_receipt( 2761 receipt: RadrootsRelayPublishRelayReceipt, 2762 targets: &[ResolvedPublishRelay], 2763 latency_ms: Option<u64>, 2764 ) -> Vec<TargetOutcome> { 2765 targets 2766 .iter() 2767 .filter(|target| target.url.as_str() == receipt.relay_url.as_str()) 2768 .map(|target| TargetOutcome { 2769 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 2770 endpoint_uri: receipt.relay_url.clone(), 2771 target_scope: target.target_scope.clone(), 2772 target_label: target.target_label.clone(), 2773 source: target.source, 2774 attempted: receipt.attempted, 2775 outcome_kind: publish_outcome_kind(receipt.outcome.kind), 2776 message: receipt.outcome.message.clone(), 2777 latency_ms, 2778 }) 2779 .collect() 2780 } 2781 2782 fn publish_outcome_kind(kind: RadrootsRelayOutcomeKind) -> OutcomeKind { 2783 match kind { 2784 RadrootsRelayOutcomeKind::Accepted => OutcomeKind::Accepted, 2785 RadrootsRelayOutcomeKind::DuplicateAccepted => OutcomeKind::DuplicateAccepted, 2786 RadrootsRelayOutcomeKind::Blocked => OutcomeKind::Blocked, 2787 RadrootsRelayOutcomeKind::RateLimited => OutcomeKind::RateLimited, 2788 RadrootsRelayOutcomeKind::Invalid => OutcomeKind::Invalid, 2789 RadrootsRelayOutcomeKind::PowRequired => OutcomeKind::PowRequired, 2790 RadrootsRelayOutcomeKind::Restricted => OutcomeKind::Restricted, 2791 RadrootsRelayOutcomeKind::AuthRequired => OutcomeKind::AuthRequired, 2792 RadrootsRelayOutcomeKind::Muted => OutcomeKind::Muted, 2793 RadrootsRelayOutcomeKind::Unsupported => OutcomeKind::Unsupported, 2794 RadrootsRelayOutcomeKind::PaymentRequired => OutcomeKind::PaymentRequired, 2795 RadrootsRelayOutcomeKind::Error => OutcomeKind::Error, 2796 RadrootsRelayOutcomeKind::Timeout => OutcomeKind::Timeout, 2797 RadrootsRelayOutcomeKind::ConnectionFailed => OutcomeKind::ConnectionFailed, 2798 RadrootsRelayOutcomeKind::Unknown => OutcomeKind::Unknown, 2799 } 2800 } 2801 2802 fn target_outcome_fingerprint( 2803 target: &TargetOutcome, 2804 index: usize, 2805 ) -> Result<TargetFingerprint, TransportPublishError> { 2806 let transport_kind = 2807 TransportId::parse_canonical(target.transport_kind.as_str()).map_err(|error| { 2808 TransportPublishError::InvalidPublishJobState(format!( 2809 "target outcome {index} has invalid transport kind: {error}" 2810 )) 2811 })?; 2812 let scope = target 2813 .target_scope 2814 .as_deref() 2815 .map(TargetScope::parse) 2816 .transpose() 2817 .map_err(|error| { 2818 TransportPublishError::InvalidPublishJobState(format!( 2819 "target outcome {index} has invalid target scope: {error}" 2820 )) 2821 })?; 2822 let label = target 2823 .target_label 2824 .as_deref() 2825 .map(TargetLabel::parse) 2826 .transpose() 2827 .map_err(|error| { 2828 TransportPublishError::InvalidPublishJobState(format!( 2829 "target outcome {index} has invalid target label: {error}" 2830 )) 2831 })?; 2832 let target = transport_target_from_outcome_parts( 2833 transport_kind, 2834 target.endpoint_uri.as_str(), 2835 scope, 2836 label, 2837 ) 2838 .map_err(|error| { 2839 TransportPublishError::InvalidPublishJobState(format!( 2840 "target outcome {index} fingerprint failed: {error}" 2841 )) 2842 })?; 2843 Ok(target.fingerprint().clone()) 2844 } 2845 2846 fn transport_target_from_outcome_parts( 2847 transport_kind: TransportId, 2848 endpoint_uri: &str, 2849 scope: Option<TargetScope>, 2850 label: Option<TargetLabel>, 2851 ) -> Result<TransportTarget, TransportError> { 2852 match transport_kind { 2853 TransportId::NOSTR => { 2854 TransportTarget::nostr_relay_with_metadata(endpoint_uri, scope, label) 2855 } 2856 TransportId::RETICULUM => { 2857 if endpoint_uri != RADROOTS_RETICULUM_ENDPOINT_URI { 2858 return Err(TransportError::InvalidTargetUri); 2859 } 2860 let scope = match scope { 2861 Some(scope) => scope, 2862 None => TargetScope::parse("local")?, 2863 }; 2864 TransportTarget::new_with_metadata( 2865 TransportId::RETICULUM, 2866 endpoint_uri, 2867 Some(scope), 2868 label, 2869 ) 2870 } 2871 TransportId::LOCAL => TransportTarget::local_with_metadata(endpoint_uri, scope, label), 2872 _ => Err(TransportError::InvalidTargetUri), 2873 } 2874 } 2875 2876 fn validate_delivery_policy_for_resolution( 2877 delivery_policy: &DeliveryPolicy, 2878 resolution: &PublishRelayResolution, 2879 ) -> Result<(), TransportPublishError> { 2880 let DeliveryPolicy::RequiredTargets { targets } = delivery_policy else { 2881 return Ok(()); 2882 }; 2883 let target_fingerprints = resolution.target_fingerprints()?; 2884 if targets.iter().any(|required| { 2885 !target_fingerprints 2886 .iter() 2887 .any(|actual| actual.as_str() == required.as_str()) 2888 }) { 2889 return Err(TransportPublishError::InvalidSignedEvent( 2890 "publish request requires a target outside the resolved set".to_owned(), 2891 )); 2892 } 2893 Ok(()) 2894 } 2895 2896 fn required_outcomes_for_policy<'a>( 2897 required_targets: &[ProtocolTargetFingerprint], 2898 outcomes: &'a [TargetOutcome], 2899 ) -> Vec<&'a TargetOutcome> { 2900 required_targets 2901 .iter() 2902 .filter_map(|required| { 2903 outcomes.iter().enumerate().find_map(|(index, outcome)| { 2904 target_outcome_fingerprint(outcome, index) 2905 .ok() 2906 .filter(|fingerprint| fingerprint.as_str() == required.as_str()) 2907 .map(|_| outcome) 2908 }) 2909 }) 2910 .collect() 2911 } 2912 2913 fn delivery_status( 2914 delivery_policy: &DeliveryPolicy, 2915 target_count: usize, 2916 outcomes: &[TargetOutcome], 2917 ) -> JobStatus { 2918 let (satisfied, status_outcomes) = match delivery_policy { 2919 DeliveryPolicy::RequiredTargets { targets } => { 2920 let required_outcomes = required_outcomes_for_policy(targets, outcomes); 2921 let satisfied = required_outcomes.len() == targets.len() 2922 && required_outcomes 2923 .iter() 2924 .all(|outcome| outcome.outcome_kind.counts_toward_accepted_delivery()); 2925 (satisfied, required_outcomes) 2926 } 2927 DeliveryPolicy::Any | DeliveryPolicy::All | DeliveryPolicy::Quorum { .. } => { 2928 let required = delivery_policy.required_target_count(target_count); 2929 let acknowledged = outcomes 2930 .iter() 2931 .filter(|outcome| outcome.outcome_kind.counts_toward_accepted_delivery()) 2932 .count(); 2933 ( 2934 acknowledged >= required, 2935 outcomes.iter().collect::<Vec<_>>(), 2936 ) 2937 } 2938 }; 2939 if satisfied { 2940 return JobStatus::DeliverySatisfied; 2941 } 2942 if status_outcomes 2943 .iter() 2944 .any(|outcome| outcome.outcome_kind.is_retryable()) 2945 { 2946 JobStatus::DeliveryUnsatisfiedRetryable 2947 } else if status_outcomes 2948 .iter() 2949 .any(|outcome| outcome.outcome_kind == OutcomeKind::DeferredUntilImplemented) 2950 && status_outcomes 2951 .iter() 2952 .all(|outcome| !outcome.outcome_kind.is_terminal_failure()) 2953 { 2954 JobStatus::DeliveryDeferredUntilImplemented 2955 } else { 2956 JobStatus::DeliveryUnsatisfiedTerminal 2957 } 2958 } 2959 2960 fn last_error_for_status(status: JobStatus) -> Option<&'static str> { 2961 match status { 2962 JobStatus::DeliverySatisfied => None, 2963 JobStatus::Rejected => Some("no_transport_publish_targets"), 2964 JobStatus::DeliveryDeferred => Some("delivery_deferred_until_implemented"), 2965 JobStatus::DeliveryDeferredUntilImplemented => Some("delivery_deferred_until_implemented"), 2966 JobStatus::Accepted 2967 | JobStatus::Publishing 2968 | JobStatus::DeliveryUnsatisfiedRetryable 2969 | JobStatus::DeliveryUnsatisfiedTerminal => Some("delivery_unsatisfied"), 2970 } 2971 } 2972 2973 fn timeout_receipts(targets: &[RadrootsRelayUrl]) -> Vec<RadrootsRelayPublishRelayReceipt> { 2974 targets 2975 .iter() 2976 .map(|target| { 2977 RadrootsRelayPublishRelayReceipt::attempted( 2978 target.as_str(), 2979 RadrootsRelayOutcome::timeout("timeout: publish attempt exceeded daemon bound"), 2980 ) 2981 }) 2982 .collect() 2983 } 2984 2985 fn transport_error_receipts( 2986 targets: &[RadrootsRelayUrl], 2987 error: TransportPublishError, 2988 ) -> Vec<RadrootsRelayPublishRelayReceipt> { 2989 let message = format!("error: {error}"); 2990 targets 2991 .iter() 2992 .map(|target| { 2993 RadrootsRelayPublishRelayReceipt::attempted( 2994 target.as_str(), 2995 RadrootsRelayOutcome::connection_failed(message.clone()), 2996 ) 2997 }) 2998 .collect() 2999 } 3000 3001 pub fn write_token_file(path: &Path, token: &str) -> Result<(), TransportPublishError> { 3002 if let Some(parent) = path 3003 .parent() 3004 .filter(|parent| !parent.as_os_str().is_empty()) 3005 { 3006 std::fs::create_dir_all(parent)?; 3007 } 3008 let mut options = std::fs::OpenOptions::new(); 3009 options.write(true).create_new(true); 3010 #[cfg(unix)] 3011 { 3012 use std::os::unix::fs::OpenOptionsExt; 3013 options.mode(0o600); 3014 } 3015 use std::io::Write; 3016 let mut file = options.open(path)?; 3017 file.write_all(token.as_bytes())?; 3018 file.write_all(b"\n")?; 3019 Ok(()) 3020 } 3021 3022 fn ensure_lower_hex( 3023 field: &str, 3024 value: &str, 3025 expected_len: usize, 3026 ) -> Result<(), TransportPublishError> { 3027 if value.len() == expected_len 3028 && value 3029 .bytes() 3030 .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f')) 3031 { 3032 Ok(()) 3033 } else { 3034 Err(TransportPublishError::InvalidScope(format!( 3035 "{field} must be {expected_len} lowercase hex characters" 3036 ))) 3037 } 3038 } 3039 3040 fn json_column<T: for<'de> Deserialize<'de>>( 3041 row: &SqliteRow, 3042 index: usize, 3043 field: &'static str, 3044 ) -> Result<T, TransportPublishError> { 3045 let value: String = row.try_get(index)?; 3046 serde_json::from_str(value.as_str()).map_err(|error| persisted_decode_error(field, error)) 3047 } 3048 3049 fn json_text<T: for<'de> Deserialize<'de>>( 3050 row: &SqliteRow, 3051 index: usize, 3052 field: &'static str, 3053 ) -> Result<T, TransportPublishError> { 3054 let value: String = row.try_get(index)?; 3055 serde_json::from_str(value.as_str()).map_err(|error| persisted_decode_error(field, error)) 3056 } 3057 3058 #[derive(Debug, Error)] 3059 #[error("{field} integer value {value} is outside {target} range")] 3060 struct TransportPublishStorageIntegerRangeError { 3061 field: &'static str, 3062 value: i64, 3063 target: &'static str, 3064 } 3065 3066 fn checked_event_kind_column(row: &SqliteRow, index: usize) -> Result<u32, TransportPublishError> { 3067 let value = row.try_get::<i64, _>(index)?; 3068 if !(0..=i64::from(u32::MAX)).contains(&value) { 3069 return Err(persisted_decode_error( 3070 "event_kind", 3071 TransportPublishStorageIntegerRangeError { 3072 field: "event_kind", 3073 value, 3074 target: "u32", 3075 }, 3076 )); 3077 } 3078 u32::try_from(value).map_err(|_| { 3079 persisted_decode_error( 3080 "event_kind", 3081 TransportPublishStorageIntegerRangeError { 3082 field: "event_kind", 3083 value, 3084 target: "u32", 3085 }, 3086 ) 3087 }) 3088 } 3089 3090 fn checked_usize_column( 3091 row: &SqliteRow, 3092 index: usize, 3093 field: &'static str, 3094 ) -> Result<usize, TransportPublishError> { 3095 let value = row.try_get::<i64, _>(index)?; 3096 usize::try_from(value).map_err(|_| { 3097 persisted_decode_error( 3098 field, 3099 TransportPublishStorageIntegerRangeError { 3100 field, 3101 value, 3102 target: "usize", 3103 }, 3104 ) 3105 }) 3106 } 3107 3108 fn storage_count_i64(value: usize, field: &'static str) -> Result<i64, TransportPublishError> { 3109 i64::try_from(value).map_err(|_| { 3110 TransportPublishError::InvalidPublishJobState(format!( 3111 "{field} value exceeds i64 storage range" 3112 )) 3113 }) 3114 } 3115 3116 fn checked_optional_u64_column( 3117 row: &SqliteRow, 3118 index: usize, 3119 field: &'static str, 3120 ) -> Result<Option<u64>, TransportPublishError> { 3121 row.try_get::<Option<i64>, _>(index)? 3122 .map(|value| { 3123 u64::try_from(value).map_err(|_| { 3124 persisted_decode_error( 3125 field, 3126 TransportPublishStorageIntegerRangeError { 3127 field, 3128 value, 3129 target: "u64", 3130 }, 3131 ) 3132 }) 3133 }) 3134 .transpose() 3135 } 3136 3137 fn storage_target_scope(target_scope: Option<&str>) -> &str { 3138 target_scope.unwrap_or("") 3139 } 3140 3141 fn storage_target_scope_to_protocol(target_scope: String) -> Option<String> { 3142 (!target_scope.is_empty()).then_some(target_scope) 3143 } 3144 3145 fn persisted_decode_error<E>(field: &'static str, error: E) -> TransportPublishError 3146 where 3147 E: std::error::Error, 3148 { 3149 TransportPublishError::InvalidPublishJobState(format!( 3150 "{field} persisted value could not be decoded: {error}" 3151 )) 3152 } 3153 3154 fn current_unix_secs() -> i64 { 3155 SystemTime::now() 3156 .duration_since(UNIX_EPOCH) 3157 .map(|duration| duration.as_secs() as i64) 3158 .unwrap_or_default() 3159 } 3160 3161 fn current_unix_millis() -> i64 { 3162 SystemTime::now() 3163 .duration_since(UNIX_EPOCH) 3164 .map(|duration| duration.as_millis() as i64) 3165 .unwrap_or_default() 3166 } 3167 3168 #[cfg(test)] 3169 mod tests { 3170 use super::{ 3171 PublishEventMetadata, PublishJobInsert, PublishJobVisibility, PublishPrincipal, 3172 PublishPrincipalInit, SCHEMA_VERSION, TRANSPORT_KIND_NOSTR, TRANSPORT_KIND_RETICULUM, 3173 TRANSPORT_PUBLISH_SCHEMA_SQL, TransportPublish, TransportPublishError, 3174 TransportPublishStore, generate_bearer_token, hash_bearer_token, parse_nostr_source_policy, 3175 }; 3176 use crate::app::config::{ 3177 NostrRelayUrlPolicy, TransportPublishConfig, TransportPublishNostrConfig, 3178 }; 3179 use crate::app::identity_storage::DaemonIdentity; 3180 use crate::host_nostr::Timestamp; 3181 use crate::transport::relay_publish::{ 3182 MockRelayPublishAdapter as RadrootsMockRelayPublishAdapter, 3183 RelayOutcome as RadrootsRelayOutcome, 3184 }; 3185 use nostr::JsonUtil; 3186 use nostr::{EventBuilder, Kind, Tag}; 3187 use radroots_protocol::radrootsd::transport_publish::v5::{ 3188 DeliveryPolicy, EventRequest, JobStatus, NostrTargetSourcePolicy, OutcomeKind, 3189 ReticulumBehavior, Target, TargetFingerprint as ProtocolTargetFingerprint, TargetOutcome, 3190 TargetPolicy, TargetPolicyName, TargetSource, 3191 }; 3192 use radroots_protocol::radrootsd::transport_publish::v5::{ 3193 RETICULUM_ENDPOINT_URI as RADROOTS_RETICULUM_ENDPOINT_URI, 3194 RETICULUM_UNAVAILABLE_MESSAGE as RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE, 3195 }; 3196 use radroots_transport::Target as TransportTarget; 3197 use sqlx::Row; 3198 use sqlx::sqlite::{SqliteConnectOptions, SqliteConnection}; 3199 use std::collections::BTreeMap; 3200 use std::net::{IpAddr, Ipv4Addr}; 3201 use std::sync::Arc; 3202 3203 const RELAY_PRIMARY: &str = "wss://relay.example.com"; 3204 const RELAY_SECONDARY: &str = "wss://relay-2.example.com"; 3205 const RELAY_FORBIDDEN: &str = "wss://forbidden-relay.example.com"; 3206 3207 fn event_metadata(pubkey: &str, kind: u32) -> PublishEventMetadata { 3208 PublishEventMetadata { 3209 event_id: "0".repeat(64), 3210 pubkey: pubkey.to_owned(), 3211 kind, 3212 } 3213 } 3214 3215 fn raw_event_json(pubkey: &str, kind: u32) -> String { 3216 format!( 3217 r#"{{"id":"{}","pubkey":"{}","created_at":1700000000,"kind":{},"tags":[["d","listing-1"]],"content":"{{}}","sig":"{}"}}"#, 3218 "0".repeat(64), 3219 pubkey, 3220 kind, 3221 "1".repeat(128) 3222 ) 3223 } 3224 3225 fn request(pubkey: &str, kind: u32) -> EventRequest { 3226 EventRequest { 3227 raw_event_json: raw_event_json(pubkey, kind), 3228 target_policy: TargetPolicy::nostr( 3229 NostrTargetSourcePolicy::DaemonDefaultOnly, 3230 Vec::new(), 3231 ), 3232 delivery_policy: DeliveryPolicy::Any, 3233 idempotency_key: Some("idem-1".to_owned()), 3234 timeout_ms: None, 3235 } 3236 } 3237 3238 fn schema_with_replacement(fragment: &str, replacement: &str) -> String { 3239 assert!(TRANSPORT_PUBLISH_SCHEMA_SQL.contains(fragment)); 3240 TRANSPORT_PUBLISH_SCHEMA_SQL.replace(fragment, replacement) 3241 } 3242 3243 fn create_existing_schema( 3244 database_path: &std::path::Path, 3245 schema_sql: &str, 3246 user_version: Option<i64>, 3247 ) { 3248 let mut connection = open_test_database(database_path); 3249 super::execute_raw_sql(&mut connection, schema_sql).expect("create schema"); 3250 if let Some(version) = user_version { 3251 super::execute_sql( 3252 &mut connection, 3253 format!("PRAGMA user_version = {version}").as_str(), 3254 ) 3255 .expect("set schema version"); 3256 } 3257 } 3258 3259 fn open_test_database(database_path: &std::path::Path) -> SqliteConnection { 3260 super::connect_sqlite( 3261 SqliteConnectOptions::new() 3262 .filename(database_path) 3263 .create_if_missing(true), 3264 ) 3265 .expect("open schema database") 3266 } 3267 3268 fn open_schema_error(database_path: &std::path::Path) -> TransportPublishError { 3269 match TransportPublishStore::open(database_path.to_path_buf()) { 3270 Ok(_) => panic!("malformed schema opened"), 3271 Err(error) => error, 3272 } 3273 } 3274 3275 fn assert_schema_error( 3276 error: TransportPublishError, 3277 expected_table: &'static str, 3278 expected_detail: &str, 3279 ) { 3280 match error { 3281 TransportPublishError::Schema { table, detail } => { 3282 assert_eq!(table, expected_table); 3283 assert!( 3284 detail.contains(expected_detail), 3285 "schema detail `{detail}` did not contain `{expected_detail}`" 3286 ); 3287 } 3288 error => panic!("unexpected error: {error}"), 3289 } 3290 } 3291 3292 fn database_user_version(database_path: &std::path::Path) -> i64 { 3293 let mut connection = open_test_database(database_path); 3294 super::transport_publish_schema_version(&mut connection).expect("user version") 3295 } 3296 3297 fn test_query_column_names(connection: &mut SqliteConnection, table: &str) -> Vec<String> { 3298 let sql = format!("PRAGMA table_info({table})"); 3299 super::fetch_all_sql(connection, sql.as_str()) 3300 .expect("query schema") 3301 .iter() 3302 .map(|row| row.try_get::<String, _>(1).expect("column name")) 3303 .collect() 3304 } 3305 3306 fn signed_event(identity: &DaemonIdentity, content: &str) -> String { 3307 // Transport tests require an already-signed wire fixture; they do not 3308 // exercise a Radroots product-authoring boundary. 3309 let event = EventBuilder::new(Kind::Custom(30_402), content) 3310 .tag(Tag::identifier("listing-1")) 3311 .custom_created_at(Timestamp::from_secs(1_700_000_000)) 3312 .sign_with_keys(identity.keys()) 3313 .expect("signed event"); 3314 event.as_json() 3315 } 3316 3317 fn raw_event_with_field( 3318 raw_event_json: String, 3319 field: &str, 3320 value: serde_json::Value, 3321 ) -> String { 3322 let mut event: serde_json::Value = 3323 serde_json::from_str(raw_event_json.as_str()).expect("raw event json"); 3324 event[field] = value; 3325 serde_json::to_string(&event).expect("mutated raw event") 3326 } 3327 3328 fn publish_request( 3329 raw_event_json: String, 3330 relays: Vec<String>, 3331 source_policy: NostrTargetSourcePolicy, 3332 delivery_policy: DeliveryPolicy, 3333 idempotency_key: Option<&str>, 3334 ) -> EventRequest { 3335 EventRequest { 3336 raw_event_json, 3337 target_policy: TargetPolicy::nostr(source_policy, relays), 3338 delivery_policy, 3339 idempotency_key: idempotency_key.map(str::to_owned), 3340 timeout_ms: Some(5_000), 3341 } 3342 } 3343 3344 fn reticulum_publish_request( 3345 raw_event_json: String, 3346 behavior: ReticulumBehavior, 3347 ) -> EventRequest { 3348 EventRequest { 3349 raw_event_json, 3350 target_policy: TargetPolicy::explicit_targets(vec![Target::reticulum(behavior)]), 3351 delivery_policy: DeliveryPolicy::Any, 3352 idempotency_key: None, 3353 timeout_ms: Some(5_000), 3354 } 3355 } 3356 3357 fn interrupted_target_snapshot(endpoint_uri: &str, source: TargetSource) -> TargetOutcome { 3358 TargetOutcome { 3359 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 3360 endpoint_uri: endpoint_uri.to_owned(), 3361 target_scope: None, 3362 target_label: None, 3363 source, 3364 attempted: true, 3365 outcome_kind: OutcomeKind::ConnectionFailed, 3366 message: Some("publish_attempt_interrupted".to_owned()), 3367 latency_ms: None, 3368 } 3369 } 3370 3371 fn accepted_target_outcome(endpoint_uri: &str, source: TargetSource) -> TargetOutcome { 3372 TargetOutcome { 3373 transport_kind: TRANSPORT_KIND_NOSTR.to_owned(), 3374 endpoint_uri: endpoint_uri.to_owned(), 3375 target_scope: None, 3376 target_label: None, 3377 source, 3378 attempted: true, 3379 outcome_kind: OutcomeKind::Accepted, 3380 message: None, 3381 latency_ms: Some(12), 3382 } 3383 } 3384 3385 fn scoped_target_outcome( 3386 mut outcome: TargetOutcome, 3387 target_scope: &str, 3388 target_label: Option<&str>, 3389 ) -> TargetOutcome { 3390 outcome.target_scope = Some(target_scope.to_owned()); 3391 outcome.target_label = target_label.map(str::to_owned); 3392 outcome 3393 } 3394 3395 fn store_principal(store: &TransportPublishStore, pubkey: &str) -> PublishPrincipal { 3396 store 3397 .create_principal(PublishPrincipalInit { 3398 label: "tester".to_owned(), 3399 token_hash: hash_bearer_token(generate_bearer_token().as_str()), 3400 allowed_pubkeys: vec![pubkey.to_owned()], 3401 allowed_kinds: vec![30_402], 3402 allowed_target_policies: vec![TargetPolicyName::Nostr], 3403 allowed_explicit_transport_kinds: Vec::new(), 3404 allowed_nostr_source_policies: vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 3405 allow_request_targets: false, 3406 job_visibility: PublishJobVisibility::Own, 3407 expires_at_unix: None, 3408 }) 3409 .expect("principal") 3410 } 3411 3412 fn assert_invalid_job_state(error: TransportPublishError, expected: &str) { 3413 match error { 3414 TransportPublishError::InvalidPublishJobState(message) => { 3415 assert!(message.contains(expected), "{message}"); 3416 assert!(!message.contains("rrd_tp_")); 3417 assert!(!message.contains("token")); 3418 } 3419 error => panic!("unexpected error: {error}"), 3420 } 3421 } 3422 3423 fn assert_storage_integer_range_error(error: TransportPublishError, expected: &str) { 3424 match error { 3425 TransportPublishError::InvalidPublishJobState(message) => { 3426 assert!(message.contains(expected), "{message}"); 3427 assert!(!message.contains("rrd_tp_")); 3428 assert!(!message.contains("token")); 3429 } 3430 error => panic!("unexpected error: {error}"), 3431 } 3432 } 3433 3434 fn transport_publish( 3435 config: TransportPublishConfig, 3436 ) -> (TransportPublish, RadrootsMockRelayPublishAdapter) { 3437 transport_publish_with_resolver(config, Arc::new(StaticPublishRelayResolver::new())) 3438 } 3439 3440 fn transport_publish_with_resolver( 3441 config: TransportPublishConfig, 3442 resolver: Arc<dyn super::PublishRelayResolver>, 3443 ) -> (TransportPublish, RadrootsMockRelayPublishAdapter) { 3444 let adapter = RadrootsMockRelayPublishAdapter::new(); 3445 let proxy = TransportPublish::memory(config) 3446 .expect("proxy") 3447 .with_relay_resolver(resolver) 3448 .with_publisher(Arc::new(adapter.clone())); 3449 (proxy, adapter) 3450 } 3451 3452 fn principal( 3453 proxy: &TransportPublish, 3454 pubkey: String, 3455 nostr_source_policies: Vec<NostrTargetSourcePolicy>, 3456 allow_request_targets: bool, 3457 visibility: PublishJobVisibility, 3458 ) -> PublishPrincipal { 3459 proxy 3460 .store 3461 .create_principal(PublishPrincipalInit { 3462 label: "tester".to_owned(), 3463 token_hash: hash_bearer_token(generate_bearer_token().as_str()), 3464 allowed_pubkeys: vec![pubkey], 3465 allowed_kinds: vec![30_402], 3466 allowed_target_policies: vec![TargetPolicyName::Nostr], 3467 allowed_explicit_transport_kinds: Vec::new(), 3468 allowed_nostr_source_policies: nostr_source_policies, 3469 allow_request_targets, 3470 job_visibility: visibility, 3471 expires_at_unix: None, 3472 }) 3473 .expect("principal") 3474 } 3475 3476 fn explicit_target_principal( 3477 proxy: &TransportPublish, 3478 pubkey: String, 3479 visibility: PublishJobVisibility, 3480 ) -> PublishPrincipal { 3481 explicit_target_principal_with_kinds( 3482 proxy, 3483 pubkey, 3484 vec![ 3485 TRANSPORT_KIND_NOSTR.to_owned(), 3486 TRANSPORT_KIND_RETICULUM.to_owned(), 3487 ], 3488 visibility, 3489 ) 3490 } 3491 3492 fn explicit_target_principal_with_kinds( 3493 proxy: &TransportPublish, 3494 pubkey: String, 3495 allowed_explicit_transport_kinds: Vec<String>, 3496 visibility: PublishJobVisibility, 3497 ) -> PublishPrincipal { 3498 proxy 3499 .store 3500 .create_principal(PublishPrincipalInit { 3501 label: "explicit-target-tester".to_owned(), 3502 token_hash: hash_bearer_token(generate_bearer_token().as_str()), 3503 allowed_pubkeys: vec![pubkey], 3504 allowed_kinds: vec![30_402], 3505 allowed_target_policies: vec![TargetPolicyName::ExplicitTargets], 3506 allowed_explicit_transport_kinds, 3507 allowed_nostr_source_policies: Vec::new(), 3508 allow_request_targets: true, 3509 job_visibility: visibility, 3510 expires_at_unix: None, 3511 }) 3512 .expect("principal") 3513 } 3514 3515 fn config_with_defaults(relays: Vec<&str>) -> TransportPublishConfig { 3516 TransportPublishConfig { 3517 nostr: TransportPublishNostrConfig { 3518 daemon_default_relays: relays.into_iter().map(str::to_owned).collect(), 3519 ..TransportPublishNostrConfig::default() 3520 }, 3521 ..TransportPublishConfig::default() 3522 } 3523 } 3524 3525 #[test] 3526 fn explicit_target_principals_require_canonical_allowed_transport_kinds() { 3527 let store = TransportPublishStore::memory().expect("store"); 3528 let base = PublishPrincipalInit { 3529 label: "explicit-target-tester".to_owned(), 3530 token_hash: hash_bearer_token(generate_bearer_token().as_str()), 3531 allowed_pubkeys: vec!["a".repeat(64)], 3532 allowed_kinds: vec![30_402], 3533 allowed_target_policies: vec![TargetPolicyName::ExplicitTargets], 3534 allowed_explicit_transport_kinds: Vec::new(), 3535 allowed_nostr_source_policies: Vec::new(), 3536 allow_request_targets: true, 3537 job_visibility: PublishJobVisibility::Own, 3538 expires_at_unix: None, 3539 }; 3540 3541 assert!(matches!( 3542 store.create_principal(base.clone()), 3543 Err(TransportPublishError::InvalidScope(message)) 3544 if message.contains("allowed explicit transport kind") 3545 )); 3546 3547 let mut uppercase = base.clone(); 3548 uppercase.allowed_explicit_transport_kinds = vec!["Nostr".to_owned()]; 3549 assert!(matches!( 3550 store.create_principal(uppercase), 3551 Err(TransportPublishError::InvalidScope(message)) 3552 if message.contains("explicit transport kind") 3553 )); 3554 3555 let mut duplicate = base.clone(); 3556 duplicate.allowed_explicit_transport_kinds = vec![ 3557 TRANSPORT_KIND_NOSTR.to_owned(), 3558 TRANSPORT_KIND_NOSTR.to_owned(), 3559 ]; 3560 assert!(matches!( 3561 store.create_principal(duplicate), 3562 Err(TransportPublishError::InvalidScope(message)) 3563 if message.contains("canonical and unique") 3564 )); 3565 3566 let mut removed_execution_kind = base.clone(); 3567 removed_execution_kind.allowed_explicit_transport_kinds = 3568 vec![removed_proxy_transport_kind_string()]; 3569 assert!(matches!( 3570 store.create_principal(removed_execution_kind), 3571 Err(TransportPublishError::InvalidScope(message)) 3572 if message.contains("unknown explicit transport kind") 3573 )); 3574 3575 let mut nostr_policy_with_explicit_kinds = base; 3576 nostr_policy_with_explicit_kinds.allowed_target_policies = vec![TargetPolicyName::Nostr]; 3577 nostr_policy_with_explicit_kinds.allowed_explicit_transport_kinds = 3578 vec![TRANSPORT_KIND_NOSTR.to_owned()]; 3579 nostr_policy_with_explicit_kinds.allowed_nostr_source_policies = 3580 vec![NostrTargetSourcePolicy::DaemonDefaultOnly]; 3581 assert!(matches!( 3582 store.create_principal(nostr_policy_with_explicit_kinds), 3583 Err(TransportPublishError::InvalidScope(message)) 3584 if message.contains("without explicit target policy") 3585 )); 3586 } 3587 3588 #[derive(Default)] 3589 struct StaticPublishRelayResolver { 3590 results: BTreeMap<String, Result<Vec<IpAddr>, String>>, 3591 } 3592 3593 impl StaticPublishRelayResolver { 3594 fn new() -> Self { 3595 Self::default() 3596 } 3597 3598 fn with_addresses(mut self, url: &str, addresses: Vec<IpAddr>) -> Self { 3599 self.results.insert(url.to_owned(), Ok(addresses)); 3600 self 3601 } 3602 3603 fn with_failure(mut self, url: &str, error: &str) -> Self { 3604 self.results.insert(url.to_owned(), Err(error.to_owned())); 3605 self 3606 } 3607 } 3608 3609 impl super::PublishRelayResolver for StaticPublishRelayResolver { 3610 fn resolve<'a>( 3611 &'a self, 3612 url: &'a super::RadrootsRelayUrl, 3613 ) -> super::PublishRelayResolveFuture<'a> { 3614 Box::pin(async move { 3615 match self.results.get(url.as_str()) { 3616 Some(Ok(addresses)) => Ok(addresses.clone()), 3617 Some(Err(error)) => Err(std::io::Error::other(error.clone())), 3618 None => Ok(vec![IpAddr::V4(Ipv4Addr::new(93, 184, 216, 34))]), 3619 } 3620 }) 3621 } 3622 } 3623 3624 struct StaticPublishAuthorRelayDiscovery { 3625 relays: Vec<String>, 3626 } 3627 3628 impl StaticPublishAuthorRelayDiscovery { 3629 fn new(relays: Vec<&str>) -> Self { 3630 Self { 3631 relays: relays.into_iter().map(str::to_owned).collect(), 3632 } 3633 } 3634 } 3635 3636 impl super::PublishAuthorRelayDiscovery for StaticPublishAuthorRelayDiscovery { 3637 fn fetch_author_write_relays<'a>( 3638 &'a self, 3639 _pubkey: &'a str, 3640 _discovery_targets: Vec<super::ResolvedPublishRelay>, 3641 _connect_timeout_secs: u64, 3642 ) -> super::PublishAuthorRelayDiscoveryFuture<'a> { 3643 let relays = self.relays.clone(); 3644 Box::pin(async move { Ok(relays) }) 3645 } 3646 } 3647 3648 #[test] 3649 fn token_generation_and_hashing_do_not_store_plaintext() { 3650 let token = generate_bearer_token(); 3651 assert!(token.starts_with("rrd_tp_")); 3652 let hash = hash_bearer_token(token.as_str()); 3653 assert!(hash.starts_with("sha256:")); 3654 assert!(!hash.contains(token.as_str())); 3655 } 3656 3657 #[test] 3658 fn nostr_source_policy_parser_accepts_contract_values() { 3659 assert_eq!( 3660 parse_nostr_source_policy("explicit_only").expect("policy"), 3661 NostrTargetSourcePolicy::ExplicitOnly 3662 ); 3663 assert!(parse_nostr_source_policy("unknown").is_err()); 3664 } 3665 3666 #[test] 3667 fn storage_authenticates_hashed_tokens_and_scopes_jobs() { 3668 let store = TransportPublishStore::memory().expect("store"); 3669 let token = generate_bearer_token(); 3670 let token_hash = hash_bearer_token(token.as_str()); 3671 let accepted_identity = DaemonIdentity::generate(); 3672 let denied_identity = DaemonIdentity::generate(); 3673 let principal = store 3674 .create_principal(PublishPrincipalInit { 3675 label: "tester".to_owned(), 3676 token_hash: token_hash.clone(), 3677 allowed_pubkeys: vec![accepted_identity.public_key_hex()], 3678 allowed_kinds: vec![30_402], 3679 allowed_target_policies: vec![TargetPolicyName::Nostr], 3680 allowed_explicit_transport_kinds: Vec::new(), 3681 allowed_nostr_source_policies: vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 3682 allow_request_targets: false, 3683 job_visibility: PublishJobVisibility::Own, 3684 expires_at_unix: None, 3685 }) 3686 .expect("principal"); 3687 assert_eq!( 3688 store 3689 .principal_for_token_hash(token_hash.as_str()) 3690 .expect("lookup") 3691 .expect("principal") 3692 .principal_id, 3693 principal.principal_id 3694 ); 3695 let denied = publish_request( 3696 signed_event(&denied_identity, "{}"), 3697 Vec::new(), 3698 NostrTargetSourcePolicy::DaemonDefaultOnly, 3699 DeliveryPolicy::Any, 3700 None, 3701 ); 3702 let denied_signed = 3703 super::signed_event_from_raw_json(denied.raw_event_json.as_str()).expect("denied raw"); 3704 assert!(principal.allows_event(&denied_signed, &denied).is_err()); 3705 3706 let accepted = publish_request( 3707 signed_event(&accepted_identity, "{}"), 3708 Vec::new(), 3709 NostrTargetSourcePolicy::DaemonDefaultOnly, 3710 DeliveryPolicy::Any, 3711 None, 3712 ); 3713 let accepted_signed = super::signed_event_from_raw_json(accepted.raw_event_json.as_str()) 3714 .expect("accepted raw"); 3715 principal 3716 .allows_event(&accepted_signed, &accepted) 3717 .expect("scope"); 3718 let response = store 3719 .record_publish_job(PublishJobInsert { 3720 principal_id: principal.principal_id.clone(), 3721 idempotency_key: Some("idem-1".to_owned()), 3722 event: PublishEventMetadata::from_signed_event(&accepted_signed), 3723 request: accepted.clone(), 3724 request_fingerprint: "fingerprint-1".to_owned(), 3725 effective_target_count: 1, 3726 target_snapshots: vec![interrupted_target_snapshot( 3727 RELAY_PRIMARY, 3728 TargetSource::DaemonDefault, 3729 )], 3730 }) 3731 .expect("record job"); 3732 assert!(!response.deduplicated); 3733 let duplicate = store 3734 .record_publish_job(PublishJobInsert { 3735 principal_id: principal.principal_id.clone(), 3736 idempotency_key: Some("idem-1".to_owned()), 3737 event: PublishEventMetadata::from_signed_event(&accepted_signed), 3738 request: accepted, 3739 request_fingerprint: "fingerprint-1".to_owned(), 3740 effective_target_count: 1, 3741 target_snapshots: vec![interrupted_target_snapshot( 3742 RELAY_PRIMARY, 3743 TargetSource::DaemonDefault, 3744 )], 3745 }) 3746 .expect("dedupe"); 3747 assert!(duplicate.deduplicated); 3748 assert_eq!(duplicate.job.job_id, response.job.job_id); 3749 assert_eq!( 3750 store 3751 .list_jobs_for_principal(&principal, 50) 3752 .expect("jobs") 3753 .len(), 3754 1 3755 ); 3756 } 3757 3758 #[test] 3759 fn store_egress_rejects_malformed_target_counts_for_get_list_and_dedupe() { 3760 let store = TransportPublishStore::memory().expect("store"); 3761 let pubkey = "a".repeat(64); 3762 let principal = store_principal(&store, pubkey.as_str()); 3763 let request = request(pubkey.as_str(), 30_402); 3764 let response = store 3765 .record_publish_job(PublishJobInsert { 3766 principal_id: principal.principal_id.clone(), 3767 idempotency_key: Some("idem-invalid-target-count".to_owned()), 3768 event: event_metadata(pubkey.as_str(), 30_402), 3769 request: request.clone(), 3770 request_fingerprint: "fingerprint-invalid-target-count".to_owned(), 3771 effective_target_count: 1, 3772 target_snapshots: vec![accepted_target_outcome( 3773 RELAY_PRIMARY, 3774 TargetSource::DaemonDefault, 3775 )], 3776 }) 3777 .expect("record job"); 3778 store 3779 .complete_publish_job( 3780 response.job.job_id.as_str(), 3781 JobStatus::DeliverySatisfied, 3782 vec![accepted_target_outcome( 3783 RELAY_PRIMARY, 3784 TargetSource::DaemonDefault, 3785 )], 3786 None, 3787 ) 3788 .expect("complete job"); 3789 { 3790 let mut connection = store 3791 .inner 3792 .lock() 3793 .unwrap_or_else(std::sync::PoisonError::into_inner); 3794 super::block_on_sqlite( 3795 sqlx::query( 3796 "UPDATE transport_publish_jobs SET effective_target_count = 2 WHERE job_id = ?1", 3797 ) 3798 .bind(response.job.job_id.as_str()) 3799 .execute(&mut *connection), 3800 ) 3801 .expect("corrupt target count"); 3802 } 3803 3804 assert_invalid_job_state( 3805 store 3806 .job_by_id(response.job.job_id.as_str()) 3807 .expect_err("invalid get"), 3808 "job target_count 2 does not match 1 target outcomes", 3809 ); 3810 assert_invalid_job_state( 3811 store 3812 .list_jobs_for_principal(&principal, 50) 3813 .expect_err("invalid list"), 3814 "job target_count 2 does not match 1 target outcomes", 3815 ); 3816 assert_invalid_job_state( 3817 store 3818 .record_publish_job(PublishJobInsert { 3819 principal_id: principal.principal_id.clone(), 3820 idempotency_key: Some("idem-invalid-target-count".to_owned()), 3821 event: event_metadata(pubkey.as_str(), 30_402), 3822 request, 3823 request_fingerprint: "fingerprint-invalid-target-count".to_owned(), 3824 effective_target_count: 1, 3825 target_snapshots: vec![accepted_target_outcome( 3826 RELAY_PRIMARY, 3827 TargetSource::DaemonDefault, 3828 )], 3829 }) 3830 .expect_err("invalid dedupe"), 3831 "job target_count 2 does not match 1 target outcomes", 3832 ); 3833 } 3834 3835 #[test] 3836 fn store_egress_rejects_impossible_persisted_event_kind_values() { 3837 for invalid_kind in [-1_i64, i64::from(u32::MAX) + 1] { 3838 let store = TransportPublishStore::memory().expect("store"); 3839 let pubkey = "a".repeat(64); 3840 let principal = store_principal(&store, pubkey.as_str()); 3841 let request = request(pubkey.as_str(), 30_402); 3842 let response = store 3843 .record_publish_job(PublishJobInsert { 3844 principal_id: principal.principal_id.clone(), 3845 idempotency_key: Some(format!("idem-invalid-kind-{invalid_kind}")), 3846 event: event_metadata(pubkey.as_str(), 30_402), 3847 request: request.clone(), 3848 request_fingerprint: format!("fingerprint-invalid-kind-{invalid_kind}"), 3849 effective_target_count: 1, 3850 target_snapshots: vec![accepted_target_outcome( 3851 RELAY_PRIMARY, 3852 TargetSource::DaemonDefault, 3853 )], 3854 }) 3855 .expect("record job"); 3856 store 3857 .complete_publish_job( 3858 response.job.job_id.as_str(), 3859 JobStatus::DeliverySatisfied, 3860 vec![accepted_target_outcome( 3861 RELAY_PRIMARY, 3862 TargetSource::DaemonDefault, 3863 )], 3864 None, 3865 ) 3866 .expect("complete job"); 3867 { 3868 let mut connection = store 3869 .inner 3870 .lock() 3871 .unwrap_or_else(std::sync::PoisonError::into_inner); 3872 super::block_on_sqlite( 3873 sqlx::query( 3874 "UPDATE transport_publish_jobs SET event_kind = ?2 WHERE job_id = ?1", 3875 ) 3876 .bind(response.job.job_id.as_str()) 3877 .bind(invalid_kind) 3878 .execute(&mut *connection), 3879 ) 3880 .expect("corrupt event kind"); 3881 } 3882 3883 assert_storage_integer_range_error( 3884 store 3885 .job_by_id(response.job.job_id.as_str()) 3886 .expect_err("invalid get"), 3887 "event_kind integer value", 3888 ); 3889 assert_storage_integer_range_error( 3890 store 3891 .list_jobs_for_principal(&principal, 50) 3892 .expect_err("invalid list"), 3893 "event_kind integer value", 3894 ); 3895 assert_storage_integer_range_error( 3896 store 3897 .record_publish_job(PublishJobInsert { 3898 principal_id: principal.principal_id.clone(), 3899 idempotency_key: Some(format!("idem-invalid-kind-{invalid_kind}")), 3900 event: event_metadata(pubkey.as_str(), 30_402), 3901 request, 3902 request_fingerprint: format!("fingerprint-invalid-kind-{invalid_kind}"), 3903 effective_target_count: 1, 3904 target_snapshots: vec![accepted_target_outcome( 3905 RELAY_PRIMARY, 3906 TargetSource::DaemonDefault, 3907 )], 3908 }) 3909 .expect_err("invalid dedupe"), 3910 "event_kind integer value", 3911 ); 3912 } 3913 } 3914 3915 #[test] 3916 fn store_egress_rejects_negative_persisted_effective_target_count() { 3917 let store = TransportPublishStore::memory().expect("store"); 3918 let pubkey = "a".repeat(64); 3919 let principal = store_principal(&store, pubkey.as_str()); 3920 let request = request(pubkey.as_str(), 30_402); 3921 let response = store 3922 .record_publish_job(PublishJobInsert { 3923 principal_id: principal.principal_id.clone(), 3924 idempotency_key: Some("idem-negative-target-count".to_owned()), 3925 event: event_metadata(pubkey.as_str(), 30_402), 3926 request: request.clone(), 3927 request_fingerprint: "fingerprint-negative-target-count".to_owned(), 3928 effective_target_count: 1, 3929 target_snapshots: vec![accepted_target_outcome( 3930 RELAY_PRIMARY, 3931 TargetSource::DaemonDefault, 3932 )], 3933 }) 3934 .expect("record job"); 3935 store 3936 .complete_publish_job( 3937 response.job.job_id.as_str(), 3938 JobStatus::DeliverySatisfied, 3939 vec![accepted_target_outcome( 3940 RELAY_PRIMARY, 3941 TargetSource::DaemonDefault, 3942 )], 3943 None, 3944 ) 3945 .expect("complete job"); 3946 { 3947 let mut connection = store 3948 .inner 3949 .lock() 3950 .unwrap_or_else(std::sync::PoisonError::into_inner); 3951 super::block_on_sqlite( 3952 sqlx::query( 3953 "UPDATE transport_publish_jobs SET effective_target_count = -1 WHERE job_id = ?1", 3954 ) 3955 .bind(response.job.job_id.as_str()) 3956 .execute(&mut *connection), 3957 ) 3958 .expect("corrupt target count"); 3959 } 3960 3961 assert_storage_integer_range_error( 3962 store 3963 .job_by_id(response.job.job_id.as_str()) 3964 .expect_err("invalid get"), 3965 "effective_target_count integer value -1 is outside usize range", 3966 ); 3967 assert_storage_integer_range_error( 3968 store 3969 .list_jobs_for_principal(&principal, 50) 3970 .expect_err("invalid list"), 3971 "effective_target_count integer value -1 is outside usize range", 3972 ); 3973 assert_storage_integer_range_error( 3974 store 3975 .record_publish_job(PublishJobInsert { 3976 principal_id: principal.principal_id.clone(), 3977 idempotency_key: Some("idem-negative-target-count".to_owned()), 3978 event: event_metadata(pubkey.as_str(), 30_402), 3979 request, 3980 request_fingerprint: "fingerprint-negative-target-count".to_owned(), 3981 effective_target_count: 1, 3982 target_snapshots: vec![accepted_target_outcome( 3983 RELAY_PRIMARY, 3984 TargetSource::DaemonDefault, 3985 )], 3986 }) 3987 .expect_err("invalid dedupe"), 3988 "effective_target_count integer value -1 is outside usize range", 3989 ); 3990 } 3991 3992 #[test] 3993 fn store_egress_rejects_negative_persisted_target_latency() { 3994 let store = TransportPublishStore::memory().expect("store"); 3995 let pubkey = "a".repeat(64); 3996 let principal = store_principal(&store, pubkey.as_str()); 3997 let request = request(pubkey.as_str(), 30_402); 3998 let response = store 3999 .record_publish_job(PublishJobInsert { 4000 principal_id: principal.principal_id.clone(), 4001 idempotency_key: Some("idem-negative-latency".to_owned()), 4002 event: event_metadata(pubkey.as_str(), 30_402), 4003 request: request.clone(), 4004 request_fingerprint: "fingerprint-negative-latency".to_owned(), 4005 effective_target_count: 1, 4006 target_snapshots: vec![accepted_target_outcome( 4007 RELAY_PRIMARY, 4008 TargetSource::DaemonDefault, 4009 )], 4010 }) 4011 .expect("record job"); 4012 store 4013 .complete_publish_job( 4014 response.job.job_id.as_str(), 4015 JobStatus::DeliverySatisfied, 4016 vec![accepted_target_outcome( 4017 RELAY_PRIMARY, 4018 TargetSource::DaemonDefault, 4019 )], 4020 None, 4021 ) 4022 .expect("complete job"); 4023 { 4024 let mut connection = store 4025 .inner 4026 .lock() 4027 .unwrap_or_else(std::sync::PoisonError::into_inner); 4028 super::block_on_sqlite( 4029 sqlx::query( 4030 "UPDATE transport_publish_target_results SET latency_ms = -5 WHERE job_id = ?1", 4031 ) 4032 .bind(response.job.job_id.as_str()) 4033 .execute(&mut *connection), 4034 ) 4035 .expect("corrupt latency"); 4036 } 4037 4038 assert_storage_integer_range_error( 4039 store 4040 .job_by_id(response.job.job_id.as_str()) 4041 .expect_err("invalid get"), 4042 "latency_ms integer value -5 is outside u64 range", 4043 ); 4044 assert_storage_integer_range_error( 4045 store 4046 .list_jobs_for_principal(&principal, 50) 4047 .expect_err("invalid list"), 4048 "latency_ms integer value -5 is outside u64 range", 4049 ); 4050 assert_storage_integer_range_error( 4051 store 4052 .record_publish_job(PublishJobInsert { 4053 principal_id: principal.principal_id.clone(), 4054 idempotency_key: Some("idem-negative-latency".to_owned()), 4055 event: event_metadata(pubkey.as_str(), 30_402), 4056 request, 4057 request_fingerprint: "fingerprint-negative-latency".to_owned(), 4058 effective_target_count: 1, 4059 target_snapshots: vec![accepted_target_outcome( 4060 RELAY_PRIMARY, 4061 TargetSource::DaemonDefault, 4062 )], 4063 }) 4064 .expect_err("invalid dedupe"), 4065 "latency_ms integer value -5 is outside u64 range", 4066 ); 4067 } 4068 4069 #[test] 4070 fn store_egress_rejects_explicit_target_outcome_drift() { 4071 let store = TransportPublishStore::memory().expect("store"); 4072 let pubkey = "a".repeat(64); 4073 let principal = store_principal(&store, pubkey.as_str()); 4074 let mut request = request(pubkey.as_str(), 30_402); 4075 request.target_policy = TargetPolicy::explicit_targets(vec![Target::nostr(RELAY_PRIMARY)]); 4076 let response = store 4077 .record_publish_job(PublishJobInsert { 4078 principal_id: principal.principal_id.clone(), 4079 idempotency_key: Some("idem-explicit-drift".to_owned()), 4080 event: event_metadata(pubkey.as_str(), 30_402), 4081 request, 4082 request_fingerprint: "fingerprint-explicit-drift".to_owned(), 4083 effective_target_count: 1, 4084 target_snapshots: vec![accepted_target_outcome( 4085 RELAY_PRIMARY, 4086 TargetSource::Request, 4087 )], 4088 }) 4089 .expect("record job"); 4090 store 4091 .complete_publish_job( 4092 response.job.job_id.as_str(), 4093 JobStatus::DeliverySatisfied, 4094 vec![accepted_target_outcome( 4095 RELAY_SECONDARY, 4096 TargetSource::Request, 4097 )], 4098 None, 4099 ) 4100 .expect("complete drifted job"); 4101 4102 assert_invalid_job_state( 4103 store 4104 .job_by_id(response.job.job_id.as_str()) 4105 .expect_err("invalid explicit drift"), 4106 "transport target outcome 0 does not match explicit target policy", 4107 ); 4108 } 4109 4110 #[test] 4111 fn store_egress_rejects_explicit_target_scope_drift() { 4112 let store = TransportPublishStore::memory().expect("store"); 4113 let pubkey = "a".repeat(64); 4114 let principal = store_principal(&store, pubkey.as_str()); 4115 let mut request = request(pubkey.as_str(), 30_402); 4116 request.target_policy = TargetPolicy::explicit_targets(vec![ 4117 Target::nostr(RELAY_PRIMARY) 4118 .with_scope("farm.local") 4119 .with_label("Farm relay"), 4120 ]); 4121 let response = store 4122 .record_publish_job(PublishJobInsert { 4123 principal_id: principal.principal_id.clone(), 4124 idempotency_key: Some("idem-explicit-scope-drift".to_owned()), 4125 event: event_metadata(pubkey.as_str(), 30_402), 4126 request, 4127 request_fingerprint: "fingerprint-explicit-scope-drift".to_owned(), 4128 effective_target_count: 1, 4129 target_snapshots: vec![scoped_target_outcome( 4130 accepted_target_outcome(RELAY_PRIMARY, TargetSource::Request), 4131 "farm.local", 4132 Some("Farm relay"), 4133 )], 4134 }) 4135 .expect("record job"); 4136 store 4137 .complete_publish_job( 4138 response.job.job_id.as_str(), 4139 JobStatus::DeliverySatisfied, 4140 vec![scoped_target_outcome( 4141 accepted_target_outcome(RELAY_PRIMARY, TargetSource::Request), 4142 "farm.remote", 4143 Some("Farm relay"), 4144 )], 4145 None, 4146 ) 4147 .expect("complete drifted job"); 4148 4149 assert_invalid_job_state( 4150 store 4151 .job_by_id(response.job.job_id.as_str()) 4152 .expect_err("invalid explicit scope drift"), 4153 "transport target outcome 0 does not match explicit target policy", 4154 ); 4155 } 4156 4157 #[test] 4158 fn store_open_recovers_interrupted_publishing_jobs() { 4159 let directory = tempfile::tempdir().expect("tempdir"); 4160 let database_path = directory.path().join("publish-proxy.sqlite"); 4161 let token_hash = hash_bearer_token(generate_bearer_token().as_str()); 4162 let pubkey = "a".repeat(64); 4163 let request = request(pubkey.as_str(), 30_402); 4164 let (job_id, principal) = { 4165 let store = TransportPublishStore::open(database_path.clone()).expect("store"); 4166 let principal = store 4167 .create_principal(PublishPrincipalInit { 4168 label: "tester".to_owned(), 4169 token_hash, 4170 allowed_pubkeys: vec![pubkey.clone()], 4171 allowed_kinds: vec![30_402], 4172 allowed_target_policies: vec![TargetPolicyName::Nostr], 4173 allowed_explicit_transport_kinds: Vec::new(), 4174 allowed_nostr_source_policies: vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 4175 allow_request_targets: false, 4176 job_visibility: PublishJobVisibility::Own, 4177 expires_at_unix: None, 4178 }) 4179 .expect("principal"); 4180 let response = store 4181 .record_publish_job(PublishJobInsert { 4182 principal_id: principal.principal_id.clone(), 4183 idempotency_key: Some("idem-interrupted".to_owned()), 4184 event: event_metadata(pubkey.as_str(), 30_402), 4185 request, 4186 request_fingerprint: "fingerprint-interrupted".to_owned(), 4187 effective_target_count: 1, 4188 target_snapshots: vec![interrupted_target_snapshot( 4189 RELAY_PRIMARY, 4190 TargetSource::DaemonDefault, 4191 )], 4192 }) 4193 .expect("record job"); 4194 assert_eq!(response.job.status, JobStatus::Publishing); 4195 (response.job.job_id, principal) 4196 }; 4197 4198 let reopened = TransportPublishStore::open(database_path).expect("reopen store"); 4199 let recovered = reopened.job_by_id(job_id.as_str()).expect("recovered job"); 4200 assert_eq!(recovered.status, JobStatus::DeliveryUnsatisfiedRetryable); 4201 assert_eq!( 4202 recovered.last_error.as_deref(), 4203 Some("publish_attempt_interrupted") 4204 ); 4205 assert!(recovered.completed_at_ms.is_some()); 4206 assert_eq!(recovered.targets.len(), 1); 4207 assert_eq!( 4208 recovered.targets[0].outcome_kind, 4209 OutcomeKind::ConnectionFailed 4210 ); 4211 recovered.validate().expect("valid recovered job"); 4212 let listed = reopened 4213 .list_jobs_for_principal(&principal, 50) 4214 .expect("listed jobs"); 4215 assert_eq!(listed.len(), 1); 4216 listed[0].validate().expect("valid listed recovered job"); 4217 } 4218 4219 #[test] 4220 fn store_open_recovers_interrupted_scoped_explicit_target_metadata() { 4221 let directory = tempfile::tempdir().expect("tempdir"); 4222 let database_path = directory.path().join("publish-proxy-scoped.sqlite"); 4223 let pubkey = "a".repeat(64); 4224 let (job_id, principal) = { 4225 let store = TransportPublishStore::open(database_path.clone()).expect("store"); 4226 let principal = store_principal(&store, pubkey.as_str()); 4227 let mut request = request(pubkey.as_str(), 30_402); 4228 request.target_policy = TargetPolicy::explicit_targets(vec![ 4229 Target::nostr(RELAY_PRIMARY) 4230 .with_scope("farm.local") 4231 .with_label("Farm relay"), 4232 ]); 4233 let response = store 4234 .record_publish_job(PublishJobInsert { 4235 principal_id: principal.principal_id.clone(), 4236 idempotency_key: Some("idem-interrupted-scoped".to_owned()), 4237 event: event_metadata(pubkey.as_str(), 30_402), 4238 request, 4239 request_fingerprint: "fingerprint-interrupted-scoped".to_owned(), 4240 effective_target_count: 1, 4241 target_snapshots: vec![scoped_target_outcome( 4242 interrupted_target_snapshot(RELAY_PRIMARY, TargetSource::Request), 4243 "farm.local", 4244 Some("Farm relay"), 4245 )], 4246 }) 4247 .expect("record job"); 4248 (response.job.job_id, principal) 4249 }; 4250 4251 let reopened = TransportPublishStore::open(database_path).expect("reopen store"); 4252 let recovered = reopened.job_by_id(job_id.as_str()).expect("recovered job"); 4253 assert_eq!(recovered.status, JobStatus::DeliveryUnsatisfiedRetryable); 4254 assert_eq!(recovered.targets.len(), 1); 4255 assert_eq!( 4256 recovered.targets[0].target_scope.as_deref(), 4257 Some("farm.local") 4258 ); 4259 assert_eq!( 4260 recovered.targets[0].target_label.as_deref(), 4261 Some("Farm relay") 4262 ); 4263 recovered.validate().expect("valid recovered job"); 4264 let listed = reopened 4265 .list_jobs_for_principal(&principal, 50) 4266 .expect("listed jobs"); 4267 assert_eq!( 4268 listed[0].targets[0].target_scope.as_deref(), 4269 Some("farm.local") 4270 ); 4271 assert_eq!( 4272 listed[0].targets[0].target_label.as_deref(), 4273 Some("Farm relay") 4274 ); 4275 } 4276 4277 #[test] 4278 fn store_egress_rejects_recovered_explicit_target_snapshot_drift() { 4279 let directory = tempfile::tempdir().expect("tempdir"); 4280 let database_path = directory.path().join("publish-proxy-drift.sqlite"); 4281 let pubkey = "a".repeat(64); 4282 let (job_id, principal) = { 4283 let store = TransportPublishStore::open(database_path.clone()).expect("store"); 4284 let principal = store_principal(&store, pubkey.as_str()); 4285 let mut request = request(pubkey.as_str(), 30_402); 4286 request.target_policy = 4287 TargetPolicy::explicit_targets(vec![Target::nostr(RELAY_PRIMARY)]); 4288 let response = store 4289 .record_publish_job(PublishJobInsert { 4290 principal_id: principal.principal_id.clone(), 4291 idempotency_key: Some("idem-recovered-drift".to_owned()), 4292 event: event_metadata(pubkey.as_str(), 30_402), 4293 request, 4294 request_fingerprint: "fingerprint-recovered-drift".to_owned(), 4295 effective_target_count: 1, 4296 target_snapshots: vec![accepted_target_outcome( 4297 RELAY_SECONDARY, 4298 TargetSource::Request, 4299 )], 4300 }) 4301 .expect("record job"); 4302 (response.job.job_id, principal) 4303 }; 4304 4305 let reopened = TransportPublishStore::open(database_path).expect("reopen store"); 4306 assert_invalid_job_state( 4307 reopened 4308 .job_by_id(job_id.as_str()) 4309 .expect_err("invalid get"), 4310 "transport target outcome 0 does not match explicit target policy", 4311 ); 4312 assert_invalid_job_state( 4313 reopened 4314 .list_jobs_for_principal(&principal, 50) 4315 .expect_err("invalid list"), 4316 "transport target outcome 0 does not match explicit target policy", 4317 ); 4318 } 4319 4320 #[test] 4321 fn store_open_rejects_interrupted_jobs_without_target_snapshots() { 4322 let directory = tempfile::tempdir().expect("tempdir"); 4323 let database_path = directory 4324 .path() 4325 .join("publish-proxy-missing-snapshot.sqlite"); 4326 let token_hash = hash_bearer_token(generate_bearer_token().as_str()); 4327 let pubkey = "a".repeat(64); 4328 let request = request(pubkey.as_str(), 30_402); 4329 { 4330 let store = TransportPublishStore::open(database_path.clone()).expect("store"); 4331 let principal = store 4332 .create_principal(PublishPrincipalInit { 4333 label: "tester".to_owned(), 4334 token_hash, 4335 allowed_pubkeys: vec![pubkey.clone()], 4336 allowed_kinds: vec![30_402], 4337 allowed_target_policies: vec![TargetPolicyName::Nostr], 4338 allowed_explicit_transport_kinds: Vec::new(), 4339 allowed_nostr_source_policies: vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 4340 allow_request_targets: false, 4341 job_visibility: PublishJobVisibility::Own, 4342 expires_at_unix: None, 4343 }) 4344 .expect("principal"); 4345 let now = super::current_unix_millis(); 4346 let mut connection = store 4347 .inner 4348 .lock() 4349 .unwrap_or_else(std::sync::PoisonError::into_inner); 4350 super::block_on_sqlite( 4351 sqlx::query( 4352 r#" 4353 INSERT INTO transport_publish_jobs ( 4354 job_id, 4355 principal_id, 4356 idempotency_key, 4357 request_fingerprint, 4358 status, 4359 event_id, 4360 event_pubkey, 4361 event_kind, 4362 target_policy_json, 4363 delivery_policy_json, 4364 requested_target_count, 4365 effective_target_count, 4366 request_json, 4367 requested_at_ms, 4368 updated_at_ms 4369 ) 4370 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15) 4371 "#, 4372 ) 4373 .bind("job-missing-snapshot") 4374 .bind(principal.principal_id.as_str()) 4375 .bind("idem-missing-snapshot") 4376 .bind("fingerprint-missing-snapshot") 4377 .bind(serde_json::to_string(&JobStatus::Publishing).expect("status")) 4378 .bind("0".repeat(64)) 4379 .bind(pubkey.as_str()) 4380 .bind(30_402_i64) 4381 .bind(serde_json::to_string(&request.target_policy).expect("target policy")) 4382 .bind(serde_json::to_string(&request.delivery_policy).expect("delivery policy")) 4383 .bind( 4384 super::storage_count_i64( 4385 request.target_policy.request_target_count(), 4386 "requested_target_count", 4387 ) 4388 .expect("requested target count"), 4389 ) 4390 .bind(1_i64) 4391 .bind(serde_json::to_string(&request).expect("request")) 4392 .bind(now) 4393 .bind(now) 4394 .execute(&mut *connection), 4395 ) 4396 .expect("insert historical job"); 4397 } 4398 4399 let reopened = TransportPublishStore::open(database_path).expect("reopen store"); 4400 let recovered = reopened 4401 .job_by_id("job-missing-snapshot") 4402 .expect("recovered job"); 4403 assert_eq!(recovered.status, JobStatus::Rejected); 4404 assert_eq!( 4405 recovered.last_error.as_deref(), 4406 Some("publish_attempt_interrupted_missing_target_snapshot") 4407 ); 4408 assert_eq!(recovered.target_count, 0); 4409 assert!(recovered.targets.is_empty()); 4410 recovered.validate().expect("valid rejected recovered job"); 4411 } 4412 4413 #[test] 4414 fn transport_store_open_validates_current_principal_schema() { 4415 let directory = tempfile::tempdir().expect("tempdir"); 4416 let database_path = directory.path().join("publish-proxy-current.sqlite"); 4417 let store = TransportPublishStore::open(database_path).expect("open current schema"); 4418 let mut connection = store 4419 .inner 4420 .lock() 4421 .unwrap_or_else(std::sync::PoisonError::into_inner); 4422 let columns = test_query_column_names(&mut connection, "transport_publish_principals"); 4423 assert!( 4424 columns 4425 .iter() 4426 .any(|column| column == "allowed_explicit_transport_kinds_json") 4427 ); 4428 let version = 4429 super::transport_publish_schema_version(&mut connection).expect("user version"); 4430 assert_eq!(version, SCHEMA_VERSION); 4431 } 4432 4433 #[test] 4434 fn transport_store_open_rejects_legacy_principal_schema_without_explicit_kind_allowlist() { 4435 let directory = tempfile::tempdir().expect("tempdir"); 4436 let database_path = directory.path().join("publish-proxy-v1.sqlite"); 4437 let token_hash = hash_bearer_token(generate_bearer_token().as_str()); 4438 { 4439 let mut connection = open_test_database(database_path.as_path()); 4440 super::execute_raw_sql( 4441 &mut connection, 4442 r#" 4443 CREATE TABLE transport_publish_principals ( 4444 principal_id TEXT PRIMARY KEY NOT NULL, 4445 label TEXT NOT NULL, 4446 token_hash TEXT NOT NULL UNIQUE, 4447 allowed_pubkeys_json TEXT NOT NULL, 4448 allowed_kinds_json TEXT NOT NULL, 4449 allowed_target_policies_json TEXT NOT NULL, 4450 allowed_nostr_source_policies_json TEXT NOT NULL, 4451 allow_request_targets INTEGER NOT NULL, 4452 job_visibility TEXT NOT NULL, 4453 expires_at_unix INTEGER, 4454 revoked_at_unix INTEGER, 4455 created_at_unix INTEGER NOT NULL 4456 ); 4457 "#, 4458 ) 4459 .expect("schema"); 4460 super::execute_sql( 4461 &mut connection, 4462 format!("PRAGMA user_version = {SCHEMA_VERSION}").as_str(), 4463 ) 4464 .expect("set schema version"); 4465 super::block_on_sqlite( 4466 sqlx::query( 4467 r#" 4468 INSERT INTO transport_publish_principals ( 4469 principal_id, 4470 label, 4471 token_hash, 4472 allowed_pubkeys_json, 4473 allowed_kinds_json, 4474 allowed_target_policies_json, 4475 allowed_nostr_source_policies_json, 4476 allow_request_targets, 4477 job_visibility, 4478 expires_at_unix, 4479 revoked_at_unix, 4480 created_at_unix 4481 ) 4482 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, NULL, NULL, ?10) 4483 "#, 4484 ) 4485 .bind("principal-v1") 4486 .bind("v1") 4487 .bind(token_hash.as_str()) 4488 .bind(serde_json::to_string(&vec!["a".repeat(64)]).expect("pubkeys")) 4489 .bind(serde_json::to_string(&vec![30_402]).expect("kinds")) 4490 .bind(serde_json::to_string(&vec![TargetPolicyName::Nostr]).expect("policies")) 4491 .bind( 4492 serde_json::to_string(&vec![NostrTargetSourcePolicy::DaemonDefaultOnly]) 4493 .expect("source policies"), 4494 ) 4495 .bind(false) 4496 .bind(PublishJobVisibility::Own.to_string()) 4497 .bind(1_i64) 4498 .execute(&mut connection), 4499 ) 4500 .expect("principal"); 4501 } 4502 4503 let error = match TransportPublishStore::open(database_path.clone()) { 4504 Ok(_) => panic!("legacy schema opened"), 4505 Err(error) => error, 4506 }; 4507 match error { 4508 TransportPublishError::Schema { table, detail } => { 4509 assert_eq!(table, "transport_publish_principals"); 4510 assert!(detail.contains("allowed_explicit_transport_kinds_json")); 4511 } 4512 error => panic!("unexpected error: {error}"), 4513 } 4514 4515 let mut connection = open_test_database(database_path.as_path()); 4516 let columns = test_query_column_names(&mut connection, "transport_publish_principals"); 4517 assert!( 4518 !columns 4519 .iter() 4520 .any(|column| column == "allowed_explicit_transport_kinds_json") 4521 ); 4522 assert_eq!( 4523 database_user_version(database_path.as_path()), 4524 SCHEMA_VERSION 4525 ); 4526 } 4527 4528 #[test] 4529 fn transport_store_open_rejects_legacy_schema_version() { 4530 let directory = tempfile::tempdir().expect("tempdir"); 4531 let database_path = directory.path().join("publish-proxy-v3.sqlite"); 4532 create_existing_schema( 4533 database_path.as_path(), 4534 TRANSPORT_PUBLISH_SCHEMA_SQL, 4535 Some(3), 4536 ); 4537 4538 assert_schema_error( 4539 open_schema_error(database_path.as_path()), 4540 "transport_publish_schema", 4541 "user_version", 4542 ); 4543 assert_eq!(database_user_version(database_path.as_path()), 3); 4544 } 4545 4546 #[test] 4547 fn transport_store_open_rejects_current_schema_missing_user_version() { 4548 let directory = tempfile::tempdir().expect("tempdir"); 4549 let database_path = directory 4550 .path() 4551 .join("publish-proxy-missing-version.sqlite"); 4552 create_existing_schema(database_path.as_path(), TRANSPORT_PUBLISH_SCHEMA_SQL, None); 4553 4554 assert_schema_error( 4555 open_schema_error(database_path.as_path()), 4556 "transport_publish_schema", 4557 "user_version", 4558 ); 4559 assert_eq!(database_user_version(database_path.as_path()), 0); 4560 } 4561 4562 #[test] 4563 fn transport_store_open_rejects_current_schema_without_token_hash_unique_index() { 4564 let directory = tempfile::tempdir().expect("tempdir"); 4565 let database_path = directory 4566 .path() 4567 .join("publish-proxy-missing-token-unique.sqlite"); 4568 let schema = schema_with_replacement( 4569 "token_hash TEXT NOT NULL UNIQUE", 4570 "token_hash TEXT NOT NULL", 4571 ); 4572 create_existing_schema( 4573 database_path.as_path(), 4574 schema.as_str(), 4575 Some(SCHEMA_VERSION), 4576 ); 4577 4578 assert_schema_error( 4579 open_schema_error(database_path.as_path()), 4580 "transport_publish_principals", 4581 "token_hash", 4582 ); 4583 assert_eq!( 4584 database_user_version(database_path.as_path()), 4585 SCHEMA_VERSION 4586 ); 4587 } 4588 4589 #[test] 4590 fn transport_store_open_rejects_current_schema_without_idempotency_unique_index() { 4591 let directory = tempfile::tempdir().expect("tempdir"); 4592 let database_path = directory 4593 .path() 4594 .join("publish-proxy-missing-idempotency-index.sqlite"); 4595 let schema = schema_with_replacement( 4596 r#"CREATE UNIQUE INDEX IF NOT EXISTS transport_publish_jobs_principal_idempotency_idx 4597 ON transport_publish_jobs(principal_id, idempotency_key) 4598 WHERE idempotency_key IS NOT NULL;"#, 4599 "", 4600 ); 4601 create_existing_schema( 4602 database_path.as_path(), 4603 schema.as_str(), 4604 Some(SCHEMA_VERSION), 4605 ); 4606 4607 assert_schema_error( 4608 open_schema_error(database_path.as_path()), 4609 "transport_publish_jobs", 4610 "principal_id, idempotency_key", 4611 ); 4612 assert_eq!( 4613 database_user_version(database_path.as_path()), 4614 SCHEMA_VERSION 4615 ); 4616 } 4617 4618 #[test] 4619 fn transport_store_open_rejects_current_target_results_schema_without_scope_not_null() { 4620 let directory = tempfile::tempdir().expect("tempdir"); 4621 let database_path = directory 4622 .path() 4623 .join("publish-proxy-target-scope-nullable.sqlite"); 4624 let schema = schema_with_replacement("target_scope TEXT NOT NULL", "target_scope TEXT"); 4625 create_existing_schema( 4626 database_path.as_path(), 4627 schema.as_str(), 4628 Some(SCHEMA_VERSION), 4629 ); 4630 4631 assert_schema_error( 4632 open_schema_error(database_path.as_path()), 4633 "transport_publish_target_results", 4634 "target_scope", 4635 ); 4636 assert_eq!( 4637 database_user_version(database_path.as_path()), 4638 SCHEMA_VERSION 4639 ); 4640 } 4641 4642 #[test] 4643 fn transport_store_open_rejects_current_target_results_schema_without_scoped_primary_key() { 4644 let directory = tempfile::tempdir().expect("tempdir"); 4645 let database_path = directory 4646 .path() 4647 .join("publish-proxy-target-results-unscoped-pk.sqlite"); 4648 let schema = schema_with_replacement( 4649 "PRIMARY KEY(job_id, transport_kind, endpoint_uri, target_scope)", 4650 "PRIMARY KEY(job_id, transport_kind, endpoint_uri)", 4651 ); 4652 create_existing_schema( 4653 database_path.as_path(), 4654 schema.as_str(), 4655 Some(SCHEMA_VERSION), 4656 ); 4657 4658 assert_schema_error( 4659 open_schema_error(database_path.as_path()), 4660 "transport_publish_target_results", 4661 "primary key", 4662 ); 4663 assert_eq!( 4664 database_user_version(database_path.as_path()), 4665 SCHEMA_VERSION 4666 ); 4667 } 4668 4669 #[test] 4670 fn transport_store_open_rejects_current_schema_without_job_primary_key() { 4671 let directory = tempfile::tempdir().expect("tempdir"); 4672 let database_path = directory.path().join("publish-proxy-missing-job-pk.sqlite"); 4673 let schema = 4674 schema_with_replacement("job_id TEXT PRIMARY KEY NOT NULL", "job_id TEXT NOT NULL"); 4675 create_existing_schema( 4676 database_path.as_path(), 4677 schema.as_str(), 4678 Some(SCHEMA_VERSION), 4679 ); 4680 4681 assert_schema_error( 4682 open_schema_error(database_path.as_path()), 4683 "transport_publish_jobs", 4684 "primary key", 4685 ); 4686 assert_eq!( 4687 database_user_version(database_path.as_path()), 4688 SCHEMA_VERSION 4689 ); 4690 } 4691 4692 #[test] 4693 fn transport_store_open_rejects_current_schema_without_job_foreign_key() { 4694 let directory = tempfile::tempdir().expect("tempdir"); 4695 let database_path = directory.path().join("publish-proxy-missing-job-fk.sqlite"); 4696 let schema = schema_with_replacement( 4697 ", 4698 FOREIGN KEY(principal_id) REFERENCES transport_publish_principals(principal_id)", 4699 "", 4700 ); 4701 create_existing_schema( 4702 database_path.as_path(), 4703 schema.as_str(), 4704 Some(SCHEMA_VERSION), 4705 ); 4706 4707 assert_schema_error( 4708 open_schema_error(database_path.as_path()), 4709 "transport_publish_jobs", 4710 "foreign key", 4711 ); 4712 assert_eq!( 4713 database_user_version(database_path.as_path()), 4714 SCHEMA_VERSION 4715 ); 4716 } 4717 4718 #[tokio::test] 4719 async fn publish_event_verifies_and_records_daemon_default_outcome() { 4720 let identity = DaemonIdentity::generate(); 4721 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 4722 let principal = principal( 4723 &proxy, 4724 identity.public_key_hex(), 4725 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 4726 false, 4727 PublishJobVisibility::Own, 4728 ); 4729 let event = signed_event(&identity, "{}"); 4730 let raw_event = event.clone(); 4731 let response = proxy 4732 .publish_event( 4733 &principal, 4734 publish_request( 4735 event, 4736 Vec::new(), 4737 NostrTargetSourcePolicy::DaemonDefaultOnly, 4738 DeliveryPolicy::Any, 4739 Some("idem-valid"), 4740 ), 4741 ) 4742 .await 4743 .expect("publish"); 4744 4745 assert!(!response.deduplicated); 4746 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 4747 assert_eq!(response.job.target_count, 1); 4748 assert_eq!(response.job.acknowledged_count, 1); 4749 assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); 4750 assert_eq!(response.job.targets[0].source, TargetSource::DaemonDefault); 4751 assert_eq!(adapter.captured_raw_events(), vec![raw_event]); 4752 } 4753 4754 #[tokio::test] 4755 async fn publish_event_rejects_tampered_content_before_publish() { 4756 let identity = DaemonIdentity::generate(); 4757 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 4758 let principal = principal( 4759 &proxy, 4760 identity.public_key_hex(), 4761 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 4762 false, 4763 PublishJobVisibility::Own, 4764 ); 4765 let event = raw_event_with_field( 4766 signed_event(&identity, "trusted"), 4767 "content", 4768 serde_json::Value::String("tampered".to_owned()), 4769 ); 4770 let error = proxy 4771 .publish_event( 4772 &principal, 4773 publish_request( 4774 event, 4775 Vec::new(), 4776 NostrTargetSourcePolicy::DaemonDefaultOnly, 4777 DeliveryPolicy::Any, 4778 None, 4779 ), 4780 ) 4781 .await 4782 .expect_err("tampered event should fail"); 4783 4784 assert!(matches!(error, TransportPublishError::EventWire(_))); 4785 assert!(adapter.captured_raw_events().is_empty()); 4786 } 4787 4788 #[tokio::test] 4789 async fn publish_event_rejects_wrong_signature_before_publish() { 4790 let identity = DaemonIdentity::generate(); 4791 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 4792 let principal = principal( 4793 &proxy, 4794 identity.public_key_hex(), 4795 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 4796 false, 4797 PublishJobVisibility::Own, 4798 ); 4799 let event = signed_event(&identity, "{}"); 4800 let parsed: serde_json::Value = serde_json::from_str(event.as_str()).expect("event json"); 4801 let sig = parsed["sig"].as_str().expect("sig"); 4802 let replacement = if sig.starts_with('0') { "1" } else { "0" }; 4803 let mut tampered_sig = sig.to_owned(); 4804 tampered_sig.replace_range(0..1, replacement); 4805 let event = raw_event_with_field(event, "sig", serde_json::Value::String(tampered_sig)); 4806 let error = proxy 4807 .publish_event( 4808 &principal, 4809 publish_request( 4810 event, 4811 Vec::new(), 4812 NostrTargetSourcePolicy::DaemonDefaultOnly, 4813 DeliveryPolicy::Any, 4814 None, 4815 ), 4816 ) 4817 .await 4818 .expect_err("wrong signature should fail"); 4819 4820 assert!(matches!( 4821 error, 4822 TransportPublishError::InvalidSignedEvent(_) 4823 )); 4824 assert!(adapter.captured_raw_events().is_empty()); 4825 } 4826 4827 #[tokio::test] 4828 async fn publish_event_rejects_malformed_wire_fields() { 4829 let identity = DaemonIdentity::generate(); 4830 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 4831 let principal = principal( 4832 &proxy, 4833 identity.public_key_hex(), 4834 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 4835 false, 4836 PublishJobVisibility::Own, 4837 ); 4838 let event = signed_event(&identity, "{}"); 4839 let parsed: serde_json::Value = serde_json::from_str(event.as_str()).expect("event json"); 4840 let event_id = parsed["id"].as_str().expect("event id").to_uppercase(); 4841 let event = raw_event_with_field(event, "id", serde_json::Value::String(event_id)); 4842 let error = proxy 4843 .publish_event( 4844 &principal, 4845 publish_request( 4846 event, 4847 Vec::new(), 4848 NostrTargetSourcePolicy::DaemonDefaultOnly, 4849 DeliveryPolicy::Any, 4850 None, 4851 ), 4852 ) 4853 .await 4854 .expect_err("malformed field should fail"); 4855 4856 assert!(matches!(error, TransportPublishError::EventWire(_))); 4857 assert!(adapter.captured_raw_events().is_empty()); 4858 } 4859 4860 #[tokio::test] 4861 async fn publish_event_uses_explicit_request_relays_when_allowed() { 4862 let identity = DaemonIdentity::generate(); 4863 let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_SECONDARY])); 4864 let principal = principal( 4865 &proxy, 4866 identity.public_key_hex(), 4867 vec![NostrTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault], 4868 true, 4869 PublishJobVisibility::Own, 4870 ); 4871 let response = proxy 4872 .publish_event( 4873 &principal, 4874 publish_request( 4875 signed_event(&identity, "{}"), 4876 vec![RELAY_PRIMARY.to_owned()], 4877 NostrTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 4878 DeliveryPolicy::Any, 4879 None, 4880 ), 4881 ) 4882 .await 4883 .expect("publish"); 4884 4885 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 4886 assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); 4887 assert_eq!(response.job.targets[0].source, TargetSource::Request); 4888 } 4889 4890 #[tokio::test] 4891 async fn publish_event_uses_cached_nip65_author_write_before_defaults() { 4892 let identity = DaemonIdentity::generate(); 4893 let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_SECONDARY])); 4894 proxy 4895 .store 4896 .cache_author_write_relays( 4897 identity.public_key_hex().as_str(), 4898 &[RELAY_PRIMARY.to_owned()], 4899 ) 4900 .expect("cache author relays"); 4901 let principal = principal( 4902 &proxy, 4903 identity.public_key_hex(), 4904 vec![NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault], 4905 false, 4906 PublishJobVisibility::Own, 4907 ); 4908 let response = proxy 4909 .publish_event( 4910 &principal, 4911 publish_request( 4912 signed_event(&identity, "{}"), 4913 Vec::new(), 4914 NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault, 4915 DeliveryPolicy::Any, 4916 None, 4917 ), 4918 ) 4919 .await 4920 .expect("publish"); 4921 4922 assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); 4923 assert_eq!( 4924 response.job.targets[0].source, 4925 TargetSource::NostrAuthorWrite 4926 ); 4927 } 4928 4929 #[tokio::test] 4930 async fn publish_event_discards_invalid_cached_author_write_relay() { 4931 let identity = DaemonIdentity::generate(); 4932 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_SECONDARY])); 4933 proxy 4934 .store 4935 .cache_author_write_relays( 4936 identity.public_key_hex().as_str(), 4937 &[RELAY_PRIMARY.to_owned(), "not a cached relay".to_owned()], 4938 ) 4939 .expect("cache author relays"); 4940 let principal = principal( 4941 &proxy, 4942 identity.public_key_hex(), 4943 vec![NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault], 4944 false, 4945 PublishJobVisibility::Own, 4946 ); 4947 let response = proxy 4948 .publish_event( 4949 &principal, 4950 publish_request( 4951 signed_event(&identity, "{}"), 4952 Vec::new(), 4953 NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault, 4954 DeliveryPolicy::Any, 4955 None, 4956 ), 4957 ) 4958 .await 4959 .expect("publish"); 4960 4961 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 4962 let accepted = response 4963 .job 4964 .targets 4965 .iter() 4966 .find(|relay| relay.endpoint_uri == RELAY_PRIMARY) 4967 .expect("accepted author relay"); 4968 assert_eq!(accepted.source, TargetSource::NostrAuthorWrite); 4969 assert!(accepted.attempted); 4970 assert!( 4971 response 4972 .job 4973 .targets 4974 .iter() 4975 .all(|relay| relay.endpoint_uri != "not a cached relay") 4976 ); 4977 assert_eq!(adapter.captured_raw_events().len(), 1); 4978 } 4979 4980 #[tokio::test] 4981 async fn publish_event_discards_invalid_author_and_discovery_relays() { 4982 let identity = DaemonIdentity::generate(); 4983 let mut config = config_with_defaults(vec![RELAY_SECONDARY]); 4984 config.nostr.author_relay_discovery_relays = vec!["not a discovery relay".to_owned()]; 4985 let (proxy, adapter) = transport_publish(config); 4986 proxy 4987 .store 4988 .cache_author_write_relays( 4989 identity.public_key_hex().as_str(), 4990 &["not a cached relay".to_owned()], 4991 ) 4992 .expect("cache author relays"); 4993 let principal = principal( 4994 &proxy, 4995 identity.public_key_hex(), 4996 vec![NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault], 4997 false, 4998 PublishJobVisibility::Own, 4999 ); 5000 let response = proxy 5001 .publish_event( 5002 &principal, 5003 publish_request( 5004 signed_event(&identity, "{}"), 5005 Vec::new(), 5006 NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault, 5007 DeliveryPolicy::Any, 5008 None, 5009 ), 5010 ) 5011 .await 5012 .expect("publish"); 5013 5014 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5015 let daemon_default = response 5016 .job 5017 .targets 5018 .iter() 5019 .find(|relay| relay.endpoint_uri == RELAY_SECONDARY) 5020 .expect("daemon default relay"); 5021 assert_eq!(daemon_default.source, TargetSource::DaemonDefault); 5022 assert!(daemon_default.attempted); 5023 assert!(response.job.targets.iter().all(|relay| { 5024 relay.endpoint_uri != "not a cached relay" 5025 && relay.endpoint_uri != "not a discovery relay" 5026 })); 5027 assert_eq!(adapter.captured_raw_events().len(), 1); 5028 } 5029 5030 #[tokio::test] 5031 async fn publish_event_preserves_valid_discovery_rejections_only() { 5032 let identity = DaemonIdentity::generate(); 5033 let mut config = config_with_defaults(vec![RELAY_PRIMARY]); 5034 config.nostr.author_relay_discovery_relays = 5035 vec![RELAY_PRIMARY.to_owned(), RELAY_FORBIDDEN.to_owned()]; 5036 let resolver = StaticPublishRelayResolver::new().with_addresses( 5037 RELAY_FORBIDDEN, 5038 vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))], 5039 ); 5040 let adapter = RadrootsMockRelayPublishAdapter::new(); 5041 let proxy = TransportPublish::memory(config) 5042 .expect("proxy") 5043 .with_relay_resolver(Arc::new(resolver)) 5044 .with_author_relay_discovery(Arc::new(StaticPublishAuthorRelayDiscovery::new(vec![ 5045 "not a discovered author relay", 5046 RELAY_SECONDARY, 5047 ]))) 5048 .with_publisher(Arc::new(adapter.clone())); 5049 let principal = principal( 5050 &proxy, 5051 identity.public_key_hex(), 5052 vec![NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault], 5053 false, 5054 PublishJobVisibility::Own, 5055 ); 5056 let response = proxy 5057 .publish_event( 5058 &principal, 5059 publish_request( 5060 signed_event(&identity, "{}"), 5061 Vec::new(), 5062 NostrTargetSourcePolicy::AuthorWriteThenDaemonDefault, 5063 DeliveryPolicy::Any, 5064 None, 5065 ), 5066 ) 5067 .await 5068 .expect("publish"); 5069 5070 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5071 let accepted = response 5072 .job 5073 .targets 5074 .iter() 5075 .find(|relay| relay.endpoint_uri == RELAY_SECONDARY) 5076 .expect("discovered author relay"); 5077 assert_eq!(accepted.source, TargetSource::NostrAuthorWrite); 5078 assert!(accepted.attempted); 5079 assert!( 5080 response 5081 .job 5082 .targets 5083 .iter() 5084 .all(|relay| relay.endpoint_uri != "not a discovered author relay") 5085 ); 5086 let discovery = response 5087 .job 5088 .targets 5089 .iter() 5090 .find(|relay| relay.endpoint_uri == RELAY_FORBIDDEN) 5091 .expect("discovery relay rejection"); 5092 assert_eq!(discovery.source, TargetSource::DaemonDefault); 5093 assert_eq!(discovery.outcome_kind, OutcomeKind::TargetRejected); 5094 assert!(!discovery.attempted); 5095 assert_eq!(adapter.captured_raw_events().len(), 1); 5096 } 5097 5098 #[tokio::test] 5099 async fn publish_event_records_no_transport_publish_targets_failure() { 5100 let identity = DaemonIdentity::generate(); 5101 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5102 let principal = principal( 5103 &proxy, 5104 identity.public_key_hex(), 5105 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5106 false, 5107 PublishJobVisibility::Own, 5108 ); 5109 let response = proxy 5110 .publish_event( 5111 &principal, 5112 publish_request( 5113 signed_event(&identity, "{}"), 5114 Vec::new(), 5115 NostrTargetSourcePolicy::DaemonDefaultOnly, 5116 DeliveryPolicy::Any, 5117 None, 5118 ), 5119 ) 5120 .await 5121 .expect("publish"); 5122 5123 assert_eq!(response.job.status, JobStatus::Rejected); 5124 assert_eq!( 5125 response.job.last_error.as_deref(), 5126 Some("no_transport_publish_targets") 5127 ); 5128 assert!(response.job.targets.is_empty()); 5129 assert!(adapter.captured_raw_events().is_empty()); 5130 } 5131 5132 #[tokio::test] 5133 async fn publish_event_records_reticulum_unavailable_as_deferred_until_implemented() { 5134 let identity = DaemonIdentity::generate(); 5135 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5136 let principal = 5137 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5138 let response = proxy 5139 .publish_event( 5140 &principal, 5141 reticulum_publish_request( 5142 signed_event(&identity, "{}"), 5143 ReticulumBehavior::RejectDeliveryAttempts, 5144 ), 5145 ) 5146 .await 5147 .expect("publish"); 5148 5149 assert_eq!( 5150 response.job.status, 5151 JobStatus::DeliveryDeferredUntilImplemented 5152 ); 5153 assert!(response.job.terminal); 5154 assert!(!response.job.delivery_satisfied); 5155 assert_eq!(response.job.terminal_count, 0); 5156 assert!(response.job.completed_at_ms.is_some()); 5157 assert_eq!( 5158 response.job.last_error.as_deref(), 5159 Some("delivery_deferred_until_implemented") 5160 ); 5161 assert_eq!(response.job.targets.len(), 1); 5162 assert_eq!( 5163 response.job.targets[0].outcome_kind, 5164 OutcomeKind::DeferredUntilImplemented 5165 ); 5166 assert_eq!( 5167 response.job.targets[0].message.as_deref(), 5168 Some(RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE) 5169 ); 5170 assert!(!response.job.targets[0].attempted); 5171 assert!(adapter.captured_raw_events().is_empty()); 5172 } 5173 5174 #[tokio::test] 5175 async fn publish_event_records_explicit_nostr_target_when_kind_allowed() { 5176 let identity = DaemonIdentity::generate(); 5177 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5178 let principal = 5179 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5180 let event = signed_event(&identity, "{}"); 5181 let raw_event = event.clone(); 5182 let request = EventRequest { 5183 raw_event_json: event, 5184 target_policy: TargetPolicy::explicit_targets(vec![Target::nostr(RELAY_PRIMARY)]), 5185 delivery_policy: DeliveryPolicy::Any, 5186 idempotency_key: None, 5187 timeout_ms: Some(5_000), 5188 }; 5189 5190 let response = proxy 5191 .publish_event(&principal, request) 5192 .await 5193 .expect("publish"); 5194 5195 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5196 assert_eq!(response.job.targets.len(), 1); 5197 assert_eq!(response.job.targets[0].source, TargetSource::Request); 5198 assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); 5199 assert_eq!(response.job.targets[0].target_scope, None); 5200 assert_eq!(response.job.targets[0].target_label, None); 5201 assert_eq!(adapter.captured_raw_events(), vec![raw_event]); 5202 } 5203 5204 #[tokio::test] 5205 async fn publish_event_preserves_explicit_nostr_target_metadata_when_kind_allowed() { 5206 let identity = DaemonIdentity::generate(); 5207 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5208 let principal = 5209 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5210 let event = signed_event(&identity, "{}"); 5211 let raw_event = event.clone(); 5212 let request = EventRequest { 5213 raw_event_json: event, 5214 target_policy: TargetPolicy::explicit_targets(vec![ 5215 Target::nostr(RELAY_PRIMARY) 5216 .with_scope("farm.local") 5217 .with_label("Farm relay"), 5218 ]), 5219 delivery_policy: DeliveryPolicy::Any, 5220 idempotency_key: None, 5221 timeout_ms: Some(5_000), 5222 }; 5223 5224 let response = proxy 5225 .publish_event(&principal, request) 5226 .await 5227 .expect("publish"); 5228 5229 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5230 assert_eq!(response.job.targets.len(), 1); 5231 assert_eq!(response.job.targets[0].endpoint_uri, RELAY_PRIMARY); 5232 assert_eq!( 5233 response.job.targets[0].target_scope.as_deref(), 5234 Some("farm.local") 5235 ); 5236 assert_eq!( 5237 response.job.targets[0].target_label.as_deref(), 5238 Some("Farm relay") 5239 ); 5240 assert_eq!(response.job.targets[0].source, TargetSource::Request); 5241 assert_eq!(adapter.captured_raw_events(), vec![raw_event]); 5242 response.job.validate().expect("valid scoped job"); 5243 } 5244 5245 #[tokio::test] 5246 async fn publish_event_records_scoped_targets_with_shared_relay_url() { 5247 let identity = DaemonIdentity::generate(); 5248 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5249 let principal = 5250 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5251 let event = signed_event(&identity, "{}"); 5252 let raw_event = event.clone(); 5253 let request = EventRequest { 5254 raw_event_json: event, 5255 target_policy: TargetPolicy::explicit_targets(vec![ 5256 Target::nostr(RELAY_PRIMARY) 5257 .with_scope("farm.a") 5258 .with_label("Farm A"), 5259 Target::nostr(RELAY_PRIMARY) 5260 .with_scope("farm.b") 5261 .with_label("Farm B"), 5262 ]), 5263 delivery_policy: DeliveryPolicy::All, 5264 idempotency_key: None, 5265 timeout_ms: Some(5_000), 5266 }; 5267 5268 let response = proxy 5269 .publish_event(&principal, request) 5270 .await 5271 .expect("publish"); 5272 5273 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5274 assert_eq!(response.job.targets.len(), 2); 5275 assert!( 5276 response 5277 .job 5278 .targets 5279 .iter() 5280 .all(|target| target.endpoint_uri == RELAY_PRIMARY) 5281 ); 5282 assert_eq!( 5283 response 5284 .job 5285 .targets 5286 .iter() 5287 .map(|target| ( 5288 target.target_scope.as_deref(), 5289 target.target_label.as_deref() 5290 )) 5291 .collect::<Vec<_>>(), 5292 vec![ 5293 (Some("farm.a"), Some("Farm A")), 5294 (Some("farm.b"), Some("Farm B")), 5295 ] 5296 ); 5297 assert_eq!(adapter.captured_raw_events(), vec![raw_event]); 5298 response 5299 .job 5300 .validate() 5301 .expect("valid scoped shared relay job"); 5302 } 5303 5304 #[tokio::test] 5305 async fn publish_event_required_targets_do_not_count_optional_success() { 5306 let identity = DaemonIdentity::generate(); 5307 let adapter = RadrootsMockRelayPublishAdapter::new() 5308 .with_outcome( 5309 RELAY_PRIMARY, 5310 RadrootsRelayOutcome::classify("restricted: required relay rejected"), 5311 ) 5312 .with_outcome(RELAY_SECONDARY, RadrootsRelayOutcome::accepted()); 5313 let (proxy, _) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5314 let proxy = proxy.with_publisher(Arc::new(adapter.clone())); 5315 let principal = 5316 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5317 let required_target = TransportTarget::nostr_relay(RELAY_PRIMARY).expect("required target"); 5318 let request = EventRequest { 5319 raw_event_json: signed_event(&identity, "{}"), 5320 target_policy: TargetPolicy::explicit_targets(vec![ 5321 Target::nostr(RELAY_PRIMARY), 5322 Target::nostr(RELAY_SECONDARY), 5323 ]), 5324 delivery_policy: DeliveryPolicy::required_targets(vec![ 5325 ProtocolTargetFingerprint::parse(required_target.fingerprint().as_str()) 5326 .expect("protocol target fingerprint"), 5327 ]) 5328 .expect("required targets"), 5329 idempotency_key: None, 5330 timeout_ms: Some(5_000), 5331 }; 5332 5333 let response = proxy 5334 .publish_event(&principal, request) 5335 .await 5336 .expect("publish"); 5337 5338 assert_eq!(response.job.status, JobStatus::DeliveryUnsatisfiedTerminal); 5339 assert!(!response.job.delivery_satisfied); 5340 assert_eq!(response.job.acknowledged_count, 1); 5341 assert_eq!(adapter.captured_raw_events().len(), 1); 5342 response 5343 .job 5344 .validate() 5345 .expect("valid required target terminal job"); 5346 } 5347 5348 #[tokio::test] 5349 async fn publish_event_required_targets_ignore_optional_retryable_failures() { 5350 let identity = DaemonIdentity::generate(); 5351 let adapter = RadrootsMockRelayPublishAdapter::new() 5352 .with_outcome(RELAY_PRIMARY, RadrootsRelayOutcome::accepted()) 5353 .with_outcome( 5354 RELAY_SECONDARY, 5355 RadrootsRelayOutcome::timeout("optional timeout"), 5356 ); 5357 let (proxy, _) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5358 let proxy = proxy.with_publisher(Arc::new(adapter.clone())); 5359 let principal = 5360 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5361 let required_target = TransportTarget::nostr_relay(RELAY_PRIMARY).expect("required target"); 5362 let request = EventRequest { 5363 raw_event_json: signed_event(&identity, "{}"), 5364 target_policy: TargetPolicy::explicit_targets(vec![ 5365 Target::nostr(RELAY_PRIMARY), 5366 Target::nostr(RELAY_SECONDARY), 5367 ]), 5368 delivery_policy: DeliveryPolicy::required_targets(vec![ 5369 ProtocolTargetFingerprint::parse(required_target.fingerprint().as_str()) 5370 .expect("protocol target fingerprint"), 5371 ]) 5372 .expect("required targets"), 5373 idempotency_key: None, 5374 timeout_ms: Some(5_000), 5375 }; 5376 5377 let response = proxy 5378 .publish_event(&principal, request) 5379 .await 5380 .expect("publish"); 5381 5382 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5383 assert!(response.job.delivery_satisfied); 5384 assert_eq!(response.job.acknowledged_count, 1); 5385 assert_eq!(response.job.retryable_count, 1); 5386 assert_eq!(adapter.captured_raw_events().len(), 1); 5387 response 5388 .job 5389 .validate() 5390 .expect("valid required target satisfied job"); 5391 } 5392 5393 #[tokio::test] 5394 async fn publish_event_rejects_required_target_not_in_resolved_set() { 5395 let identity = DaemonIdentity::generate(); 5396 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5397 let principal = 5398 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5399 let stale_target = TransportTarget::nostr_relay(RELAY_SECONDARY).expect("stale target"); 5400 let request = EventRequest { 5401 raw_event_json: signed_event(&identity, "{}"), 5402 target_policy: TargetPolicy::explicit_targets(vec![Target::nostr(RELAY_PRIMARY)]), 5403 delivery_policy: DeliveryPolicy::required_targets(vec![ 5404 ProtocolTargetFingerprint::parse(stale_target.fingerprint().as_str()) 5405 .expect("protocol target fingerprint"), 5406 ]) 5407 .expect("required targets"), 5408 idempotency_key: None, 5409 timeout_ms: Some(5_000), 5410 }; 5411 5412 let err = proxy 5413 .publish_event(&principal, request) 5414 .await 5415 .expect_err("stale required target"); 5416 5417 assert!( 5418 matches!( 5419 err, 5420 TransportPublishError::InvalidSignedEvent(ref message) 5421 if message.contains("requires a target") 5422 ), 5423 "{err:?}" 5424 ); 5425 assert!(adapter.captured_raw_events().is_empty()); 5426 assert!( 5427 proxy 5428 .store 5429 .list_jobs_for_principal(&principal, 10) 5430 .expect("jobs") 5431 .is_empty() 5432 ); 5433 } 5434 5435 #[tokio::test] 5436 async fn publish_event_rejects_duplicate_explicit_targets_before_recording_job() { 5437 let identity = DaemonIdentity::generate(); 5438 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5439 let principal = 5440 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5441 let request = EventRequest { 5442 raw_event_json: signed_event(&identity, "{}"), 5443 target_policy: TargetPolicy::explicit_targets(vec![ 5444 Target::nostr(RELAY_PRIMARY), 5445 Target::nostr(RELAY_PRIMARY), 5446 ]), 5447 delivery_policy: DeliveryPolicy::Any, 5448 idempotency_key: None, 5449 timeout_ms: Some(5_000), 5450 }; 5451 5452 let err = proxy 5453 .publish_event(&principal, request) 5454 .await 5455 .expect_err("duplicate explicit targets"); 5456 5457 assert!(matches!( 5458 err, 5459 TransportPublishError::InvalidSignedEvent(ref message) 5460 if message.contains("duplicates an earlier target") 5461 )); 5462 assert!(adapter.captured_raw_events().is_empty()); 5463 assert!( 5464 proxy 5465 .store 5466 .list_jobs_for_principal(&principal, 10) 5467 .expect("jobs") 5468 .is_empty() 5469 ); 5470 } 5471 5472 #[tokio::test] 5473 async fn publish_event_rejects_explicit_target_kind_not_allowed_before_recording_job() { 5474 let identity = DaemonIdentity::generate(); 5475 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5476 let principal = explicit_target_principal_with_kinds( 5477 &proxy, 5478 identity.public_key_hex(), 5479 vec![TRANSPORT_KIND_NOSTR.to_owned()], 5480 PublishJobVisibility::Own, 5481 ); 5482 5483 let err = proxy 5484 .publish_event( 5485 &principal, 5486 reticulum_publish_request( 5487 signed_event(&identity, "{}"), 5488 ReticulumBehavior::RejectDeliveryAttempts, 5489 ), 5490 ) 5491 .await 5492 .expect_err("disallowed explicit target kind"); 5493 5494 assert!(matches!(err, TransportPublishError::InvalidScope(_))); 5495 assert!(adapter.captured_raw_events().is_empty()); 5496 assert!( 5497 proxy 5498 .store 5499 .list_jobs_for_principal(&principal, 10) 5500 .expect("jobs") 5501 .is_empty() 5502 ); 5503 } 5504 5505 #[tokio::test] 5506 async fn publish_event_records_reticulum_deferred_as_terminal_nonfailure() { 5507 let identity = DaemonIdentity::generate(); 5508 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5509 let principal = 5510 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5511 let response = proxy 5512 .publish_event( 5513 &principal, 5514 reticulum_publish_request( 5515 signed_event(&identity, "{}"), 5516 ReticulumBehavior::DeferDeliveryPlans, 5517 ), 5518 ) 5519 .await 5520 .expect("publish"); 5521 5522 assert_eq!( 5523 response.job.status, 5524 JobStatus::DeliveryDeferredUntilImplemented 5525 ); 5526 assert!(response.job.terminal); 5527 assert!(!response.job.delivery_satisfied); 5528 assert_eq!(response.job.terminal_count, 0); 5529 assert!(response.job.completed_at_ms.is_some()); 5530 assert_eq!( 5531 response.job.last_error.as_deref(), 5532 Some("delivery_deferred_until_implemented") 5533 ); 5534 assert_eq!(response.job.targets.len(), 1); 5535 assert_eq!( 5536 response.job.targets[0].outcome_kind, 5537 OutcomeKind::DeferredUntilImplemented 5538 ); 5539 assert!(!response.job.targets[0].attempted); 5540 assert!(adapter.captured_raw_events().is_empty()); 5541 } 5542 5543 #[tokio::test] 5544 async fn publish_event_rejects_noncanonical_reticulum_endpoint_before_recording_job() { 5545 let identity = DaemonIdentity::generate(); 5546 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5547 let principal = 5548 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5549 let mut request = reticulum_publish_request( 5550 signed_event(&identity, "{}"), 5551 ReticulumBehavior::RejectDeliveryAttempts, 5552 ); 5553 request.target_policy = TargetPolicy::explicit_targets(vec![Target { 5554 transport_kind: "reticulum".to_owned(), 5555 endpoint_uri: "reticulum:unavailable-alt".to_owned(), 5556 target_scope: None, 5557 target_label: None, 5558 reticulum_behavior: Some(ReticulumBehavior::RejectDeliveryAttempts), 5559 }]); 5560 5561 let err = proxy 5562 .publish_event(&principal, request) 5563 .await 5564 .expect_err("noncanonical Reticulum endpoint"); 5565 5566 assert!(matches!(err, TransportPublishError::InvalidSignedEvent(_))); 5567 assert!(adapter.captured_raw_events().is_empty()); 5568 assert!( 5569 proxy 5570 .store 5571 .list_jobs_for_principal(&principal, 10) 5572 .expect("jobs") 5573 .is_empty() 5574 ); 5575 } 5576 5577 #[tokio::test] 5578 async fn publish_event_rejects_reticulum_behavior_on_non_reticulum_before_recording_job() { 5579 let identity = DaemonIdentity::generate(); 5580 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5581 let principal = 5582 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5583 let mut request = publish_request( 5584 signed_event(&identity, "{}"), 5585 vec![RELAY_PRIMARY.to_owned()], 5586 NostrTargetSourcePolicy::ExplicitOnly, 5587 DeliveryPolicy::Any, 5588 None, 5589 ); 5590 request.target_policy = TargetPolicy::explicit_targets(vec![Target { 5591 transport_kind: "nostr".to_owned(), 5592 endpoint_uri: RELAY_PRIMARY.to_owned(), 5593 target_scope: None, 5594 target_label: None, 5595 reticulum_behavior: Some(ReticulumBehavior::RejectDeliveryAttempts), 5596 }]); 5597 5598 let err = proxy 5599 .publish_event(&principal, request) 5600 .await 5601 .expect_err("non-Reticulum reticulum behavior"); 5602 5603 assert!(matches!(err, TransportPublishError::InvalidSignedEvent(_))); 5604 assert!(adapter.captured_raw_events().is_empty()); 5605 assert!( 5606 proxy 5607 .store 5608 .list_jobs_for_principal(&principal, 10) 5609 .expect("jobs") 5610 .is_empty() 5611 ); 5612 } 5613 5614 #[tokio::test] 5615 async fn publish_event_rejects_noncanonical_reticulum_kind_before_recording_job() { 5616 let identity = DaemonIdentity::generate(); 5617 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5618 let principal = 5619 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5620 let mut request = reticulum_publish_request( 5621 signed_event(&identity, "{}"), 5622 ReticulumBehavior::RejectDeliveryAttempts, 5623 ); 5624 request.target_policy = TargetPolicy::explicit_targets(vec![Target { 5625 transport_kind: "Reticulum".to_owned(), 5626 endpoint_uri: RADROOTS_RETICULUM_ENDPOINT_URI.to_owned(), 5627 target_scope: None, 5628 target_label: None, 5629 reticulum_behavior: Some(ReticulumBehavior::RejectDeliveryAttempts), 5630 }]); 5631 5632 let err = proxy 5633 .publish_event(&principal, request) 5634 .await 5635 .expect_err("noncanonical Reticulum kind"); 5636 5637 assert!(matches!(err, TransportPublishError::InvalidSignedEvent(_))); 5638 assert!(adapter.captured_raw_events().is_empty()); 5639 assert!( 5640 proxy 5641 .store 5642 .list_jobs_for_principal(&principal, 10) 5643 .expect("jobs") 5644 .is_empty() 5645 ); 5646 } 5647 5648 #[tokio::test] 5649 async fn publish_event_rejects_removed_execution_kind_before_recording_job() { 5650 let identity = DaemonIdentity::generate(); 5651 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5652 let principal = 5653 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5654 let mut request = reticulum_publish_request( 5655 signed_event(&identity, "{}"), 5656 ReticulumBehavior::RejectDeliveryAttempts, 5657 ); 5658 request.target_policy = TargetPolicy::explicit_targets(vec![Target { 5659 transport_kind: removed_execution_kind_string(), 5660 endpoint_uri: removed_execution_endpoint_uri(), 5661 target_scope: None, 5662 target_label: None, 5663 reticulum_behavior: None, 5664 }]); 5665 5666 let err = proxy 5667 .publish_event(&principal, request) 5668 .await 5669 .expect_err("removed execution kind"); 5670 5671 assert!(matches!(err, TransportPublishError::InvalidSignedEvent(_))); 5672 assert!(adapter.captured_raw_events().is_empty()); 5673 assert!( 5674 proxy 5675 .store 5676 .list_jobs_for_principal(&principal, 10) 5677 .expect("jobs") 5678 .is_empty() 5679 ); 5680 } 5681 5682 #[tokio::test] 5683 async fn publish_event_rejects_removed_execution_target_before_recording_job() { 5684 let identity = DaemonIdentity::generate(); 5685 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5686 let principal = 5687 explicit_target_principal(&proxy, identity.public_key_hex(), PublishJobVisibility::Own); 5688 let mut request = reticulum_publish_request( 5689 signed_event(&identity, "{}"), 5690 ReticulumBehavior::RejectDeliveryAttempts, 5691 ); 5692 request.target_policy = TargetPolicy::explicit_targets(vec![Target { 5693 transport_kind: removed_proxy_transport_kind_string(), 5694 endpoint_uri: removed_execution_endpoint_uri(), 5695 target_scope: None, 5696 target_label: None, 5697 reticulum_behavior: None, 5698 }]); 5699 5700 let err = proxy 5701 .publish_event(&principal, request) 5702 .await 5703 .expect_err("removed execution target"); 5704 5705 assert!(matches!(err, TransportPublishError::InvalidSignedEvent(_))); 5706 assert!(adapter.captured_raw_events().is_empty()); 5707 assert!( 5708 proxy 5709 .store 5710 .list_jobs_for_principal(&principal, 10) 5711 .expect("jobs") 5712 .is_empty() 5713 ); 5714 } 5715 5716 fn removed_execution_kind_string() -> String { 5717 ["radrootsd", "_proxy"].concat() 5718 } 5719 5720 fn removed_proxy_transport_kind_string() -> String { 5721 ["pro", "xy"].concat() 5722 } 5723 5724 fn removed_execution_endpoint_uri() -> String { 5725 ["radrootsd-", "pro", "xy:publish"].concat() 5726 } 5727 5728 #[tokio::test] 5729 async fn publish_event_records_unsafe_request_relay_rejection() { 5730 let identity = DaemonIdentity::generate(); 5731 let (proxy, adapter) = transport_publish(TransportPublishConfig::default()); 5732 let principal = principal( 5733 &proxy, 5734 identity.public_key_hex(), 5735 vec![NostrTargetSourcePolicy::ExplicitOnly], 5736 true, 5737 PublishJobVisibility::Own, 5738 ); 5739 let response = proxy 5740 .publish_event( 5741 &principal, 5742 publish_request( 5743 signed_event(&identity, "{}"), 5744 vec!["wss://127.0.0.1:7777".to_owned()], 5745 NostrTargetSourcePolicy::ExplicitOnly, 5746 DeliveryPolicy::Any, 5747 None, 5748 ), 5749 ) 5750 .await 5751 .expect("publish"); 5752 5753 assert_eq!(response.job.status, JobStatus::DeliveryUnsatisfiedTerminal); 5754 assert_eq!(response.job.targets.len(), 1); 5755 assert_eq!( 5756 response.job.targets[0].outcome_kind, 5757 OutcomeKind::TargetRejected 5758 ); 5759 assert!(!response.job.targets[0].attempted); 5760 assert!(adapter.captured_raw_events().is_empty()); 5761 } 5762 5763 #[tokio::test] 5764 async fn publish_event_rejects_forbidden_public_dns_destination_before_publish() { 5765 let identity = DaemonIdentity::generate(); 5766 let resolver = StaticPublishRelayResolver::new() 5767 .with_addresses(RELAY_PRIMARY, vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))]); 5768 let (proxy, adapter) = transport_publish_with_resolver( 5769 config_with_defaults(vec![RELAY_PRIMARY]), 5770 Arc::new(resolver), 5771 ); 5772 let principal = principal( 5773 &proxy, 5774 identity.public_key_hex(), 5775 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5776 false, 5777 PublishJobVisibility::Own, 5778 ); 5779 let response = proxy 5780 .publish_event( 5781 &principal, 5782 publish_request( 5783 signed_event(&identity, "{}"), 5784 Vec::new(), 5785 NostrTargetSourcePolicy::DaemonDefaultOnly, 5786 DeliveryPolicy::Any, 5787 None, 5788 ), 5789 ) 5790 .await 5791 .expect("publish"); 5792 5793 assert_eq!(response.job.status, JobStatus::DeliveryUnsatisfiedTerminal); 5794 assert_eq!(response.job.targets.len(), 1); 5795 assert_eq!( 5796 response.job.targets[0].outcome_kind, 5797 OutcomeKind::TargetRejected 5798 ); 5799 assert!(!response.job.targets[0].attempted); 5800 assert!(adapter.captured_raw_events().is_empty()); 5801 } 5802 5803 #[tokio::test] 5804 async fn publish_event_records_dns_failure_as_unattempted_retryable_outcome() { 5805 let identity = DaemonIdentity::generate(); 5806 let resolver = StaticPublishRelayResolver::new().with_failure(RELAY_PRIMARY, "no records"); 5807 let (proxy, adapter) = transport_publish_with_resolver( 5808 config_with_defaults(vec![RELAY_PRIMARY]), 5809 Arc::new(resolver), 5810 ); 5811 let principal = principal( 5812 &proxy, 5813 identity.public_key_hex(), 5814 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5815 false, 5816 PublishJobVisibility::Own, 5817 ); 5818 let response = proxy 5819 .publish_event( 5820 &principal, 5821 publish_request( 5822 signed_event(&identity, "{}"), 5823 Vec::new(), 5824 NostrTargetSourcePolicy::DaemonDefaultOnly, 5825 DeliveryPolicy::Any, 5826 None, 5827 ), 5828 ) 5829 .await 5830 .expect("publish"); 5831 5832 assert_eq!(response.job.status, JobStatus::DeliveryUnsatisfiedRetryable); 5833 assert_eq!( 5834 response.job.last_error.as_deref(), 5835 Some("delivery_unsatisfied") 5836 ); 5837 assert_eq!(response.job.targets.len(), 1); 5838 assert_eq!( 5839 response.job.targets[0].outcome_kind, 5840 OutcomeKind::ConnectionFailed 5841 ); 5842 assert!(!response.job.targets[0].attempted); 5843 assert!(adapter.captured_raw_events().is_empty()); 5844 } 5845 5846 #[tokio::test] 5847 async fn publish_event_localhost_policy_skips_public_dns_guard() { 5848 let identity = DaemonIdentity::generate(); 5849 let mut config = config_with_defaults(vec!["ws://localhost:7777"]); 5850 config.nostr.relay_url_policy = NostrRelayUrlPolicy::Localhost; 5851 let resolver = StaticPublishRelayResolver::new() 5852 .with_failure("ws://localhost:7777", "localhost resolution should not run"); 5853 let (proxy, adapter) = transport_publish_with_resolver(config, Arc::new(resolver)); 5854 let principal = principal( 5855 &proxy, 5856 identity.public_key_hex(), 5857 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5858 false, 5859 PublishJobVisibility::Own, 5860 ); 5861 let response = proxy 5862 .publish_event( 5863 &principal, 5864 publish_request( 5865 signed_event(&identity, "{}"), 5866 Vec::new(), 5867 NostrTargetSourcePolicy::DaemonDefaultOnly, 5868 DeliveryPolicy::Any, 5869 None, 5870 ), 5871 ) 5872 .await 5873 .expect("publish"); 5874 5875 assert_eq!(response.job.status, JobStatus::DeliverySatisfied); 5876 assert_eq!(response.job.targets[0].endpoint_uri, "ws://localhost:7777"); 5877 assert!(!adapter.captured_raw_events().is_empty()); 5878 } 5879 5880 #[tokio::test] 5881 async fn publish_event_deduplicates_same_intent_and_conflicts_different_intent() { 5882 let identity = DaemonIdentity::generate(); 5883 let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5884 let principal = principal( 5885 &proxy, 5886 identity.public_key_hex(), 5887 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5888 false, 5889 PublishJobVisibility::Own, 5890 ); 5891 let request = publish_request( 5892 signed_event(&identity, "{}"), 5893 Vec::new(), 5894 NostrTargetSourcePolicy::DaemonDefaultOnly, 5895 DeliveryPolicy::Any, 5896 Some("idem-conflict"), 5897 ); 5898 let first = proxy 5899 .publish_event(&principal, request.clone()) 5900 .await 5901 .expect("first"); 5902 let duplicate = proxy 5903 .publish_event(&principal, request) 5904 .await 5905 .expect("duplicate"); 5906 5907 assert!(!first.deduplicated); 5908 assert!(duplicate.deduplicated); 5909 assert_eq!(duplicate.job.job_id, first.job.job_id); 5910 5911 let conflict = proxy 5912 .publish_event( 5913 &principal, 5914 publish_request( 5915 signed_event(&identity, "changed"), 5916 Vec::new(), 5917 NostrTargetSourcePolicy::DaemonDefaultOnly, 5918 DeliveryPolicy::Any, 5919 Some("idem-conflict"), 5920 ), 5921 ) 5922 .await 5923 .expect_err("conflict"); 5924 assert!(matches!( 5925 conflict, 5926 TransportPublishError::IdempotencyConflict(_) 5927 )); 5928 } 5929 5930 #[tokio::test] 5931 async fn publish_event_rejects_zero_and_excessive_timeout_before_job_creation() { 5932 let identity = DaemonIdentity::generate(); 5933 let (proxy, adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5934 let principal = principal( 5935 &proxy, 5936 identity.public_key_hex(), 5937 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5938 false, 5939 PublishJobVisibility::Own, 5940 ); 5941 let mut zero = publish_request( 5942 signed_event(&identity, "{}"), 5943 Vec::new(), 5944 NostrTargetSourcePolicy::DaemonDefaultOnly, 5945 DeliveryPolicy::Any, 5946 Some("idem-zero-timeout"), 5947 ); 5948 zero.timeout_ms = Some(0); 5949 let zero_error = proxy 5950 .publish_event(&principal, zero) 5951 .await 5952 .expect_err("zero timeout should fail"); 5953 assert!(matches!( 5954 zero_error, 5955 TransportPublishError::InvalidSignedEvent(_) 5956 )); 5957 5958 let mut excessive = publish_request( 5959 signed_event(&identity, "changed"), 5960 Vec::new(), 5961 NostrTargetSourcePolicy::DaemonDefaultOnly, 5962 DeliveryPolicy::Any, 5963 Some("idem-excessive-timeout"), 5964 ); 5965 excessive.timeout_ms = Some(10_001); 5966 let excessive_error = proxy 5967 .publish_event(&principal, excessive) 5968 .await 5969 .expect_err("excessive timeout should fail"); 5970 assert!(matches!( 5971 excessive_error, 5972 TransportPublishError::InvalidSignedEvent(_) 5973 )); 5974 assert!( 5975 proxy 5976 .store 5977 .list_jobs_for_principal(&principal, 50) 5978 .expect("jobs") 5979 .is_empty() 5980 ); 5981 assert!(adapter.captured_raw_events().is_empty()); 5982 } 5983 5984 #[tokio::test] 5985 async fn publish_event_default_timeout_fingerprints_as_effective_timeout() { 5986 let identity = DaemonIdentity::generate(); 5987 let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 5988 let principal = principal( 5989 &proxy, 5990 identity.public_key_hex(), 5991 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 5992 false, 5993 PublishJobVisibility::Own, 5994 ); 5995 let event = signed_event(&identity, "{}"); 5996 let mut default_timeout = publish_request( 5997 event.clone(), 5998 Vec::new(), 5999 NostrTargetSourcePolicy::DaemonDefaultOnly, 6000 DeliveryPolicy::Any, 6001 Some("idem-default-timeout"), 6002 ); 6003 default_timeout.timeout_ms = None; 6004 let mut explicit_default = publish_request( 6005 event, 6006 Vec::new(), 6007 NostrTargetSourcePolicy::DaemonDefaultOnly, 6008 DeliveryPolicy::Any, 6009 Some("idem-default-timeout"), 6010 ); 6011 explicit_default.timeout_ms = Some(10_000); 6012 6013 let first = proxy 6014 .publish_event(&principal, default_timeout) 6015 .await 6016 .expect("first"); 6017 let duplicate = proxy 6018 .publish_event(&principal, explicit_default) 6019 .await 6020 .expect("duplicate"); 6021 assert!(!first.deduplicated); 6022 assert!(duplicate.deduplicated); 6023 assert_eq!(duplicate.job.job_id, first.job.job_id); 6024 } 6025 6026 #[tokio::test] 6027 async fn publish_event_fingerprint_conflicts_on_different_effective_timeout() { 6028 let identity = DaemonIdentity::generate(); 6029 let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 6030 let principal = principal( 6031 &proxy, 6032 identity.public_key_hex(), 6033 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 6034 false, 6035 PublishJobVisibility::Own, 6036 ); 6037 let event = signed_event(&identity, "{}"); 6038 let first = publish_request( 6039 event.clone(), 6040 Vec::new(), 6041 NostrTargetSourcePolicy::DaemonDefaultOnly, 6042 DeliveryPolicy::Any, 6043 Some("idem-timeout-conflict"), 6044 ); 6045 let mut conflict = publish_request( 6046 event, 6047 Vec::new(), 6048 NostrTargetSourcePolicy::DaemonDefaultOnly, 6049 DeliveryPolicy::Any, 6050 Some("idem-timeout-conflict"), 6051 ); 6052 conflict.timeout_ms = Some(6_000); 6053 6054 proxy.publish_event(&principal, first).await.expect("first"); 6055 let error = proxy 6056 .publish_event(&principal, conflict) 6057 .await 6058 .expect_err("timeout conflict"); 6059 assert!(matches!( 6060 error, 6061 TransportPublishError::IdempotencyConflict(_) 6062 )); 6063 } 6064 6065 #[tokio::test] 6066 async fn publish_event_concurrency_limit_rejects_without_job_creation() { 6067 let identity = DaemonIdentity::generate(); 6068 let mut config = config_with_defaults(vec![RELAY_PRIMARY]); 6069 config.max_concurrent_publish_jobs = 1; 6070 let (proxy, adapter) = transport_publish(config); 6071 let principal = principal( 6072 &proxy, 6073 identity.public_key_hex(), 6074 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 6075 false, 6076 PublishJobVisibility::Own, 6077 ); 6078 let _permit = proxy.acquire_publish_permit().expect("permit"); 6079 let error = proxy 6080 .publish_event( 6081 &principal, 6082 publish_request( 6083 signed_event(&identity, "{}"), 6084 Vec::new(), 6085 NostrTargetSourcePolicy::DaemonDefaultOnly, 6086 DeliveryPolicy::Any, 6087 Some("idem-concurrency"), 6088 ), 6089 ) 6090 .await 6091 .expect_err("concurrency limit"); 6092 assert!(matches!(error, TransportPublishError::ConcurrencyLimit)); 6093 assert!( 6094 proxy 6095 .store 6096 .list_jobs_for_principal(&principal, 50) 6097 .expect("jobs") 6098 .is_empty() 6099 ); 6100 assert!(adapter.captured_raw_events().is_empty()); 6101 } 6102 6103 #[tokio::test] 6104 async fn publish_jobs_respect_own_and_admin_visibility() { 6105 let identity = DaemonIdentity::generate(); 6106 let other_identity = DaemonIdentity::generate(); 6107 let (proxy, _adapter) = transport_publish(config_with_defaults(vec![RELAY_PRIMARY])); 6108 let owner = principal( 6109 &proxy, 6110 identity.public_key_hex(), 6111 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 6112 false, 6113 PublishJobVisibility::Own, 6114 ); 6115 let other = principal( 6116 &proxy, 6117 other_identity.public_key_hex(), 6118 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 6119 false, 6120 PublishJobVisibility::Own, 6121 ); 6122 let admin = principal( 6123 &proxy, 6124 other_identity.public_key_hex(), 6125 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 6126 false, 6127 PublishJobVisibility::Admin, 6128 ); 6129 let response = proxy 6130 .publish_event( 6131 &owner, 6132 publish_request( 6133 signed_event(&identity, "{}"), 6134 Vec::new(), 6135 NostrTargetSourcePolicy::DaemonDefaultOnly, 6136 DeliveryPolicy::Any, 6137 None, 6138 ), 6139 ) 6140 .await 6141 .expect("publish"); 6142 6143 assert!( 6144 proxy 6145 .store 6146 .job_by_id_for_principal(response.job.job_id.as_str(), &other) 6147 .expect("other read") 6148 .is_none() 6149 ); 6150 assert!( 6151 proxy 6152 .store 6153 .job_by_id_for_principal(response.job.job_id.as_str(), &admin) 6154 .expect("admin read") 6155 .is_some() 6156 ); 6157 } 6158 6159 #[tokio::test] 6160 async fn publish_event_records_retryable_relay_failures() { 6161 let identity = DaemonIdentity::generate(); 6162 let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( 6163 RELAY_PRIMARY, 6164 RadrootsRelayOutcome::connection_failed("error: unavailable"), 6165 ); 6166 let proxy = TransportPublish::memory(config_with_defaults(vec![RELAY_PRIMARY])) 6167 .expect("proxy") 6168 .with_publisher(Arc::new(adapter)); 6169 let principal = principal( 6170 &proxy, 6171 identity.public_key_hex(), 6172 vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 6173 false, 6174 PublishJobVisibility::Own, 6175 ); 6176 let response = proxy 6177 .publish_event( 6178 &principal, 6179 publish_request( 6180 signed_event(&identity, "{}"), 6181 Vec::new(), 6182 NostrTargetSourcePolicy::DaemonDefaultOnly, 6183 DeliveryPolicy::Any, 6184 None, 6185 ), 6186 ) 6187 .await 6188 .expect("publish"); 6189 6190 assert_eq!(response.job.status, JobStatus::DeliveryUnsatisfiedRetryable); 6191 assert_eq!(response.job.retryable_count, 1); 6192 } 6193 }