state_discovery.rs (58510B)
1 //! Durable discovery desired/current state and exact committed publication bytes. 2 3 use core::fmt; 4 use std::{ 5 collections::{BTreeMap, BTreeSet}, 6 error::Error, 7 }; 8 9 use radroots_service_sqlite::{ 10 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 11 }; 12 use serde::{Deserialize, Serialize}; 13 use sha2::{Digest, Sha256}; 14 use sqlx::Row; 15 16 #[cfg(any(target_os = "linux", target_os = "macos"))] 17 use crate::state_admin::AdminJournalOperationError; 18 use crate::state_delivery::{ 19 DeliveryOperationError, MycDeliveryArtifactDigest, MycDeliveryJobId, MycDeliveryJobRecord, 20 MycDeliveryJobStatus, MycDeliverySource, MycDeliveryTimeUnixMs, create_job, 21 }; 22 use crate::state_repository::{ 23 MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, 24 RepositoryOperationError, require_expected_metadata, 25 }; 26 use crate::{MycExpectedIdentities, MycStateMetadata}; 27 use nostr::RelayUrl as RadrootsNostrRelayUrl; 28 use radroots_nostr::event::{ 29 ApplicationHandlerSpec as RadrootsNostrApplicationHandlerSpec, Event as RadrootsNostrEvent, 30 Kind as RadrootsNostrKind, Metadata as RadrootsNostrMetadata, 31 Timestamp as RadrootsNostrTimestamp, 32 build_application_handler as radroots_nostr_build_application_handler_event, 33 }; 34 35 /// Maximum exact signed-event bytes admitted from the configured event bound. 36 pub const MYC_DISCOVERY_DOCUMENT_MAX_BYTES: usize = 524_288; 37 /// Maximum deterministic NIP-05 projection input bytes. 38 pub const MYC_NIP05_PROJECTION_MAX_BYTES: usize = 524_288; 39 /// Maximum deterministic offline NIP-05 document bytes. 40 pub const MYC_NIP05_DOCUMENT_MAX_BYTES: usize = 524_288; 41 42 const NIP46_RPC_KIND: u32 = 24_133; 43 const NIP89_HANDLER_KIND: u16 = 31_990; 44 const DESIRED_DIGEST_DOMAIN: &[u8] = b"radroots.myc.discovery_desired.v1\0"; 45 const GENERATION_ID_DOMAIN: &[u8] = b"radroots.myc.discovery_generation.v1\0"; 46 47 const READ_STATE_SQL: &str = r#"SELECT 48 CASE WHEN typeof(desired_generation_id) = 'blob' AND length(desired_generation_id) = 32 49 THEN desired_generation_id ELSE NULL END AS desired_generation_id, 50 CASE WHEN typeof(desired_job_id) = 'blob' AND length(desired_job_id) = 32 51 THEN desired_job_id ELSE NULL END AS desired_job_id, 52 typeof(current_generation_id) AS current_generation_id_type, 53 CASE WHEN typeof(current_generation_id) = 'blob' AND length(current_generation_id) = 32 54 THEN current_generation_id ELSE NULL END AS current_generation_id, 55 typeof(current_job_id) AS current_job_id_type, 56 CASE WHEN typeof(current_job_id) = 'blob' AND length(current_job_id) = 32 57 THEN current_job_id ELSE NULL END AS current_job_id, 58 updated_at_unix_ms 59 FROM discovery_publication_state 60 WHERE singleton = 1 61 LIMIT 2"#; 62 63 const READ_DOCUMENT_BY_GENERATION_SQL: &str = r#"SELECT 64 CASE WHEN typeof(generation_id) = 'blob' AND length(generation_id) = 32 65 THEN generation_id ELSE NULL END AS generation_id, 66 CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32 67 THEN event_id ELSE NULL END AS event_id, 68 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 69 THEN event_sha256 ELSE NULL END AS event_sha256, 70 CASE WHEN typeof(event_bytes) = 'blob' AND length(event_bytes) BETWEEN 1 AND 524288 71 THEN event_bytes ELSE NULL END AS event_bytes, 72 CASE WHEN typeof(nip05_projection_sha256) = 'blob' 73 AND length(nip05_projection_sha256) = 32 74 THEN nip05_projection_sha256 ELSE NULL END AS nip05_projection_sha256, 75 CASE WHEN typeof(nip05_projection_bytes) = 'blob' 76 AND length(nip05_projection_bytes) BETWEEN 1 AND 524288 77 THEN nip05_projection_bytes ELSE NULL END AS nip05_projection_bytes 78 FROM discovery_documents 79 WHERE generation_id = ? 80 LIMIT 2"#; 81 82 const READ_DESIRED_SQL: &str = r#"SELECT 83 CASE WHEN typeof(generation_id) = 'blob' AND length(generation_id) = 32 84 THEN generation_id ELSE NULL END AS generation_id, 85 CASE WHEN typeof(normalized_config_sha256) = 'blob' 86 AND length(normalized_config_sha256) = 32 87 THEN normalized_config_sha256 ELSE NULL END AS normalized_config_sha256, 88 CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 89 THEN desired_sha256 ELSE NULL END AS desired_sha256, 90 created_at_unix_ms 91 FROM discovery_desired_state 92 WHERE generation_id = ? 93 LIMIT 2"#; 94 95 const READ_JOB_SQL: &str = r#"SELECT status 96 FROM delivery_jobs 97 WHERE job_id = ? AND source_kind = 'discovery_handler' AND source_id = ? 98 LIMIT 2"#; 99 100 const INSERT_DESIRED_SQL: &str = r#"INSERT INTO discovery_desired_state ( 101 generation_id, normalized_config_sha256, desired_sha256, created_at_unix_ms 102 ) VALUES (?, ?, ?, ?)"#; 103 104 const INSERT_DOCUMENT_SQL: &str = r#"INSERT INTO discovery_documents ( 105 generation_id, event_id, event_sha256, event_bytes, 106 nip05_projection_sha256, nip05_projection_bytes 107 ) VALUES (?, ?, ?, ?, ?, ?)"#; 108 109 const INSERT_STATE_SQL: &str = r#"INSERT INTO discovery_publication_state ( 110 singleton, desired_generation_id, desired_job_id, current_generation_id, 111 current_job_id, updated_at_unix_ms 112 ) VALUES (1, ?, ?, NULL, NULL, ?)"#; 113 114 const REPLACE_DESIRED_SQL: &str = r#"UPDATE discovery_publication_state 115 SET desired_generation_id = ?, desired_job_id = ?, updated_at_unix_ms = ? 116 WHERE singleton = 1 AND updated_at_unix_ms <= ? 117 AND desired_generation_id = ? AND desired_job_id = ?"#; 118 119 const PROMOTE_CURRENT_SQL: &str = r#"UPDATE discovery_publication_state 120 SET current_generation_id = desired_generation_id, 121 current_job_id = desired_job_id, 122 updated_at_unix_ms = ? 123 WHERE singleton = 1 AND desired_generation_id = ? AND desired_job_id = ? 124 AND updated_at_unix_ms <= ? 125 AND (current_generation_id IS NULL OR current_generation_id != desired_generation_id 126 OR current_job_id IS NULL OR current_job_id != desired_job_id)"#; 127 128 /// Stable source-free construction and admission failure classes. 129 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 130 pub enum MycDiscoveryStateErrorKind { 131 Disabled, 132 TooLarge, 133 InvalidProjection, 134 InvalidEvent, 135 IdentityMismatch, 136 } 137 138 impl MycDiscoveryStateErrorKind { 139 /// Returns the stable machine-readable code. 140 #[must_use] 141 pub const fn code(self) -> &'static str { 142 match self { 143 Self::Disabled => "discovery_state_disabled", 144 Self::TooLarge => "discovery_state_too_large", 145 Self::InvalidProjection => "discovery_projection_invalid", 146 Self::InvalidEvent => "discovery_event_invalid", 147 Self::IdentityMismatch => "discovery_identity_mismatch", 148 } 149 } 150 } 151 152 /// Source-free failure at the discovery-state boundary. 153 #[derive(Clone, Copy, PartialEq, Eq)] 154 pub struct MycDiscoveryStateError { 155 kind: MycDiscoveryStateErrorKind, 156 } 157 158 impl MycDiscoveryStateError { 159 const fn new(kind: MycDiscoveryStateErrorKind) -> Self { 160 Self { kind } 161 } 162 163 /// Returns the stable failure class. 164 #[must_use] 165 pub const fn kind(self) -> MycDiscoveryStateErrorKind { 166 self.kind 167 } 168 169 /// Returns the stable machine-readable code. 170 #[must_use] 171 pub const fn code(self) -> &'static str { 172 self.kind.code() 173 } 174 } 175 176 impl fmt::Display for MycDiscoveryStateError { 177 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 178 formatter.write_str(match self.kind { 179 MycDiscoveryStateErrorKind::Disabled => "discovery state is not enabled", 180 MycDiscoveryStateErrorKind::TooLarge => "discovery state exceeds its configured bound", 181 MycDiscoveryStateErrorKind::InvalidProjection => { 182 "discovery projection inputs are invalid" 183 } 184 MycDiscoveryStateErrorKind::InvalidEvent => "discovery event is invalid", 185 MycDiscoveryStateErrorKind::IdentityMismatch => { 186 "discovery event identity does not match configuration" 187 } 188 }) 189 } 190 } 191 192 impl fmt::Debug for MycDiscoveryStateError { 193 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 194 formatter 195 .debug_struct("MycDiscoveryStateError") 196 .field("kind", &self.kind) 197 .finish() 198 } 199 } 200 201 impl Error for MycDiscoveryStateError {} 202 203 macro_rules! redacted_id { 204 ($name:ident, $debug:literal) => { 205 #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] 206 pub struct $name([u8; 32]); 207 208 impl $name { 209 /// Returns the exact identity bytes. 210 #[must_use] 211 pub const fn as_bytes(&self) -> &[u8; 32] { 212 &self.0 213 } 214 } 215 216 impl fmt::Debug for $name { 217 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 218 formatter.write_str($debug) 219 } 220 } 221 }; 222 } 223 224 redacted_id!( 225 MycDiscoveryGenerationId, 226 "MycDiscoveryGenerationId([redacted])" 227 ); 228 redacted_id!( 229 MycDiscoveryDocumentDigest, 230 "MycDiscoveryDocumentDigest([redacted])" 231 ); 232 redacted_id!( 233 MycNip05ProjectionDigest, 234 "MycNip05ProjectionDigest([redacted])" 235 ); 236 redacted_id!(MycNip05DocumentDigest, "MycNip05DocumentDigest([redacted])"); 237 238 /// Explicit discovery generation selected for an offline NIP-05 export. 239 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 240 pub enum MycNip05ExportSelection { 241 Desired, 242 Current, 243 } 244 245 /// Canonical offline NIP-05/NIP-46 document derived from verified state. 246 #[derive(Clone, PartialEq, Eq)] 247 pub struct MycNip05Document { 248 selection: MycNip05ExportSelection, 249 domain: Box<str>, 250 digest: MycNip05DocumentDigest, 251 bytes: Box<[u8]>, 252 } 253 254 impl MycNip05Document { 255 #[must_use] 256 pub const fn selection(&self) -> MycNip05ExportSelection { 257 self.selection 258 } 259 260 /// Returns the configured domain associated with this explicit export. 261 #[must_use] 262 pub fn domain(&self) -> &str { 263 &self.domain 264 } 265 266 /// Returns the exact compact UTF-8 JSON bytes to export. 267 #[must_use] 268 pub fn bytes(&self) -> &[u8] { 269 &self.bytes 270 } 271 272 /// Returns the SHA-256 identity of the exact export bytes. 273 #[must_use] 274 pub const fn digest(&self) -> MycNip05DocumentDigest { 275 self.digest 276 } 277 } 278 279 impl fmt::Debug for MycNip05Document { 280 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 281 formatter 282 .debug_struct("MycNip05Document") 283 .field("selection", &self.selection) 284 .field("domain", &"[redacted]") 285 .field("bytes", &"[redacted]") 286 .field("digest", &"[redacted]") 287 .finish() 288 } 289 } 290 291 #[derive(Clone, PartialEq, Eq)] 292 pub(crate) struct MycDiscoveryPolicies { 293 domain: Box<str>, 294 handler_identifier: Box<str>, 295 author_public_key: Box<str>, 296 public_relays: Box<[Box<str>]>, 297 nostrconnect_url: Option<Box<str>>, 298 metadata_json: Box<str>, 299 event_max_bytes: usize, 300 } 301 302 impl MycDiscoveryPolicies { 303 pub(crate) fn from_normalized( 304 normalized: &serde_json::Value, 305 identities: &MycExpectedIdentities, 306 ) -> Result<Option<Self>, MycDiscoveryStateError> { 307 let discovery = normalized 308 .pointer("/discovery") 309 .and_then(serde_json::Value::as_object) 310 .ok_or_else(|| { 311 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 312 })?; 313 let enabled = discovery 314 .get("enabled") 315 .and_then(serde_json::Value::as_bool) 316 .ok_or_else(|| { 317 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 318 })?; 319 if !enabled { 320 return Ok(None); 321 } 322 let author = identities.discovery().ok_or_else(|| { 323 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::IdentityMismatch) 324 })?; 325 let string = |name: &str| { 326 discovery 327 .get(name) 328 .and_then(serde_json::Value::as_str) 329 .filter(|value| !value.is_empty()) 330 .ok_or_else(|| { 331 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 332 }) 333 }; 334 let public_ids = discovery 335 .get("public_relay_ids") 336 .and_then(serde_json::Value::as_array) 337 .ok_or_else(|| { 338 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 339 })?; 340 let relay_map = normalized 341 .pointer("/relays") 342 .and_then(serde_json::Value::as_array) 343 .ok_or_else(|| { 344 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 345 })? 346 .iter() 347 .map(|relay| { 348 let id = relay.pointer("/id").and_then(serde_json::Value::as_str)?; 349 let url = relay.pointer("/url").and_then(serde_json::Value::as_str)?; 350 Some((id, url)) 351 }) 352 .collect::<Option<BTreeMap<_, _>>>() 353 .ok_or_else(|| { 354 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 355 })?; 356 let public_relays = public_ids 357 .iter() 358 .map(|id| { 359 let id = id.as_str()?; 360 let url = *relay_map.get(id)?; 361 RadrootsNostrRelayUrl::parse(url).ok()?; 362 Some(Box::<str>::from(url)) 363 }) 364 .collect::<Option<Vec<_>>>() 365 .filter(|relays| !relays.is_empty()) 366 .ok_or_else(|| { 367 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 368 })? 369 .into_boxed_slice(); 370 let metadata = discovery 371 .get("metadata") 372 .and_then(serde_json::Value::as_object) 373 .ok_or_else(|| { 374 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 375 })?; 376 let optional = |name: &str| { 377 metadata 378 .get(name) 379 .and_then(serde_json::Value::as_str) 380 .map(str::trim) 381 .filter(|value| !value.is_empty()) 382 .map(ToOwned::to_owned) 383 }; 384 let metadata = RadrootsNostrMetadata { 385 name: optional("name"), 386 display_name: optional("display_name"), 387 about: optional("about"), 388 website: optional("website"), 389 picture: optional("picture"), 390 ..RadrootsNostrMetadata::default() 391 }; 392 let metadata_json = if metadata.name.is_none() 393 && metadata.display_name.is_none() 394 && metadata.about.is_none() 395 && metadata.website.is_none() 396 && metadata.picture.is_none() 397 { 398 Box::<str>::from("") 399 } else { 400 serde_json::to_string(&metadata) 401 .map_err(|_| { 402 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 403 })? 404 .into_boxed_str() 405 }; 406 let template = discovery 407 .get("nostrconnect_url_template") 408 .and_then(serde_json::Value::as_str) 409 .filter(|value| !value.is_empty()); 410 let signer = identities.user().as_hex(); 411 let nostrconnect_url = template 412 .map(|template| render_nostrconnect_url(template, signer, &public_relays)) 413 .transpose()? 414 .map(String::into_boxed_str); 415 let event_max_bytes = normalized 416 .pointer("/resource_limits/events/wire_bytes") 417 .and_then(serde_json::Value::as_u64) 418 .and_then(|value| usize::try_from(value).ok()) 419 .filter(|value| (1..=MYC_DISCOVERY_DOCUMENT_MAX_BYTES).contains(value)) 420 .ok_or_else(|| { 421 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 422 })?; 423 Ok(Some(Self { 424 domain: string("domain")?.into(), 425 handler_identifier: string("handler_identifier")?.into(), 426 author_public_key: author.as_hex().into(), 427 public_relays, 428 nostrconnect_url, 429 metadata_json, 430 event_max_bytes, 431 })) 432 } 433 } 434 435 pub(crate) fn prepare_discovery_signing_bytes( 436 metadata: &MycStateMetadata, 437 created_at_unix_seconds: u64, 438 ) -> Result<Box<[u8]>, MycDiscoveryStateError> { 439 let policy = metadata 440 .discovery_policies() 441 .ok_or_else(|| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::Disabled))?; 442 let metadata = if policy.metadata_json.is_empty() { 443 RadrootsNostrMetadata::default() 444 } else { 445 serde_json::from_str(&policy.metadata_json).map_err(|_| { 446 MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) 447 })? 448 }; 449 let mut spec = RadrootsNostrApplicationHandlerSpec::new(vec![NIP46_RPC_KIND]) 450 .with_identifier(policy.handler_identifier.to_string()) 451 .with_relays( 452 policy 453 .public_relays 454 .iter() 455 .map(ToString::to_string) 456 .collect(), 457 ) 458 .with_metadata(metadata); 459 if let Some(url) = &policy.nostrconnect_url { 460 spec = spec.with_nostr_connect_url(url.to_string()); 461 } 462 let public_key = policy 463 .author_public_key 464 .parse() 465 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::IdentityMismatch))?; 466 let request = radroots_nostr_build_application_handler_event(&spec) 467 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))? 468 .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at_unix_seconds)) 469 .into_external_signing_request(public_key) 470 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))?; 471 let bytes = serde_json::to_vec(&request) 472 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))?; 473 if bytes.is_empty() || bytes.len() > policy.event_max_bytes { 474 return Err(MycDiscoveryStateError::new( 475 MycDiscoveryStateErrorKind::TooLarge, 476 )); 477 } 478 Ok(bytes.into_boxed_slice()) 479 } 480 481 impl fmt::Debug for MycDiscoveryPolicies { 482 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 483 formatter 484 .debug_struct("MycDiscoveryPolicies") 485 .field("relay_count", &self.public_relays.len()) 486 .field("values", &"[redacted]") 487 .finish() 488 } 489 } 490 491 /// Exact signature-verified discovery generation prepared outside SQLite. 492 pub struct MycDiscoveryCommitRequest { 493 generation_id: MycDiscoveryGenerationId, 494 desired_digest: MycDiscoveryDocumentDigest, 495 configuration_digest: [u8; 32], 496 event_id: [u8; 32], 497 event_digest: MycDeliveryArtifactDigest, 498 event_bytes: Box<[u8]>, 499 projection_digest: MycNip05ProjectionDigest, 500 projection_bytes: Box<[u8]>, 501 created_at: MycDeliveryTimeUnixMs, 502 } 503 504 impl MycDiscoveryCommitRequest { 505 /// Validates exact canonical event bytes, signature, identity, NIP-89 semantics, and bounds. 506 pub fn new( 507 metadata: &MycStateMetadata, 508 event_bytes: &[u8], 509 created_at: MycDeliveryTimeUnixMs, 510 ) -> Result<Self, MycDiscoveryStateError> { 511 let policy = metadata 512 .discovery_policies() 513 .ok_or_else(|| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::Disabled))?; 514 if event_bytes.is_empty() || event_bytes.len() > policy.event_max_bytes { 515 return Err(MycDiscoveryStateError::new( 516 MycDiscoveryStateErrorKind::TooLarge, 517 )); 518 } 519 let event: RadrootsNostrEvent = serde_json::from_slice(event_bytes) 520 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidEvent))?; 521 let canonical = serde_json::to_vec(&event) 522 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidEvent))?; 523 if canonical != event_bytes || event.verify().is_err() { 524 return Err(MycDiscoveryStateError::new( 525 MycDiscoveryStateErrorKind::InvalidEvent, 526 )); 527 } 528 validate_event(&event, policy)?; 529 let projection_bytes = projection_bytes(policy)?; 530 if projection_bytes.is_empty() || projection_bytes.len() > MYC_NIP05_PROJECTION_MAX_BYTES { 531 return Err(MycDiscoveryStateError::new( 532 MycDiscoveryStateErrorKind::TooLarge, 533 )); 534 } 535 let event_digest_bytes: [u8; 32] = Sha256::digest(event_bytes).into(); 536 let projection_digest_bytes: [u8; 32] = Sha256::digest(&projection_bytes).into(); 537 let configuration_digest = *metadata.configuration_digest().as_bytes(); 538 let mut desired_hasher = Sha256::new(); 539 desired_hasher.update(DESIRED_DIGEST_DOMAIN); 540 desired_hasher.update(configuration_digest); 541 desired_hasher.update(event_digest_bytes); 542 desired_hasher.update(projection_digest_bytes); 543 let desired_digest_bytes: [u8; 32] = desired_hasher.finalize().into(); 544 let mut generation_hasher = Sha256::new(); 545 generation_hasher.update(GENERATION_ID_DOMAIN); 546 generation_hasher.update(desired_digest_bytes); 547 let generation_id = MycDiscoveryGenerationId(generation_hasher.finalize().into()); 548 Ok(Self { 549 generation_id, 550 desired_digest: MycDiscoveryDocumentDigest(desired_digest_bytes), 551 configuration_digest, 552 event_id: *event.id.as_bytes(), 553 event_digest: MycDeliveryArtifactDigest::from_bytes(event_digest_bytes), 554 event_bytes: event_bytes.into(), 555 projection_digest: MycNip05ProjectionDigest(projection_digest_bytes), 556 projection_bytes: projection_bytes.into_boxed_slice(), 557 created_at, 558 }) 559 } 560 561 /// Returns the deterministic desired-state generation identity. 562 #[must_use] 563 pub const fn generation_id(&self) -> MycDiscoveryGenerationId { 564 self.generation_id 565 } 566 567 /// Returns the exact digest binding normalized configuration, event, and projection inputs. 568 #[must_use] 569 pub const fn desired_digest(&self) -> MycDiscoveryDocumentDigest { 570 self.desired_digest 571 } 572 573 pub(crate) const fn event_digest(&self) -> MycDeliveryArtifactDigest { 574 self.event_digest 575 } 576 577 pub(crate) const fn projection_digest(&self) -> MycNip05ProjectionDigest { 578 self.projection_digest 579 } 580 581 fn owned(&self) -> Self { 582 Self { 583 generation_id: self.generation_id, 584 desired_digest: self.desired_digest, 585 configuration_digest: self.configuration_digest, 586 event_id: self.event_id, 587 event_digest: self.event_digest, 588 event_bytes: self.event_bytes.clone(), 589 projection_digest: self.projection_digest, 590 projection_bytes: self.projection_bytes.clone(), 591 created_at: self.created_at, 592 } 593 } 594 } 595 596 impl fmt::Debug for MycDiscoveryCommitRequest { 597 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 598 formatter.write_str("MycDiscoveryCommitRequest([redacted])") 599 } 600 } 601 602 /// Immutable exact document bytes retained for a discovery generation. 603 #[derive(Clone, PartialEq, Eq)] 604 pub struct MycDiscoveryDocumentRecord { 605 generation_id: MycDiscoveryGenerationId, 606 event_id: [u8; 32], 607 event_digest: MycDeliveryArtifactDigest, 608 event_bytes: Box<[u8]>, 609 projection_digest: MycNip05ProjectionDigest, 610 projection_bytes: Box<[u8]>, 611 } 612 613 impl MycDiscoveryDocumentRecord { 614 #[must_use] 615 pub const fn generation_id(&self) -> MycDiscoveryGenerationId { 616 self.generation_id 617 } 618 619 /// Returns the sole exact publishable event bytes. 620 #[must_use] 621 pub fn event_bytes(&self) -> &[u8] { 622 &self.event_bytes 623 } 624 625 /// Returns deterministic NIP-05 projection inputs, not a hosted response. 626 #[must_use] 627 pub fn nip05_projection_bytes(&self) -> &[u8] { 628 &self.projection_bytes 629 } 630 631 #[must_use] 632 pub const fn event_digest(&self) -> MycDeliveryArtifactDigest { 633 self.event_digest 634 } 635 636 /// Returns the verified NIP-01 event identity. 637 #[must_use] 638 pub const fn event_id(&self) -> &[u8; 32] { 639 &self.event_id 640 } 641 642 /// Returns the exact digest of the deterministic NIP-05 projection inputs. 643 #[must_use] 644 pub const fn nip05_projection_digest(&self) -> MycNip05ProjectionDigest { 645 self.projection_digest 646 } 647 } 648 649 impl fmt::Debug for MycDiscoveryDocumentRecord { 650 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 651 formatter 652 .debug_struct("MycDiscoveryDocumentRecord") 653 .field("event_bytes", &"[redacted]") 654 .field("nip05_projection_bytes", &"[redacted]") 655 .finish() 656 } 657 } 658 659 /// Current desired and proven-delivered discovery generations. 660 #[derive(Clone, PartialEq, Eq)] 661 pub struct MycDiscoveryPublicationState { 662 desired_generation_id: MycDiscoveryGenerationId, 663 desired_job_id: MycDeliveryJobId, 664 current_generation_id: Option<MycDiscoveryGenerationId>, 665 current_job_id: Option<MycDeliveryJobId>, 666 updated_at: MycDeliveryTimeUnixMs, 667 } 668 669 impl MycDiscoveryPublicationState { 670 #[must_use] 671 pub const fn desired_generation_id(&self) -> MycDiscoveryGenerationId { 672 self.desired_generation_id 673 } 674 #[must_use] 675 pub const fn desired_job_id(&self) -> MycDeliveryJobId { 676 self.desired_job_id 677 } 678 #[must_use] 679 pub const fn current_generation_id(&self) -> Option<MycDiscoveryGenerationId> { 680 self.current_generation_id 681 } 682 #[must_use] 683 pub const fn current_job_id(&self) -> Option<MycDeliveryJobId> { 684 self.current_job_id 685 } 686 } 687 688 impl fmt::Debug for MycDiscoveryPublicationState { 689 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 690 formatter 691 .debug_struct("MycDiscoveryPublicationState") 692 .field("has_current", &self.current_generation_id.is_some()) 693 .field("updated_at", &self.updated_at) 694 .finish() 695 } 696 } 697 698 /// Atomic result of one discovery desired-state commit. 699 #[derive(Clone, PartialEq, Eq)] 700 pub struct MycDiscoveryCommitRecord { 701 state: MycDiscoveryPublicationState, 702 document: MycDiscoveryDocumentRecord, 703 job: MycDeliveryJobRecord, 704 } 705 706 impl MycDiscoveryCommitRecord { 707 #[must_use] 708 pub const fn state(&self) -> &MycDiscoveryPublicationState { 709 &self.state 710 } 711 #[must_use] 712 pub const fn document(&self) -> &MycDiscoveryDocumentRecord { 713 &self.document 714 } 715 #[must_use] 716 pub const fn job(&self) -> &MycDeliveryJobRecord { 717 &self.job 718 } 719 } 720 721 impl fmt::Debug for MycDiscoveryCommitRecord { 722 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 723 formatter.write_str("MycDiscoveryCommitRecord([redacted])") 724 } 725 } 726 727 /// Created or exact replay outcome of one atomic discovery commit. 728 #[derive(Clone, PartialEq, Eq)] 729 pub enum MycDiscoveryCommitAdmission { 730 Created(MycDiscoveryCommitRecord), 731 ExactReplay(MycDiscoveryCommitRecord), 732 } 733 734 impl MycDiscoveryCommitAdmission { 735 #[must_use] 736 pub const fn record(&self) -> &MycDiscoveryCommitRecord { 737 match self { 738 Self::Created(record) | Self::ExactReplay(record) => record, 739 } 740 } 741 } 742 743 impl fmt::Debug for MycDiscoveryCommitAdmission { 744 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 745 formatter.write_str(match self { 746 Self::Created(_) => "MycDiscoveryCommitAdmission::Created([redacted])", 747 Self::ExactReplay(_) => "MycDiscoveryCommitAdmission::ExactReplay([redacted])", 748 }) 749 } 750 } 751 752 impl MycStateRepository<'_> { 753 /// Atomically commits desired state, exact documents, targets, and initial delivery evidence. 754 pub async fn commit_discovery_desired_state( 755 &self, 756 request: &MycDiscoveryCommitRequest, 757 ) -> Result<MycDiscoveryCommitAdmission, MycStateRepositoryError> { 758 let request = request.owned(); 759 let expected = PersistedMetadata::from(self.expected()); 760 let delivery_policy = self.expected().delivery_policies().clone(); 761 let discovery_policy = self.expected().discovery_policies().cloned(); 762 self.host() 763 .transaction(move |transaction| { 764 Box::pin(async move { 765 verify_metadata(transaction, &expected).await?; 766 let discovery_policy = discovery_policy 767 .as_ref() 768 .ok_or(DiscoveryOperationError::Binding)?; 769 commit_desired(transaction, &request, &delivery_policy, discovery_policy).await 770 }) 771 }) 772 .await 773 .map_err(map_transaction_error) 774 } 775 776 /// Reads the bounded desired/current discovery state, when configured and committed. 777 pub async fn read_discovery_publication_state( 778 &self, 779 ) -> Result<Option<MycDiscoveryPublicationState>, MycStateRepositoryError> { 780 let expected = PersistedMetadata::from(self.expected()); 781 self.host() 782 .transaction(move |transaction| { 783 Box::pin(async move { 784 verify_metadata(transaction, &expected).await?; 785 read_state(transaction).await 786 }) 787 }) 788 .await 789 .map_err(map_transaction_error) 790 } 791 792 /// Reads exact committed bytes for one discovery delivery job. 793 pub async fn read_discovery_document_for_job( 794 &self, 795 job_id: MycDeliveryJobId, 796 ) -> Result<Option<MycDiscoveryDocumentRecord>, MycStateRepositoryError> { 797 let expected = PersistedMetadata::from(self.expected()); 798 let policy = self.expected().discovery_policies().cloned(); 799 self.host() 800 .transaction(move |transaction| { 801 Box::pin(async move { 802 verify_metadata(transaction, &expected).await?; 803 let generation = discovery_generation_for_job(transaction, job_id).await?; 804 match (generation, policy.as_ref()) { 805 (Some(generation), Some(policy)) => { 806 read_document(transaction, generation, policy).await 807 } 808 (Some(_), None) => Err(DiscoveryOperationError::Binding), 809 (None, _) => Ok(None), 810 } 811 }) 812 }) 813 .await 814 .map_err(map_transaction_error) 815 } 816 817 /// Advances current state only from the current desired job's proven-delivered evidence. 818 pub async fn promote_delivered_discovery_state( 819 &self, 820 job_id: MycDeliveryJobId, 821 observed_at: MycDeliveryTimeUnixMs, 822 ) -> Result<MycDiscoveryPublicationState, MycStateRepositoryError> { 823 let expected = PersistedMetadata::from(self.expected()); 824 self.host() 825 .transaction(move |transaction| { 826 Box::pin(async move { 827 verify_metadata(transaction, &expected).await?; 828 promote_current(transaction, job_id, observed_at).await 829 }) 830 }) 831 .await 832 .map_err(map_transaction_error) 833 } 834 835 /// Renders an explicit desired or proven-current NIP-05 document offline. 836 /// 837 /// This performs only verified state reads. It does not host, write, or 838 /// publish the returned bytes and is valid through an inspection host. 839 pub async fn render_offline_nip05( 840 &self, 841 selection: MycNip05ExportSelection, 842 ) -> Result<MycNip05Document, MycStateRepositoryError> { 843 let expected = PersistedMetadata::from(self.expected()); 844 let policy = self.expected().discovery_policies().cloned(); 845 self.host() 846 .transaction(move |transaction| { 847 Box::pin(async move { 848 verify_metadata(transaction, &expected).await?; 849 let policy = policy.as_ref().ok_or(DiscoveryOperationError::Binding)?; 850 let state = read_state(transaction) 851 .await? 852 .ok_or(DiscoveryOperationError::Binding)?; 853 let generation = match selection { 854 MycNip05ExportSelection::Desired => state.desired_generation_id, 855 MycNip05ExportSelection::Current => state 856 .current_generation_id 857 .ok_or(DiscoveryOperationError::Binding)?, 858 }; 859 let document = read_document(transaction, generation, policy) 860 .await? 861 .ok_or(DiscoveryOperationError::Binding)?; 862 render_nip05(selection, &document) 863 }) 864 }) 865 .await 866 .map_err(map_transaction_error) 867 } 868 } 869 870 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 871 pub(crate) enum DiscoveryOperationError { 872 Binding, 873 Storage, 874 } 875 876 #[cfg(any(target_os = "linux", target_os = "macos"))] 877 impl From<DiscoveryOperationError> for AdminJournalOperationError { 878 fn from(error: DiscoveryOperationError) -> Self { 879 match error { 880 DiscoveryOperationError::Binding => Self::Binding, 881 DiscoveryOperationError::Storage => Self::Storage, 882 } 883 } 884 } 885 886 #[cfg(any(target_os = "linux", target_os = "macos"))] 887 pub(crate) async fn apply_admin_discovery_publish( 888 transaction: &mut ServiceSqliteTransaction<'_>, 889 request: &MycDiscoveryCommitRequest, 890 delivery_policy: &crate::state_delivery::MycDeliveryPolicies, 891 discovery_policy: &MycDiscoveryPolicies, 892 ) -> Result<MycDiscoveryCommitAdmission, AdminJournalOperationError> { 893 commit_desired(transaction, request, delivery_policy, discovery_policy) 894 .await 895 .map_err(Into::into) 896 } 897 898 impl From<DeliveryOperationError> for DiscoveryOperationError { 899 fn from(error: DeliveryOperationError) -> Self { 900 match error { 901 DeliveryOperationError::Binding => Self::Binding, 902 DeliveryOperationError::Storage => Self::Storage, 903 } 904 } 905 } 906 907 async fn verify_metadata( 908 transaction: &mut ServiceSqliteTransaction<'_>, 909 expected: &PersistedMetadata, 910 ) -> Result<(), DiscoveryOperationError> { 911 require_expected_metadata(transaction, expected) 912 .await 913 .map_err(|error| match error { 914 RepositoryOperationError::Binding => DiscoveryOperationError::Binding, 915 RepositoryOperationError::Storage => DiscoveryOperationError::Storage, 916 }) 917 } 918 919 async fn commit_desired( 920 transaction: &mut ServiceSqliteTransaction<'_>, 921 request: &MycDiscoveryCommitRequest, 922 delivery_policy: &crate::state_delivery::MycDeliveryPolicies, 923 discovery_policy: &MycDiscoveryPolicies, 924 ) -> Result<MycDiscoveryCommitAdmission, DiscoveryOperationError> { 925 let existing_desired = read_desired(transaction, request.generation_id).await?; 926 let existing_document = 927 read_document(transaction, request.generation_id, discovery_policy).await?; 928 let exact_existing = match (&existing_desired, &existing_document) { 929 (Some(desired), Some(document)) => { 930 desired.configuration_digest == request.configuration_digest 931 && desired.desired_digest == request.desired_digest 932 && desired.created_at == request.created_at 933 && exact_document(document, request) 934 } 935 (None, None) => false, 936 _ => return Err(DiscoveryOperationError::Binding), 937 }; 938 if existing_desired.is_some() && !exact_existing { 939 return Err(DiscoveryOperationError::Binding); 940 } 941 if existing_desired.is_none() { 942 insert_desired(transaction, request).await?; 943 insert_document(transaction, request).await?; 944 } 945 let source = MycDeliverySource::discovery_handler(*request.generation_id.as_bytes()); 946 let job = create_job( 947 transaction, 948 source, 949 request.event_digest, 950 request.created_at, 951 delivery_policy, 952 ) 953 .await?; 954 let job_record = job.record().clone(); 955 let prior = read_state(transaction).await?; 956 let exact_pointer = prior.as_ref().is_some_and(|state| { 957 state.desired_generation_id == request.generation_id 958 && state.desired_job_id == job_record.id() 959 }); 960 let created = existing_desired.is_none(); 961 match prior { 962 None => { 963 let result = sqlx::query(INSERT_STATE_SQL) 964 .bind(request.generation_id.as_bytes().as_slice()) 965 .bind(job_record.id().as_bytes().as_slice()) 966 .bind(request.created_at.sqlite_value()) 967 .execute(&mut *transaction) 968 .await 969 .map_err(|_| DiscoveryOperationError::Storage)?; 970 require_one(result.rows_affected())?; 971 } 972 Some(state) if exact_pointer => { 973 if !exact_existing { 974 return Err(DiscoveryOperationError::Binding); 975 } 976 let document = existing_document.ok_or(DiscoveryOperationError::Binding)?; 977 return Ok(MycDiscoveryCommitAdmission::ExactReplay( 978 MycDiscoveryCommitRecord { 979 state, 980 document, 981 job: job_record, 982 }, 983 )); 984 } 985 Some(state) => { 986 if request.created_at <= state.updated_at { 987 return Err(DiscoveryOperationError::Binding); 988 } 989 let result = sqlx::query(REPLACE_DESIRED_SQL) 990 .bind(request.generation_id.as_bytes().as_slice()) 991 .bind(job_record.id().as_bytes().as_slice()) 992 .bind(request.created_at.sqlite_value()) 993 .bind(request.created_at.sqlite_value()) 994 .bind(state.desired_generation_id.as_bytes().as_slice()) 995 .bind(state.desired_job_id.as_bytes().as_slice()) 996 .execute(&mut *transaction) 997 .await 998 .map_err(|_| DiscoveryOperationError::Storage)?; 999 require_one(result.rows_affected())?; 1000 } 1001 } 1002 let state = read_state(transaction) 1003 .await? 1004 .ok_or(DiscoveryOperationError::Binding)?; 1005 let document = read_document(transaction, request.generation_id, discovery_policy) 1006 .await? 1007 .ok_or(DiscoveryOperationError::Binding)?; 1008 let record = MycDiscoveryCommitRecord { 1009 state, 1010 document, 1011 job: job_record, 1012 }; 1013 if created { 1014 Ok(MycDiscoveryCommitAdmission::Created(record)) 1015 } else { 1016 Err(DiscoveryOperationError::Binding) 1017 } 1018 } 1019 1020 #[derive(Clone, PartialEq, Eq)] 1021 struct DesiredRecord { 1022 configuration_digest: [u8; 32], 1023 desired_digest: MycDiscoveryDocumentDigest, 1024 created_at: MycDeliveryTimeUnixMs, 1025 } 1026 1027 async fn insert_desired( 1028 transaction: &mut ServiceSqliteTransaction<'_>, 1029 request: &MycDiscoveryCommitRequest, 1030 ) -> Result<(), DiscoveryOperationError> { 1031 let result = sqlx::query(INSERT_DESIRED_SQL) 1032 .bind(request.generation_id.as_bytes().as_slice()) 1033 .bind(request.configuration_digest.as_slice()) 1034 .bind(request.desired_digest.as_bytes().as_slice()) 1035 .bind(request.created_at.sqlite_value()) 1036 .execute(&mut *transaction) 1037 .await 1038 .map_err(|_| DiscoveryOperationError::Storage)?; 1039 require_one(result.rows_affected()) 1040 } 1041 1042 async fn insert_document( 1043 transaction: &mut ServiceSqliteTransaction<'_>, 1044 request: &MycDiscoveryCommitRequest, 1045 ) -> Result<(), DiscoveryOperationError> { 1046 let result = sqlx::query(INSERT_DOCUMENT_SQL) 1047 .bind(request.generation_id.as_bytes().as_slice()) 1048 .bind(request.event_id.as_slice()) 1049 .bind(request.event_digest.as_bytes().as_slice()) 1050 .bind(request.event_bytes.as_ref()) 1051 .bind(request.projection_digest.as_bytes().as_slice()) 1052 .bind(request.projection_bytes.as_ref()) 1053 .execute(&mut *transaction) 1054 .await 1055 .map_err(|_| DiscoveryOperationError::Storage)?; 1056 require_one(result.rows_affected()) 1057 } 1058 1059 async fn read_desired( 1060 transaction: &mut ServiceSqliteTransaction<'_>, 1061 generation: MycDiscoveryGenerationId, 1062 ) -> Result<Option<DesiredRecord>, DiscoveryOperationError> { 1063 let rows = sqlx::query(READ_DESIRED_SQL) 1064 .bind(generation.as_bytes().as_slice()) 1065 .fetch_all(&mut *transaction) 1066 .await 1067 .map_err(|_| DiscoveryOperationError::Storage)?; 1068 if rows.len() > 1 { 1069 return Err(DiscoveryOperationError::Binding); 1070 } 1071 rows.first() 1072 .map(|row| { 1073 if blob32(row, "generation_id")? != *generation.as_bytes() { 1074 return Err(DiscoveryOperationError::Binding); 1075 } 1076 Ok(DesiredRecord { 1077 configuration_digest: blob32(row, "normalized_config_sha256")?, 1078 desired_digest: MycDiscoveryDocumentDigest(blob32(row, "desired_sha256")?), 1079 created_at: time(row, "created_at_unix_ms")?, 1080 }) 1081 }) 1082 .transpose() 1083 } 1084 1085 async fn read_document( 1086 transaction: &mut ServiceSqliteTransaction<'_>, 1087 generation: MycDiscoveryGenerationId, 1088 policy: &MycDiscoveryPolicies, 1089 ) -> Result<Option<MycDiscoveryDocumentRecord>, DiscoveryOperationError> { 1090 let rows = sqlx::query(READ_DOCUMENT_BY_GENERATION_SQL) 1091 .bind(generation.as_bytes().as_slice()) 1092 .fetch_all(&mut *transaction) 1093 .await 1094 .map_err(|_| DiscoveryOperationError::Storage)?; 1095 if rows.len() > 1 { 1096 return Err(DiscoveryOperationError::Binding); 1097 } 1098 rows.first() 1099 .map(|row| parse_document(row, policy)) 1100 .transpose() 1101 } 1102 1103 pub(crate) async fn verify_document_for_delivery_job( 1104 transaction: &mut ServiceSqliteTransaction<'_>, 1105 job: &MycDeliveryJobRecord, 1106 policy: Option<&MycDiscoveryPolicies>, 1107 ) -> Result<(), DeliveryOperationError> { 1108 if job.source_kind() != crate::state_delivery::MycDeliverySourceKind::DiscoveryHandler { 1109 return Err(DeliveryOperationError::Binding); 1110 } 1111 let generation = discovery_generation_for_job(transaction, job.id()) 1112 .await 1113 .map_err(|error| match error { 1114 DiscoveryOperationError::Binding => DeliveryOperationError::Binding, 1115 DiscoveryOperationError::Storage => DeliveryOperationError::Storage, 1116 })? 1117 .ok_or(DeliveryOperationError::Binding)?; 1118 if generation.as_bytes() != job.source_id() { 1119 return Err(DeliveryOperationError::Binding); 1120 } 1121 let document = read_document( 1122 transaction, 1123 generation, 1124 policy.ok_or(DeliveryOperationError::Binding)?, 1125 ) 1126 .await 1127 .map_err(|error| match error { 1128 DiscoveryOperationError::Binding => DeliveryOperationError::Binding, 1129 DiscoveryOperationError::Storage => DeliveryOperationError::Storage, 1130 })? 1131 .ok_or(DeliveryOperationError::Binding)?; 1132 (document.event_digest() == job.artifact_digest()) 1133 .then_some(()) 1134 .ok_or(DeliveryOperationError::Binding) 1135 } 1136 1137 fn parse_document( 1138 row: &sqlx::sqlite::SqliteRow, 1139 policy: &MycDiscoveryPolicies, 1140 ) -> Result<MycDiscoveryDocumentRecord, DiscoveryOperationError> { 1141 let generation_id = MycDiscoveryGenerationId(blob32(row, "generation_id")?); 1142 let event_id = blob32(row, "event_id")?; 1143 let event_digest = MycDeliveryArtifactDigest::from_bytes(blob32(row, "event_sha256")?); 1144 let event_bytes = bounded_blob(row, "event_bytes", MYC_DISCOVERY_DOCUMENT_MAX_BYTES)?; 1145 let projection_digest = MycNip05ProjectionDigest(blob32(row, "nip05_projection_sha256")?); 1146 let stored_projection_bytes = bounded_blob( 1147 row, 1148 "nip05_projection_bytes", 1149 MYC_NIP05_PROJECTION_MAX_BYTES, 1150 )?; 1151 let actual_event_digest: [u8; 32] = Sha256::digest(&event_bytes).into(); 1152 let actual_projection_digest: [u8; 32] = Sha256::digest(&stored_projection_bytes).into(); 1153 let event: RadrootsNostrEvent = 1154 serde_json::from_slice(&event_bytes).map_err(|_| DiscoveryOperationError::Binding)?; 1155 let canonical = serde_json::to_vec(&event).map_err(|_| DiscoveryOperationError::Binding)?; 1156 let expected_projection = 1157 projection_bytes(policy).map_err(|_| DiscoveryOperationError::Binding)?; 1158 if actual_event_digest != *event_digest.as_bytes() 1159 || actual_projection_digest != *projection_digest.as_bytes() 1160 || canonical.as_slice() != event_bytes.as_ref() 1161 || event.verify().is_err() 1162 || event.id.as_bytes() != &event_id 1163 || validate_event(&event, policy).is_err() 1164 || expected_projection.as_slice() != stored_projection_bytes.as_ref() 1165 { 1166 return Err(DiscoveryOperationError::Binding); 1167 } 1168 Ok(MycDiscoveryDocumentRecord { 1169 generation_id, 1170 event_id, 1171 event_digest, 1172 event_bytes, 1173 projection_digest, 1174 projection_bytes: stored_projection_bytes, 1175 }) 1176 } 1177 1178 async fn read_state( 1179 transaction: &mut ServiceSqliteTransaction<'_>, 1180 ) -> Result<Option<MycDiscoveryPublicationState>, DiscoveryOperationError> { 1181 let rows = sqlx::query(READ_STATE_SQL) 1182 .fetch_all(&mut *transaction) 1183 .await 1184 .map_err(|_| DiscoveryOperationError::Storage)?; 1185 if rows.len() > 1 { 1186 return Err(DiscoveryOperationError::Binding); 1187 } 1188 rows.first().map(parse_state).transpose() 1189 } 1190 1191 fn parse_state( 1192 row: &sqlx::sqlite::SqliteRow, 1193 ) -> Result<MycDiscoveryPublicationState, DiscoveryOperationError> { 1194 let current_generation_id = 1195 optional_blob32(row, "current_generation_id", "current_generation_id_type")? 1196 .map(MycDiscoveryGenerationId); 1197 let current_job_id = optional_blob32(row, "current_job_id", "current_job_id_type")? 1198 .map(MycDeliveryJobId::from_persisted); 1199 if current_generation_id.is_some() != current_job_id.is_some() { 1200 return Err(DiscoveryOperationError::Binding); 1201 } 1202 Ok(MycDiscoveryPublicationState { 1203 desired_generation_id: MycDiscoveryGenerationId(blob32(row, "desired_generation_id")?), 1204 desired_job_id: MycDeliveryJobId::from_persisted(blob32(row, "desired_job_id")?), 1205 current_generation_id, 1206 current_job_id, 1207 updated_at: time(row, "updated_at_unix_ms")?, 1208 }) 1209 } 1210 1211 async fn discovery_generation_for_job( 1212 transaction: &mut ServiceSqliteTransaction<'_>, 1213 job_id: MycDeliveryJobId, 1214 ) -> Result<Option<MycDiscoveryGenerationId>, DiscoveryOperationError> { 1215 let rows = sqlx::query( 1216 "SELECT source_id FROM delivery_jobs \ 1217 WHERE job_id = ? AND source_kind = 'discovery_handler' LIMIT 2", 1218 ) 1219 .bind(job_id.as_bytes().as_slice()) 1220 .fetch_all(&mut *transaction) 1221 .await 1222 .map_err(|_| DiscoveryOperationError::Storage)?; 1223 if rows.len() > 1 { 1224 return Err(DiscoveryOperationError::Binding); 1225 } 1226 rows.first() 1227 .map(|row| blob32(row, "source_id").map(MycDiscoveryGenerationId)) 1228 .transpose() 1229 } 1230 1231 async fn promote_current( 1232 transaction: &mut ServiceSqliteTransaction<'_>, 1233 job_id: MycDeliveryJobId, 1234 observed_at: MycDeliveryTimeUnixMs, 1235 ) -> Result<MycDiscoveryPublicationState, DiscoveryOperationError> { 1236 let state = read_state(transaction) 1237 .await? 1238 .ok_or(DiscoveryOperationError::Binding)?; 1239 if state.desired_job_id != job_id || observed_at < state.updated_at { 1240 return Err(DiscoveryOperationError::Binding); 1241 } 1242 let rows = sqlx::query(READ_JOB_SQL) 1243 .bind(job_id.as_bytes().as_slice()) 1244 .bind(state.desired_generation_id.as_bytes().as_slice()) 1245 .fetch_all(&mut *transaction) 1246 .await 1247 .map_err(|_| DiscoveryOperationError::Storage)?; 1248 if rows.len() != 1 1249 || rows[0] 1250 .try_get::<&str, _>("status") 1251 .ok() 1252 .and_then(MycDeliveryJobStatus::parse) 1253 != Some(MycDeliveryJobStatus::Delivered) 1254 { 1255 return Err(DiscoveryOperationError::Binding); 1256 } 1257 if state.current_generation_id == Some(state.desired_generation_id) 1258 && state.current_job_id == Some(state.desired_job_id) 1259 { 1260 return Ok(state); 1261 } 1262 let result = sqlx::query(PROMOTE_CURRENT_SQL) 1263 .bind(observed_at.sqlite_value()) 1264 .bind(state.desired_generation_id.as_bytes().as_slice()) 1265 .bind(job_id.as_bytes().as_slice()) 1266 .bind(observed_at.sqlite_value()) 1267 .execute(&mut *transaction) 1268 .await 1269 .map_err(|_| DiscoveryOperationError::Storage)?; 1270 require_one(result.rows_affected())?; 1271 read_state(transaction) 1272 .await? 1273 .ok_or(DiscoveryOperationError::Binding) 1274 } 1275 1276 pub(crate) async fn promote_current_if_desired( 1277 transaction: &mut ServiceSqliteTransaction<'_>, 1278 job_id: MycDeliveryJobId, 1279 observed_at: MycDeliveryTimeUnixMs, 1280 ) -> Result<bool, DeliveryOperationError> { 1281 let Some(state) = read_state(transaction).await.map_err(|error| match error { 1282 DiscoveryOperationError::Binding => DeliveryOperationError::Binding, 1283 DiscoveryOperationError::Storage => DeliveryOperationError::Storage, 1284 })? 1285 else { 1286 return Err(DeliveryOperationError::Binding); 1287 }; 1288 if state.desired_job_id != job_id { 1289 return Ok(false); 1290 } 1291 promote_current(transaction, job_id, observed_at) 1292 .await 1293 .map_err(|error| match error { 1294 DiscoveryOperationError::Binding => DeliveryOperationError::Binding, 1295 DiscoveryOperationError::Storage => DeliveryOperationError::Storage, 1296 })?; 1297 Ok(true) 1298 } 1299 1300 fn exact_document( 1301 record: &MycDiscoveryDocumentRecord, 1302 request: &MycDiscoveryCommitRequest, 1303 ) -> bool { 1304 record.generation_id == request.generation_id 1305 && record.event_id == request.event_id 1306 && record.event_digest == request.event_digest 1307 && record.event_bytes.as_ref() == request.event_bytes.as_ref() 1308 && record.projection_digest == request.projection_digest 1309 && record.projection_bytes.as_ref() == request.projection_bytes.as_ref() 1310 } 1311 1312 fn validate_event( 1313 event: &RadrootsNostrEvent, 1314 policy: &MycDiscoveryPolicies, 1315 ) -> Result<(), MycDiscoveryStateError> { 1316 if event.pubkey.to_hex() != policy.author_public_key.as_ref() { 1317 return Err(MycDiscoveryStateError::new( 1318 MycDiscoveryStateErrorKind::IdentityMismatch, 1319 )); 1320 } 1321 if event.kind != RadrootsNostrKind::Custom(NIP89_HANDLER_KIND) { 1322 return Err(MycDiscoveryStateError::new( 1323 MycDiscoveryStateErrorKind::InvalidEvent, 1324 )); 1325 } 1326 let mut expected_tags = vec![ 1327 vec!["d".to_owned(), policy.handler_identifier.to_string()], 1328 vec!["k".to_owned(), NIP46_RPC_KIND.to_string()], 1329 ]; 1330 expected_tags.extend( 1331 policy 1332 .public_relays 1333 .iter() 1334 .map(|relay| vec!["relay".to_owned(), relay.to_string()]), 1335 ); 1336 if let Some(url) = policy.nostrconnect_url.as_ref() { 1337 expected_tags.push(vec!["nostrconnect_url".to_owned(), url.to_string()]); 1338 } 1339 let actual_tags = event 1340 .tags 1341 .iter() 1342 .map(|tag| tag.as_slice().to_vec()) 1343 .collect::<Vec<_>>(); 1344 if actual_tags != expected_tags || event.content != policy.metadata_json.as_ref() { 1345 return Err(MycDiscoveryStateError::new( 1346 MycDiscoveryStateErrorKind::InvalidEvent, 1347 )); 1348 } 1349 Ok(()) 1350 } 1351 1352 #[derive(Serialize)] 1353 struct Nip05Projection<'a> { 1354 schema: &'static str, 1355 schema_version: u32, 1356 domain: &'a str, 1357 name: &'static str, 1358 public_key: &'a str, 1359 relays: Vec<&'a str>, 1360 #[serde(skip_serializing_if = "Option::is_none")] 1361 nostrconnect_url: Option<&'a str>, 1362 } 1363 1364 #[derive(Deserialize, Serialize)] 1365 #[serde(deny_unknown_fields)] 1366 struct Nip05ProjectionInput { 1367 schema: Box<str>, 1368 schema_version: u32, 1369 domain: Box<str>, 1370 name: Box<str>, 1371 public_key: Box<str>, 1372 relays: Box<[Box<str>]>, 1373 #[serde(skip_serializing_if = "Option::is_none")] 1374 nostrconnect_url: Option<Box<str>>, 1375 } 1376 1377 #[derive(Serialize)] 1378 struct Nip05Names<'a> { 1379 #[serde(rename = "_")] 1380 root: &'a str, 1381 } 1382 1383 #[derive(Serialize)] 1384 struct Nip46Discovery<'a> { 1385 relays: &'a [Box<str>], 1386 #[serde(skip_serializing_if = "Option::is_none")] 1387 nostrconnect_url: Option<&'a str>, 1388 } 1389 1390 #[derive(Serialize)] 1391 struct Nip05Output<'a> { 1392 names: Nip05Names<'a>, 1393 nip46: Nip46Discovery<'a>, 1394 } 1395 1396 fn render_nip05( 1397 selection: MycNip05ExportSelection, 1398 document: &MycDiscoveryDocumentRecord, 1399 ) -> Result<MycNip05Document, DiscoveryOperationError> { 1400 let input: Nip05ProjectionInput = serde_json::from_slice(document.nip05_projection_bytes()) 1401 .map_err(|_| DiscoveryOperationError::Binding)?; 1402 let canonical = serde_json::to_vec(&input).map_err(|_| DiscoveryOperationError::Binding)?; 1403 let valid_public_key = input.public_key.len() == 64 1404 && input 1405 .public_key 1406 .as_bytes() 1407 .iter() 1408 .all(u8::is_ascii_hexdigit) 1409 && !input 1410 .public_key 1411 .as_bytes() 1412 .iter() 1413 .any(u8::is_ascii_uppercase); 1414 let unique_relays = input 1415 .relays 1416 .iter() 1417 .map(Box::as_ref) 1418 .collect::<BTreeSet<_>>(); 1419 let valid_relays = (1..=32).contains(&input.relays.len()) 1420 && input 1421 .relays 1422 .iter() 1423 .all(|relay| RadrootsNostrRelayUrl::parse(relay).is_ok()) 1424 && unique_relays.len() == input.relays.len(); 1425 let valid_url = input 1426 .nostrconnect_url 1427 .as_deref() 1428 .is_none_or(|url| nostr::Url::parse(url).is_ok()); 1429 if canonical.as_slice() != document.nip05_projection_bytes() 1430 || input.schema.as_ref() != "radroots.myc.nip05-projection-input.v1" 1431 || input.schema_version != 1 1432 || input.domain.is_empty() 1433 || input.domain.len() > 253 1434 || input.name.as_ref() != "_" 1435 || !valid_public_key 1436 || !valid_relays 1437 || !valid_url 1438 { 1439 return Err(DiscoveryOperationError::Binding); 1440 } 1441 let bytes = serde_json::to_vec(&Nip05Output { 1442 names: Nip05Names { 1443 root: &input.public_key, 1444 }, 1445 nip46: Nip46Discovery { 1446 relays: &input.relays, 1447 nostrconnect_url: input.nostrconnect_url.as_deref(), 1448 }, 1449 }) 1450 .map_err(|_| DiscoveryOperationError::Binding)?; 1451 if bytes.is_empty() || bytes.len() > MYC_NIP05_DOCUMENT_MAX_BYTES { 1452 return Err(DiscoveryOperationError::Binding); 1453 } 1454 Ok(MycNip05Document { 1455 selection, 1456 domain: input.domain, 1457 digest: MycNip05DocumentDigest(Sha256::digest(&bytes).into()), 1458 bytes: bytes.into_boxed_slice(), 1459 }) 1460 } 1461 1462 fn projection_bytes(policy: &MycDiscoveryPolicies) -> Result<Vec<u8>, MycDiscoveryStateError> { 1463 serde_json::to_vec(&Nip05Projection { 1464 schema: "radroots.myc.nip05-projection-input.v1", 1465 schema_version: 1, 1466 domain: &policy.domain, 1467 name: "_", 1468 public_key: &policy.author_public_key, 1469 relays: policy.public_relays.iter().map(AsRef::as_ref).collect(), 1470 nostrconnect_url: policy.nostrconnect_url.as_deref(), 1471 }) 1472 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection)) 1473 } 1474 1475 fn render_nostrconnect_url( 1476 template: &str, 1477 signer_public_key: &str, 1478 public_relays: &[Box<str>], 1479 ) -> Result<String, MycDiscoveryStateError> { 1480 let mut serializer = url::form_urlencoded::Serializer::new(String::new()); 1481 for relay in public_relays { 1482 serializer.append_pair("relay", relay); 1483 } 1484 let bunker_uri = format!("bunker://{signer_public_key}?{}", serializer.finish()); 1485 let bunker_uri = radroots_nostr_connect::uri::Uri::parse(&bunker_uri) 1486 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))? 1487 .to_string(); 1488 let encoded: String = url::form_urlencoded::byte_serialize(bunker_uri.as_bytes()).collect(); 1489 let rendered = template.replace("<nostrconnect>", &encoded); 1490 nostr::Url::parse(&rendered) 1491 .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))?; 1492 Ok(rendered) 1493 } 1494 1495 fn bounded_blob( 1496 row: &sqlx::sqlite::SqliteRow, 1497 column: &str, 1498 maximum: usize, 1499 ) -> Result<Box<[u8]>, DiscoveryOperationError> { 1500 let value = row 1501 .try_get::<Option<Vec<u8>>, _>(column) 1502 .map_err(|_| DiscoveryOperationError::Binding)? 1503 .ok_or(DiscoveryOperationError::Binding)?; 1504 if value.is_empty() || value.len() > maximum { 1505 return Err(DiscoveryOperationError::Binding); 1506 } 1507 Ok(value.into_boxed_slice()) 1508 } 1509 1510 fn blob32( 1511 row: &sqlx::sqlite::SqliteRow, 1512 column: &str, 1513 ) -> Result<[u8; 32], DiscoveryOperationError> { 1514 row.try_get::<Option<Vec<u8>>, _>(column) 1515 .map_err(|_| DiscoveryOperationError::Binding)? 1516 .ok_or(DiscoveryOperationError::Binding)? 1517 .try_into() 1518 .map_err(|_| DiscoveryOperationError::Binding) 1519 } 1520 1521 fn optional_blob32( 1522 row: &sqlx::sqlite::SqliteRow, 1523 column: &str, 1524 type_column: &str, 1525 ) -> Result<Option<[u8; 32]>, DiscoveryOperationError> { 1526 match row 1527 .try_get::<&str, _>(type_column) 1528 .map_err(|_| DiscoveryOperationError::Binding)? 1529 { 1530 "null" => Ok(None), 1531 "blob" => blob32(row, column).map(Some), 1532 _ => Err(DiscoveryOperationError::Binding), 1533 } 1534 } 1535 1536 fn time( 1537 row: &sqlx::sqlite::SqliteRow, 1538 column: &str, 1539 ) -> Result<MycDeliveryTimeUnixMs, DiscoveryOperationError> { 1540 let value = row 1541 .try_get::<i64, _>(column) 1542 .map_err(|_| DiscoveryOperationError::Binding)?; 1543 u64::try_from(value) 1544 .ok() 1545 .and_then(|value| MycDeliveryTimeUnixMs::new(value).ok()) 1546 .ok_or(DiscoveryOperationError::Binding) 1547 } 1548 1549 fn require_one(rows: u64) -> Result<(), DiscoveryOperationError> { 1550 (rows == 1) 1551 .then_some(()) 1552 .ok_or(DiscoveryOperationError::Storage) 1553 } 1554 1555 fn map_transaction_error( 1556 error: ServiceSqliteTransactionError<DiscoveryOperationError>, 1557 ) -> MycStateRepositoryError { 1558 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 1559 return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); 1560 } 1561 let kind = match error.operation_error() { 1562 Some(DiscoveryOperationError::Binding) => MycStateRepositoryErrorKind::Binding, 1563 Some(DiscoveryOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, 1564 }; 1565 MycStateRepositoryError::new(kind) 1566 }