presence_publication.rs (116530B)
1 //! Exact signed presence documents and durable target delivery. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use nostr::{EventBuilder, Kind, PublicKey as NostrPublicKey, Tag, Timestamp}; 7 use radroots_event::{ 8 GenericEventDraft, 9 envelope::{EventEnvelope, kind::KIND_APPLICATION_HANDLER}, 10 profile::AuthoredProfile, 11 wire::{EventWireLimits, Nip01EventWire}, 12 }; 13 use radroots_event_codec::authoring::AuthoredEventPlan; 14 use radroots_nostr::event::{ 15 ApplicationHandlerSpec, Verification, build_application_handler, verify, verify_id, 16 }; 17 use radroots_service_host::{EntropySource, UnixTimeSeconds}; 18 use radroots_service_sqlite::{ 19 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 20 }; 21 use radroots_transport::BoxFuture; 22 use sha2::{Digest as _, Sha256}; 23 use sqlx::Row as _; 24 use zeroize::Zeroizing; 25 26 use crate::{ 27 RHI_PRESENCE_DESIRED_MAX_TARGETS, RhiDecryptedIdentity, RhiJitterBoundMilliseconds, 28 RhiPresenceDesiredAuthority, RhiPresenceDesiredCommitOutcome, RhiPresenceDesiredMode, 29 RhiPresenceDesiredState, RhiPresenceDocumentKind, RhiPresenceOutboxRepository, 30 RhiPresenceTarget, RhiRuntimeAdapterErrorKind, RhiStateHostMode, RhiTimeEntropyAdapters, 31 presence_desired::presence_authority_matches_state, 32 }; 33 34 /// Exact version of the signed presence and delivery contract. 35 pub const RHI_PRESENCE_PUBLICATION_CONTRACT_VERSION: u32 = 1; 36 /// Maximum canonical bytes retained for one signed presence event. 37 pub const RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES: usize = 32 * 1024; 38 /// Maximum durable attempts for one presence target. 39 pub const RHI_PRESENCE_MAX_ATTEMPTS: u16 = 100; 40 41 const PROFILE_NAME: &str = "rhi"; 42 const PROFILE_DISPLAY_NAME: &str = "Radroots RHI"; 43 const PROFILE_ABOUT: &str = "Radroots evidence reconciliation and attestation service"; 44 const APPLICATION_HANDLER_IDENTIFIER: &str = "rhi"; 45 const APPLICATION_HANDLER_KINDS: [u32; 1] = [3441]; 46 const INITIAL_BACKOFF_MILLISECONDS: u64 = 250; 47 const MAXIMUM_BACKOFF_MILLISECONDS: u64 = 30_000; 48 const ATTEMPT_DEADLINE_MILLISECONDS: u64 = 15_000; 49 const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64; 50 const OUTBOX_ID_DOMAIN: &[u8] = b"radroots.rhi.presence_outbox.v1\0"; 51 const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.rhi.presence_attempt.v1\0"; 52 53 const READ_DESIRED_BINDING_SQL: &str = r#"SELECT singleton, generation, enabled, 54 profile, application_handler, 55 CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 56 THEN target_set_sha256 ELSE NULL END AS target_set_sha256, 57 target_count, required_target_count, queue_capacity, 58 CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 59 THEN desired_sha256 ELSE NULL END AS desired_sha256 60 FROM presence_desired_state 61 LIMIT 2"#; 62 63 const READ_GENERATION_OUTBOX_SQL: &str = r#"SELECT 64 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 65 THEN outbox_id ELSE NULL END AS outbox_id, 66 desired_generation, 67 length(CAST(document_kind AS BLOB)) AS document_kind_bytes, 68 substr(document_kind, 1, 20) AS document_kind, 69 CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 70 THEN desired_sha256 ELSE NULL END AS desired_sha256, 71 CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 72 THEN target_set_sha256 ELSE NULL END AS target_set_sha256, 73 CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32 74 THEN event_id ELSE NULL END AS event_id, 75 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 76 THEN event_sha256 ELSE NULL END AS event_sha256, 77 length(exact_signed_event_bytes) AS exact_signed_event_bytes_length, 78 substr(exact_signed_event_bytes, 1, 32769) AS exact_signed_event_bytes, 79 authored_at_unix_s, 80 length(CAST(service_public_key AS BLOB)) AS service_public_key_bytes, 81 substr(service_public_key, 1, 65) AS service_public_key, 82 target_count, required_target_count, max_attempts, 83 initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 84 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 11) AS state, 85 revision, next_attempt_unix_ms, 86 CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, 87 CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, 88 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 89 FROM presence_outbox 90 WHERE desired_generation = ? 91 ORDER BY CASE document_kind 92 WHEN 'service_profile' THEN 0 93 WHEN 'application_handler' THEN 1 94 ELSE 2 END 95 LIMIT 3"#; 96 97 const READ_OUTBOX_SQL: &str = r#"SELECT 98 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 99 THEN outbox_id ELSE NULL END AS outbox_id, 100 desired_generation, 101 length(CAST(document_kind AS BLOB)) AS document_kind_bytes, 102 substr(document_kind, 1, 20) AS document_kind, 103 CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 104 THEN desired_sha256 ELSE NULL END AS desired_sha256, 105 CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 106 THEN target_set_sha256 ELSE NULL END AS target_set_sha256, 107 CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32 108 THEN event_id ELSE NULL END AS event_id, 109 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 110 THEN event_sha256 ELSE NULL END AS event_sha256, 111 length(exact_signed_event_bytes) AS exact_signed_event_bytes_length, 112 substr(exact_signed_event_bytes, 1, 32769) AS exact_signed_event_bytes, 113 authored_at_unix_s, 114 length(CAST(service_public_key AS BLOB)) AS service_public_key_bytes, 115 substr(service_public_key, 1, 65) AS service_public_key, 116 target_count, required_target_count, max_attempts, 117 initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 118 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 11) AS state, 119 revision, next_attempt_unix_ms, 120 CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, 121 CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, 122 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 123 FROM presence_outbox 124 WHERE outbox_id = ? 125 LIMIT 2"#; 126 127 const READ_CLAIMABLE_OUTBOX_SQL: &str = r#"SELECT 128 CASE WHEN typeof(outbox.outbox_id) = 'blob' AND length(outbox.outbox_id) = 32 129 THEN outbox.outbox_id ELSE NULL END AS outbox_id 130 FROM presence_outbox AS outbox 131 JOIN presence_desired_state AS desired 132 ON desired.singleton = 1 133 AND desired.generation = outbox.desired_generation 134 AND desired.desired_sha256 = outbox.desired_sha256 135 WHERE outbox.state = 'pending' AND outbox.next_attempt_unix_ms <= ? 136 ORDER BY outbox.next_attempt_unix_ms, outbox.created_at_unix_ms, 137 CASE outbox.document_kind 138 WHEN 'service_profile' THEN 0 139 WHEN 'application_handler' THEN 1 140 ELSE 2 END, 141 outbox.outbox_id 142 LIMIT 1"#; 143 144 const READ_EXPIRED_OUTBOX_SQL: &str = r#"SELECT 145 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 146 THEN outbox_id ELSE NULL END AS outbox_id 147 FROM presence_outbox 148 WHERE state = 'leased' AND lease_expires_unix_ms <= ? 149 ORDER BY lease_expires_unix_ms, created_at_unix_ms, document_kind, outbox_id 150 LIMIT 1"#; 151 152 const READ_TARGETS_SQL: &str = r#"SELECT target_ordinal, 153 length(CAST(relay_id AS BLOB)) AS relay_id_bytes, 154 substr(relay_id, 1, 65) AS relay_id, required, 155 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 14) AS state, 156 revision, attempt_count, next_attempt_unix_ms, 157 CASE WHEN last_attempt_id IS NULL THEN NULL ELSE length(last_attempt_id) END 158 AS last_attempt_id_bytes, 159 CASE WHEN last_attempt_id IS NULL THEN NULL ELSE substr(last_attempt_id, 1, 33) END 160 AS last_attempt_id, 161 updated_at_unix_ms 162 FROM presence_targets 163 WHERE outbox_id = ? 164 ORDER BY target_ordinal 165 LIMIT 33"#; 166 167 const INSERT_OUTBOX_SQL: &str = r#"INSERT INTO presence_outbox ( 168 outbox_id, desired_generation, document_kind, desired_sha256, 169 target_set_sha256, event_id, event_sha256, exact_signed_event_bytes, 170 authored_at_unix_s, service_public_key, target_count, required_target_count, 171 max_attempts, initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 172 state, revision, next_attempt_unix_ms, lease_owner, lease_expires_unix_ms, 173 created_at_unix_ms, updated_at_unix_ms 174 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 100, 250, 30000, 15000, 175 'pending', 1, ?, NULL, NULL, ?, ?)"#; 176 177 const INSERT_TARGET_SQL: &str = r#"INSERT INTO presence_targets ( 178 outbox_id, target_ordinal, relay_id, required, state, revision, 179 attempt_count, next_attempt_unix_ms, last_attempt_id, updated_at_unix_ms 180 ) VALUES (?, ?, ?, ?, 'pending', 1, 0, ?, NULL, ?)"#; 181 182 const SUPERSEDE_OUTBOX_SQL: &str = r#"UPDATE presence_outbox 183 SET state = 'superseded', revision = revision + 1, 184 next_attempt_unix_ms = NULL, lease_owner = NULL, 185 lease_expires_unix_ms = NULL, updated_at_unix_ms = ? 186 WHERE desired_generation != ? AND state IN ('pending', 'blocked') 187 AND updated_at_unix_ms <= ?"#; 188 189 const COUNT_STALE_LEASES_SQL: &str = r#"SELECT COUNT(*) AS lease_count 190 FROM presence_outbox 191 WHERE desired_generation != ? AND state = 'leased'"#; 192 193 const READ_MAX_AUTHORED_SQL: &str = r#"SELECT MAX(authored_at_unix_s) AS authored_at_unix_s 194 FROM presence_outbox WHERE document_kind = ?"#; 195 196 const CLAIM_OUTBOX_SQL: &str = r#"UPDATE presence_outbox 197 SET state = 'leased', revision = revision + 1, 198 next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?, 199 updated_at_unix_ms = ? 200 WHERE outbox_id = ? AND revision = ? AND state = 'pending' 201 AND next_attempt_unix_ms <= ? AND updated_at_unix_ms <= ?"#; 202 203 const PREPARE_TARGET_SQL: &str = r#"UPDATE presence_targets 204 SET state = 'submitted', revision = revision + 1, 205 attempt_count = attempt_count + 1, next_attempt_unix_ms = NULL, 206 last_attempt_id = ?, updated_at_unix_ms = ? 207 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ? 208 AND state IN ('pending', 'failed', 'rate_limited', 'unknown') 209 AND attempt_count < ? AND next_attempt_unix_ms <= ? 210 AND updated_at_unix_ms <= ?"#; 211 212 const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO presence_attempts ( 213 attempt_id, outbox_id, target_ordinal, attempt_number, event_sha256, 214 lease_owner, started_at_unix_ms, finished_at_unix_ms, outcome, result_code 215 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 216 217 const READ_ATTEMPT_SQL: &str = r#"SELECT 218 CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 219 THEN attempt_id ELSE NULL END AS attempt_id, 220 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 221 THEN outbox_id ELSE NULL END AS outbox_id, 222 target_ordinal, attempt_number, 223 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 224 THEN event_sha256 ELSE NULL END AS event_sha256, 225 CASE WHEN typeof(lease_owner) = 'blob' AND length(lease_owner) = 16 226 THEN lease_owner ELSE NULL END AS lease_owner, 227 started_at_unix_ms, finished_at_unix_ms, 228 length(CAST(outcome AS BLOB)) AS outcome_bytes, substr(outcome, 1, 14) AS outcome, 229 length(CAST(result_code AS BLOB)) AS result_code_bytes, 230 substr(result_code, 1, 65) AS result_code 231 FROM presence_attempts 232 WHERE attempt_id = ? 233 LIMIT 2"#; 234 235 const UPDATE_TARGET_OUTCOME_SQL: &str = r#"UPDATE presence_targets 236 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, 237 updated_at_unix_ms = ? 238 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ? 239 AND state = 'submitted' AND attempt_count = ? AND last_attempt_id = ?"#; 240 241 const UPDATE_OUTBOX_AFTER_ATTEMPT_SQL: &str = r#"UPDATE presence_outbox 242 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, 243 lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? 244 WHERE outbox_id = ? AND revision = ? AND state = 'leased' 245 AND lease_owner = ? AND lease_expires_unix_ms = ? 246 AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#; 247 248 const UPDATE_OUTBOX_RECOVERY_SQL: &str = r#"UPDATE presence_outbox 249 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, 250 lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? 251 WHERE outbox_id = ? AND revision = ? AND state = 'leased' 252 AND lease_owner = ? AND lease_expires_unix_ms = ? 253 AND lease_expires_unix_ms <= ? AND updated_at_unix_ms <= ?"#; 254 255 /// Stable source-free signed-presence and delivery failure class. 256 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 257 pub enum RhiPresencePublicationErrorKind { 258 InvalidMode, 259 InvalidInput, 260 DesiredStateMismatch, 261 IdentityMismatch, 262 EntropyUnavailable, 263 RenderingFailed, 264 VerificationFailed, 265 NotReady, 266 LeaseLost, 267 Invariant, 268 ClockUnavailable, 269 Storage, 270 CommitOutcomeUnknown, 271 } 272 273 impl RhiPresencePublicationErrorKind { 274 /// Returns the stable machine-readable failure code. 275 #[must_use] 276 pub const fn code(self) -> &'static str { 277 match self { 278 Self::InvalidMode => "presence_publication_mode_invalid", 279 Self::InvalidInput => "presence_publication_input_invalid", 280 Self::DesiredStateMismatch => "presence_publication_desired_state_mismatch", 281 Self::IdentityMismatch => "presence_publication_identity_mismatch", 282 Self::EntropyUnavailable => "presence_publication_entropy_unavailable", 283 Self::RenderingFailed => "presence_publication_rendering_failed", 284 Self::VerificationFailed => "presence_publication_verification_failed", 285 Self::NotReady => "presence_publication_not_ready", 286 Self::LeaseLost => "presence_publication_lease_lost", 287 Self::Invariant => "presence_publication_invariant_failed", 288 Self::ClockUnavailable => "presence_publication_clock_unavailable", 289 Self::Storage => "presence_publication_storage_failed", 290 Self::CommitOutcomeUnknown => "presence_publication_commit_outcome_unknown", 291 } 292 } 293 } 294 295 /// Redacted source-free signed-presence and delivery failure. 296 #[derive(Clone, Copy, PartialEq, Eq)] 297 pub struct RhiPresencePublicationError { 298 kind: RhiPresencePublicationErrorKind, 299 } 300 301 impl RhiPresencePublicationError { 302 /// Returns the stable failure class. 303 #[must_use] 304 pub const fn kind(self) -> RhiPresencePublicationErrorKind { 305 self.kind 306 } 307 308 /// Returns the stable machine-readable failure code. 309 #[must_use] 310 pub const fn code(self) -> &'static str { 311 self.kind.code() 312 } 313 } 314 315 impl fmt::Display for RhiPresencePublicationError { 316 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 317 formatter.write_str(match self.kind { 318 RhiPresencePublicationErrorKind::InvalidMode => { 319 "RHI presence publication requires writable state" 320 } 321 RhiPresencePublicationErrorKind::InvalidInput => { 322 "RHI presence publication input is invalid" 323 } 324 RhiPresencePublicationErrorKind::DesiredStateMismatch => { 325 "RHI presence publication desired state does not match" 326 } 327 RhiPresencePublicationErrorKind::IdentityMismatch => { 328 "RHI presence publication identity does not match" 329 } 330 RhiPresencePublicationErrorKind::EntropyUnavailable => { 331 "RHI presence publication entropy is unavailable" 332 } 333 RhiPresencePublicationErrorKind::RenderingFailed => { 334 "RHI presence publication rendering failed" 335 } 336 RhiPresencePublicationErrorKind::VerificationFailed => { 337 "RHI presence publication verification failed" 338 } 339 RhiPresencePublicationErrorKind::NotReady => "RHI presence work is not ready", 340 RhiPresencePublicationErrorKind::LeaseLost => { 341 "RHI presence publication lease is no longer authoritative" 342 } 343 RhiPresencePublicationErrorKind::Invariant => { 344 "RHI presence publication state invariant failed" 345 } 346 RhiPresencePublicationErrorKind::ClockUnavailable => { 347 "RHI presence publication clock is unavailable" 348 } 349 RhiPresencePublicationErrorKind::Storage => { 350 "RHI presence publication transaction failed" 351 } 352 RhiPresencePublicationErrorKind::CommitOutcomeUnknown => { 353 "RHI presence publication commit outcome is unknown" 354 } 355 }) 356 } 357 } 358 359 impl fmt::Debug for RhiPresencePublicationError { 360 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 361 formatter 362 .debug_struct("RhiPresencePublicationError") 363 .field("kind", &self.kind) 364 .finish() 365 } 366 } 367 368 impl Error for RhiPresencePublicationError {} 369 370 /// One independently verified exact signed presence document. 371 pub struct RhiSignedPresenceDocument { 372 kind: RhiPresenceDocumentKind, 373 event_id: [u8; 32], 374 signed_event_sha256: [u8; 32], 375 signed_event_bytes: Box<[u8]>, 376 created_at_unix_seconds: u64, 377 } 378 379 impl RhiSignedPresenceDocument { 380 /// Returns the closed document kind. 381 #[must_use] 382 pub const fn kind(&self) -> RhiPresenceDocumentKind { 383 self.kind 384 } 385 386 /// Returns the independently verified NIP-01 event identifier. 387 #[must_use] 388 pub const fn event_id(&self) -> &[u8; 32] { 389 &self.event_id 390 } 391 392 /// Returns the digest of the exact retained bytes. 393 #[must_use] 394 pub const fn signed_event_sha256(&self) -> &[u8; 32] { 395 &self.signed_event_sha256 396 } 397 398 /// Returns the exact canonical signed bytes. 399 #[must_use] 400 pub fn signed_event_bytes(&self) -> &[u8] { 401 &self.signed_event_bytes 402 } 403 404 /// Returns the injected authored timestamp. 405 #[must_use] 406 pub const fn created_at_unix_seconds(&self) -> u64 { 407 self.created_at_unix_seconds 408 } 409 } 410 411 impl fmt::Debug for RhiSignedPresenceDocument { 412 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 413 formatter 414 .debug_struct("RhiSignedPresenceDocument") 415 .field("kind", &self.kind) 416 .field("signed_event_bytes", &self.signed_event_bytes.len()) 417 .finish_non_exhaustive() 418 } 419 } 420 421 /// Sealed exact document set bound to one committed desired-state generation. 422 pub struct RhiSignedPresenceDocuments { 423 desired_state: RhiPresenceDesiredState, 424 service_public_key: Box<str>, 425 targets: Box<[RhiPresenceTarget]>, 426 documents: Box<[RhiSignedPresenceDocument]>, 427 } 428 429 impl RhiSignedPresenceDocuments { 430 /// Returns the exact committed desired state. 431 #[must_use] 432 pub const fn desired_state(&self) -> RhiPresenceDesiredState { 433 self.desired_state 434 } 435 436 /// Returns the exact ordered signed-document inventory. 437 #[must_use] 438 pub fn documents(&self) -> &[RhiSignedPresenceDocument] { 439 &self.documents 440 } 441 442 /// Returns the stable target count without exposing endpoints. 443 #[must_use] 444 pub fn target_count(&self) -> usize { 445 self.targets.len() 446 } 447 448 fn retained_copy(&self) -> Self { 449 Self { 450 desired_state: self.desired_state, 451 service_public_key: self.service_public_key.clone(), 452 targets: self.targets.clone(), 453 documents: self 454 .documents 455 .iter() 456 .map(|document| RhiSignedPresenceDocument { 457 kind: document.kind, 458 event_id: document.event_id, 459 signed_event_sha256: document.signed_event_sha256, 460 signed_event_bytes: document.signed_event_bytes.clone(), 461 created_at_unix_seconds: document.created_at_unix_seconds, 462 }) 463 .collect(), 464 } 465 } 466 } 467 468 impl fmt::Debug for RhiSignedPresenceDocuments { 469 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 470 formatter 471 .debug_struct("RhiSignedPresenceDocuments") 472 .field("desired_generation", &self.desired_state.generation()) 473 .field("document_count", &self.documents.len()) 474 .field("target_count", &self.targets.len()) 475 .finish_non_exhaustive() 476 } 477 } 478 479 /// Opaque identity of one exact durable presence outbox. 480 #[derive(Clone, Copy, PartialEq, Eq, Hash)] 481 pub struct RhiPresenceOutboxId([u8; 32]); 482 483 impl RhiPresenceOutboxId { 484 /// Returns the exact identity bytes. 485 #[must_use] 486 pub const fn as_bytes(&self) -> &[u8; 32] { 487 &self.0 488 } 489 } 490 491 impl fmt::Debug for RhiPresenceOutboxId { 492 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 493 formatter.write_str("RhiPresenceOutboxId([redacted])") 494 } 495 } 496 497 /// Opaque identity of one durable presence attempt. 498 #[derive(Clone, Copy, PartialEq, Eq, Hash)] 499 pub struct RhiPresenceAttemptId([u8; 32]); 500 501 impl RhiPresenceAttemptId { 502 /// Returns the exact identity bytes. 503 #[must_use] 504 pub const fn as_bytes(&self) -> &[u8; 32] { 505 &self.0 506 } 507 } 508 509 impl fmt::Debug for RhiPresenceAttemptId { 510 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 511 formatter.write_str("RhiPresenceAttemptId([redacted])") 512 } 513 } 514 515 /// Bounded wall-clock instant in whole Unix milliseconds. 516 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 517 pub struct RhiPresenceUnixMilliseconds(u64); 518 519 impl RhiPresenceUnixMilliseconds { 520 /// Validates one injected instant against SQLite's signed representation. 521 pub fn new(value: u64) -> Result<Self, RhiPresencePublicationError> { 522 if value > MAX_UNIX_MILLISECONDS { 523 return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); 524 } 525 Ok(Self(value)) 526 } 527 528 /// Returns the exact instant. 529 #[must_use] 530 pub const fn get(self) -> u64 { 531 self.0 532 } 533 } 534 535 /// Bounded caller-injected retry delay in whole milliseconds. 536 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 537 pub struct RhiPresenceRetryDelayMilliseconds(u64); 538 539 impl RhiPresenceRetryDelayMilliseconds { 540 /// Validates one delay against the absolute retry ceiling. 541 pub fn new(value: u64) -> Result<Self, RhiPresencePublicationError> { 542 if value > MAXIMUM_BACKOFF_MILLISECONDS { 543 return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); 544 } 545 Ok(Self(value)) 546 } 547 548 /// Returns the exact delay. 549 #[must_use] 550 pub const fn get(self) -> u64 { 551 self.0 552 } 553 } 554 555 /// Stable process-local owner token for one compare-and-swap lease. 556 #[derive(Clone, Copy, PartialEq, Eq, Hash)] 557 pub struct RhiPresenceLeaseOwner([u8; 16]); 558 559 impl RhiPresenceLeaseOwner { 560 /// Validates one injected nonzero owner token. 561 pub fn from_bytes(bytes: [u8; 16]) -> Result<Self, RhiPresencePublicationError> { 562 if bytes.iter().all(|byte| *byte == 0) { 563 return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); 564 } 565 Ok(Self(bytes)) 566 } 567 } 568 569 impl fmt::Debug for RhiPresenceLeaseOwner { 570 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 571 formatter.write_str("RhiPresenceLeaseOwner([redacted])") 572 } 573 } 574 575 /// Stable durable presence-outbox lifecycle state. 576 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 577 pub enum RhiPresenceOutboxState { 578 Pending, 579 Leased, 580 Complete, 581 Blocked, 582 Superseded, 583 } 584 585 impl RhiPresenceOutboxState { 586 /// Returns the exact machine-contract spelling. 587 #[must_use] 588 pub const fn code(self) -> &'static str { 589 match self { 590 Self::Pending => "pending", 591 Self::Leased => "leased", 592 Self::Complete => "complete", 593 Self::Blocked => "blocked", 594 Self::Superseded => "superseded", 595 } 596 } 597 } 598 599 /// Stable durable per-target presence state. 600 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 601 pub enum RhiPresenceTargetState { 602 Pending, 603 Submitted, 604 Accepted, 605 Rejected, 606 RateLimited, 607 AuthRequired, 608 Failed, 609 Unknown, 610 } 611 612 impl RhiPresenceTargetState { 613 /// Returns the exact machine-contract spelling. 614 #[must_use] 615 pub const fn code(self) -> &'static str { 616 match self { 617 Self::Pending => "pending", 618 Self::Submitted => "submitted", 619 Self::Accepted => "accepted", 620 Self::Rejected => "rejected", 621 Self::RateLimited => "rate_limited", 622 Self::AuthRequired => "auth_required", 623 Self::Failed => "failed", 624 Self::Unknown => "unknown", 625 } 626 } 627 } 628 629 /// Closed result observed from one exact-byte relay submission. 630 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 631 pub enum RhiPresenceAttemptOutcome { 632 Submitted, 633 Accepted, 634 Rejected, 635 RateLimited, 636 AuthRequired, 637 Failed, 638 Unknown, 639 } 640 641 impl RhiPresenceAttemptOutcome { 642 /// Returns the exact machine-contract spelling. 643 #[must_use] 644 pub const fn code(self) -> &'static str { 645 match self { 646 Self::Submitted => "submitted", 647 Self::Accepted => "accepted", 648 Self::Rejected => "rejected", 649 Self::RateLimited => "rate_limited", 650 Self::AuthRequired => "auth_required", 651 Self::Failed => "failed", 652 Self::Unknown => "unknown", 653 } 654 } 655 } 656 657 /// Confirmed result of committing one complete signed presence generation. 658 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 659 pub struct RhiPresenceCommitOutcome { 660 desired_state: RhiPresenceDesiredState, 661 document_count: u8, 662 changed: bool, 663 } 664 665 impl RhiPresenceCommitOutcome { 666 /// Returns the exact committed desired-state identity. 667 #[must_use] 668 pub const fn desired_state(self) -> RhiPresenceDesiredState { 669 self.desired_state 670 } 671 672 /// Returns the exact committed document count. 673 #[must_use] 674 pub const fn document_count(self) -> u8 { 675 self.document_count 676 } 677 678 /// Reports whether any durable presence workflow state changed. 679 #[must_use] 680 pub const fn changed(self) -> bool { 681 self.changed 682 } 683 } 684 685 /// Non-forgeable compare-and-swap authority for one claimed presence outbox. 686 #[derive(PartialEq, Eq)] 687 pub struct RhiPresenceLease { 688 outbox: PresenceOutboxRecord, 689 owner: RhiPresenceLeaseOwner, 690 expires_at: RhiPresenceUnixMilliseconds, 691 } 692 693 impl RhiPresenceLease { 694 /// Returns the exact claimed outbox identity. 695 #[must_use] 696 pub const fn outbox_id(&self) -> RhiPresenceOutboxId { 697 self.outbox.id 698 } 699 700 /// Returns the exact lease expiry. 701 #[must_use] 702 pub const fn expires_at(&self) -> RhiPresenceUnixMilliseconds { 703 self.expires_at 704 } 705 } 706 707 impl fmt::Debug for RhiPresenceLease { 708 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 709 formatter 710 .debug_struct("RhiPresenceLease") 711 .field("identity", &"[redacted]") 712 .field("revision", &self.outbox.revision) 713 .field("expires_at", &self.expires_at) 714 .finish() 715 } 716 } 717 718 /// Exact signed bytes exposed only after durable Submitted state exists. 719 #[must_use = "prepared presence must be executed exactly or left for unknown recovery"] 720 pub struct RhiPreparedPresenceAttempt { 721 lease: RhiPresenceLease, 722 target: PresenceTargetRecord, 723 attempt_id: RhiPresenceAttemptId, 724 started_at: RhiPresenceUnixMilliseconds, 725 deadline_at: RhiPresenceUnixMilliseconds, 726 exact_signed_event_bytes: Box<[u8]>, 727 } 728 729 impl RhiPreparedPresenceAttempt { 730 /// Returns the exact attempt identity. 731 #[must_use] 732 pub const fn attempt_id(&self) -> RhiPresenceAttemptId { 733 self.attempt_id 734 } 735 736 /// Returns the stable relay identity without an endpoint or secret. 737 #[must_use] 738 pub fn relay_id(&self) -> &str { 739 &self.target.relay_id 740 } 741 742 /// Returns the original committed bytes with no transformation. 743 #[must_use] 744 pub const fn exact_signed_event_bytes(&self) -> &[u8] { 745 &self.exact_signed_event_bytes 746 } 747 748 /// Returns the absolute attempt deadline. 749 #[must_use] 750 pub const fn deadline_at(&self) -> RhiPresenceUnixMilliseconds { 751 self.deadline_at 752 } 753 754 /// Returns the one-based durable attempt number. 755 #[must_use] 756 pub const fn attempt_number(&self) -> u16 { 757 self.target.attempt_count 758 } 759 760 fn retry_upper_bound(&self) -> u64 { 761 retry_upper_bound(self.target.attempt_count) 762 } 763 } 764 765 impl fmt::Debug for RhiPreparedPresenceAttempt { 766 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 767 formatter 768 .debug_struct("RhiPreparedPresenceAttempt") 769 .field("identity", &"[redacted]") 770 .field("target_ordinal", &self.target.ordinal) 771 .field("attempt_number", &self.target.attempt_count) 772 .field("deadline_at", &self.deadline_at) 773 .finish() 774 } 775 } 776 777 /// Closed exact-byte transport boundary for durable presence delivery. 778 pub trait RhiExactPresenceSink: Send + Sync { 779 /// Submits the exact retained payload and returns one closed observation. 780 fn submit_exact<'a>( 781 &'a self, 782 attempt: &'a RhiPreparedPresenceAttempt, 783 ) -> BoxFuture<'a, RhiPresenceAttemptOutcome>; 784 } 785 786 /// Confirmed durable result for one exact presence attempt. 787 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 788 pub struct RhiPresenceAttemptCommit { 789 outbox_id: RhiPresenceOutboxId, 790 attempt_id: RhiPresenceAttemptId, 791 target_ordinal: u8, 792 attempt_number: u16, 793 outcome: RhiPresenceAttemptOutcome, 794 target_state: RhiPresenceTargetState, 795 outbox_state: RhiPresenceOutboxState, 796 } 797 798 impl RhiPresenceAttemptCommit { 799 #[must_use] 800 pub const fn outbox_id(self) -> RhiPresenceOutboxId { 801 self.outbox_id 802 } 803 804 #[must_use] 805 pub const fn attempt_id(self) -> RhiPresenceAttemptId { 806 self.attempt_id 807 } 808 809 #[must_use] 810 pub const fn target_ordinal(self) -> u8 { 811 self.target_ordinal 812 } 813 814 #[must_use] 815 pub const fn attempt_number(self) -> u16 { 816 self.attempt_number 817 } 818 819 #[must_use] 820 pub const fn outcome(self) -> RhiPresenceAttemptOutcome { 821 self.outcome 822 } 823 824 #[must_use] 825 pub const fn target_state(self) -> RhiPresenceTargetState { 826 self.target_state 827 } 828 829 #[must_use] 830 pub const fn outbox_state(self) -> RhiPresenceOutboxState { 831 self.outbox_state 832 } 833 } 834 835 /// Builds, signs, and independently revalidates one complete desired document set. 836 /// 837 /// Exactly 32 entropy bytes are consumed per enabled document. No clock, 838 /// persistence, relay, task, or network authority is acquired implicitly. 839 pub fn build_rhi_signed_presence_documents( 840 desired: RhiPresenceDesiredCommitOutcome, 841 authority: &RhiPresenceDesiredAuthority, 842 identity: &RhiDecryptedIdentity, 843 created_at: UnixTimeSeconds, 844 entropy: &dyn EntropySource, 845 ) -> Result<RhiSignedPresenceDocuments, RhiPresencePublicationError> { 846 let state = desired.state(); 847 if created_at.get() > i64::MAX as u64 { 848 return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); 849 } 850 if !presence_authority_matches_state(state, authority) { 851 return Err(failure( 852 RhiPresencePublicationErrorKind::DesiredStateMismatch, 853 )); 854 } 855 if identity.public_identity().as_hex() != authority.service_public_key() { 856 return Err(failure(RhiPresencePublicationErrorKind::IdentityMismatch)); 857 } 858 let mut documents = Vec::with_capacity(authority.document_kinds().len()); 859 for kind in authority.document_kinds() { 860 let plan = presence_plan(*kind, identity.public_identity().as_hex(), created_at.get())?; 861 let mut auxiliary = Zeroizing::new([0_u8; 32]); 862 entropy 863 .fill_bytes(&mut auxiliary[..]) 864 .map_err(|_| failure(RhiPresencePublicationErrorKind::EntropyUnavailable))?; 865 let event = sign_plan(identity, &plan, &auxiliary)?; 866 let bytes = serde_json::to_vec(&event) 867 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 868 let verified = validate_signed_document( 869 *kind, 870 identity.public_identity().as_hex(), 871 created_at.get(), 872 &bytes, 873 )?; 874 documents.push(RhiSignedPresenceDocument { 875 kind: *kind, 876 event_id: *verified.id().as_bytes(), 877 signed_event_sha256: Sha256::digest(&bytes).into(), 878 signed_event_bytes: bytes.into_boxed_slice(), 879 created_at_unix_seconds: created_at.get(), 880 }); 881 } 882 if (state.mode() == RhiPresenceDesiredMode::Disabled && !documents.is_empty()) 883 || (state.mode() == RhiPresenceDesiredMode::Enabled 884 && documents.len() != authority.document_kinds().len()) 885 { 886 return Err(failure(RhiPresencePublicationErrorKind::Invariant)); 887 } 888 Ok(RhiSignedPresenceDocuments { 889 desired_state: state, 890 service_public_key: authority.service_public_key().into(), 891 targets: authority.targets().to_vec().into_boxed_slice(), 892 documents: documents.into_boxed_slice(), 893 }) 894 } 895 896 /// Independently verifies every exact signed document against its bound intent. 897 pub fn validate_rhi_signed_presence_documents( 898 documents: &RhiSignedPresenceDocuments, 899 authority: &RhiPresenceDesiredAuthority, 900 ) -> Result<(), RhiPresencePublicationError> { 901 if !presence_authority_matches_state(documents.desired_state, authority) 902 || documents.service_public_key.as_ref() != authority.service_public_key() 903 || documents.targets.as_ref() != authority.targets() 904 || documents.documents.len() != authority.document_kinds().len() 905 { 906 return Err(failure( 907 RhiPresencePublicationErrorKind::DesiredStateMismatch, 908 )); 909 } 910 for (document, kind) in documents.documents.iter().zip(authority.document_kinds()) { 911 if document.kind != *kind 912 || Sha256::digest(document.signed_event_bytes.as_ref()).as_slice() 913 != document.signed_event_sha256 914 { 915 return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); 916 } 917 let verified = validate_signed_document( 918 document.kind, 919 authority.service_public_key(), 920 document.created_at_unix_seconds, 921 &document.signed_event_bytes, 922 )?; 923 if verified.id().as_bytes() != &document.event_id { 924 return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); 925 } 926 } 927 Ok(()) 928 } 929 930 impl RhiPresenceOutboxRepository<'_> { 931 /// Atomically commits exact signed bytes and the complete immutable target set. 932 /// 933 /// A successful return means every exact payload and initial target state is 934 /// durable before any relay adapter can observe those bytes. The caller 935 /// retains the sealed exact-byte capability so an unknown commit result can 936 /// be reconciled by replaying this same value. Exact replay is idempotent. 937 pub async fn commit_signed_presence( 938 &self, 939 documents: &RhiSignedPresenceDocuments, 940 committed_at: RhiPresenceUnixMilliseconds, 941 ) -> Result<RhiPresenceCommitOutcome, RhiPresencePublicationError> { 942 require_writable(self)?; 943 let documents = documents.retained_copy(); 944 self.host() 945 .sqlite_host() 946 .transaction(move |transaction| { 947 Box::pin(async move { commit_signed(transaction, documents, committed_at).await }) 948 }) 949 .await 950 .map_err(map_transaction_error) 951 } 952 953 /// Claims the oldest currently desired, due presence outbox. 954 pub async fn claim_next_presence( 955 &self, 956 owner: RhiPresenceLeaseOwner, 957 now: RhiPresenceUnixMilliseconds, 958 ) -> Result<Option<RhiPresenceLease>, RhiPresencePublicationError> { 959 require_writable(self)?; 960 self.host() 961 .sqlite_host() 962 .transaction(move |transaction| { 963 Box::pin(async move { claim_next(transaction, owner, now).await }) 964 }) 965 .await 966 .map_err(map_transaction_error) 967 } 968 969 /// Persists Submitted before exposing the committed exact bytes. 970 pub async fn prepare_next_presence_target( 971 &self, 972 lease: RhiPresenceLease, 973 started_at: RhiPresenceUnixMilliseconds, 974 ) -> Result<RhiPreparedPresenceAttempt, RhiPresencePublicationError> { 975 require_writable(self)?; 976 self.host() 977 .sqlite_host() 978 .transaction(move |transaction| { 979 Box::pin(async move { prepare_next(transaction, lease, started_at).await }) 980 }) 981 .await 982 .map_err(map_transaction_error) 983 } 984 985 /// Commits one closed observed outcome by exact compare-and-swap. 986 pub async fn record_presence_outcome( 987 &self, 988 prepared: &RhiPreparedPresenceAttempt, 989 finished_at: RhiPresenceUnixMilliseconds, 990 outcome: RhiPresenceAttemptOutcome, 991 retry_delay: RhiPresenceRetryDelayMilliseconds, 992 ) -> Result<RhiPresenceAttemptCommit, RhiPresencePublicationError> { 993 require_writable(self)?; 994 let input = 995 PresenceRecordInput::from_prepared(prepared, finished_at, outcome, retry_delay)?; 996 self.host() 997 .sqlite_host() 998 .transaction(move |transaction| { 999 Box::pin(async move { record_outcome(transaction, input).await }) 1000 }) 1001 .await 1002 .map_err(map_transaction_error) 1003 } 1004 1005 /// Recovers at most one expired lease and records Unknown for submitted work. 1006 pub async fn recover_one_expired_presence( 1007 &self, 1008 adapters: &RhiTimeEntropyAdapters, 1009 now: RhiPresenceUnixMilliseconds, 1010 ) -> Result<bool, RhiPresencePublicationError> { 1011 require_writable(self)?; 1012 let candidate = self 1013 .host() 1014 .sqlite_host() 1015 .transaction(move |transaction| { 1016 Box::pin(async move { read_expired_candidate(transaction, now).await }) 1017 }) 1018 .await 1019 .map_err(map_transaction_error)?; 1020 let Some(candidate) = candidate else { 1021 return Ok(false); 1022 }; 1023 let delay = sample_retry_delay(adapters, candidate.retry_upper_bound)?; 1024 self.host() 1025 .sqlite_host() 1026 .transaction(move |transaction| { 1027 Box::pin(async move { recover_expired(transaction, candidate, now, delay).await }) 1028 }) 1029 .await 1030 .map_err(map_transaction_error) 1031 } 1032 1033 /// Executes at most one exact-byte presence attempt through an injected sink. 1034 /// 1035 /// Cancellation after durable preparation leaves Submitted evidence. Lease 1036 /// expiry recovery records Unknown before retrying the same retained bytes. 1037 pub async fn execute_next_presence( 1038 &self, 1039 owner: RhiPresenceLeaseOwner, 1040 adapters: &RhiTimeEntropyAdapters, 1041 sink: &dyn RhiExactPresenceSink, 1042 ) -> Result<Option<RhiPresenceAttemptCommit>, RhiPresencePublicationError> { 1043 let now = presence_now(adapters)?; 1044 self.recover_one_expired_presence(adapters, now).await?; 1045 let Some(lease) = self.claim_next_presence(owner, now).await? else { 1046 return Ok(None); 1047 }; 1048 let prepared = self.prepare_next_presence_target(lease, now).await?; 1049 let outcome = match sink.submit_exact(&prepared).await { 1050 RhiPresenceAttemptOutcome::Submitted => RhiPresenceAttemptOutcome::Unknown, 1051 outcome => outcome, 1052 }; 1053 let finished_at = presence_now(adapters)?; 1054 let retry_delay = if retryable_outcome(outcome) { 1055 sample_retry_delay(adapters, prepared.retry_upper_bound())? 1056 } else { 1057 RhiPresenceRetryDelayMilliseconds(0) 1058 }; 1059 self.record_presence_outcome(&prepared, finished_at, outcome, retry_delay) 1060 .await 1061 .map(Some) 1062 } 1063 } 1064 1065 fn require_writable( 1066 repository: &RhiPresenceOutboxRepository<'_>, 1067 ) -> Result<(), RhiPresencePublicationError> { 1068 if repository.host().mode() == RhiStateHostMode::ReadWriteExisting { 1069 Ok(()) 1070 } else { 1071 Err(failure(RhiPresencePublicationErrorKind::InvalidMode)) 1072 } 1073 } 1074 1075 fn presence_now( 1076 adapters: &RhiTimeEntropyAdapters, 1077 ) -> Result<RhiPresenceUnixMilliseconds, RhiPresencePublicationError> { 1078 let value = adapters.now_utc_milliseconds().map_err(|error| { 1079 failure(match error.kind() { 1080 RhiRuntimeAdapterErrorKind::WallClockUnavailable => { 1081 RhiPresencePublicationErrorKind::ClockUnavailable 1082 } 1083 _ => RhiPresencePublicationErrorKind::ClockUnavailable, 1084 }) 1085 })?; 1086 RhiPresenceUnixMilliseconds::new(value) 1087 .map_err(|_| failure(RhiPresencePublicationErrorKind::ClockUnavailable)) 1088 } 1089 1090 fn sample_retry_delay( 1091 adapters: &RhiTimeEntropyAdapters, 1092 cap: u64, 1093 ) -> Result<RhiPresenceRetryDelayMilliseconds, RhiPresencePublicationError> { 1094 if cap == 0 { 1095 return Ok(RhiPresenceRetryDelayMilliseconds(0)); 1096 } 1097 let bound = RhiJitterBoundMilliseconds::new(cap) 1098 .map_err(|_| failure(RhiPresencePublicationErrorKind::InvalidInput))?; 1099 adapters 1100 .sample_full_jitter(bound) 1101 .map(|delay| RhiPresenceRetryDelayMilliseconds(delay.get())) 1102 .map_err(|_| failure(RhiPresencePublicationErrorKind::EntropyUnavailable)) 1103 } 1104 1105 async fn commit_signed( 1106 transaction: &mut ServiceSqliteTransaction<'_>, 1107 documents: RhiSignedPresenceDocuments, 1108 committed_at: RhiPresenceUnixMilliseconds, 1109 ) -> Result<RhiPresenceCommitOutcome, PresenceOperationError> { 1110 let desired = read_desired_binding(transaction).await?; 1111 validate_document_set(&documents, desired)?; 1112 let current = read_generation_outboxes(transaction, desired.generation).await?; 1113 if !current.is_empty() && generation_matches(&documents, ¤t)? { 1114 return Ok(RhiPresenceCommitOutcome { 1115 desired_state: documents.desired_state, 1116 document_count: u8::try_from(documents.documents.len()) 1117 .map_err(|_| PresenceOperationError::Invariant)?, 1118 changed: false, 1119 }); 1120 } 1121 if !current.is_empty() { 1122 return Err(PresenceOperationError::Invariant); 1123 } 1124 if documents.desired_state.mode() == RhiPresenceDesiredMode::Enabled { 1125 for document in &documents.documents { 1126 let rows = sqlx::query(READ_MAX_AUTHORED_SQL) 1127 .bind(document.kind.code()) 1128 .fetch_all(&mut *transaction) 1129 .await 1130 .map_err(|_| PresenceOperationError::Storage)?; 1131 let [row] = rows.as_slice() else { 1132 return Err(PresenceOperationError::Invariant); 1133 }; 1134 if let Some(previous) = row 1135 .try_get::<Option<i64>, _>("authored_at_unix_s") 1136 .map_err(|_| PresenceOperationError::Invariant)? 1137 { 1138 let previous = 1139 u64::try_from(previous).map_err(|_| PresenceOperationError::Invariant)?; 1140 if document.created_at_unix_seconds <= previous { 1141 return Err(PresenceOperationError::InvalidInput); 1142 } 1143 } 1144 } 1145 } 1146 let lease_rows = sqlx::query(COUNT_STALE_LEASES_SQL) 1147 .bind(i64_value(desired.generation)?) 1148 .fetch_all(&mut *transaction) 1149 .await 1150 .map_err(|_| PresenceOperationError::Storage)?; 1151 let [lease_row] = lease_rows.as_slice() else { 1152 return Err(PresenceOperationError::Invariant); 1153 }; 1154 if lease_row 1155 .try_get::<i64, _>("lease_count") 1156 .ok() 1157 .filter(|value| *value >= 0) 1158 .ok_or(PresenceOperationError::Invariant)? 1159 != 0 1160 { 1161 return Err(PresenceOperationError::NotReady); 1162 } 1163 let superseded = sqlx::query(SUPERSEDE_OUTBOX_SQL) 1164 .bind(i64_value(committed_at.get())?) 1165 .bind(i64_value(desired.generation)?) 1166 .bind(i64_value(committed_at.get())?) 1167 .execute(&mut *transaction) 1168 .await 1169 .map_err(|_| PresenceOperationError::Storage)? 1170 .rows_affected(); 1171 1172 for document in &documents.documents { 1173 let outbox_id = derive_outbox_id( 1174 desired.generation, 1175 document.kind, 1176 desired.desired_sha256, 1177 document.signed_event_sha256, 1178 ); 1179 let inserted = sqlx::query(INSERT_OUTBOX_SQL) 1180 .bind(outbox_id.as_bytes().as_slice()) 1181 .bind(i64_value(desired.generation)?) 1182 .bind(document.kind.code()) 1183 .bind(desired.desired_sha256.as_slice()) 1184 .bind(desired.target_set_sha256.as_slice()) 1185 .bind(document.event_id.as_slice()) 1186 .bind(document.signed_event_sha256.as_slice()) 1187 .bind(document.signed_event_bytes.as_ref()) 1188 .bind(i64_value(document.created_at_unix_seconds)?) 1189 .bind(documents.service_public_key.as_ref()) 1190 .bind(i64::from(desired.target_count)) 1191 .bind(i64::from(desired.required_target_count)) 1192 .bind(i64_value(committed_at.get())?) 1193 .bind(i64_value(committed_at.get())?) 1194 .bind(i64_value(committed_at.get())?) 1195 .execute(&mut *transaction) 1196 .await 1197 .map_err(|_| PresenceOperationError::Storage)?; 1198 if inserted.rows_affected() != 1 { 1199 return Err(PresenceOperationError::Invariant); 1200 } 1201 for target in &documents.targets { 1202 let inserted = sqlx::query(INSERT_TARGET_SQL) 1203 .bind(outbox_id.as_bytes().as_slice()) 1204 .bind(i64::from(target.ordinal())) 1205 .bind(target.relay_id()) 1206 .bind(i64::from(target.required())) 1207 .bind(i64_value(committed_at.get())?) 1208 .bind(i64_value(committed_at.get())?) 1209 .bind(i64_value(committed_at.get())?) 1210 .execute(&mut *transaction) 1211 .await 1212 .map_err(|_| PresenceOperationError::Storage)?; 1213 if inserted.rows_affected() != 1 { 1214 return Err(PresenceOperationError::Invariant); 1215 } 1216 } 1217 } 1218 let committed = read_generation_outboxes(transaction, desired.generation).await?; 1219 if !generation_matches(&documents, &committed)? { 1220 return Err(PresenceOperationError::Invariant); 1221 } 1222 Ok(RhiPresenceCommitOutcome { 1223 desired_state: documents.desired_state, 1224 document_count: u8::try_from(documents.documents.len()) 1225 .map_err(|_| PresenceOperationError::Invariant)?, 1226 changed: !documents.documents.is_empty() || superseded != 0, 1227 }) 1228 } 1229 1230 fn validate_document_set( 1231 documents: &RhiSignedPresenceDocuments, 1232 desired: PresenceDesiredBinding, 1233 ) -> Result<(), PresenceOperationError> { 1234 let state = documents.desired_state; 1235 if state.generation() != desired.generation 1236 || state.mode() != desired.mode 1237 || state.profile() != desired.profile 1238 || state.application_handler() != desired.application_handler 1239 || state.target_set_sha256() != &desired.target_set_sha256 1240 || state.target_count() != desired.target_count 1241 || state.required_target_count() != desired.required_target_count 1242 || state.queue_capacity() != desired.queue_capacity 1243 || state.desired_sha256() != &desired.desired_sha256 1244 || documents.targets.len() != usize::from(desired.target_count) 1245 || documents.documents.len() != desired.document_count() 1246 || target_set_digest(&documents.targets)? != desired.target_set_sha256 1247 || !valid_public_key(&documents.service_public_key) 1248 { 1249 return Err(PresenceOperationError::DesiredStateMismatch); 1250 } 1251 let expected = desired.document_kinds(); 1252 for (document, kind) in documents.documents.iter().zip(expected) { 1253 if document.kind != kind 1254 || document.signed_event_bytes.is_empty() 1255 || document.signed_event_bytes.len() > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES 1256 || <[u8; 32]>::from(Sha256::digest(&document.signed_event_bytes)) 1257 != document.signed_event_sha256 1258 { 1259 return Err(PresenceOperationError::VerificationFailed); 1260 } 1261 let event = validate_signed_document( 1262 document.kind, 1263 &documents.service_public_key, 1264 document.created_at_unix_seconds, 1265 &document.signed_event_bytes, 1266 ) 1267 .map_err(|_| PresenceOperationError::VerificationFailed)?; 1268 if event.id().as_bytes() != &document.event_id { 1269 return Err(PresenceOperationError::VerificationFailed); 1270 } 1271 } 1272 Ok(()) 1273 } 1274 1275 fn generation_matches( 1276 documents: &RhiSignedPresenceDocuments, 1277 current: &[PresenceOutboxRecord], 1278 ) -> Result<bool, PresenceOperationError> { 1279 if documents.documents.len() != current.len() { 1280 return Ok(false); 1281 } 1282 for (document, outbox) in documents.documents.iter().zip(current) { 1283 if outbox.desired_generation != documents.desired_state.generation() 1284 || outbox.desired_sha256 != *documents.desired_state.desired_sha256() 1285 || outbox.target_set_sha256 != *documents.desired_state.target_set_sha256() 1286 || outbox.target_count != documents.desired_state.target_count() 1287 || outbox.required_target_count != documents.desired_state.required_target_count() 1288 || outbox.id 1289 != derive_outbox_id( 1290 outbox.desired_generation, 1291 document.kind, 1292 outbox.desired_sha256, 1293 document.signed_event_sha256, 1294 ) 1295 || document.kind != outbox.document_kind 1296 || document.event_id != outbox.event_id 1297 || document.signed_event_sha256 != outbox.event_sha256 1298 || document.signed_event_bytes.as_ref() != outbox.exact_signed_event_bytes.as_ref() 1299 || document.created_at_unix_seconds != outbox.authored_at_unix_s 1300 || documents.service_public_key.as_ref() != outbox.service_public_key.as_ref() 1301 || documents 1302 .targets 1303 .iter() 1304 .zip(outbox.targets.iter()) 1305 .any(|(expected, actual)| { 1306 expected.ordinal() != actual.ordinal 1307 || expected.relay_id() != actual.relay_id.as_ref() 1308 || expected.required() != actual.required 1309 }) 1310 { 1311 return Ok(false); 1312 } 1313 validate_targets(outbox, &outbox.targets)?; 1314 } 1315 Ok(true) 1316 } 1317 1318 fn derive_outbox_id( 1319 generation: u64, 1320 kind: RhiPresenceDocumentKind, 1321 desired_sha256: [u8; 32], 1322 event_sha256: [u8; 32], 1323 ) -> RhiPresenceOutboxId { 1324 let mut digest = Sha256::new(); 1325 digest.update(OUTBOX_ID_DOMAIN); 1326 digest.update(generation.to_be_bytes()); 1327 digest.update([document_kind_tag(kind)]); 1328 digest.update(desired_sha256); 1329 digest.update(event_sha256); 1330 RhiPresenceOutboxId(digest.finalize().into()) 1331 } 1332 1333 fn derive_attempt_id( 1334 outbox: RhiPresenceOutboxId, 1335 event_sha256: [u8; 32], 1336 ordinal: u8, 1337 attempt_number: u16, 1338 ) -> RhiPresenceAttemptId { 1339 let mut digest = Sha256::new(); 1340 digest.update(ATTEMPT_ID_DOMAIN); 1341 digest.update(outbox.0); 1342 digest.update(event_sha256); 1343 digest.update([ordinal]); 1344 digest.update(attempt_number.to_be_bytes()); 1345 RhiPresenceAttemptId(digest.finalize().into()) 1346 } 1347 1348 const fn document_kind_tag(kind: RhiPresenceDocumentKind) -> u8 { 1349 match kind { 1350 RhiPresenceDocumentKind::ServiceProfile => 0, 1351 RhiPresenceDocumentKind::ApplicationHandler => 1, 1352 } 1353 } 1354 1355 fn target_set_digest(targets: &[RhiPresenceTarget]) -> Result<[u8; 32], PresenceOperationError> { 1356 const DOMAIN: &[u8] = b"radroots.rhi.presence_target_set.v1\0"; 1357 let mut digest = Sha256::new(); 1358 digest.update(DOMAIN); 1359 digest.update( 1360 u32::try_from(targets.len()) 1361 .map_err(|_| PresenceOperationError::InvalidInput)? 1362 .to_be_bytes(), 1363 ); 1364 for target in targets { 1365 digest.update(u32::from(target.ordinal()).to_be_bytes()); 1366 digest.update( 1367 u64::try_from(target.relay_id().len()) 1368 .map_err(|_| PresenceOperationError::InvalidInput)? 1369 .to_be_bytes(), 1370 ); 1371 digest.update(target.relay_id().as_bytes()); 1372 digest.update([u8::from(target.required())]); 1373 } 1374 Ok(digest.finalize().into()) 1375 } 1376 1377 async fn read_desired_binding( 1378 transaction: &mut ServiceSqliteTransaction<'_>, 1379 ) -> Result<PresenceDesiredBinding, PresenceOperationError> { 1380 let rows = sqlx::query(READ_DESIRED_BINDING_SQL) 1381 .fetch_all(&mut *transaction) 1382 .await 1383 .map_err(|_| PresenceOperationError::Storage)?; 1384 let [row] = rows.as_slice() else { 1385 return Err(PresenceOperationError::DesiredStateMismatch); 1386 }; 1387 if row.try_get::<i64, _>("singleton").ok() != Some(1) { 1388 return Err(PresenceOperationError::Invariant); 1389 } 1390 let generation = positive_u64(row, "generation")?; 1391 let enabled = boolean_i64(row, "enabled")?; 1392 let profile = boolean_i64(row, "profile")?; 1393 let application_handler = boolean_i64(row, "application_handler")?; 1394 let target_set_sha256 = blob::<32>(row, "target_set_sha256")?; 1395 let target_count = bounded_u8(row, "target_count", RHI_PRESENCE_DESIRED_MAX_TARGETS)?; 1396 let required_target_count = 1397 bounded_u8(row, "required_target_count", usize::from(target_count))?; 1398 let queue_capacity = row 1399 .try_get::<i64, _>("queue_capacity") 1400 .ok() 1401 .and_then(|value| u32::try_from(value).ok()) 1402 .filter(|value| *value <= 4_096) 1403 .ok_or(PresenceOperationError::Invariant)?; 1404 let desired_sha256 = blob::<32>(row, "desired_sha256")?; 1405 let mode = if enabled { 1406 RhiPresenceDesiredMode::Enabled 1407 } else { 1408 RhiPresenceDesiredMode::Disabled 1409 }; 1410 let valid = if enabled { 1411 (profile || application_handler) && target_count > 0 && queue_capacity > 0 1412 } else { 1413 !profile 1414 && !application_handler 1415 && target_count == 0 1416 && required_target_count == 0 1417 && queue_capacity == 0 1418 }; 1419 if !valid { 1420 return Err(PresenceOperationError::Invariant); 1421 } 1422 Ok(PresenceDesiredBinding { 1423 generation, 1424 mode, 1425 profile, 1426 application_handler, 1427 target_set_sha256, 1428 target_count, 1429 required_target_count, 1430 queue_capacity, 1431 desired_sha256, 1432 }) 1433 } 1434 1435 async fn read_generation_outboxes( 1436 transaction: &mut ServiceSqliteTransaction<'_>, 1437 generation: u64, 1438 ) -> Result<Vec<PresenceOutboxRecord>, PresenceOperationError> { 1439 let rows = sqlx::query(READ_GENERATION_OUTBOX_SQL) 1440 .bind(i64_value(generation)?) 1441 .fetch_all(&mut *transaction) 1442 .await 1443 .map_err(|_| PresenceOperationError::Storage)?; 1444 if rows.len() > 2 { 1445 return Err(PresenceOperationError::Invariant); 1446 } 1447 let mut records = Vec::with_capacity(rows.len()); 1448 for row in rows { 1449 let mut outbox = decode_outbox(row)?; 1450 outbox.targets = read_targets(transaction, outbox.id) 1451 .await? 1452 .into_boxed_slice(); 1453 records.push(outbox); 1454 } 1455 Ok(records) 1456 } 1457 1458 async fn read_outbox( 1459 transaction: &mut ServiceSqliteTransaction<'_>, 1460 id: RhiPresenceOutboxId, 1461 ) -> Result<Option<PresenceOutboxRecord>, PresenceOperationError> { 1462 let rows = sqlx::query(READ_OUTBOX_SQL) 1463 .bind(id.as_bytes().as_slice()) 1464 .fetch_all(&mut *transaction) 1465 .await 1466 .map_err(|_| PresenceOperationError::Storage)?; 1467 let Some(row) = exactly_zero_or_one(rows)? else { 1468 return Ok(None); 1469 }; 1470 let mut outbox = decode_outbox(row)?; 1471 outbox.targets = read_targets(transaction, outbox.id) 1472 .await? 1473 .into_boxed_slice(); 1474 Ok(Some(outbox)) 1475 } 1476 1477 async fn read_targets( 1478 transaction: &mut ServiceSqliteTransaction<'_>, 1479 outbox_id: RhiPresenceOutboxId, 1480 ) -> Result<Vec<PresenceTargetRecord>, PresenceOperationError> { 1481 let rows = sqlx::query(READ_TARGETS_SQL) 1482 .bind(outbox_id.as_bytes().as_slice()) 1483 .fetch_all(&mut *transaction) 1484 .await 1485 .map_err(|_| PresenceOperationError::Storage)?; 1486 if rows.len() > RHI_PRESENCE_DESIRED_MAX_TARGETS { 1487 return Err(PresenceOperationError::Invariant); 1488 } 1489 rows.into_iter().map(decode_target).collect() 1490 } 1491 1492 fn decode_outbox( 1493 row: sqlx::sqlite::SqliteRow, 1494 ) -> Result<PresenceOutboxRecord, PresenceOperationError> { 1495 let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); 1496 let desired_generation = positive_u64(&row, "desired_generation")?; 1497 let document_kind = decode_document_kind(&bounded_text( 1498 &row, 1499 "document_kind", 1500 "document_kind_bytes", 1501 19, 1502 )?)?; 1503 let desired_sha256 = blob::<32>(&row, "desired_sha256")?; 1504 let target_set_sha256 = blob::<32>(&row, "target_set_sha256")?; 1505 let event_id = blob::<32>(&row, "event_id")?; 1506 let event_sha256 = blob::<32>(&row, "event_sha256")?; 1507 let exact_length = positive_usize(&row, "exact_signed_event_bytes_length")?; 1508 if exact_length > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES { 1509 return Err(PresenceOperationError::Invariant); 1510 } 1511 let exact_signed_event_bytes = row 1512 .try_get::<Vec<u8>, _>("exact_signed_event_bytes") 1513 .map_err(|_| PresenceOperationError::Invariant)?; 1514 let authored_at_unix_s = nonnegative_u64(&row, "authored_at_unix_s")?; 1515 let service_public_key = 1516 bounded_text(&row, "service_public_key", "service_public_key_bytes", 64)?; 1517 let target_count = bounded_u8(&row, "target_count", RHI_PRESENCE_DESIRED_MAX_TARGETS)?; 1518 let required_target_count = 1519 bounded_u8(&row, "required_target_count", usize::from(target_count))?; 1520 let max_attempts = bounded_u16(&row, "max_attempts", RHI_PRESENCE_MAX_ATTEMPTS)?; 1521 let initial_backoff_ms = bounded_u64(&row, "initial_backoff_ms", INITIAL_BACKOFF_MILLISECONDS)?; 1522 let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", MAXIMUM_BACKOFF_MILLISECONDS)?; 1523 let attempt_deadline_ms = 1524 bounded_u64(&row, "attempt_deadline_ms", ATTEMPT_DEADLINE_MILLISECONDS)?; 1525 let state = decode_outbox_state(&bounded_text(&row, "state", "state_bytes", 10)?)?; 1526 let revision = positive_u64(&row, "revision")?; 1527 let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?; 1528 let lease_owner = optional_owner(&row)?; 1529 let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?; 1530 let created_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "created_at_unix_ms")?); 1531 let updated_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?); 1532 if exact_signed_event_bytes.len() != exact_length 1533 || exact_signed_event_bytes.is_empty() 1534 || <[u8; 32]>::from(Sha256::digest(&exact_signed_event_bytes)) != event_sha256 1535 || !valid_public_key(&service_public_key) 1536 || max_attempts != RHI_PRESENCE_MAX_ATTEMPTS 1537 || initial_backoff_ms != INITIAL_BACKOFF_MILLISECONDS 1538 || maximum_backoff_ms != MAXIMUM_BACKOFF_MILLISECONDS 1539 || attempt_deadline_ms != ATTEMPT_DEADLINE_MILLISECONDS 1540 || initial_backoff_ms > maximum_backoff_ms 1541 || updated_at < created_at 1542 || !valid_outbox_shape(state, next_attempt, lease_owner, lease_expires) 1543 { 1544 return Err(PresenceOperationError::Invariant); 1545 } 1546 Ok(PresenceOutboxRecord { 1547 id, 1548 desired_generation, 1549 document_kind, 1550 desired_sha256, 1551 target_set_sha256, 1552 event_id, 1553 event_sha256, 1554 exact_signed_event_bytes: exact_signed_event_bytes.into_boxed_slice(), 1555 authored_at_unix_s, 1556 service_public_key: service_public_key.into_boxed_str(), 1557 target_count, 1558 required_target_count, 1559 max_attempts, 1560 state, 1561 revision, 1562 next_attempt, 1563 lease_owner, 1564 lease_expires, 1565 created_at, 1566 updated_at, 1567 targets: Box::new([]), 1568 }) 1569 } 1570 1571 fn decode_target( 1572 row: sqlx::sqlite::SqliteRow, 1573 ) -> Result<PresenceTargetRecord, PresenceOperationError> { 1574 let ordinal = bounded_u8(&row, "target_ordinal", RHI_PRESENCE_DESIRED_MAX_TARGETS - 1)?; 1575 let relay_id = bounded_text(&row, "relay_id", "relay_id_bytes", 64)?; 1576 let required = boolean_i64(&row, "required")?; 1577 let state = decode_target_state(&bounded_text(&row, "state", "state_bytes", 13)?)?; 1578 let revision = positive_u64(&row, "revision")?; 1579 let attempt_count = bounded_u16(&row, "attempt_count", RHI_PRESENCE_MAX_ATTEMPTS)?; 1580 let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?; 1581 let last_attempt_id = optional_attempt_id(&row)?; 1582 let updated_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?); 1583 if !valid_relay_id(&relay_id) 1584 || !valid_target_shape(state, attempt_count, next_attempt, last_attempt_id) 1585 { 1586 return Err(PresenceOperationError::Invariant); 1587 } 1588 Ok(PresenceTargetRecord { 1589 ordinal, 1590 relay_id: relay_id.into_boxed_str(), 1591 required, 1592 state, 1593 revision, 1594 attempt_count, 1595 next_attempt, 1596 last_attempt_id, 1597 updated_at, 1598 }) 1599 } 1600 1601 async fn claim_next( 1602 transaction: &mut ServiceSqliteTransaction<'_>, 1603 owner: RhiPresenceLeaseOwner, 1604 now: RhiPresenceUnixMilliseconds, 1605 ) -> Result<Option<RhiPresenceLease>, PresenceOperationError> { 1606 let rows = sqlx::query(READ_CLAIMABLE_OUTBOX_SQL) 1607 .bind(i64_value(now.get())?) 1608 .fetch_all(&mut *transaction) 1609 .await 1610 .map_err(|_| PresenceOperationError::Storage)?; 1611 let Some(row) = exactly_zero_or_one(rows)? else { 1612 return Ok(None); 1613 }; 1614 let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); 1615 let outbox = read_outbox(transaction, id) 1616 .await? 1617 .ok_or(PresenceOperationError::Invariant)?; 1618 validate_current_desired(transaction, &outbox).await?; 1619 validate_targets(&outbox, &outbox.targets)?; 1620 let expires = now 1621 .get() 1622 .checked_add(ATTEMPT_DEADLINE_MILLISECONDS) 1623 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 1624 .ok_or(PresenceOperationError::InvalidInput)?; 1625 let result = sqlx::query(CLAIM_OUTBOX_SQL) 1626 .bind(owner.0.as_slice()) 1627 .bind(i64_value(expires)?) 1628 .bind(i64_value(now.get())?) 1629 .bind(outbox.id.as_bytes().as_slice()) 1630 .bind(i64_value(outbox.revision)?) 1631 .bind(i64_value(now.get())?) 1632 .bind(i64_value(now.get())?) 1633 .execute(&mut *transaction) 1634 .await 1635 .map_err(|_| PresenceOperationError::Storage)?; 1636 if result.rows_affected() != 1 { 1637 return Err(PresenceOperationError::LeaseLost); 1638 } 1639 let leased = read_outbox(transaction, id) 1640 .await? 1641 .filter(|record| { 1642 record.state == RhiPresenceOutboxState::Leased 1643 && record.lease_owner == Some(owner) 1644 && record.lease_expires == Some(RhiPresenceUnixMilliseconds(expires)) 1645 }) 1646 .ok_or(PresenceOperationError::Invariant)?; 1647 Ok(Some(RhiPresenceLease { 1648 outbox: leased, 1649 owner, 1650 expires_at: RhiPresenceUnixMilliseconds(expires), 1651 })) 1652 } 1653 1654 async fn prepare_next( 1655 transaction: &mut ServiceSqliteTransaction<'_>, 1656 lease: RhiPresenceLease, 1657 started_at: RhiPresenceUnixMilliseconds, 1658 ) -> Result<RhiPreparedPresenceAttempt, PresenceOperationError> { 1659 validate_lease(transaction, &lease, started_at, true).await?; 1660 let candidate = lease 1661 .outbox 1662 .targets 1663 .iter() 1664 .find(|target| target_is_due(target, lease.outbox.max_attempts, started_at)) 1665 .ok_or(PresenceOperationError::NotReady)?; 1666 let candidate_ordinal = candidate.ordinal; 1667 let candidate_revision = candidate.revision; 1668 let attempt_number = candidate 1669 .attempt_count 1670 .checked_add(1) 1671 .filter(|value| *value <= lease.outbox.max_attempts) 1672 .ok_or(PresenceOperationError::Invariant)?; 1673 let attempt_id = derive_attempt_id( 1674 lease.outbox.id, 1675 lease.outbox.event_sha256, 1676 candidate_ordinal, 1677 attempt_number, 1678 ); 1679 let result = sqlx::query(PREPARE_TARGET_SQL) 1680 .bind(attempt_id.as_bytes().as_slice()) 1681 .bind(i64_value(started_at.get())?) 1682 .bind(lease.outbox.id.as_bytes().as_slice()) 1683 .bind(i64::from(candidate_ordinal)) 1684 .bind(i64_value(candidate_revision)?) 1685 .bind(i64::from(lease.outbox.max_attempts)) 1686 .bind(i64_value(started_at.get())?) 1687 .bind(i64_value(started_at.get())?) 1688 .execute(&mut *transaction) 1689 .await 1690 .map_err(|_| PresenceOperationError::Storage)?; 1691 if result.rows_affected() != 1 { 1692 return Err(PresenceOperationError::LeaseLost); 1693 } 1694 let current = read_outbox(transaction, lease.outbox.id) 1695 .await? 1696 .ok_or(PresenceOperationError::Invariant)?; 1697 if !same_outbox_identity(¤t, &lease.outbox) { 1698 return Err(PresenceOperationError::LeaseLost); 1699 } 1700 let target = current 1701 .targets 1702 .into_vec() 1703 .into_iter() 1704 .find(|target| target.ordinal == candidate_ordinal) 1705 .filter(|target| { 1706 target.state == RhiPresenceTargetState::Submitted 1707 && target.attempt_count == attempt_number 1708 && target.last_attempt_id == Some(attempt_id) 1709 && target.updated_at == started_at 1710 }) 1711 .ok_or(PresenceOperationError::Invariant)?; 1712 validate_signed_document( 1713 lease.outbox.document_kind, 1714 &lease.outbox.service_public_key, 1715 lease.outbox.authored_at_unix_s, 1716 &lease.outbox.exact_signed_event_bytes, 1717 ) 1718 .map_err(|_| PresenceOperationError::Invariant)?; 1719 let deadline_at = lease.expires_at; 1720 let exact_signed_event_bytes = lease.outbox.exact_signed_event_bytes.clone(); 1721 Ok(RhiPreparedPresenceAttempt { 1722 lease, 1723 target, 1724 attempt_id, 1725 started_at, 1726 deadline_at, 1727 exact_signed_event_bytes, 1728 }) 1729 } 1730 1731 async fn record_outcome( 1732 transaction: &mut ServiceSqliteTransaction<'_>, 1733 input: PresenceRecordInput, 1734 ) -> Result<RhiPresenceAttemptCommit, PresenceOperationError> { 1735 if let Some(existing) = read_attempt(transaction, input.attempt_id).await? { 1736 return reconcile_recorded(transaction, input, existing).await; 1737 } 1738 validate_input_lease(transaction, input, true).await?; 1739 let outbox = read_outbox(transaction, input.outbox_id) 1740 .await? 1741 .ok_or(PresenceOperationError::Invariant)?; 1742 let target = outbox 1743 .targets 1744 .iter() 1745 .find(|target| target.ordinal == input.target_ordinal) 1746 .ok_or(PresenceOperationError::Invariant)?; 1747 if target.state != RhiPresenceTargetState::Submitted 1748 || target.revision != input.target_revision 1749 || target.attempt_count != input.attempt_number 1750 || target.last_attempt_id != Some(input.attempt_id) 1751 || target.updated_at != input.started_at 1752 { 1753 return Err(PresenceOperationError::LeaseLost); 1754 } 1755 insert_attempt(transaction, input).await?; 1756 let (next_state, next_attempt) = target_outcome_schedule(input)?; 1757 let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL) 1758 .bind(next_state.code()) 1759 .bind(optional_i64(next_attempt)?) 1760 .bind(i64_value(input.finished_at.get())?) 1761 .bind(input.outbox_id.as_bytes().as_slice()) 1762 .bind(i64::from(input.target_ordinal)) 1763 .bind(i64_value(input.target_revision)?) 1764 .bind(i64::from(input.attempt_number)) 1765 .bind(input.attempt_id.as_bytes().as_slice()) 1766 .execute(&mut *transaction) 1767 .await 1768 .map_err(|_| PresenceOperationError::Storage)?; 1769 if result.rows_affected() != 1 { 1770 return Err(PresenceOperationError::LeaseLost); 1771 } 1772 let targets = read_targets(transaction, input.outbox_id).await?; 1773 let disposition = disposition(&outbox, &targets)?; 1774 update_outbox_after_attempt(transaction, input, disposition).await?; 1775 Ok(RhiPresenceAttemptCommit { 1776 outbox_id: input.outbox_id, 1777 attempt_id: input.attempt_id, 1778 target_ordinal: input.target_ordinal, 1779 attempt_number: input.attempt_number, 1780 outcome: input.outcome, 1781 target_state: next_state, 1782 outbox_state: disposition.state, 1783 }) 1784 } 1785 1786 async fn read_expired_candidate( 1787 transaction: &mut ServiceSqliteTransaction<'_>, 1788 now: RhiPresenceUnixMilliseconds, 1789 ) -> Result<Option<PresenceRecoveryCandidate>, PresenceOperationError> { 1790 let rows = sqlx::query(READ_EXPIRED_OUTBOX_SQL) 1791 .bind(i64_value(now.get())?) 1792 .fetch_all(&mut *transaction) 1793 .await 1794 .map_err(|_| PresenceOperationError::Storage)?; 1795 let Some(row) = exactly_zero_or_one(rows)? else { 1796 return Ok(None); 1797 }; 1798 let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); 1799 let outbox = read_outbox(transaction, id) 1800 .await? 1801 .ok_or(PresenceOperationError::Invariant)?; 1802 let desired = read_desired_binding(transaction).await?; 1803 let current_desired = desired_matches_outbox(desired, &outbox); 1804 validate_targets(&outbox, &outbox.targets)?; 1805 let submitted: Vec<_> = outbox 1806 .targets 1807 .iter() 1808 .filter(|target| target.state == RhiPresenceTargetState::Submitted) 1809 .map(|target| target.attempt_count) 1810 .collect(); 1811 if submitted.len() > 1 { 1812 return Err(PresenceOperationError::Invariant); 1813 } 1814 let retry_upper_bound = if current_desired { 1815 submitted 1816 .first() 1817 .copied() 1818 .filter(|attempt| *attempt < outbox.max_attempts) 1819 .map_or(0, retry_upper_bound) 1820 } else { 1821 0 1822 }; 1823 Ok(Some(PresenceRecoveryCandidate { 1824 outbox, 1825 retry_upper_bound, 1826 current_desired, 1827 })) 1828 } 1829 1830 async fn recover_expired( 1831 transaction: &mut ServiceSqliteTransaction<'_>, 1832 candidate: PresenceRecoveryCandidate, 1833 now: RhiPresenceUnixMilliseconds, 1834 delay: RhiPresenceRetryDelayMilliseconds, 1835 ) -> Result<bool, PresenceOperationError> { 1836 let current = read_outbox(transaction, candidate.outbox.id) 1837 .await? 1838 .filter(|current| same_outbox(current, &candidate.outbox)) 1839 .ok_or(PresenceOperationError::LeaseLost)?; 1840 if current.state != RhiPresenceOutboxState::Leased 1841 || current.lease_expires.is_none_or(|expires| expires > now) 1842 { 1843 return Err(PresenceOperationError::LeaseLost); 1844 } 1845 let owner = current 1846 .lease_owner 1847 .ok_or(PresenceOperationError::Invariant)?; 1848 let expires = current 1849 .lease_expires 1850 .ok_or(PresenceOperationError::Invariant)?; 1851 for target in current 1852 .targets 1853 .iter() 1854 .filter(|target| target.state == RhiPresenceTargetState::Submitted) 1855 { 1856 let attempt_id = target 1857 .last_attempt_id 1858 .ok_or(PresenceOperationError::Invariant)?; 1859 if read_attempt(transaction, attempt_id).await?.is_some() { 1860 return Err(PresenceOperationError::Invariant); 1861 } 1862 let input = PresenceRecordInput { 1863 outbox_id: current.id, 1864 outbox_revision: current.revision, 1865 event_sha256: current.event_sha256, 1866 owner, 1867 lease_expires: expires, 1868 target_ordinal: target.ordinal, 1869 target_revision: target.revision, 1870 attempt_number: target.attempt_count, 1871 attempt_id, 1872 started_at: target.updated_at, 1873 finished_at: now, 1874 outcome: RhiPresenceAttemptOutcome::Unknown, 1875 retry_delay: delay, 1876 }; 1877 insert_attempt(transaction, input).await?; 1878 let (state, next) = target_outcome_schedule(input)?; 1879 let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL) 1880 .bind(state.code()) 1881 .bind(optional_i64(next)?) 1882 .bind(i64_value(now.get())?) 1883 .bind(current.id.as_bytes().as_slice()) 1884 .bind(i64::from(target.ordinal)) 1885 .bind(i64_value(target.revision)?) 1886 .bind(i64::from(target.attempt_count)) 1887 .bind(attempt_id.as_bytes().as_slice()) 1888 .execute(&mut *transaction) 1889 .await 1890 .map_err(|_| PresenceOperationError::Storage)?; 1891 if result.rows_affected() != 1 { 1892 return Err(PresenceOperationError::LeaseLost); 1893 } 1894 } 1895 let targets = read_targets(transaction, current.id).await?; 1896 let disposition = if candidate.current_desired { 1897 disposition(¤t, &targets)? 1898 } else { 1899 PresenceOutboxDisposition { 1900 state: RhiPresenceOutboxState::Superseded, 1901 next_attempt: None, 1902 } 1903 }; 1904 let result = sqlx::query(UPDATE_OUTBOX_RECOVERY_SQL) 1905 .bind(disposition.state.code()) 1906 .bind(optional_i64(disposition.next_attempt)?) 1907 .bind(i64_value(now.get())?) 1908 .bind(current.id.as_bytes().as_slice()) 1909 .bind(i64_value(current.revision)?) 1910 .bind(owner.0.as_slice()) 1911 .bind(i64_value(expires.get())?) 1912 .bind(i64_value(now.get())?) 1913 .bind(i64_value(now.get())?) 1914 .execute(&mut *transaction) 1915 .await 1916 .map_err(|_| PresenceOperationError::Storage)?; 1917 if result.rows_affected() != 1 { 1918 return Err(PresenceOperationError::LeaseLost); 1919 } 1920 Ok(true) 1921 } 1922 1923 async fn validate_current_desired( 1924 transaction: &mut ServiceSqliteTransaction<'_>, 1925 outbox: &PresenceOutboxRecord, 1926 ) -> Result<(), PresenceOperationError> { 1927 let desired = read_desired_binding(transaction).await?; 1928 if !desired_matches_outbox(desired, outbox) { 1929 return Err(PresenceOperationError::DesiredStateMismatch); 1930 } 1931 Ok(()) 1932 } 1933 1934 fn desired_matches_outbox(desired: PresenceDesiredBinding, outbox: &PresenceOutboxRecord) -> bool { 1935 desired.mode == RhiPresenceDesiredMode::Enabled 1936 && desired.generation == outbox.desired_generation 1937 && desired.desired_sha256 == outbox.desired_sha256 1938 && desired.target_set_sha256 == outbox.target_set_sha256 1939 && desired.target_count == outbox.target_count 1940 && desired.required_target_count == outbox.required_target_count 1941 && desired 1942 .document_kinds() 1943 .any(|kind| kind == outbox.document_kind) 1944 } 1945 1946 async fn validate_lease( 1947 transaction: &mut ServiceSqliteTransaction<'_>, 1948 lease: &RhiPresenceLease, 1949 now: RhiPresenceUnixMilliseconds, 1950 require_unexpired: bool, 1951 ) -> Result<(), PresenceOperationError> { 1952 let current = read_outbox(transaction, lease.outbox.id) 1953 .await? 1954 .filter(|current| same_outbox(current, &lease.outbox)) 1955 .ok_or(PresenceOperationError::LeaseLost)?; 1956 validate_current_desired(transaction, ¤t).await?; 1957 if current.state != RhiPresenceOutboxState::Leased 1958 || current.lease_owner != Some(lease.owner) 1959 || current.lease_expires != Some(lease.expires_at) 1960 || (require_unexpired && lease.expires_at <= now) 1961 { 1962 return Err(PresenceOperationError::LeaseLost); 1963 } 1964 Ok(()) 1965 } 1966 1967 async fn validate_input_lease( 1968 transaction: &mut ServiceSqliteTransaction<'_>, 1969 input: PresenceRecordInput, 1970 require_unexpired: bool, 1971 ) -> Result<(), PresenceOperationError> { 1972 let current = read_outbox(transaction, input.outbox_id) 1973 .await? 1974 .ok_or(PresenceOperationError::LeaseLost)?; 1975 validate_current_desired(transaction, ¤t).await?; 1976 if current.revision != input.outbox_revision 1977 || current.event_sha256 != input.event_sha256 1978 || current.state != RhiPresenceOutboxState::Leased 1979 || current.lease_owner != Some(input.owner) 1980 || current.lease_expires != Some(input.lease_expires) 1981 || (require_unexpired && input.lease_expires <= input.finished_at) 1982 { 1983 return Err(PresenceOperationError::LeaseLost); 1984 } 1985 Ok(()) 1986 } 1987 1988 async fn insert_attempt( 1989 transaction: &mut ServiceSqliteTransaction<'_>, 1990 input: PresenceRecordInput, 1991 ) -> Result<(), PresenceOperationError> { 1992 let result = sqlx::query(INSERT_ATTEMPT_SQL) 1993 .bind(input.attempt_id.as_bytes().as_slice()) 1994 .bind(input.outbox_id.as_bytes().as_slice()) 1995 .bind(i64::from(input.target_ordinal)) 1996 .bind(i64::from(input.attempt_number)) 1997 .bind(input.event_sha256.as_slice()) 1998 .bind(input.owner.0.as_slice()) 1999 .bind(i64_value(input.started_at.get())?) 2000 .bind(i64_value(input.finished_at.get())?) 2001 .bind(input.outcome.code()) 2002 .bind(input.outcome.code()) 2003 .execute(&mut *transaction) 2004 .await 2005 .map_err(|_| PresenceOperationError::Storage)?; 2006 if result.rows_affected() != 1 { 2007 return Err(PresenceOperationError::Storage); 2008 } 2009 Ok(()) 2010 } 2011 2012 async fn update_outbox_after_attempt( 2013 transaction: &mut ServiceSqliteTransaction<'_>, 2014 input: PresenceRecordInput, 2015 disposition: PresenceOutboxDisposition, 2016 ) -> Result<(), PresenceOperationError> { 2017 let result = sqlx::query(UPDATE_OUTBOX_AFTER_ATTEMPT_SQL) 2018 .bind(disposition.state.code()) 2019 .bind(optional_i64(disposition.next_attempt)?) 2020 .bind(i64_value(input.finished_at.get())?) 2021 .bind(input.outbox_id.as_bytes().as_slice()) 2022 .bind(i64_value(input.outbox_revision)?) 2023 .bind(input.owner.0.as_slice()) 2024 .bind(i64_value(input.lease_expires.get())?) 2025 .bind(i64_value(input.finished_at.get())?) 2026 .bind(i64_value(input.finished_at.get())?) 2027 .execute(&mut *transaction) 2028 .await 2029 .map_err(|_| PresenceOperationError::Storage)?; 2030 if result.rows_affected() != 1 { 2031 return Err(PresenceOperationError::LeaseLost); 2032 } 2033 Ok(()) 2034 } 2035 2036 async fn read_attempt( 2037 transaction: &mut ServiceSqliteTransaction<'_>, 2038 attempt_id: RhiPresenceAttemptId, 2039 ) -> Result<Option<PresenceAttemptRecord>, PresenceOperationError> { 2040 let rows = sqlx::query(READ_ATTEMPT_SQL) 2041 .bind(attempt_id.as_bytes().as_slice()) 2042 .fetch_all(&mut *transaction) 2043 .await 2044 .map_err(|_| PresenceOperationError::Storage)?; 2045 exactly_zero_or_one(rows)?.map(decode_attempt).transpose() 2046 } 2047 2048 fn decode_attempt( 2049 row: sqlx::sqlite::SqliteRow, 2050 ) -> Result<PresenceAttemptRecord, PresenceOperationError> { 2051 let attempt_id = RhiPresenceAttemptId(blob::<32>(&row, "attempt_id")?); 2052 let outbox_id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); 2053 let target_ordinal = bounded_u8(&row, "target_ordinal", RHI_PRESENCE_DESIRED_MAX_TARGETS - 1)?; 2054 let attempt_number = bounded_u16(&row, "attempt_number", RHI_PRESENCE_MAX_ATTEMPTS)?; 2055 let event_sha256 = blob::<32>(&row, "event_sha256")?; 2056 let owner = RhiPresenceLeaseOwner(blob::<16>(&row, "lease_owner")?); 2057 let started_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "started_at_unix_ms")?); 2058 let finished_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "finished_at_unix_ms")?); 2059 let outcome = decode_attempt_outcome(&bounded_text(&row, "outcome", "outcome_bytes", 13)?)?; 2060 let result_code = bounded_text(&row, "result_code", "result_code_bytes", 64)?; 2061 if owner.0.iter().all(|byte| *byte == 0) 2062 || started_at > finished_at 2063 || result_code != outcome.code() 2064 { 2065 return Err(PresenceOperationError::Invariant); 2066 } 2067 Ok(PresenceAttemptRecord { 2068 attempt_id, 2069 outbox_id, 2070 target_ordinal, 2071 attempt_number, 2072 event_sha256, 2073 owner, 2074 started_at, 2075 finished_at, 2076 outcome, 2077 }) 2078 } 2079 2080 async fn reconcile_recorded( 2081 transaction: &mut ServiceSqliteTransaction<'_>, 2082 input: PresenceRecordInput, 2083 existing: PresenceAttemptRecord, 2084 ) -> Result<RhiPresenceAttemptCommit, PresenceOperationError> { 2085 if existing != PresenceAttemptRecord::from_input(input) { 2086 return Err(PresenceOperationError::Invariant); 2087 } 2088 let outbox = read_outbox(transaction, input.outbox_id) 2089 .await? 2090 .ok_or(PresenceOperationError::Invariant)?; 2091 if outbox.revision 2092 != input 2093 .outbox_revision 2094 .checked_add(1) 2095 .ok_or(PresenceOperationError::Invariant)? 2096 || outbox.updated_at != input.finished_at 2097 || outbox.lease_owner.is_some() 2098 || outbox.lease_expires.is_some() 2099 || outbox.event_sha256 != input.event_sha256 2100 { 2101 return Err(PresenceOperationError::Invariant); 2102 } 2103 validate_targets(&outbox, &outbox.targets)?; 2104 let (expected_state, expected_next) = target_outcome_schedule(input)?; 2105 let expected_target_revision = input 2106 .target_revision 2107 .checked_add(1) 2108 .ok_or(PresenceOperationError::Invariant)?; 2109 let target = outbox 2110 .targets 2111 .iter() 2112 .find(|target| target.ordinal == input.target_ordinal) 2113 .filter(|target| { 2114 target.revision == expected_target_revision 2115 && target.attempt_count == input.attempt_number 2116 && target.last_attempt_id == Some(input.attempt_id) 2117 && target.state == expected_state 2118 && target.next_attempt == expected_next 2119 && target.updated_at == input.finished_at 2120 }) 2121 .ok_or(PresenceOperationError::Invariant)?; 2122 let expected_disposition = disposition(&outbox, &outbox.targets)?; 2123 if outbox.state != expected_disposition.state 2124 || outbox.next_attempt != expected_disposition.next_attempt 2125 { 2126 return Err(PresenceOperationError::Invariant); 2127 } 2128 Ok(RhiPresenceAttemptCommit { 2129 outbox_id: outbox.id, 2130 attempt_id: input.attempt_id, 2131 target_ordinal: input.target_ordinal, 2132 attempt_number: input.attempt_number, 2133 outcome: input.outcome, 2134 target_state: target.state, 2135 outbox_state: outbox.state, 2136 }) 2137 } 2138 2139 fn disposition( 2140 outbox: &PresenceOutboxRecord, 2141 targets: &[PresenceTargetRecord], 2142 ) -> Result<PresenceOutboxDisposition, PresenceOperationError> { 2143 validate_targets(outbox, targets)?; 2144 let required = targets.iter().filter(|target| target.required); 2145 if required 2146 .clone() 2147 .all(|target| target.state == RhiPresenceTargetState::Accepted) 2148 { 2149 return Ok(PresenceOutboxDisposition { 2150 state: RhiPresenceOutboxState::Complete, 2151 next_attempt: None, 2152 }); 2153 } 2154 if required 2155 .clone() 2156 .any(|target| target_is_blocking(target, outbox.max_attempts)) 2157 { 2158 return Ok(PresenceOutboxDisposition { 2159 state: RhiPresenceOutboxState::Blocked, 2160 next_attempt: None, 2161 }); 2162 } 2163 let next_attempt = required 2164 .filter_map(|target| target.next_attempt) 2165 .min() 2166 .ok_or(PresenceOperationError::Invariant)?; 2167 Ok(PresenceOutboxDisposition { 2168 state: RhiPresenceOutboxState::Pending, 2169 next_attempt: Some(next_attempt), 2170 }) 2171 } 2172 2173 fn validate_targets( 2174 outbox: &PresenceOutboxRecord, 2175 targets: &[PresenceTargetRecord], 2176 ) -> Result<(), PresenceOperationError> { 2177 if targets.len() != usize::from(outbox.target_count) 2178 || targets.is_empty() 2179 || targets.len() > RHI_PRESENCE_DESIRED_MAX_TARGETS 2180 || targets.iter().filter(|target| target.required).count() 2181 != usize::from(outbox.required_target_count) 2182 || targets 2183 .iter() 2184 .enumerate() 2185 .any(|(ordinal, target)| usize::from(target.ordinal) != ordinal) 2186 || targets.iter().enumerate().any(|(index, target)| { 2187 targets[index + 1..] 2188 .iter() 2189 .any(|later| later.relay_id == target.relay_id) 2190 }) 2191 { 2192 return Err(PresenceOperationError::Invariant); 2193 } 2194 let mut digest = Sha256::new(); 2195 digest.update(b"radroots.rhi.presence_target_set.v1\0"); 2196 digest.update( 2197 u32::try_from(targets.len()) 2198 .map_err(|_| PresenceOperationError::Invariant)? 2199 .to_be_bytes(), 2200 ); 2201 for target in targets { 2202 digest.update(u32::from(target.ordinal).to_be_bytes()); 2203 digest.update( 2204 u64::try_from(target.relay_id.len()) 2205 .map_err(|_| PresenceOperationError::Invariant)? 2206 .to_be_bytes(), 2207 ); 2208 digest.update(target.relay_id.as_bytes()); 2209 digest.update([u8::from(target.required)]); 2210 } 2211 if <[u8; 32]>::from(digest.finalize()) != outbox.target_set_sha256 { 2212 return Err(PresenceOperationError::Invariant); 2213 } 2214 Ok(()) 2215 } 2216 2217 fn target_outcome_schedule( 2218 input: PresenceRecordInput, 2219 ) -> Result<(RhiPresenceTargetState, Option<RhiPresenceUnixMilliseconds>), PresenceOperationError> { 2220 let state = target_state_for_outcome(input.outcome)?; 2221 let next = 2222 if retryable_outcome(input.outcome) && input.attempt_number < RHI_PRESENCE_MAX_ATTEMPTS { 2223 if input.retry_delay.get() > retry_upper_bound(input.attempt_number) { 2224 return Err(PresenceOperationError::InvalidInput); 2225 } 2226 Some(RhiPresenceUnixMilliseconds( 2227 input 2228 .finished_at 2229 .get() 2230 .checked_add(input.retry_delay.get()) 2231 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 2232 .ok_or(PresenceOperationError::InvalidInput)?, 2233 )) 2234 } else { 2235 if input.retry_delay.get() != 0 { 2236 return Err(PresenceOperationError::InvalidInput); 2237 } 2238 None 2239 }; 2240 Ok((state, next)) 2241 } 2242 2243 fn target_is_blocking(target: &PresenceTargetRecord, max_attempts: u16) -> bool { 2244 matches!( 2245 target.state, 2246 RhiPresenceTargetState::Rejected | RhiPresenceTargetState::AuthRequired 2247 ) || (retryable_state(target.state) 2248 && target.attempt_count >= max_attempts 2249 && target.next_attempt.is_none()) 2250 } 2251 2252 fn target_is_due( 2253 target: &PresenceTargetRecord, 2254 max_attempts: u16, 2255 now: RhiPresenceUnixMilliseconds, 2256 ) -> bool { 2257 retryable_state(target.state) 2258 && target.attempt_count < max_attempts 2259 && target.next_attempt.is_some_and(|next| next <= now) 2260 } 2261 2262 const fn retryable_state(state: RhiPresenceTargetState) -> bool { 2263 matches!( 2264 state, 2265 RhiPresenceTargetState::Pending 2266 | RhiPresenceTargetState::Failed 2267 | RhiPresenceTargetState::RateLimited 2268 | RhiPresenceTargetState::Unknown 2269 ) 2270 } 2271 2272 const fn retryable_outcome(outcome: RhiPresenceAttemptOutcome) -> bool { 2273 matches!( 2274 outcome, 2275 RhiPresenceAttemptOutcome::RateLimited 2276 | RhiPresenceAttemptOutcome::Failed 2277 | RhiPresenceAttemptOutcome::Unknown 2278 ) 2279 } 2280 2281 const fn target_state_for_outcome( 2282 outcome: RhiPresenceAttemptOutcome, 2283 ) -> Result<RhiPresenceTargetState, PresenceOperationError> { 2284 match outcome { 2285 RhiPresenceAttemptOutcome::Submitted => Err(PresenceOperationError::InvalidInput), 2286 RhiPresenceAttemptOutcome::Accepted => Ok(RhiPresenceTargetState::Accepted), 2287 RhiPresenceAttemptOutcome::Rejected => Ok(RhiPresenceTargetState::Rejected), 2288 RhiPresenceAttemptOutcome::RateLimited => Ok(RhiPresenceTargetState::RateLimited), 2289 RhiPresenceAttemptOutcome::AuthRequired => Ok(RhiPresenceTargetState::AuthRequired), 2290 RhiPresenceAttemptOutcome::Failed => Ok(RhiPresenceTargetState::Failed), 2291 RhiPresenceAttemptOutcome::Unknown => Ok(RhiPresenceTargetState::Unknown), 2292 } 2293 } 2294 2295 fn retry_upper_bound(attempt_number: u16) -> u64 { 2296 let mut bound = INITIAL_BACKOFF_MILLISECONDS; 2297 for _ in 1..attempt_number { 2298 bound = bound.saturating_mul(2).min(MAXIMUM_BACKOFF_MILLISECONDS); 2299 } 2300 bound.min(MAXIMUM_BACKOFF_MILLISECONDS) 2301 } 2302 2303 fn same_outbox(current: &PresenceOutboxRecord, prior: &PresenceOutboxRecord) -> bool { 2304 same_outbox_identity(current, prior) 2305 && current.state == prior.state 2306 && current.revision == prior.revision 2307 && current.next_attempt == prior.next_attempt 2308 && current.lease_owner == prior.lease_owner 2309 && current.lease_expires == prior.lease_expires 2310 && current.updated_at == prior.updated_at 2311 && current.targets == prior.targets 2312 } 2313 2314 fn same_outbox_identity(current: &PresenceOutboxRecord, prior: &PresenceOutboxRecord) -> bool { 2315 current.id == prior.id 2316 && current.desired_generation == prior.desired_generation 2317 && current.document_kind == prior.document_kind 2318 && current.desired_sha256 == prior.desired_sha256 2319 && current.target_set_sha256 == prior.target_set_sha256 2320 && current.event_id == prior.event_id 2321 && current.event_sha256 == prior.event_sha256 2322 && current.exact_signed_event_bytes == prior.exact_signed_event_bytes 2323 && current.authored_at_unix_s == prior.authored_at_unix_s 2324 && current.service_public_key == prior.service_public_key 2325 && current.target_count == prior.target_count 2326 && current.required_target_count == prior.required_target_count 2327 && current.max_attempts == prior.max_attempts 2328 && current.created_at == prior.created_at 2329 } 2330 2331 fn valid_outbox_shape( 2332 state: RhiPresenceOutboxState, 2333 next_attempt: Option<RhiPresenceUnixMilliseconds>, 2334 lease_owner: Option<RhiPresenceLeaseOwner>, 2335 lease_expires: Option<RhiPresenceUnixMilliseconds>, 2336 ) -> bool { 2337 match state { 2338 RhiPresenceOutboxState::Pending => { 2339 next_attempt.is_some() && lease_owner.is_none() && lease_expires.is_none() 2340 } 2341 RhiPresenceOutboxState::Leased => { 2342 next_attempt.is_none() && lease_owner.is_some() && lease_expires.is_some() 2343 } 2344 RhiPresenceOutboxState::Complete 2345 | RhiPresenceOutboxState::Blocked 2346 | RhiPresenceOutboxState::Superseded => { 2347 next_attempt.is_none() && lease_owner.is_none() && lease_expires.is_none() 2348 } 2349 } 2350 } 2351 2352 fn valid_target_shape( 2353 state: RhiPresenceTargetState, 2354 attempt_count: u16, 2355 next_attempt: Option<RhiPresenceUnixMilliseconds>, 2356 last_attempt_id: Option<RhiPresenceAttemptId>, 2357 ) -> bool { 2358 match state { 2359 RhiPresenceTargetState::Pending => { 2360 attempt_count == 0 && next_attempt.is_some() && last_attempt_id.is_none() 2361 } 2362 RhiPresenceTargetState::Submitted => { 2363 attempt_count > 0 && next_attempt.is_none() && last_attempt_id.is_some() 2364 } 2365 RhiPresenceTargetState::Accepted 2366 | RhiPresenceTargetState::Rejected 2367 | RhiPresenceTargetState::AuthRequired => { 2368 attempt_count > 0 && next_attempt.is_none() && last_attempt_id.is_some() 2369 } 2370 RhiPresenceTargetState::RateLimited 2371 | RhiPresenceTargetState::Failed 2372 | RhiPresenceTargetState::Unknown => { 2373 attempt_count > 0 2374 && last_attempt_id.is_some() 2375 && if attempt_count < RHI_PRESENCE_MAX_ATTEMPTS { 2376 next_attempt.is_some() 2377 } else { 2378 next_attempt.is_none() 2379 } 2380 } 2381 } 2382 } 2383 2384 fn decode_document_kind(value: &str) -> Result<RhiPresenceDocumentKind, PresenceOperationError> { 2385 match value { 2386 "service_profile" => Ok(RhiPresenceDocumentKind::ServiceProfile), 2387 "application_handler" => Ok(RhiPresenceDocumentKind::ApplicationHandler), 2388 _ => Err(PresenceOperationError::Invariant), 2389 } 2390 } 2391 2392 fn decode_outbox_state(value: &str) -> Result<RhiPresenceOutboxState, PresenceOperationError> { 2393 match value { 2394 "pending" => Ok(RhiPresenceOutboxState::Pending), 2395 "leased" => Ok(RhiPresenceOutboxState::Leased), 2396 "complete" => Ok(RhiPresenceOutboxState::Complete), 2397 "blocked" => Ok(RhiPresenceOutboxState::Blocked), 2398 "superseded" => Ok(RhiPresenceOutboxState::Superseded), 2399 _ => Err(PresenceOperationError::Invariant), 2400 } 2401 } 2402 2403 fn decode_target_state(value: &str) -> Result<RhiPresenceTargetState, PresenceOperationError> { 2404 match value { 2405 "pending" => Ok(RhiPresenceTargetState::Pending), 2406 "submitted" => Ok(RhiPresenceTargetState::Submitted), 2407 "accepted" => Ok(RhiPresenceTargetState::Accepted), 2408 "rejected" => Ok(RhiPresenceTargetState::Rejected), 2409 "rate_limited" => Ok(RhiPresenceTargetState::RateLimited), 2410 "auth_required" => Ok(RhiPresenceTargetState::AuthRequired), 2411 "failed" => Ok(RhiPresenceTargetState::Failed), 2412 "unknown" => Ok(RhiPresenceTargetState::Unknown), 2413 _ => Err(PresenceOperationError::Invariant), 2414 } 2415 } 2416 2417 fn decode_attempt_outcome( 2418 value: &str, 2419 ) -> Result<RhiPresenceAttemptOutcome, PresenceOperationError> { 2420 match value { 2421 "submitted" => Ok(RhiPresenceAttemptOutcome::Submitted), 2422 "accepted" => Ok(RhiPresenceAttemptOutcome::Accepted), 2423 "rejected" => Ok(RhiPresenceAttemptOutcome::Rejected), 2424 "rate_limited" => Ok(RhiPresenceAttemptOutcome::RateLimited), 2425 "auth_required" => Ok(RhiPresenceAttemptOutcome::AuthRequired), 2426 "failed" => Ok(RhiPresenceAttemptOutcome::Failed), 2427 "unknown" => Ok(RhiPresenceAttemptOutcome::Unknown), 2428 _ => Err(PresenceOperationError::Invariant), 2429 } 2430 } 2431 2432 fn exactly_zero_or_one( 2433 rows: Vec<sqlx::sqlite::SqliteRow>, 2434 ) -> Result<Option<sqlx::sqlite::SqliteRow>, PresenceOperationError> { 2435 match rows.as_slice() { 2436 [] => Ok(None), 2437 [_] => Ok(rows.into_iter().next()), 2438 _ => Err(PresenceOperationError::Invariant), 2439 } 2440 } 2441 2442 fn blob<const N: usize>( 2443 row: &sqlx::sqlite::SqliteRow, 2444 column: &str, 2445 ) -> Result<[u8; N], PresenceOperationError> { 2446 row.try_get::<Option<Vec<u8>>, _>(column) 2447 .map_err(|_| PresenceOperationError::Invariant)? 2448 .ok_or(PresenceOperationError::Invariant)? 2449 .try_into() 2450 .map_err(|_| PresenceOperationError::Invariant) 2451 } 2452 2453 fn bounded_text( 2454 row: &sqlx::sqlite::SqliteRow, 2455 value_column: &str, 2456 length_column: &str, 2457 maximum: usize, 2458 ) -> Result<String, PresenceOperationError> { 2459 let length = row 2460 .try_get::<i64, _>(length_column) 2461 .ok() 2462 .and_then(|value| usize::try_from(value).ok()) 2463 .filter(|value| *value > 0 && *value <= maximum) 2464 .ok_or(PresenceOperationError::Invariant)?; 2465 let value = row 2466 .try_get::<String, _>(value_column) 2467 .map_err(|_| PresenceOperationError::Invariant)?; 2468 if value.len() != length { 2469 return Err(PresenceOperationError::Invariant); 2470 } 2471 Ok(value) 2472 } 2473 2474 fn bounded_u8( 2475 row: &sqlx::sqlite::SqliteRow, 2476 column: &str, 2477 maximum: usize, 2478 ) -> Result<u8, PresenceOperationError> { 2479 row.try_get::<i64, _>(column) 2480 .ok() 2481 .and_then(|value| u8::try_from(value).ok()) 2482 .filter(|value| usize::from(*value) <= maximum) 2483 .ok_or(PresenceOperationError::Invariant) 2484 } 2485 2486 fn bounded_u16( 2487 row: &sqlx::sqlite::SqliteRow, 2488 column: &str, 2489 maximum: u16, 2490 ) -> Result<u16, PresenceOperationError> { 2491 row.try_get::<i64, _>(column) 2492 .ok() 2493 .and_then(|value| u16::try_from(value).ok()) 2494 .filter(|value| *value <= maximum) 2495 .ok_or(PresenceOperationError::Invariant) 2496 } 2497 2498 fn bounded_u64( 2499 row: &sqlx::sqlite::SqliteRow, 2500 column: &str, 2501 maximum: u64, 2502 ) -> Result<u64, PresenceOperationError> { 2503 row.try_get::<i64, _>(column) 2504 .ok() 2505 .and_then(|value| u64::try_from(value).ok()) 2506 .filter(|value| *value <= maximum) 2507 .ok_or(PresenceOperationError::Invariant) 2508 } 2509 2510 fn positive_u64( 2511 row: &sqlx::sqlite::SqliteRow, 2512 column: &str, 2513 ) -> Result<u64, PresenceOperationError> { 2514 nonnegative_u64(row, column).and_then(|value| { 2515 (value != 0) 2516 .then_some(value) 2517 .ok_or(PresenceOperationError::Invariant) 2518 }) 2519 } 2520 2521 fn nonnegative_u64( 2522 row: &sqlx::sqlite::SqliteRow, 2523 column: &str, 2524 ) -> Result<u64, PresenceOperationError> { 2525 row.try_get::<i64, _>(column) 2526 .ok() 2527 .and_then(|value| u64::try_from(value).ok()) 2528 .ok_or(PresenceOperationError::Invariant) 2529 } 2530 2531 fn positive_usize( 2532 row: &sqlx::sqlite::SqliteRow, 2533 column: &str, 2534 ) -> Result<usize, PresenceOperationError> { 2535 row.try_get::<i64, _>(column) 2536 .ok() 2537 .and_then(|value| usize::try_from(value).ok()) 2538 .filter(|value| *value != 0) 2539 .ok_or(PresenceOperationError::Invariant) 2540 } 2541 2542 fn boolean_i64( 2543 row: &sqlx::sqlite::SqliteRow, 2544 column: &str, 2545 ) -> Result<bool, PresenceOperationError> { 2546 match row.try_get::<i64, _>(column) { 2547 Ok(0) => Ok(false), 2548 Ok(1) => Ok(true), 2549 Ok(_) | Err(_) => Err(PresenceOperationError::Invariant), 2550 } 2551 } 2552 2553 fn optional_millis( 2554 row: &sqlx::sqlite::SqliteRow, 2555 column: &str, 2556 ) -> Result<Option<RhiPresenceUnixMilliseconds>, PresenceOperationError> { 2557 row.try_get::<Option<i64>, _>(column) 2558 .map_err(|_| PresenceOperationError::Invariant)? 2559 .map(|value| { 2560 u64::try_from(value) 2561 .map(RhiPresenceUnixMilliseconds) 2562 .map_err(|_| PresenceOperationError::Invariant) 2563 }) 2564 .transpose() 2565 } 2566 2567 fn optional_owner( 2568 row: &sqlx::sqlite::SqliteRow, 2569 ) -> Result<Option<RhiPresenceLeaseOwner>, PresenceOperationError> { 2570 let length = row 2571 .try_get::<Option<i64>, _>("lease_owner_bytes") 2572 .map_err(|_| PresenceOperationError::Invariant)?; 2573 let value = row 2574 .try_get::<Option<Vec<u8>>, _>("lease_owner") 2575 .map_err(|_| PresenceOperationError::Invariant)?; 2576 match (length, value) { 2577 (None, None) => Ok(None), 2578 (Some(16), Some(value)) => { 2579 let value: [u8; 16] = value 2580 .try_into() 2581 .map_err(|_| PresenceOperationError::Invariant)?; 2582 if value.iter().all(|byte| *byte == 0) { 2583 return Err(PresenceOperationError::Invariant); 2584 } 2585 Ok(Some(RhiPresenceLeaseOwner(value))) 2586 } 2587 _ => Err(PresenceOperationError::Invariant), 2588 } 2589 } 2590 2591 fn optional_attempt_id( 2592 row: &sqlx::sqlite::SqliteRow, 2593 ) -> Result<Option<RhiPresenceAttemptId>, PresenceOperationError> { 2594 let length = row 2595 .try_get::<Option<i64>, _>("last_attempt_id_bytes") 2596 .map_err(|_| PresenceOperationError::Invariant)?; 2597 let value = row 2598 .try_get::<Option<Vec<u8>>, _>("last_attempt_id") 2599 .map_err(|_| PresenceOperationError::Invariant)?; 2600 match (length, value) { 2601 (None, None) => Ok(None), 2602 (Some(32), Some(value)) => Ok(Some(RhiPresenceAttemptId( 2603 value 2604 .try_into() 2605 .map_err(|_| PresenceOperationError::Invariant)?, 2606 ))), 2607 _ => Err(PresenceOperationError::Invariant), 2608 } 2609 } 2610 2611 fn valid_public_key(value: &str) -> bool { 2612 value.len() == 64 2613 && value 2614 .bytes() 2615 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) 2616 && NostrPublicKey::from_hex(value).is_ok() 2617 } 2618 2619 fn valid_relay_id(value: &str) -> bool { 2620 !value.is_empty() 2621 && value.len() <= 64 2622 && value.as_bytes()[0].is_ascii_lowercase() 2623 && value.bytes().all(|byte| { 2624 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-') 2625 }) 2626 } 2627 2628 fn i64_value(value: u64) -> Result<i64, PresenceOperationError> { 2629 i64::try_from(value).map_err(|_| PresenceOperationError::InvalidInput) 2630 } 2631 2632 fn optional_i64( 2633 value: Option<RhiPresenceUnixMilliseconds>, 2634 ) -> Result<Option<i64>, PresenceOperationError> { 2635 value.map(|value| i64_value(value.get())).transpose() 2636 } 2637 2638 #[derive(Clone, Copy)] 2639 struct PresenceDesiredBinding { 2640 generation: u64, 2641 mode: RhiPresenceDesiredMode, 2642 profile: bool, 2643 application_handler: bool, 2644 target_set_sha256: [u8; 32], 2645 target_count: u8, 2646 required_target_count: u8, 2647 queue_capacity: u32, 2648 desired_sha256: [u8; 32], 2649 } 2650 2651 impl PresenceDesiredBinding { 2652 const fn document_count(self) -> usize { 2653 self.profile as usize + self.application_handler as usize 2654 } 2655 2656 fn document_kinds(self) -> impl Iterator<Item = RhiPresenceDocumentKind> { 2657 [ 2658 self.profile 2659 .then_some(RhiPresenceDocumentKind::ServiceProfile), 2660 self.application_handler 2661 .then_some(RhiPresenceDocumentKind::ApplicationHandler), 2662 ] 2663 .into_iter() 2664 .flatten() 2665 } 2666 } 2667 2668 #[derive(PartialEq, Eq)] 2669 struct PresenceOutboxRecord { 2670 id: RhiPresenceOutboxId, 2671 desired_generation: u64, 2672 document_kind: RhiPresenceDocumentKind, 2673 desired_sha256: [u8; 32], 2674 target_set_sha256: [u8; 32], 2675 event_id: [u8; 32], 2676 event_sha256: [u8; 32], 2677 exact_signed_event_bytes: Box<[u8]>, 2678 authored_at_unix_s: u64, 2679 service_public_key: Box<str>, 2680 target_count: u8, 2681 required_target_count: u8, 2682 max_attempts: u16, 2683 state: RhiPresenceOutboxState, 2684 revision: u64, 2685 next_attempt: Option<RhiPresenceUnixMilliseconds>, 2686 lease_owner: Option<RhiPresenceLeaseOwner>, 2687 lease_expires: Option<RhiPresenceUnixMilliseconds>, 2688 created_at: RhiPresenceUnixMilliseconds, 2689 updated_at: RhiPresenceUnixMilliseconds, 2690 targets: Box<[PresenceTargetRecord]>, 2691 } 2692 2693 #[derive(Clone, PartialEq, Eq)] 2694 struct PresenceTargetRecord { 2695 ordinal: u8, 2696 relay_id: Box<str>, 2697 required: bool, 2698 state: RhiPresenceTargetState, 2699 revision: u64, 2700 attempt_count: u16, 2701 next_attempt: Option<RhiPresenceUnixMilliseconds>, 2702 last_attempt_id: Option<RhiPresenceAttemptId>, 2703 updated_at: RhiPresenceUnixMilliseconds, 2704 } 2705 2706 struct PresenceRecoveryCandidate { 2707 outbox: PresenceOutboxRecord, 2708 retry_upper_bound: u64, 2709 current_desired: bool, 2710 } 2711 2712 #[derive(Clone, Copy)] 2713 struct PresenceRecordInput { 2714 outbox_id: RhiPresenceOutboxId, 2715 outbox_revision: u64, 2716 event_sha256: [u8; 32], 2717 owner: RhiPresenceLeaseOwner, 2718 lease_expires: RhiPresenceUnixMilliseconds, 2719 target_ordinal: u8, 2720 target_revision: u64, 2721 attempt_number: u16, 2722 attempt_id: RhiPresenceAttemptId, 2723 started_at: RhiPresenceUnixMilliseconds, 2724 finished_at: RhiPresenceUnixMilliseconds, 2725 outcome: RhiPresenceAttemptOutcome, 2726 retry_delay: RhiPresenceRetryDelayMilliseconds, 2727 } 2728 2729 impl PresenceRecordInput { 2730 fn from_prepared( 2731 prepared: &RhiPreparedPresenceAttempt, 2732 finished_at: RhiPresenceUnixMilliseconds, 2733 outcome: RhiPresenceAttemptOutcome, 2734 retry_delay: RhiPresenceRetryDelayMilliseconds, 2735 ) -> Result<Self, RhiPresencePublicationError> { 2736 if finished_at < prepared.started_at || outcome == RhiPresenceAttemptOutcome::Submitted { 2737 return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); 2738 } 2739 let expected = derive_attempt_id( 2740 prepared.lease.outbox.id, 2741 prepared.lease.outbox.event_sha256, 2742 prepared.target.ordinal, 2743 prepared.target.attempt_count, 2744 ); 2745 if expected != prepared.attempt_id { 2746 return Err(failure(RhiPresencePublicationErrorKind::Invariant)); 2747 } 2748 Ok(Self { 2749 outbox_id: prepared.lease.outbox.id, 2750 outbox_revision: prepared.lease.outbox.revision, 2751 event_sha256: prepared.lease.outbox.event_sha256, 2752 owner: prepared.lease.owner, 2753 lease_expires: prepared.lease.expires_at, 2754 target_ordinal: prepared.target.ordinal, 2755 target_revision: prepared.target.revision, 2756 attempt_number: prepared.target.attempt_count, 2757 attempt_id: prepared.attempt_id, 2758 started_at: prepared.started_at, 2759 finished_at, 2760 outcome, 2761 retry_delay, 2762 }) 2763 } 2764 } 2765 2766 #[derive(Clone, Copy, PartialEq, Eq)] 2767 struct PresenceAttemptRecord { 2768 attempt_id: RhiPresenceAttemptId, 2769 outbox_id: RhiPresenceOutboxId, 2770 target_ordinal: u8, 2771 attempt_number: u16, 2772 event_sha256: [u8; 32], 2773 owner: RhiPresenceLeaseOwner, 2774 started_at: RhiPresenceUnixMilliseconds, 2775 finished_at: RhiPresenceUnixMilliseconds, 2776 outcome: RhiPresenceAttemptOutcome, 2777 } 2778 2779 impl PresenceAttemptRecord { 2780 const fn from_input(input: PresenceRecordInput) -> Self { 2781 Self { 2782 attempt_id: input.attempt_id, 2783 outbox_id: input.outbox_id, 2784 target_ordinal: input.target_ordinal, 2785 attempt_number: input.attempt_number, 2786 event_sha256: input.event_sha256, 2787 owner: input.owner, 2788 started_at: input.started_at, 2789 finished_at: input.finished_at, 2790 outcome: input.outcome, 2791 } 2792 } 2793 } 2794 2795 #[derive(Clone, Copy)] 2796 struct PresenceOutboxDisposition { 2797 state: RhiPresenceOutboxState, 2798 next_attempt: Option<RhiPresenceUnixMilliseconds>, 2799 } 2800 2801 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 2802 enum PresenceOperationError { 2803 InvalidInput, 2804 DesiredStateMismatch, 2805 VerificationFailed, 2806 NotReady, 2807 LeaseLost, 2808 Invariant, 2809 Storage, 2810 } 2811 2812 fn map_transaction_error( 2813 error: ServiceSqliteTransactionError<PresenceOperationError>, 2814 ) -> RhiPresencePublicationError { 2815 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 2816 return failure(RhiPresencePublicationErrorKind::CommitOutcomeUnknown); 2817 } 2818 failure(match error.operation_error().copied() { 2819 Some(PresenceOperationError::InvalidInput) => RhiPresencePublicationErrorKind::InvalidInput, 2820 Some(PresenceOperationError::DesiredStateMismatch) => { 2821 RhiPresencePublicationErrorKind::DesiredStateMismatch 2822 } 2823 Some(PresenceOperationError::VerificationFailed) => { 2824 RhiPresencePublicationErrorKind::VerificationFailed 2825 } 2826 Some(PresenceOperationError::NotReady) => RhiPresencePublicationErrorKind::NotReady, 2827 Some(PresenceOperationError::LeaseLost) => RhiPresencePublicationErrorKind::LeaseLost, 2828 Some(PresenceOperationError::Invariant) => RhiPresencePublicationErrorKind::Invariant, 2829 Some(PresenceOperationError::Storage) | None => RhiPresencePublicationErrorKind::Storage, 2830 }) 2831 } 2832 2833 fn presence_plan( 2834 kind: RhiPresenceDocumentKind, 2835 author: &str, 2836 created_at: u64, 2837 ) -> Result<AuthoredEventPlan, RhiPresencePublicationError> { 2838 match kind { 2839 RhiPresenceDocumentKind::ServiceProfile => { 2840 let profile = AuthoredProfile::new(PROFILE_NAME) 2841 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))? 2842 .with_display_name(PROFILE_DISPLAY_NAME) 2843 .with_about(PROFILE_ABOUT) 2844 .with_bot(true); 2845 AuthoredEventPlan::from_profile(&profile, created_at, author) 2846 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed)) 2847 } 2848 RhiPresenceDocumentKind::ApplicationHandler => { 2849 let metadata = nostr::Metadata::new() 2850 .name(PROFILE_DISPLAY_NAME) 2851 .about(PROFILE_ABOUT); 2852 let content = serde_json::to_string(&metadata) 2853 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 2854 let spec = ApplicationHandlerSpec::new(APPLICATION_HANDLER_KINDS.to_vec()) 2855 .with_identifier(APPLICATION_HANDLER_IDENTIFIER) 2856 .with_metadata(metadata); 2857 let builder = build_application_handler(&spec) 2858 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))? 2859 .custom_created_at(Timestamp::from_secs(created_at)); 2860 let public_key = NostrPublicKey::from_hex(author) 2861 .map_err(|_| failure(RhiPresencePublicationErrorKind::IdentityMismatch))?; 2862 let request = builder 2863 .into_external_signing_request(public_key) 2864 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 2865 let draft = GenericEventDraft::new( 2866 "radroots.application.handler.v1", 2867 KIND_APPLICATION_HANDLER, 2868 created_at, 2869 vec![ 2870 vec!["d".to_owned(), APPLICATION_HANDLER_IDENTIFIER.to_owned()], 2871 vec!["k".to_owned(), APPLICATION_HANDLER_KINDS[0].to_string()], 2872 ], 2873 content, 2874 author, 2875 ) 2876 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 2877 let plan = AuthoredEventPlan::from_generic(draft) 2878 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 2879 if request.expected_event_id().to_bytes() != *plan.expected_event_id().as_bytes() { 2880 return Err(failure(RhiPresencePublicationErrorKind::RenderingFailed)); 2881 } 2882 Ok(plan) 2883 } 2884 } 2885 } 2886 2887 fn sign_plan( 2888 identity: &RhiDecryptedIdentity, 2889 plan: &AuthoredEventPlan, 2890 auxiliary: &[u8; 32], 2891 ) -> Result<nostr::Event, RhiPresencePublicationError> { 2892 let kind = u16::try_from(plan.body().kind()) 2893 .map(Kind::Custom) 2894 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 2895 let tags = plan 2896 .body() 2897 .tags() 2898 .iter() 2899 .cloned() 2900 .map(Tag::parse) 2901 .collect::<Result<Vec<_>, _>>() 2902 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; 2903 let author = NostrPublicKey::from_hex(identity.public_identity().as_hex()) 2904 .map_err(|_| failure(RhiPresencePublicationErrorKind::IdentityMismatch))?; 2905 let unsigned = EventBuilder::new(kind, plan.body().content()) 2906 .tags(tags) 2907 .custom_created_at(Timestamp::from_secs(plan.created_at())) 2908 .build(author); 2909 if unsigned.id.as_ref().map(|id| id.to_bytes()) != Some(*plan.expected_event_id().as_bytes()) { 2910 return Err(failure(RhiPresencePublicationErrorKind::RenderingFailed)); 2911 } 2912 identity 2913 .sign_nostr_event(unsigned, auxiliary) 2914 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed)) 2915 } 2916 2917 fn validate_signed_document( 2918 kind: RhiPresenceDocumentKind, 2919 author: &str, 2920 created_at: u64, 2921 bytes: &[u8], 2922 ) -> Result<EventEnvelope, RhiPresencePublicationError> { 2923 if bytes.is_empty() || bytes.len() > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES { 2924 return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); 2925 } 2926 let source = core::str::from_utf8(bytes) 2927 .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?; 2928 let wire = Nip01EventWire::parse_json_unverified_with_limits(source, presence_wire_limits()) 2929 .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?; 2930 let event = wire 2931 .into_unverified_envelope() 2932 .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?; 2933 if verify_id(&event) != Verification::IdVerified 2934 || verify(&event) != Verification::Verified 2935 || event.author().to_hex() != author 2936 || event.created_at_u64() != created_at 2937 { 2938 return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); 2939 } 2940 let expected = presence_plan(kind, author, created_at)?; 2941 if event.id().as_bytes() != expected.expected_event_id().as_bytes() 2942 || event.kind_u32() != expected.body().kind() 2943 || event.tags_as_vec() != expected.body().tags() 2944 || event.content() != expected.body().content() 2945 { 2946 return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); 2947 } 2948 Ok(event) 2949 } 2950 2951 const fn presence_wire_limits() -> EventWireLimits { 2952 EventWireLimits { 2953 max_raw_json_bytes: RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES, 2954 max_content_bytes: 4 * 1024, 2955 max_tag_count: 4, 2956 max_total_tag_elements: 8, 2957 max_tag_element_bytes: 64, 2958 max_total_tag_bytes: 256, 2959 max_extra_fields: 0, 2960 max_total_extra_json_bytes: 0, 2961 } 2962 } 2963 2964 const fn failure(kind: RhiPresencePublicationErrorKind) -> RhiPresencePublicationError { 2965 RhiPresencePublicationError { kind } 2966 } 2967 2968 #[cfg(test)] 2969 mod tests { 2970 use std::{collections::BTreeSet, error::Error as _}; 2971 2972 use super::*; 2973 2974 #[test] 2975 fn public_vocabularies_bounds_and_diagnostics_are_closed() { 2976 let kinds = [ 2977 RhiPresencePublicationErrorKind::InvalidMode, 2978 RhiPresencePublicationErrorKind::InvalidInput, 2979 RhiPresencePublicationErrorKind::DesiredStateMismatch, 2980 RhiPresencePublicationErrorKind::IdentityMismatch, 2981 RhiPresencePublicationErrorKind::EntropyUnavailable, 2982 RhiPresencePublicationErrorKind::RenderingFailed, 2983 RhiPresencePublicationErrorKind::VerificationFailed, 2984 RhiPresencePublicationErrorKind::NotReady, 2985 RhiPresencePublicationErrorKind::LeaseLost, 2986 RhiPresencePublicationErrorKind::Invariant, 2987 RhiPresencePublicationErrorKind::ClockUnavailable, 2988 RhiPresencePublicationErrorKind::Storage, 2989 RhiPresencePublicationErrorKind::CommitOutcomeUnknown, 2990 ]; 2991 let codes: BTreeSet<_> = kinds.iter().map(|kind| kind.code()).collect(); 2992 assert_eq!(codes.len(), kinds.len()); 2993 for kind in kinds { 2994 let error = failure(kind); 2995 assert_eq!(error.kind(), kind); 2996 assert_eq!(error.code(), kind.code()); 2997 assert!(error.source().is_none()); 2998 let rendered = format!("{error} {error:?}"); 2999 for forbidden in ["relay", "wss://", "state.sqlite", "010101", "sqlx"] { 3000 assert!(!rendered.contains(forbidden)); 3001 } 3002 } 3003 3004 assert_eq!( 3005 RhiPresenceUnixMilliseconds::new(i64::MAX as u64) 3006 .expect("maximum time") 3007 .get(), 3008 i64::MAX as u64 3009 ); 3010 assert!(RhiPresenceUnixMilliseconds::new(i64::MAX as u64 + 1).is_err()); 3011 assert_eq!( 3012 RhiPresenceRetryDelayMilliseconds::new(MAXIMUM_BACKOFF_MILLISECONDS) 3013 .expect("maximum delay") 3014 .get(), 3015 MAXIMUM_BACKOFF_MILLISECONDS 3016 ); 3017 assert!(RhiPresenceRetryDelayMilliseconds::new(MAXIMUM_BACKOFF_MILLISECONDS + 1).is_err()); 3018 assert!(RhiPresenceLeaseOwner::from_bytes([0; 16]).is_err()); 3019 assert_eq!( 3020 format!("{:?}", RhiPresenceLeaseOwner::from_bytes([1; 16]).unwrap()), 3021 "RhiPresenceLeaseOwner([redacted])" 3022 ); 3023 assert_eq!( 3024 format!("{:?}", RhiPresenceOutboxId([0x5a; 32])), 3025 "RhiPresenceOutboxId([redacted])" 3026 ); 3027 assert_eq!( 3028 format!("{:?}", RhiPresenceAttemptId([0x6a; 32])), 3029 "RhiPresenceAttemptId([redacted])" 3030 ); 3031 } 3032 3033 #[test] 3034 fn state_and_outcome_codes_are_exact() { 3035 assert_eq!( 3036 [ 3037 RhiPresenceOutboxState::Pending, 3038 RhiPresenceOutboxState::Leased, 3039 RhiPresenceOutboxState::Complete, 3040 RhiPresenceOutboxState::Blocked, 3041 RhiPresenceOutboxState::Superseded, 3042 ] 3043 .map(RhiPresenceOutboxState::code), 3044 ["pending", "leased", "complete", "blocked", "superseded"] 3045 ); 3046 assert_eq!( 3047 [ 3048 RhiPresenceTargetState::Pending, 3049 RhiPresenceTargetState::Submitted, 3050 RhiPresenceTargetState::Accepted, 3051 RhiPresenceTargetState::Rejected, 3052 RhiPresenceTargetState::RateLimited, 3053 RhiPresenceTargetState::AuthRequired, 3054 RhiPresenceTargetState::Failed, 3055 RhiPresenceTargetState::Unknown, 3056 ] 3057 .map(RhiPresenceTargetState::code), 3058 [ 3059 "pending", 3060 "submitted", 3061 "accepted", 3062 "rejected", 3063 "rate_limited", 3064 "auth_required", 3065 "failed", 3066 "unknown", 3067 ] 3068 ); 3069 assert_eq!( 3070 [ 3071 RhiPresenceAttemptOutcome::Submitted, 3072 RhiPresenceAttemptOutcome::Accepted, 3073 RhiPresenceAttemptOutcome::Rejected, 3074 RhiPresenceAttemptOutcome::RateLimited, 3075 RhiPresenceAttemptOutcome::AuthRequired, 3076 RhiPresenceAttemptOutcome::Failed, 3077 RhiPresenceAttemptOutcome::Unknown, 3078 ] 3079 .map(RhiPresenceAttemptOutcome::code), 3080 [ 3081 "submitted", 3082 "accepted", 3083 "rejected", 3084 "rate_limited", 3085 "auth_required", 3086 "failed", 3087 "unknown", 3088 ] 3089 ); 3090 } 3091 }