state_connection.rs (102868B)
1 //! Durable typed Myc connection, approval, and authorization-challenge state. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use radroots_service_sqlite::{ 7 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 8 }; 9 use sha2::{Digest, Sha256}; 10 use sqlx::Row; 11 use url::{Host, Url}; 12 13 #[cfg(any(target_os = "linux", target_os = "macos"))] 14 use crate::state_admin::AdminJournalOperationError; 15 use crate::state_governance::{ 16 AuditEvidence, GovernanceOperationError, MycAuditCorrelationId, MycAuditKind, MycAuditOutcome, 17 MycAuditReasonCode, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, connection_subject, 18 global_subject, govern_rate_attempt, record_audit, relay_subject, 19 }; 20 use crate::state_repository::{ 21 MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, 22 RepositoryOperationError, require_expected_metadata, 23 }; 24 use crate::{MycNip46ClientPublicKey, MycSignerOperationId, MycSignerRequestMethod}; 25 26 /// Maximum number of independently granted permissions on one connection. 27 pub const MYC_CONNECTION_PERMISSION_MAX_COUNT: usize = 64; 28 /// Maximum canonical byte length of an operator-owned challenge URL. 29 pub const MYC_AUTHORIZATION_CHALLENGE_URL_MAX_BYTES: usize = 2_048; 30 31 const CONNECTION_ID_DOMAIN: &[u8] = b"radroots.myc.connection.v1\0"; 32 const CHALLENGE_ID_DOMAIN: &[u8] = b"radroots.myc.authorization_challenge.v1\0"; 33 const PERMISSION_SET_DOMAIN: &[u8] = b"radroots.myc.connection_permissions.v1\0"; 34 35 const READ_REQUEST_BINDING_SQL: &str = r#"SELECT 36 CASE WHEN typeof(correlation_id) = 'blob' AND length(correlation_id) = 32 37 THEN correlation_id ELSE NULL END AS correlation_id, 38 CASE WHEN typeof(client_public_key) = 'text' 39 AND length(CAST(client_public_key AS BLOB)) = 64 40 THEN client_public_key ELSE NULL END AS client_public_key, 41 CASE WHEN typeof(method) = 'text' 42 AND length(CAST(method AS BLOB)) BETWEEN 1 AND 32 43 THEN method ELSE NULL END AS method, 44 received_at_unix_ms 45 FROM nip46_requests 46 WHERE operation_id = ? 47 LIMIT 2"#; 48 49 const READ_DECISION_SQL: &str = r#"SELECT 50 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 51 THEN connection_id ELSE NULL END AS connection_id, 52 typeof(connection_id) AS connection_id_type, 53 CASE WHEN typeof(decision) = 'text' AND length(CAST(decision AS BLOB)) <= 32 54 THEN decision ELSE NULL END AS decision, 55 CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 40 56 THEN reason_code ELSE NULL END AS reason_code, 57 policy_generation, 58 CASE WHEN typeof(requested_permissions_sha256) = 'blob' 59 AND length(requested_permissions_sha256) = 32 60 THEN requested_permissions_sha256 ELSE NULL END AS requested_permissions_sha256, 61 CASE WHEN typeof(challenge_id) = 'blob' AND length(challenge_id) = 32 62 THEN challenge_id ELSE NULL END AS challenge_id, 63 typeof(challenge_id) AS challenge_id_type, 64 decided_at_unix_ms 65 FROM nip46_request_decisions 66 WHERE operation_id = ? 67 LIMIT 2"#; 68 69 const INSERT_DECISION_SQL: &str = r#"INSERT INTO nip46_request_decisions ( 70 operation_id, connection_id, decision, reason_code, policy_generation, 71 requested_permissions_sha256, challenge_id, decided_at_unix_ms 72 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)"#; 73 74 const INSERT_CONNECTION_SQL: &str = r#"INSERT INTO connections ( 75 connection_id, connection_nonce, client_public_key, requested_permissions_sha256, 76 policy_generation, status, created_at_unix_ms, updated_at_unix_ms, 77 authorized_until_unix_ms 78 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 79 80 const INSERT_PERMISSION_SQL: &str = r#"INSERT INTO connection_permissions ( 81 connection_id, permission_scope, permission_code 82 ) VALUES (?, ?, ?)"#; 83 84 const READ_CONNECTION_SQL: &str = r#"SELECT 85 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 86 THEN connection_id ELSE NULL END AS connection_id, 87 CASE WHEN typeof(connection_nonce) = 'blob' AND length(connection_nonce) = 32 88 THEN connection_nonce ELSE NULL END AS connection_nonce, 89 CASE WHEN typeof(client_public_key) = 'text' 90 AND length(CAST(client_public_key AS BLOB)) = 64 91 THEN client_public_key ELSE NULL END AS client_public_key, 92 CASE WHEN typeof(requested_permissions_sha256) = 'blob' 93 AND length(requested_permissions_sha256) = 32 94 THEN requested_permissions_sha256 ELSE NULL END AS requested_permissions_sha256, 95 policy_generation, 96 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 97 THEN status ELSE NULL END AS status, 98 created_at_unix_ms, 99 updated_at_unix_ms, 100 authorized_until_unix_ms, 101 typeof(authorized_until_unix_ms) AS authorized_until_type 102 FROM connections 103 WHERE connection_id = ? 104 LIMIT 2"#; 105 106 const READ_ACTIVE_CONNECTION_FOR_CLIENT_SQL: &str = r#"SELECT 107 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 108 THEN connection_id ELSE NULL END AS connection_id 109 FROM connections 110 WHERE client_public_key = ? AND status = 'active' 111 AND (authorized_until_unix_ms IS NULL OR authorized_until_unix_ms >= ?) 112 ORDER BY updated_at_unix_ms DESC, connection_id ASC 113 LIMIT 2"#; 114 115 const READ_RUNTIME_CONNECTION_COUNTS_SQL: &str = r#"SELECT 116 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 117 THEN status ELSE NULL END AS status, 118 COUNT(*) AS row_count 119 FROM connections 120 GROUP BY status 121 LIMIT 5"#; 122 123 const READ_ADMIN_CONNECTION_PAGE_SQL: &str = r#"SELECT 124 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 125 THEN connection_id ELSE NULL END AS connection_id, 126 updated_at_unix_ms 127 FROM connections 128 WHERE updated_at_unix_ms <= ? 129 AND (? IS NULL OR status = ?) 130 AND (? IS NULL OR updated_at_unix_ms < ? 131 OR (updated_at_unix_ms = ? AND connection_id > ?)) 132 ORDER BY updated_at_unix_ms DESC, connection_id ASC 133 LIMIT ?"#; 134 135 const READ_PERMISSIONS_SQL: &str = r#"SELECT 136 CASE WHEN typeof(permission_code) = 'text' 137 AND length(CAST(permission_code AS BLOB)) BETWEEN 1 AND 64 138 THEN permission_code ELSE NULL END AS permission_code 139 FROM connection_permissions 140 WHERE connection_id = ? AND permission_scope = ? 141 LIMIT 65"#; 142 143 const APPROVE_CONNECTION_SQL: &str = r#"UPDATE connections 144 SET status = 'active', updated_at_unix_ms = ?, authorized_until_unix_ms = ? 145 WHERE connection_id = ? AND status = 'pending' AND policy_generation = ?"#; 146 147 const DENY_CONNECTION_SQL: &str = r#"UPDATE connections 148 SET status = 'denied', updated_at_unix_ms = ?, authorized_until_unix_ms = NULL 149 WHERE connection_id = ? AND status = 'pending' AND policy_generation = ?"#; 150 151 const EXPIRE_CONNECTION_SQL: &str = r#"UPDATE connections 152 SET status = 'expired', updated_at_unix_ms = ?, authorized_until_unix_ms = NULL 153 WHERE connection_id = ? AND status = 'active' AND policy_generation = ? 154 AND authorized_until_unix_ms IS NOT NULL AND authorized_until_unix_ms < ?"#; 155 156 const REVOKE_CONNECTION_SQL: &str = r#"UPDATE connections 157 SET status = 'expired', updated_at_unix_ms = ?, authorized_until_unix_ms = NULL 158 WHERE connection_id = ? AND status = 'active' AND policy_generation = ?"#; 159 160 const READ_PENDING_DECISION_FOR_CONNECTION_SQL: &str = r#"SELECT 161 CASE WHEN typeof(operation_id) = 'blob' AND length(operation_id) = 32 162 THEN operation_id ELSE NULL END AS operation_id 163 FROM nip46_request_decisions 164 WHERE connection_id = ? AND decision = 'pending_approval' AND policy_generation = ? 165 ORDER BY decided_at_unix_ms DESC, operation_id ASC 166 LIMIT 2"#; 167 168 const UPDATE_APPROVAL_DECISION_SQL: &str = r#"UPDATE nip46_request_decisions 169 SET decision = ?, reason_code = ?, decided_at_unix_ms = ? 170 WHERE operation_id = ? AND connection_id = ? AND decision = 'pending_approval' 171 AND policy_generation = ?"#; 172 173 const INSERT_CHALLENGE_SQL: &str = r#"INSERT INTO connection_auth_challenges ( 174 challenge_id, challenge_nonce, connection_id, operation_id, policy_generation, 175 challenge_url, state, issued_at_unix_ms, expires_at_unix_ms, resolved_at_unix_ms 176 ) VALUES (?, ?, ?, ?, ?, ?, 'pending', ?, ?, NULL)"#; 177 178 const READ_CHALLENGE_SQL: &str = r#"SELECT 179 CASE WHEN typeof(challenge_id) = 'blob' AND length(challenge_id) = 32 180 THEN challenge_id ELSE NULL END AS challenge_id, 181 CASE WHEN typeof(challenge_nonce) = 'blob' AND length(challenge_nonce) = 32 182 THEN challenge_nonce ELSE NULL END AS challenge_nonce, 183 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 184 THEN connection_id ELSE NULL END AS connection_id, 185 CASE WHEN typeof(operation_id) = 'blob' AND length(operation_id) = 32 186 THEN operation_id ELSE NULL END AS operation_id, 187 policy_generation, 188 CASE WHEN typeof(challenge_url) = 'text' 189 AND length(CAST(challenge_url AS BLOB)) BETWEEN 1 AND 2048 190 THEN challenge_url ELSE NULL END AS challenge_url, 191 CASE WHEN typeof(state) = 'text' AND length(CAST(state AS BLOB)) <= 16 192 THEN state ELSE NULL END AS state, 193 issued_at_unix_ms, expires_at_unix_ms, resolved_at_unix_ms, 194 typeof(resolved_at_unix_ms) AS resolved_at_type 195 FROM connection_auth_challenges 196 WHERE operation_id = ? 197 LIMIT 2"#; 198 199 const READ_CHALLENGE_OPERATION_BY_ID_SQL: &str = r#"SELECT 200 CASE WHEN typeof(operation_id) = 'blob' AND length(operation_id) = 32 201 THEN operation_id ELSE NULL END AS operation_id 202 FROM connection_auth_challenges 203 WHERE challenge_id = ? 204 LIMIT 2"#; 205 206 const RESOLVE_CHALLENGE_SQL: &str = r#"UPDATE connection_auth_challenges 207 SET state = ?, resolved_at_unix_ms = ? 208 WHERE challenge_id = ? AND connection_id = ? AND operation_id = ? 209 AND policy_generation = ? AND state = 'pending'"#; 210 211 const UPDATE_CHALLENGE_DECISION_SQL: &str = r#"UPDATE nip46_request_decisions 212 SET decision = ?, reason_code = ?, decided_at_unix_ms = ? 213 WHERE operation_id = ? AND connection_id = ? AND challenge_id = ? 214 AND policy_generation = ? AND decision = 'challenged'"#; 215 216 /// Stable construction and validation failure classes. 217 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 218 pub enum MycConnectionStateErrorKind { 219 InvalidPermissionSet, 220 InvalidPolicyGeneration, 221 InvalidTime, 222 InvalidChallengeUrl, 223 InvalidChallengeLifetime, 224 } 225 226 impl MycConnectionStateErrorKind { 227 /// Returns the stable machine-readable error code. 228 #[must_use] 229 pub const fn code(self) -> &'static str { 230 match self { 231 Self::InvalidPermissionSet => "connection_permission_set_invalid", 232 Self::InvalidPolicyGeneration => "connection_policy_generation_invalid", 233 Self::InvalidTime => "connection_time_invalid", 234 Self::InvalidChallengeUrl => "authorization_challenge_url_invalid", 235 Self::InvalidChallengeLifetime => "authorization_challenge_lifetime_invalid", 236 } 237 } 238 } 239 240 /// Source-free connection-state validation failure. 241 #[derive(Clone, Copy, PartialEq, Eq)] 242 pub struct MycConnectionStateError { 243 kind: MycConnectionStateErrorKind, 244 } 245 246 impl MycConnectionStateError { 247 const fn new(kind: MycConnectionStateErrorKind) -> Self { 248 Self { kind } 249 } 250 251 /// Returns the stable failure class. 252 #[must_use] 253 pub const fn kind(self) -> MycConnectionStateErrorKind { 254 self.kind 255 } 256 257 /// Returns the stable machine-readable error code. 258 #[must_use] 259 pub const fn code(self) -> &'static str { 260 self.kind.code() 261 } 262 } 263 264 impl fmt::Display for MycConnectionStateError { 265 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 266 formatter.write_str(match self.kind { 267 MycConnectionStateErrorKind::InvalidPermissionSet => { 268 "connection permission set is invalid" 269 } 270 MycConnectionStateErrorKind::InvalidPolicyGeneration => { 271 "connection policy generation is invalid" 272 } 273 MycConnectionStateErrorKind::InvalidTime => "connection time is invalid", 274 MycConnectionStateErrorKind::InvalidChallengeUrl => { 275 "authorization challenge URL is invalid" 276 } 277 MycConnectionStateErrorKind::InvalidChallengeLifetime => { 278 "authorization challenge lifetime is invalid" 279 } 280 }) 281 } 282 } 283 284 impl fmt::Debug for MycConnectionStateError { 285 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 286 formatter 287 .debug_struct("MycConnectionStateError") 288 .field("kind", &self.kind) 289 .finish() 290 } 291 } 292 293 impl Error for MycConnectionStateError {} 294 295 /// One closed Myc connection permission. 296 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 297 pub enum MycConnectionPermission { 298 GetPublicKey, 299 GetSessionCapability, 300 SignEvent(u32), 301 Nip04Encrypt, 302 Nip04Decrypt, 303 Nip44Encrypt, 304 Nip44Decrypt, 305 Ping, 306 SwitchRelays, 307 Logout, 308 } 309 310 impl MycConnectionPermission { 311 pub(crate) fn code(self) -> String { 312 match self { 313 Self::GetPublicKey => "get_public_key".into(), 314 Self::GetSessionCapability => "get_session_capability".into(), 315 Self::SignEvent(kind) => format!("sign_event:kind:{kind}"), 316 Self::Nip04Encrypt => "nip04_encrypt".into(), 317 Self::Nip04Decrypt => "nip04_decrypt".into(), 318 Self::Nip44Encrypt => "nip44_encrypt".into(), 319 Self::Nip44Decrypt => "nip44_decrypt".into(), 320 Self::Ping => "ping".into(), 321 Self::SwitchRelays => "switch_relays".into(), 322 Self::Logout => "logout".into(), 323 } 324 } 325 326 pub(crate) fn parse(value: &str) -> Option<Self> { 327 match value { 328 "get_public_key" => Some(Self::GetPublicKey), 329 "get_session_capability" => Some(Self::GetSessionCapability), 330 "nip04_encrypt" => Some(Self::Nip04Encrypt), 331 "nip04_decrypt" => Some(Self::Nip04Decrypt), 332 "nip44_encrypt" => Some(Self::Nip44Encrypt), 333 "nip44_decrypt" => Some(Self::Nip44Decrypt), 334 "ping" => Some(Self::Ping), 335 "switch_relays" => Some(Self::SwitchRelays), 336 "logout" => Some(Self::Logout), 337 _ => value 338 .strip_prefix("sign_event:kind:") 339 .and_then(|kind| kind.parse::<u32>().ok().map(|parsed| (kind, parsed))) 340 .filter(|(kind, parsed)| *kind == parsed.to_string()) 341 .map(|(_, parsed)| Self::SignEvent(parsed)), 342 } 343 } 344 } 345 346 /// Immutable canonical set of connection permissions. 347 #[derive(Clone, PartialEq, Eq)] 348 pub struct MycConnectionPermissionSet { 349 permissions: Box<[MycConnectionPermission]>, 350 digest: [u8; 32], 351 } 352 353 impl MycConnectionPermissionSet { 354 /// Validates, orders, and seals one bounded permission set. 355 pub fn new(permissions: &[MycConnectionPermission]) -> Result<Self, MycConnectionStateError> { 356 if permissions.len() > MYC_CONNECTION_PERMISSION_MAX_COUNT { 357 return Err(MycConnectionStateError::new( 358 MycConnectionStateErrorKind::InvalidPermissionSet, 359 )); 360 } 361 let mut normalized = permissions.to_vec(); 362 normalized.sort_unstable(); 363 if normalized.windows(2).any(|window| window[0] == window[1]) { 364 return Err(MycConnectionStateError::new( 365 MycConnectionStateErrorKind::InvalidPermissionSet, 366 )); 367 } 368 let digest = permission_digest(&normalized); 369 Ok(Self { 370 permissions: normalized.into_boxed_slice(), 371 digest, 372 }) 373 } 374 375 /// Returns the canonical ordered permissions. 376 #[must_use] 377 pub fn permissions(&self) -> &[MycConnectionPermission] { 378 &self.permissions 379 } 380 381 fn digest(&self) -> &[u8; 32] { 382 &self.digest 383 } 384 385 pub(crate) fn is_subset_of(&self, other: &Self) -> bool { 386 self.permissions 387 .iter() 388 .all(|permission| other.permissions.binary_search(permission).is_ok()) 389 } 390 391 pub(crate) fn contains(&self, permission: MycConnectionPermission) -> bool { 392 self.permissions.binary_search(&permission).is_ok() 393 } 394 } 395 396 impl fmt::Debug for MycConnectionPermissionSet { 397 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 398 formatter 399 .debug_struct("MycConnectionPermissionSet") 400 .field("permission_count", &self.permissions.len()) 401 .finish() 402 } 403 } 404 405 /// Injected entropy used once to derive a durable connection identity. 406 pub struct MycConnectionNonce([u8; 32]); 407 408 impl MycConnectionNonce { 409 /// Wraps caller-injected entropy without reading ambient state. 410 #[must_use] 411 pub const fn from_injected_entropy(bytes: [u8; 32]) -> Self { 412 Self(bytes) 413 } 414 } 415 416 impl fmt::Debug for MycConnectionNonce { 417 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 418 formatter.write_str("MycConnectionNonce([redacted])") 419 } 420 } 421 422 /// Injected entropy used once to derive a durable challenge identity. 423 pub struct MycAuthorizationChallengeNonce([u8; 32]); 424 425 impl MycAuthorizationChallengeNonce { 426 /// Wraps caller-injected entropy without reading ambient state. 427 #[must_use] 428 pub const fn from_injected_entropy(bytes: [u8; 32]) -> Self { 429 Self(bytes) 430 } 431 } 432 433 impl fmt::Debug for MycAuthorizationChallengeNonce { 434 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 435 formatter.write_str("MycAuthorizationChallengeNonce([redacted])") 436 } 437 } 438 439 macro_rules! digest_id { 440 ($name:ident, $debug:literal) => { 441 #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] 442 pub struct $name([u8; 32]); 443 444 impl $name { 445 pub(crate) const fn from_bytes(bytes: [u8; 32]) -> Self { 446 Self(bytes) 447 } 448 449 /// Returns the exact stable identity bytes. 450 #[must_use] 451 pub const fn as_bytes(&self) -> &[u8; 32] { 452 &self.0 453 } 454 } 455 456 impl fmt::Debug for $name { 457 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 458 formatter.write_str($debug) 459 } 460 } 461 }; 462 } 463 464 digest_id!(MycConnectionId, "MycConnectionId([redacted])"); 465 digest_id!( 466 MycAuthorizationChallengeId, 467 "MycAuthorizationChallengeId([redacted])" 468 ); 469 470 /// Nonzero immutable connection-policy generation. 471 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 472 pub struct MycConnectionPolicyGeneration(u64); 473 474 impl MycConnectionPolicyGeneration { 475 /// Constructs a policy generation representable by SQLite. 476 pub fn new(value: u64) -> Result<Self, MycConnectionStateError> { 477 if value == 0 || i64::try_from(value).is_err() { 478 return Err(MycConnectionStateError::new( 479 MycConnectionStateErrorKind::InvalidPolicyGeneration, 480 )); 481 } 482 Ok(Self(value)) 483 } 484 485 /// Returns the generation value. 486 #[must_use] 487 pub const fn get(self) -> u64 { 488 self.0 489 } 490 491 fn sqlite_value(self) -> i64 { 492 i64::try_from(self.0).expect("validated generation fits SQLite") 493 } 494 } 495 496 /// Nonzero immutable millisecond UTC evidence supplied by the caller. 497 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 498 pub struct MycConnectionTimeUnixMs(u64); 499 500 impl MycConnectionTimeUnixMs { 501 /// Constructs a timestamp representable by SQLite. 502 pub fn new(value: u64) -> Result<Self, MycConnectionStateError> { 503 if value == 0 || i64::try_from(value).is_err() { 504 return Err(MycConnectionStateError::new( 505 MycConnectionStateErrorKind::InvalidTime, 506 )); 507 } 508 Ok(Self(value)) 509 } 510 511 /// Returns the injected timestamp. 512 #[must_use] 513 pub const fn get(self) -> u64 { 514 self.0 515 } 516 517 fn sqlite_value(self) -> i64 { 518 i64::try_from(self.0).expect("validated time fits SQLite") 519 } 520 } 521 522 /// Closed admission authority for a connect request. 523 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 524 pub enum MycConnectionAdmissionPolicy { 525 Trusted, 526 ExplicitApproval, 527 Denied, 528 } 529 530 impl MycConnectionAdmissionPolicy { 531 const fn decision(self) -> MycConnectionDecision { 532 match self { 533 Self::Trusted => MycConnectionDecision::Allowed, 534 Self::ExplicitApproval => MycConnectionDecision::PendingApproval, 535 Self::Denied => MycConnectionDecision::Denied, 536 } 537 } 538 } 539 540 /// Stable durable request decision. 541 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 542 pub enum MycConnectionDecision { 543 PendingApproval, 544 Challenged, 545 Allowed, 546 Denied, 547 } 548 549 impl MycConnectionDecision { 550 const fn as_str(self) -> &'static str { 551 match self { 552 Self::PendingApproval => "pending_approval", 553 Self::Challenged => "challenged", 554 Self::Allowed => "allowed", 555 Self::Denied => "denied", 556 } 557 } 558 559 const fn reason(self, policy: MycConnectionAdmissionPolicy) -> &'static str { 560 match (self, policy) { 561 (Self::PendingApproval, _) => "explicit_approval_required", 562 (Self::Allowed, _) => "trusted_client", 563 (Self::Denied, _) => "policy_denied", 564 (Self::Challenged, _) => "authorization_challenge_required", 565 } 566 } 567 568 fn parse(value: &str) -> Option<Self> { 569 match value { 570 "pending_approval" => Some(Self::PendingApproval), 571 "challenged" => Some(Self::Challenged), 572 "allowed" => Some(Self::Allowed), 573 "denied" => Some(Self::Denied), 574 _ => None, 575 } 576 } 577 } 578 579 /// Immutable connect-admission input. 580 pub struct MycConnectionAdmissionRequest { 581 operation_id: MycSignerOperationId, 582 client_public_key: MycNip46ClientPublicKey, 583 requested_permissions: MycConnectionPermissionSet, 584 policy_generation: MycConnectionPolicyGeneration, 585 nonce: MycConnectionNonce, 586 observed_at: MycConnectionTimeUnixMs, 587 authorized_until: Option<MycConnectionTimeUnixMs>, 588 policy: MycConnectionAdmissionPolicy, 589 relay_id: MycRateRelayId, 590 } 591 592 impl MycConnectionAdmissionRequest { 593 /// Constructs a request from validated protocol, policy, entropy, and time evidence. 594 #[allow(clippy::too_many_arguments)] 595 pub fn new( 596 operation_id: MycSignerOperationId, 597 client_public_key: MycNip46ClientPublicKey, 598 requested_permissions: MycConnectionPermissionSet, 599 policy_generation: MycConnectionPolicyGeneration, 600 nonce: MycConnectionNonce, 601 observed_at: MycConnectionTimeUnixMs, 602 authorized_until: Option<MycConnectionTimeUnixMs>, 603 policy: MycConnectionAdmissionPolicy, 604 relay_id: MycRateRelayId, 605 ) -> Result<Self, MycConnectionStateError> { 606 if (matches!(policy, MycConnectionAdmissionPolicy::Trusted) 607 && authorized_until.is_some_and(|until| until <= observed_at)) 608 || (!matches!(policy, MycConnectionAdmissionPolicy::Trusted) 609 && authorized_until.is_some()) 610 { 611 return Err(MycConnectionStateError::new( 612 MycConnectionStateErrorKind::InvalidTime, 613 )); 614 } 615 Ok(Self { 616 operation_id, 617 client_public_key, 618 requested_permissions, 619 policy_generation, 620 nonce, 621 observed_at, 622 authorized_until, 623 policy, 624 relay_id, 625 }) 626 } 627 628 fn owned(&self) -> Self { 629 Self { 630 operation_id: self.operation_id, 631 client_public_key: self.client_public_key.clone(), 632 requested_permissions: self.requested_permissions.clone(), 633 policy_generation: self.policy_generation, 634 nonce: MycConnectionNonce(self.nonce.0), 635 observed_at: self.observed_at, 636 authorized_until: self.authorized_until, 637 policy: self.policy, 638 relay_id: self.relay_id.clone(), 639 } 640 } 641 642 pub(crate) const fn client_public_key(&self) -> &MycNip46ClientPublicKey { 643 &self.client_public_key 644 } 645 646 pub(crate) const fn requested_permissions(&self) -> &MycConnectionPermissionSet { 647 &self.requested_permissions 648 } 649 650 pub(crate) const fn observed_at(&self) -> MycConnectionTimeUnixMs { 651 self.observed_at 652 } 653 654 pub(crate) const fn authorized_until(&self) -> Option<MycConnectionTimeUnixMs> { 655 self.authorized_until 656 } 657 658 pub(crate) const fn policy(&self) -> MycConnectionAdmissionPolicy { 659 self.policy 660 } 661 } 662 663 impl fmt::Debug for MycConnectionAdmissionRequest { 664 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 665 formatter.write_str("MycConnectionAdmissionRequest([redacted])") 666 } 667 } 668 669 /// Durable connection status. 670 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 671 pub enum MycConnectionStatus { 672 Pending, 673 Active, 674 Denied, 675 Expired, 676 } 677 678 impl MycConnectionStatus { 679 const fn as_str(self) -> &'static str { 680 match self { 681 Self::Pending => "pending", 682 Self::Active => "active", 683 Self::Denied => "denied", 684 Self::Expired => "expired", 685 } 686 } 687 688 fn parse(value: &str) -> Option<Self> { 689 match value { 690 "pending" => Some(Self::Pending), 691 "active" => Some(Self::Active), 692 "denied" => Some(Self::Denied), 693 "expired" => Some(Self::Expired), 694 _ => None, 695 } 696 } 697 698 pub(crate) const fn admin_state(self) -> &'static str { 699 match self { 700 Self::Pending => "pending", 701 Self::Active => "approved", 702 Self::Denied => "rejected", 703 Self::Expired => "revoked", 704 } 705 } 706 } 707 708 /// Validated durable connection record. 709 #[derive(Clone, PartialEq, Eq)] 710 pub struct MycConnectionRecord { 711 id: MycConnectionId, 712 client_public_key: MycNip46ClientPublicKey, 713 requested_permissions: MycConnectionPermissionSet, 714 granted_permissions: MycConnectionPermissionSet, 715 policy_generation: MycConnectionPolicyGeneration, 716 status: MycConnectionStatus, 717 created_at: MycConnectionTimeUnixMs, 718 updated_at: MycConnectionTimeUnixMs, 719 authorized_until: Option<MycConnectionTimeUnixMs>, 720 } 721 722 pub(crate) struct MycAdminConnectionPage { 723 items: Box<[MycConnectionRecord]>, 724 next: Option<(MycConnectionTimeUnixMs, MycConnectionId)>, 725 } 726 727 impl MycAdminConnectionPage { 728 pub(crate) fn items(&self) -> &[MycConnectionRecord] { 729 &self.items 730 } 731 732 pub(crate) const fn next(&self) -> Option<(MycConnectionTimeUnixMs, MycConnectionId)> { 733 self.next 734 } 735 } 736 737 impl MycConnectionRecord { 738 #[must_use] 739 pub const fn id(&self) -> MycConnectionId { 740 self.id 741 } 742 #[must_use] 743 pub const fn client_public_key(&self) -> &MycNip46ClientPublicKey { 744 &self.client_public_key 745 } 746 #[must_use] 747 pub const fn requested_permissions(&self) -> &MycConnectionPermissionSet { 748 &self.requested_permissions 749 } 750 #[must_use] 751 pub const fn granted_permissions(&self) -> &MycConnectionPermissionSet { 752 &self.granted_permissions 753 } 754 #[must_use] 755 pub const fn policy_generation(&self) -> MycConnectionPolicyGeneration { 756 self.policy_generation 757 } 758 #[must_use] 759 pub const fn status(&self) -> MycConnectionStatus { 760 self.status 761 } 762 #[must_use] 763 pub const fn created_at(&self) -> MycConnectionTimeUnixMs { 764 self.created_at 765 } 766 #[must_use] 767 pub const fn updated_at(&self) -> MycConnectionTimeUnixMs { 768 self.updated_at 769 } 770 #[must_use] 771 pub const fn authorized_until(&self) -> Option<MycConnectionTimeUnixMs> { 772 self.authorized_until 773 } 774 775 pub(crate) fn admin_permissions(&self) -> String { 776 self.granted_permissions 777 .permissions() 778 .iter() 779 .map(|permission| permission.code()) 780 .collect::<Vec<_>>() 781 .join(",") 782 } 783 784 #[cfg(test)] 785 pub(crate) fn active_for_test( 786 client_public_key: MycNip46ClientPublicKey, 787 granted_permissions: MycConnectionPermissionSet, 788 authorized_until: Option<MycConnectionTimeUnixMs>, 789 ) -> Self { 790 Self { 791 id: MycConnectionId([0x45; 32]), 792 client_public_key, 793 requested_permissions: granted_permissions.clone(), 794 granted_permissions, 795 policy_generation: MycConnectionPolicyGeneration(1), 796 status: MycConnectionStatus::Active, 797 created_at: MycConnectionTimeUnixMs(1), 798 updated_at: MycConnectionTimeUnixMs(1), 799 authorized_until, 800 } 801 } 802 } 803 804 impl fmt::Debug for MycConnectionRecord { 805 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 806 formatter 807 .debug_struct("MycConnectionRecord") 808 .field("status", &self.status) 809 .field( 810 "permission_count", 811 &self.granted_permissions.permissions.len(), 812 ) 813 .finish() 814 } 815 } 816 817 /// Durable result for initial connection admission. 818 #[derive(Clone, PartialEq, Eq)] 819 pub struct MycConnectionDecisionRecord { 820 operation_id: MycSignerOperationId, 821 connection: Option<MycConnectionRecord>, 822 decision: MycConnectionDecision, 823 policy_generation: MycConnectionPolicyGeneration, 824 decided_at: MycConnectionTimeUnixMs, 825 } 826 827 impl MycConnectionDecisionRecord { 828 #[must_use] 829 pub const fn operation_id(&self) -> MycSignerOperationId { 830 self.operation_id 831 } 832 #[must_use] 833 pub const fn connection(&self) -> Option<&MycConnectionRecord> { 834 self.connection.as_ref() 835 } 836 #[must_use] 837 pub const fn decision(&self) -> MycConnectionDecision { 838 self.decision 839 } 840 #[must_use] 841 pub const fn policy_generation(&self) -> MycConnectionPolicyGeneration { 842 self.policy_generation 843 } 844 #[must_use] 845 pub const fn decided_at(&self) -> MycConnectionTimeUnixMs { 846 self.decided_at 847 } 848 } 849 850 impl fmt::Debug for MycConnectionDecisionRecord { 851 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 852 formatter 853 .debug_struct("MycConnectionDecisionRecord") 854 .field("decision", &self.decision) 855 .finish() 856 } 857 } 858 859 /// New or replayed connection-admission result. 860 #[derive(Clone, PartialEq, Eq)] 861 pub enum MycConnectionAdmission { 862 Admitted(MycConnectionDecisionRecord), 863 ExactReplay(MycConnectionDecisionRecord), 864 RateLimited, 865 } 866 867 impl MycConnectionAdmission { 868 #[must_use] 869 pub const fn record(&self) -> Option<&MycConnectionDecisionRecord> { 870 match self { 871 Self::Admitted(record) | Self::ExactReplay(record) => Some(record), 872 Self::RateLimited => None, 873 } 874 } 875 } 876 877 impl fmt::Debug for MycConnectionAdmission { 878 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 879 formatter.write_str(match self { 880 Self::Admitted(_) => "MycConnectionAdmission::Admitted([redacted])", 881 Self::ExactReplay(_) => "MycConnectionAdmission::ExactReplay([redacted])", 882 Self::RateLimited => "MycConnectionAdmission::RateLimited", 883 }) 884 } 885 } 886 887 /// Operator decision for a pending explicit-approval connection. 888 pub enum MycConnectionOperatorDecision { 889 Approve { 890 granted_permissions: MycConnectionPermissionSet, 891 authorized_until: Option<MycConnectionTimeUnixMs>, 892 }, 893 Deny, 894 } 895 896 #[cfg(any(target_os = "linux", target_os = "macos"))] 897 pub(crate) enum MycAdminConnectionAction { 898 Approve { 899 permissions: MycConnectionPermissionSet, 900 authorized_until: Option<MycConnectionTimeUnixMs>, 901 }, 902 Reject, 903 Revoke, 904 } 905 906 impl fmt::Debug for MycConnectionOperatorDecision { 907 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 908 formatter.write_str(match self { 909 Self::Approve { .. } => "MycConnectionOperatorDecision::Approve([redacted])", 910 Self::Deny => "MycConnectionOperatorDecision::Deny", 911 }) 912 } 913 } 914 915 /// Validated operator-owned challenge URL. 916 #[derive(Clone, PartialEq, Eq)] 917 pub struct MycAuthorizationChallengeUrl(Box<str>); 918 919 impl MycAuthorizationChallengeUrl { 920 /// Admits canonical HTTPS URLs and loopback-only HTTP URLs. 921 pub fn new(value: &str) -> Result<Self, MycConnectionStateError> { 922 if value.is_empty() || value.len() > MYC_AUTHORIZATION_CHALLENGE_URL_MAX_BYTES { 923 return Err(MycConnectionStateError::new( 924 MycConnectionStateErrorKind::InvalidChallengeUrl, 925 )); 926 } 927 let parsed = Url::parse(value).map_err(|_| { 928 MycConnectionStateError::new(MycConnectionStateErrorKind::InvalidChallengeUrl) 929 })?; 930 let loopback_http = parsed.scheme() == "http" 931 && match parsed.host() { 932 Some(Host::Domain("localhost")) => true, 933 Some(Host::Ipv4(address)) => address.is_loopback(), 934 Some(Host::Ipv6(address)) => address.is_loopback(), 935 _ => false, 936 }; 937 if (parsed.scheme() != "https" && !loopback_http) 938 || !parsed.username().is_empty() 939 || parsed.password().is_some() 940 || parsed.fragment().is_some() 941 || parsed.as_str() != value 942 { 943 return Err(MycConnectionStateError::new( 944 MycConnectionStateErrorKind::InvalidChallengeUrl, 945 )); 946 } 947 Ok(Self(value.into())) 948 } 949 950 #[must_use] 951 pub fn as_str(&self) -> &str { 952 &self.0 953 } 954 } 955 956 impl fmt::Debug for MycAuthorizationChallengeUrl { 957 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 958 formatter.write_str("MycAuthorizationChallengeUrl([redacted])") 959 } 960 } 961 962 /// Immutable authorization-challenge creation input. 963 pub struct MycAuthorizationChallengeRequest { 964 operation_id: MycSignerOperationId, 965 connection_id: MycConnectionId, 966 policy_generation: MycConnectionPolicyGeneration, 967 url: MycAuthorizationChallengeUrl, 968 nonce: MycAuthorizationChallengeNonce, 969 issued_at: MycConnectionTimeUnixMs, 970 expires_at: MycConnectionTimeUnixMs, 971 } 972 973 impl MycAuthorizationChallengeRequest { 974 #[allow(clippy::too_many_arguments)] 975 pub fn new( 976 operation_id: MycSignerOperationId, 977 connection_id: MycConnectionId, 978 policy_generation: MycConnectionPolicyGeneration, 979 url: MycAuthorizationChallengeUrl, 980 nonce: MycAuthorizationChallengeNonce, 981 issued_at: MycConnectionTimeUnixMs, 982 expires_at: MycConnectionTimeUnixMs, 983 ) -> Result<Self, MycConnectionStateError> { 984 if expires_at <= issued_at { 985 return Err(MycConnectionStateError::new( 986 MycConnectionStateErrorKind::InvalidChallengeLifetime, 987 )); 988 } 989 Ok(Self { 990 operation_id, 991 connection_id, 992 policy_generation, 993 url, 994 nonce, 995 issued_at, 996 expires_at, 997 }) 998 } 999 1000 fn owned(&self) -> Self { 1001 Self { 1002 operation_id: self.operation_id, 1003 connection_id: self.connection_id, 1004 policy_generation: self.policy_generation, 1005 url: self.url.clone(), 1006 nonce: MycAuthorizationChallengeNonce(self.nonce.0), 1007 issued_at: self.issued_at, 1008 expires_at: self.expires_at, 1009 } 1010 } 1011 1012 pub(crate) const fn url(&self) -> &MycAuthorizationChallengeUrl { 1013 &self.url 1014 } 1015 1016 pub(crate) const fn issued_at(&self) -> MycConnectionTimeUnixMs { 1017 self.issued_at 1018 } 1019 1020 pub(crate) const fn expires_at(&self) -> MycConnectionTimeUnixMs { 1021 self.expires_at 1022 } 1023 } 1024 1025 impl fmt::Debug for MycAuthorizationChallengeRequest { 1026 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1027 formatter.write_str("MycAuthorizationChallengeRequest([redacted])") 1028 } 1029 } 1030 1031 /// Durable authorization-challenge state. 1032 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1033 pub enum MycAuthorizationChallengeState { 1034 Pending, 1035 Authorized, 1036 Expired, 1037 } 1038 1039 impl MycAuthorizationChallengeState { 1040 const fn as_str(self) -> &'static str { 1041 match self { 1042 Self::Pending => "pending", 1043 Self::Authorized => "authorized", 1044 Self::Expired => "expired", 1045 } 1046 } 1047 1048 fn parse(value: &str) -> Option<Self> { 1049 match value { 1050 "pending" => Some(Self::Pending), 1051 "authorized" => Some(Self::Authorized), 1052 "expired" => Some(Self::Expired), 1053 _ => None, 1054 } 1055 } 1056 } 1057 1058 /// Validated durable challenge record. 1059 #[derive(Clone, PartialEq, Eq)] 1060 pub struct MycAuthorizationChallengeRecord { 1061 id: MycAuthorizationChallengeId, 1062 operation_id: MycSignerOperationId, 1063 connection_id: MycConnectionId, 1064 policy_generation: MycConnectionPolicyGeneration, 1065 url: MycAuthorizationChallengeUrl, 1066 state: MycAuthorizationChallengeState, 1067 issued_at: MycConnectionTimeUnixMs, 1068 expires_at: MycConnectionTimeUnixMs, 1069 resolved_at: Option<MycConnectionTimeUnixMs>, 1070 } 1071 1072 impl MycAuthorizationChallengeRecord { 1073 #[must_use] 1074 pub const fn id(&self) -> MycAuthorizationChallengeId { 1075 self.id 1076 } 1077 #[must_use] 1078 pub const fn operation_id(&self) -> MycSignerOperationId { 1079 self.operation_id 1080 } 1081 #[must_use] 1082 pub const fn connection_id(&self) -> MycConnectionId { 1083 self.connection_id 1084 } 1085 #[must_use] 1086 pub const fn policy_generation(&self) -> MycConnectionPolicyGeneration { 1087 self.policy_generation 1088 } 1089 #[must_use] 1090 pub const fn url(&self) -> &MycAuthorizationChallengeUrl { 1091 &self.url 1092 } 1093 #[must_use] 1094 pub const fn state(&self) -> MycAuthorizationChallengeState { 1095 self.state 1096 } 1097 #[must_use] 1098 pub const fn issued_at(&self) -> MycConnectionTimeUnixMs { 1099 self.issued_at 1100 } 1101 #[must_use] 1102 pub const fn expires_at(&self) -> MycConnectionTimeUnixMs { 1103 self.expires_at 1104 } 1105 #[must_use] 1106 pub const fn resolved_at(&self) -> Option<MycConnectionTimeUnixMs> { 1107 self.resolved_at 1108 } 1109 } 1110 1111 impl fmt::Debug for MycAuthorizationChallengeRecord { 1112 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1113 formatter 1114 .debug_struct("MycAuthorizationChallengeRecord") 1115 .field("state", &self.state) 1116 .finish() 1117 } 1118 } 1119 1120 /// New or replayed challenge creation result. 1121 #[derive(Clone, PartialEq, Eq)] 1122 pub enum MycAuthorizationChallengeAdmission { 1123 Created(MycAuthorizationChallengeRecord), 1124 ExactReplay(MycAuthorizationChallengeRecord), 1125 RateLimited, 1126 } 1127 1128 impl MycAuthorizationChallengeAdmission { 1129 #[must_use] 1130 pub const fn record(&self) -> Option<&MycAuthorizationChallengeRecord> { 1131 match self { 1132 Self::Created(record) | Self::ExactReplay(record) => Some(record), 1133 Self::RateLimited => None, 1134 } 1135 } 1136 } 1137 1138 impl fmt::Debug for MycAuthorizationChallengeAdmission { 1139 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1140 formatter.write_str(match self { 1141 Self::Created(_) => "MycAuthorizationChallengeAdmission::Created([redacted])", 1142 Self::ExactReplay(_) => "MycAuthorizationChallengeAdmission::ExactReplay([redacted])", 1143 Self::RateLimited => "MycAuthorizationChallengeAdmission::RateLimited", 1144 }) 1145 } 1146 } 1147 1148 /// New, replayed, or rate-limited authorization result. 1149 #[derive(Clone, PartialEq, Eq)] 1150 pub enum MycAuthorizationChallengeAuthorization { 1151 Resolved(MycAuthorizationChallengeRecord), 1152 ExactReplay(MycAuthorizationChallengeRecord), 1153 RateLimited, 1154 } 1155 1156 impl MycAuthorizationChallengeAuthorization { 1157 #[must_use] 1158 pub const fn record(&self) -> Option<&MycAuthorizationChallengeRecord> { 1159 match self { 1160 Self::Resolved(record) | Self::ExactReplay(record) => Some(record), 1161 Self::RateLimited => None, 1162 } 1163 } 1164 } 1165 1166 impl fmt::Debug for MycAuthorizationChallengeAuthorization { 1167 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1168 formatter.write_str(match self { 1169 Self::Resolved(_) => "MycAuthorizationChallengeAuthorization::Resolved([redacted])", 1170 Self::ExactReplay(_) => { 1171 "MycAuthorizationChallengeAuthorization::ExactReplay([redacted])" 1172 } 1173 Self::RateLimited => "MycAuthorizationChallengeAuthorization::RateLimited", 1174 }) 1175 } 1176 } 1177 1178 impl MycStateRepository<'_> { 1179 pub(crate) async fn read_runtime_connection_counts( 1180 &self, 1181 ) -> Result<crate::MycConnectionCountsV1, MycStateRepositoryError> { 1182 let expected = PersistedMetadata::from(self.expected()); 1183 self.host() 1184 .transaction(move |transaction| { 1185 Box::pin(async move { 1186 verify_metadata(transaction, &expected).await?; 1187 let rows = sqlx::query(READ_RUNTIME_CONNECTION_COUNTS_SQL) 1188 .fetch_all(&mut *transaction) 1189 .await 1190 .map_err(|_| ConnectionOperationError::Storage)?; 1191 if rows.len() > 4 { 1192 return Err(ConnectionOperationError::Binding); 1193 } 1194 let mut counts = [None; 4]; 1195 for row in rows { 1196 let status = bounded_text(&row, "status")?; 1197 let index = match MycConnectionStatus::parse(status) { 1198 Some(MycConnectionStatus::Pending) => 0, 1199 Some(MycConnectionStatus::Active) => 1, 1200 Some(MycConnectionStatus::Denied) => 2, 1201 Some(MycConnectionStatus::Expired) => 3, 1202 None => return Err(ConnectionOperationError::Binding), 1203 }; 1204 let value = row 1205 .try_get::<i64, _>("row_count") 1206 .ok() 1207 .and_then(|value| u64::try_from(value).ok()) 1208 .ok_or(ConnectionOperationError::Binding)?; 1209 if counts[index].replace(value).is_some() { 1210 return Err(ConnectionOperationError::Binding); 1211 } 1212 } 1213 Ok(crate::MycConnectionCountsV1::new( 1214 counts[0].unwrap_or(0), 1215 counts[1].unwrap_or(0), 1216 counts[2].unwrap_or(0), 1217 counts[3].unwrap_or(0), 1218 )) 1219 }) 1220 }) 1221 .await 1222 .map_err(map_transaction_error) 1223 } 1224 1225 pub(crate) async fn read_admin_connection_page( 1226 &self, 1227 limit: u16, 1228 status: Option<MycConnectionStatus>, 1229 snapshot: MycConnectionTimeUnixMs, 1230 before: Option<(MycConnectionTimeUnixMs, MycConnectionId)>, 1231 ) -> Result<MycAdminConnectionPage, MycStateRepositoryError> { 1232 if limit == 0 || limit > 200 { 1233 return Err(MycStateRepositoryError::new( 1234 MycStateRepositoryErrorKind::Binding, 1235 )); 1236 } 1237 let expected = PersistedMetadata::from(self.expected()); 1238 self.host() 1239 .transaction(move |transaction| { 1240 Box::pin(async move { 1241 verify_metadata(transaction, &expected).await?; 1242 read_admin_connection_page(transaction, limit, status, snapshot, before).await 1243 }) 1244 }) 1245 .await 1246 .map_err(map_transaction_error) 1247 } 1248 1249 /// Resolves the sole active session for one verified transport client. 1250 /// 1251 /// Multiple simultaneously active sessions are ambiguous because NIP-46 1252 /// requests carry no server connection identifier; fail closed rather 1253 /// than silently selecting a different authorization grant. 1254 pub(crate) async fn read_active_connection_for_client( 1255 &self, 1256 client: &MycNip46ClientPublicKey, 1257 observed_at: MycConnectionTimeUnixMs, 1258 ) -> Result<Option<MycConnectionRecord>, MycStateRepositoryError> { 1259 let expected = PersistedMetadata::from(self.expected()); 1260 let client = client.clone(); 1261 self.host() 1262 .transaction(move |transaction| { 1263 Box::pin(async move { 1264 verify_metadata(transaction, &expected).await?; 1265 let rows = sqlx::query(READ_ACTIVE_CONNECTION_FOR_CLIENT_SQL) 1266 .bind(client.as_hex()) 1267 .bind(observed_at.sqlite_value()) 1268 .fetch_all(&mut *transaction) 1269 .await 1270 .map_err(|_| ConnectionOperationError::Storage)?; 1271 match rows.as_slice() { 1272 [] => Ok(None), 1273 [row] => { 1274 let id = MycConnectionId(exact_digest(row, "connection_id")?); 1275 read_connection(transaction, id).await.map(Some) 1276 } 1277 _ => Err(ConnectionOperationError::Binding), 1278 } 1279 }) 1280 }) 1281 .await 1282 .map_err(map_transaction_error) 1283 } 1284 1285 /// Reads one fully validated durable connection decision by stable operation identity. 1286 pub async fn read_connection_decision( 1287 &self, 1288 operation_id: MycSignerOperationId, 1289 ) -> Result<MycConnectionDecisionRecord, MycStateRepositoryError> { 1290 let expected = PersistedMetadata::from(self.expected()); 1291 self.host() 1292 .transaction(move |transaction| { 1293 Box::pin(async move { 1294 verify_metadata(transaction, &expected).await?; 1295 let decision = read_decision(transaction, operation_id) 1296 .await? 1297 .ok_or(ConnectionOperationError::Binding)?; 1298 decision_record(transaction, operation_id, decision).await 1299 }) 1300 }) 1301 .await 1302 .map_err(map_transaction_error) 1303 } 1304 1305 /// Atomically records trusted, explicit-approval, or direct-denial admission. 1306 pub async fn admit_connection( 1307 &self, 1308 request: &MycConnectionAdmissionRequest, 1309 ) -> Result<MycConnectionAdmission, MycStateRepositoryError> { 1310 if !self.expected().admits_rate_relay(&request.relay_id) 1311 || !self.expected().admits_connection_request(request) 1312 { 1313 return Err(MycStateRepositoryError::new( 1314 MycStateRepositoryErrorKind::Binding, 1315 )); 1316 } 1317 let request = request.owned(); 1318 let rate_policy = (request.policy != MycConnectionAdmissionPolicy::Denied).then(|| { 1319 self.expected() 1320 .governance_rate_policy(MycRateLimitClass::ConnectionAdmission) 1321 }); 1322 let expected = PersistedMetadata::from(self.expected()); 1323 self.host() 1324 .transaction(move |transaction| { 1325 Box::pin(async move { 1326 verify_metadata(transaction, &expected).await?; 1327 admit_connection(transaction, &request, rate_policy).await 1328 }) 1329 }) 1330 .await 1331 .map_err(map_transaction_error) 1332 } 1333 1334 /// Atomically approves or denies one pending connection. 1335 pub async fn decide_pending_connection( 1336 &self, 1337 operation_id: MycSignerOperationId, 1338 connection_id: MycConnectionId, 1339 policy_generation: MycConnectionPolicyGeneration, 1340 observed_at: MycConnectionTimeUnixMs, 1341 audit_correlation: MycAuditCorrelationId, 1342 decision: MycConnectionOperatorDecision, 1343 ) -> Result<MycConnectionRecord, MycStateRepositoryError> { 1344 if !self 1345 .expected() 1346 .admits_connection_operator_decision(observed_at, &decision) 1347 { 1348 return Err(MycStateRepositoryError::new( 1349 MycStateRepositoryErrorKind::Binding, 1350 )); 1351 } 1352 let expected = PersistedMetadata::from(self.expected()); 1353 self.host() 1354 .transaction(move |transaction| { 1355 Box::pin(async move { 1356 verify_metadata(transaction, &expected).await?; 1357 decide_connection( 1358 transaction, 1359 operation_id, 1360 connection_id, 1361 policy_generation, 1362 observed_at, 1363 audit_correlation, 1364 decision, 1365 ) 1366 .await 1367 }) 1368 }) 1369 .await 1370 .map_err(map_transaction_error) 1371 } 1372 1373 /// Idempotently marks an expired authorized connection. 1374 pub async fn expire_connection( 1375 &self, 1376 connection_id: MycConnectionId, 1377 policy_generation: MycConnectionPolicyGeneration, 1378 observed_at: MycConnectionTimeUnixMs, 1379 audit_correlation: MycAuditCorrelationId, 1380 ) -> Result<MycConnectionRecord, MycStateRepositoryError> { 1381 let expected = PersistedMetadata::from(self.expected()); 1382 self.host() 1383 .transaction(move |transaction| { 1384 Box::pin(async move { 1385 verify_metadata(transaction, &expected).await?; 1386 let before = read_connection(transaction, connection_id).await?; 1387 if before.policy_generation != policy_generation { 1388 return Err(ConnectionOperationError::Binding); 1389 } 1390 if before.status == MycConnectionStatus::Expired { 1391 record_connection_expiry_audit( 1392 transaction, 1393 audit_correlation, 1394 before.updated_at, 1395 ) 1396 .await?; 1397 return Ok(before); 1398 } 1399 if before.status != MycConnectionStatus::Active 1400 || before 1401 .authorized_until 1402 .is_none_or(|until| until >= observed_at) 1403 { 1404 return Err(ConnectionOperationError::Binding); 1405 } 1406 let result = sqlx::query(EXPIRE_CONNECTION_SQL) 1407 .bind(observed_at.sqlite_value()) 1408 .bind(connection_id.as_bytes().as_slice()) 1409 .bind(policy_generation.sqlite_value()) 1410 .bind(observed_at.sqlite_value()) 1411 .execute(&mut *transaction) 1412 .await 1413 .map_err(|_| ConnectionOperationError::Storage)?; 1414 require_one(result.rows_affected())?; 1415 let record = read_connection(transaction, connection_id).await?; 1416 record_connection_expiry_audit(transaction, audit_correlation, observed_at) 1417 .await?; 1418 Ok(record) 1419 }) 1420 }) 1421 .await 1422 .map_err(map_transaction_error) 1423 } 1424 1425 /// Creates or replays one request-bound authorization challenge. 1426 pub async fn issue_authorization_challenge( 1427 &self, 1428 request: &MycAuthorizationChallengeRequest, 1429 ) -> Result<MycAuthorizationChallengeAdmission, MycStateRepositoryError> { 1430 if !self 1431 .expected() 1432 .admits_authorization_challenge_request(request) 1433 { 1434 return Err(MycStateRepositoryError::new( 1435 MycStateRepositoryErrorKind::Binding, 1436 )); 1437 } 1438 let request = request.owned(); 1439 let rate_policy = self 1440 .expected() 1441 .governance_rate_policy(MycRateLimitClass::ChallengeCreation); 1442 let expected = PersistedMetadata::from(self.expected()); 1443 self.host() 1444 .transaction(move |transaction| { 1445 Box::pin(async move { 1446 verify_metadata(transaction, &expected).await?; 1447 issue_challenge(transaction, &request, rate_policy).await 1448 }) 1449 }) 1450 .await 1451 .map_err(map_transaction_error) 1452 } 1453 1454 /// Authorizes or expires an exact bound challenge, replaying terminal state safely. 1455 pub async fn authorize_challenge( 1456 &self, 1457 challenge_id: MycAuthorizationChallengeId, 1458 connection_id: MycConnectionId, 1459 operation_id: MycSignerOperationId, 1460 policy_generation: MycConnectionPolicyGeneration, 1461 observed_at: MycConnectionTimeUnixMs, 1462 ) -> Result<MycAuthorizationChallengeAuthorization, MycStateRepositoryError> { 1463 let rate_policy = self 1464 .expected() 1465 .governance_rate_policy(MycRateLimitClass::ChallengeAuthorization); 1466 let expected = PersistedMetadata::from(self.expected()); 1467 let result = self 1468 .host() 1469 .transaction(move |transaction| { 1470 Box::pin(async move { 1471 verify_metadata(transaction, &expected).await?; 1472 authorize_challenge( 1473 transaction, 1474 challenge_id, 1475 connection_id, 1476 operation_id, 1477 policy_generation, 1478 observed_at, 1479 rate_policy, 1480 ) 1481 .await 1482 }) 1483 }) 1484 .await 1485 .map_err(map_transaction_error)?; 1486 if result.record().is_some_and(|record| { 1487 record.state() == MycAuthorizationChallengeState::Authorized 1488 && !self 1489 .expected() 1490 .authorization_challenge_is_current(record, observed_at) 1491 }) { 1492 return Err(MycStateRepositoryError::new( 1493 MycStateRepositoryErrorKind::Binding, 1494 )); 1495 } 1496 Ok(result) 1497 } 1498 } 1499 1500 async fn read_admin_connection_page( 1501 transaction: &mut ServiceSqliteTransaction<'_>, 1502 limit: u16, 1503 status: Option<MycConnectionStatus>, 1504 snapshot: MycConnectionTimeUnixMs, 1505 before: Option<(MycConnectionTimeUnixMs, MycConnectionId)>, 1506 ) -> Result<MycAdminConnectionPage, ConnectionOperationError> { 1507 let status = status.map(MycConnectionStatus::as_str); 1508 let before_time = before.map(|(time, _)| time.sqlite_value()); 1509 let before_id = before.map(|(_, id)| id.as_bytes().to_vec()); 1510 let fetch_limit = i64::from(limit) + 1; 1511 let rows = sqlx::query(READ_ADMIN_CONNECTION_PAGE_SQL) 1512 .bind(snapshot.sqlite_value()) 1513 .bind(status) 1514 .bind(status) 1515 .bind(before_time) 1516 .bind(before_time) 1517 .bind(before_time) 1518 .bind(before_id) 1519 .bind(fetch_limit) 1520 .fetch_all(&mut *transaction) 1521 .await 1522 .map_err(|_| ConnectionOperationError::Storage)?; 1523 if rows.len() > usize::try_from(fetch_limit).map_err(|_| ConnectionOperationError::Binding)? { 1524 return Err(ConnectionOperationError::Binding); 1525 } 1526 let has_more = rows.len() > usize::from(limit); 1527 let mut items = Vec::with_capacity(rows.len().min(usize::from(limit))); 1528 for row in rows.iter().take(usize::from(limit)) { 1529 let id = MycConnectionId(exact_digest(row, "connection_id")?); 1530 items.push(read_connection(transaction, id).await?); 1531 } 1532 let next = has_more 1533 .then(|| items.last().map(|item| (item.updated_at, item.id))) 1534 .flatten(); 1535 Ok(MycAdminConnectionPage { 1536 items: items.into_boxed_slice(), 1537 next, 1538 }) 1539 } 1540 1541 async fn record_connection_expiry_audit( 1542 transaction: &mut ServiceSqliteTransaction<'_>, 1543 correlation: MycAuditCorrelationId, 1544 occurred_at: MycConnectionTimeUnixMs, 1545 ) -> Result<(), ConnectionOperationError> { 1546 record_audit( 1547 transaction, 1548 AuditEvidence { 1549 correlation, 1550 occurred_at, 1551 operation_id: None, 1552 }, 1553 MycAuditKind::ConnectionExpiry, 1554 MycAuditOutcome::Succeeded, 1555 MycAuditReasonCode::ConnectionExpired, 1556 ) 1557 .await 1558 .map_err(map_governance_error)?; 1559 Ok(()) 1560 } 1561 1562 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1563 enum ConnectionOperationError { 1564 Binding, 1565 Storage, 1566 } 1567 1568 #[cfg(any(target_os = "linux", target_os = "macos"))] 1569 impl From<ConnectionOperationError> for AdminJournalOperationError { 1570 fn from(error: ConnectionOperationError) -> Self { 1571 match error { 1572 ConnectionOperationError::Binding => Self::Binding, 1573 ConnectionOperationError::Storage => Self::Storage, 1574 } 1575 } 1576 } 1577 1578 #[cfg(any(target_os = "linux", target_os = "macos"))] 1579 pub(crate) async fn apply_admin_connection_mutation( 1580 transaction: &mut ServiceSqliteTransaction<'_>, 1581 connection_id: MycConnectionId, 1582 policy_generation: MycConnectionPolicyGeneration, 1583 observed_at: MycConnectionTimeUnixMs, 1584 audit_correlation: MycAuditCorrelationId, 1585 action: MycAdminConnectionAction, 1586 ) -> Result<(MycConnectionRecord, MycConnectionRecord), AdminJournalOperationError> { 1587 let before = read_connection(transaction, connection_id).await?; 1588 if before.policy_generation != policy_generation { 1589 return Err(AdminJournalOperationError::Conflict); 1590 } 1591 let after = match action { 1592 MycAdminConnectionAction::Approve { 1593 permissions, 1594 authorized_until, 1595 } => { 1596 let operation_id = 1597 read_pending_decision_operation(transaction, connection_id, policy_generation) 1598 .await?; 1599 decide_connection( 1600 transaction, 1601 operation_id, 1602 connection_id, 1603 policy_generation, 1604 observed_at, 1605 audit_correlation, 1606 MycConnectionOperatorDecision::Approve { 1607 granted_permissions: permissions, 1608 authorized_until, 1609 }, 1610 ) 1611 .await? 1612 } 1613 MycAdminConnectionAction::Reject => { 1614 let operation_id = 1615 read_pending_decision_operation(transaction, connection_id, policy_generation) 1616 .await?; 1617 decide_connection( 1618 transaction, 1619 operation_id, 1620 connection_id, 1621 policy_generation, 1622 observed_at, 1623 audit_correlation, 1624 MycConnectionOperatorDecision::Deny, 1625 ) 1626 .await? 1627 } 1628 MycAdminConnectionAction::Revoke => { 1629 if before.status == MycConnectionStatus::Expired { 1630 before.clone() 1631 } else { 1632 if before.status != MycConnectionStatus::Active || observed_at < before.updated_at { 1633 return Err(AdminJournalOperationError::Conflict); 1634 } 1635 let result = sqlx::query(REVOKE_CONNECTION_SQL) 1636 .bind(observed_at.sqlite_value()) 1637 .bind(connection_id.as_bytes().as_slice()) 1638 .bind(policy_generation.sqlite_value()) 1639 .execute(&mut *transaction) 1640 .await 1641 .map_err(|_| AdminJournalOperationError::Storage)?; 1642 require_one(result.rows_affected())?; 1643 let record = read_connection(transaction, connection_id).await?; 1644 record_operator_audit( 1645 transaction, 1646 audit_correlation, 1647 observed_at, 1648 MycAuditOutcome::Succeeded, 1649 MycAuditReasonCode::ConnectionExpired, 1650 ) 1651 .await?; 1652 record 1653 } 1654 } 1655 }; 1656 Ok((before, after)) 1657 } 1658 1659 #[cfg(any(target_os = "linux", target_os = "macos"))] 1660 pub(crate) async fn apply_admin_challenge_require( 1661 transaction: &mut ServiceSqliteTransaction<'_>, 1662 request: MycAuthorizationChallengeRequest, 1663 rate_policy: MycRateLimitPolicy, 1664 ) -> Result<MycAuthorizationChallengeRecord, AdminJournalOperationError> { 1665 match issue_challenge(transaction, &request, rate_policy).await? { 1666 MycAuthorizationChallengeAdmission::Created(record) 1667 | MycAuthorizationChallengeAdmission::ExactReplay(record) => Ok(record), 1668 MycAuthorizationChallengeAdmission::RateLimited => { 1669 Err(AdminJournalOperationError::ResourceExhausted) 1670 } 1671 } 1672 } 1673 1674 #[cfg(any(target_os = "linux", target_os = "macos"))] 1675 pub(crate) async fn apply_admin_challenge_authorize( 1676 transaction: &mut ServiceSqliteTransaction<'_>, 1677 challenge_id: MycAuthorizationChallengeId, 1678 policy_generation: MycConnectionPolicyGeneration, 1679 observed_at: MycConnectionTimeUnixMs, 1680 rate_policy: MycRateLimitPolicy, 1681 ) -> Result<MycAuthorizationChallengeRecord, AdminJournalOperationError> { 1682 let rows = sqlx::query(READ_CHALLENGE_OPERATION_BY_ID_SQL) 1683 .bind(challenge_id.as_bytes().as_slice()) 1684 .fetch_all(&mut *transaction) 1685 .await 1686 .map_err(|_| AdminJournalOperationError::Storage)?; 1687 let operation_id = match rows.as_slice() { 1688 [row] => MycSignerOperationId::from_persisted(exact_digest(row, "operation_id")?), 1689 [] | [_, ..] => return Err(AdminJournalOperationError::Binding), 1690 }; 1691 let before = read_challenge(transaction, operation_id) 1692 .await? 1693 .ok_or(AdminJournalOperationError::Binding)?; 1694 if before.policy_generation != policy_generation { 1695 return Err(AdminJournalOperationError::Conflict); 1696 } 1697 match authorize_challenge( 1698 transaction, 1699 challenge_id, 1700 before.connection_id, 1701 before.operation_id, 1702 policy_generation, 1703 observed_at, 1704 rate_policy, 1705 ) 1706 .await? 1707 { 1708 MycAuthorizationChallengeAuthorization::Resolved(record) 1709 | MycAuthorizationChallengeAuthorization::ExactReplay(record) => Ok(record), 1710 MycAuthorizationChallengeAuthorization::RateLimited => { 1711 Err(AdminJournalOperationError::ResourceExhausted) 1712 } 1713 } 1714 } 1715 1716 async fn read_pending_decision_operation( 1717 transaction: &mut ServiceSqliteTransaction<'_>, 1718 connection_id: MycConnectionId, 1719 policy_generation: MycConnectionPolicyGeneration, 1720 ) -> Result<MycSignerOperationId, ConnectionOperationError> { 1721 let rows = sqlx::query(READ_PENDING_DECISION_FOR_CONNECTION_SQL) 1722 .bind(connection_id.as_bytes().as_slice()) 1723 .bind(policy_generation.sqlite_value()) 1724 .fetch_all(&mut *transaction) 1725 .await 1726 .map_err(|_| ConnectionOperationError::Storage)?; 1727 match rows.as_slice() { 1728 [row] => Ok(MycSignerOperationId::from_persisted(exact_digest( 1729 row, 1730 "operation_id", 1731 )?)), 1732 [] | [_, ..] => Err(ConnectionOperationError::Binding), 1733 } 1734 } 1735 1736 const fn map_governance_error(error: GovernanceOperationError) -> ConnectionOperationError { 1737 match error { 1738 GovernanceOperationError::Binding => ConnectionOperationError::Binding, 1739 GovernanceOperationError::Storage => ConnectionOperationError::Storage, 1740 } 1741 } 1742 1743 async fn verify_metadata( 1744 transaction: &mut ServiceSqliteTransaction<'_>, 1745 expected: &PersistedMetadata, 1746 ) -> Result<(), ConnectionOperationError> { 1747 require_expected_metadata(transaction, expected) 1748 .await 1749 .map_err(|error| match error { 1750 RepositoryOperationError::Binding => ConnectionOperationError::Binding, 1751 RepositoryOperationError::Storage => ConnectionOperationError::Storage, 1752 }) 1753 } 1754 1755 async fn admit_connection( 1756 transaction: &mut ServiceSqliteTransaction<'_>, 1757 request: &MycConnectionAdmissionRequest, 1758 rate_policy: Option<MycRateLimitPolicy>, 1759 ) -> Result<MycConnectionAdmission, ConnectionOperationError> { 1760 let binding = read_request_binding(transaction, request.operation_id).await?; 1761 if binding.client_public_key != request.client_public_key 1762 || binding.method != MycSignerRequestMethod::Connect 1763 || request.observed_at < binding.received_at 1764 { 1765 return Err(ConnectionOperationError::Binding); 1766 } 1767 if let Some(existing) = read_decision(transaction, request.operation_id).await? { 1768 if existing.policy_generation != request.policy_generation 1769 || existing.admission_policy != Some(request.policy) 1770 || existing.requested_permissions_sha256 != *request.requested_permissions.digest() 1771 { 1772 return Err(ConnectionOperationError::Binding); 1773 } 1774 let record = decision_record(transaction, request.operation_id, existing).await?; 1775 if record 1776 .connection() 1777 .is_some_and(|connection| connection.client_public_key != request.client_public_key) 1778 { 1779 return Err(ConnectionOperationError::Binding); 1780 } 1781 return Ok(MycConnectionAdmission::ExactReplay(record)); 1782 } 1783 1784 let evidence = AuditEvidence { 1785 correlation: binding.correlation, 1786 occurred_at: request.observed_at, 1787 operation_id: Some(request.operation_id), 1788 }; 1789 if let Some(rate_policy) = rate_policy 1790 && !govern_rate_attempt( 1791 transaction, 1792 rate_policy, 1793 MycRateLimitClass::ConnectionAdmission, 1794 &[global_subject(), relay_subject(&request.relay_id)], 1795 evidence, 1796 MycAuditKind::ConnectionAdmission, 1797 ) 1798 .await 1799 .map_err(map_governance_error)? 1800 { 1801 return Ok(MycConnectionAdmission::RateLimited); 1802 } 1803 1804 let decision = request.policy.decision(); 1805 let connection_id = if request.policy == MycConnectionAdmissionPolicy::Denied { 1806 None 1807 } else { 1808 let id = derive_connection_id(&request.client_public_key, &request.nonce); 1809 insert_connection(transaction, request, id).await?; 1810 Some(id) 1811 }; 1812 insert_decision( 1813 transaction, 1814 request.operation_id, 1815 connection_id, 1816 decision, 1817 decision.reason(request.policy), 1818 request.policy_generation, 1819 request.requested_permissions.digest(), 1820 None, 1821 request.observed_at, 1822 ) 1823 .await?; 1824 let persisted = read_decision(transaction, request.operation_id) 1825 .await? 1826 .ok_or(ConnectionOperationError::Binding)?; 1827 let record = decision_record(transaction, request.operation_id, persisted).await?; 1828 let (outcome, reason) = match request.policy { 1829 MycConnectionAdmissionPolicy::Trusted => { 1830 (MycAuditOutcome::Succeeded, MycAuditReasonCode::Trusted) 1831 } 1832 MycConnectionAdmissionPolicy::ExplicitApproval => ( 1833 MycAuditOutcome::Succeeded, 1834 MycAuditReasonCode::ApprovalRequired, 1835 ), 1836 MycConnectionAdmissionPolicy::Denied => { 1837 (MycAuditOutcome::Rejected, MycAuditReasonCode::PolicyDenied) 1838 } 1839 }; 1840 record_audit( 1841 transaction, 1842 evidence, 1843 MycAuditKind::ConnectionAdmission, 1844 outcome, 1845 reason, 1846 ) 1847 .await 1848 .map_err(map_governance_error)?; 1849 Ok(MycConnectionAdmission::Admitted(record)) 1850 } 1851 1852 async fn insert_connection( 1853 transaction: &mut ServiceSqliteTransaction<'_>, 1854 request: &MycConnectionAdmissionRequest, 1855 id: MycConnectionId, 1856 ) -> Result<(), ConnectionOperationError> { 1857 let status = match request.policy { 1858 MycConnectionAdmissionPolicy::Trusted => MycConnectionStatus::Active, 1859 MycConnectionAdmissionPolicy::ExplicitApproval => MycConnectionStatus::Pending, 1860 MycConnectionAdmissionPolicy::Denied => return Err(ConnectionOperationError::Binding), 1861 }; 1862 let result = sqlx::query(INSERT_CONNECTION_SQL) 1863 .bind(id.as_bytes().as_slice()) 1864 .bind(request.nonce.0.as_slice()) 1865 .bind(request.client_public_key.as_hex()) 1866 .bind(request.requested_permissions.digest().as_slice()) 1867 .bind(request.policy_generation.sqlite_value()) 1868 .bind(status.as_str()) 1869 .bind(request.observed_at.sqlite_value()) 1870 .bind(request.observed_at.sqlite_value()) 1871 .bind( 1872 request 1873 .authorized_until 1874 .map(MycConnectionTimeUnixMs::sqlite_value), 1875 ) 1876 .execute(&mut *transaction) 1877 .await 1878 .map_err(|_| ConnectionOperationError::Storage)?; 1879 require_one(result.rows_affected())?; 1880 insert_permissions(transaction, id, "requested", &request.requested_permissions).await?; 1881 if status == MycConnectionStatus::Active { 1882 insert_permissions(transaction, id, "granted", &request.requested_permissions).await?; 1883 } 1884 Ok(()) 1885 } 1886 1887 #[allow(clippy::too_many_arguments)] 1888 async fn insert_decision( 1889 transaction: &mut ServiceSqliteTransaction<'_>, 1890 operation_id: MycSignerOperationId, 1891 connection_id: Option<MycConnectionId>, 1892 decision: MycConnectionDecision, 1893 reason: &'static str, 1894 policy_generation: MycConnectionPolicyGeneration, 1895 permissions_digest: &[u8; 32], 1896 challenge_id: Option<MycAuthorizationChallengeId>, 1897 decided_at: MycConnectionTimeUnixMs, 1898 ) -> Result<(), ConnectionOperationError> { 1899 let result = sqlx::query(INSERT_DECISION_SQL) 1900 .bind(operation_id.as_bytes().as_slice()) 1901 .bind(connection_id.map(|value| value.0.to_vec())) 1902 .bind(decision.as_str()) 1903 .bind(reason) 1904 .bind(policy_generation.sqlite_value()) 1905 .bind(permissions_digest.as_slice()) 1906 .bind(challenge_id.map(|value| value.0.to_vec())) 1907 .bind(decided_at.sqlite_value()) 1908 .execute(&mut *transaction) 1909 .await 1910 .map_err(|_| ConnectionOperationError::Storage)?; 1911 require_one(result.rows_affected()) 1912 } 1913 1914 async fn decide_connection( 1915 transaction: &mut ServiceSqliteTransaction<'_>, 1916 operation_id: MycSignerOperationId, 1917 connection_id: MycConnectionId, 1918 policy_generation: MycConnectionPolicyGeneration, 1919 observed_at: MycConnectionTimeUnixMs, 1920 audit_correlation: MycAuditCorrelationId, 1921 decision: MycConnectionOperatorDecision, 1922 ) -> Result<MycConnectionRecord, ConnectionOperationError> { 1923 let existing_decision = read_decision(transaction, operation_id) 1924 .await? 1925 .ok_or(ConnectionOperationError::Binding)?; 1926 if existing_decision.connection_id != Some(connection_id) 1927 || existing_decision.policy_generation != policy_generation 1928 { 1929 return Err(ConnectionOperationError::Binding); 1930 } 1931 let connection = read_connection(transaction, connection_id).await?; 1932 if connection.policy_generation != policy_generation { 1933 return Err(ConnectionOperationError::Binding); 1934 } 1935 1936 let (audit_outcome, audit_reason) = match decision { 1937 MycConnectionOperatorDecision::Approve { 1938 granted_permissions, 1939 authorized_until, 1940 } => { 1941 if !granted_permissions.is_subset_of(&connection.requested_permissions) 1942 || authorized_until.is_some_and(|until| until <= observed_at) 1943 { 1944 return Err(ConnectionOperationError::Binding); 1945 } 1946 if connection.status == MycConnectionStatus::Active 1947 && existing_decision.decision == MycConnectionDecision::Allowed 1948 { 1949 let replay = (connection.granted_permissions == granted_permissions 1950 && connection.authorized_until == authorized_until) 1951 .then_some(connection) 1952 .ok_or(ConnectionOperationError::Binding)?; 1953 record_operator_audit( 1954 transaction, 1955 audit_correlation, 1956 replay.updated_at, 1957 MycAuditOutcome::Succeeded, 1958 MycAuditReasonCode::OperatorApproved, 1959 ) 1960 .await?; 1961 return Ok(replay); 1962 } 1963 if connection.status != MycConnectionStatus::Pending 1964 || existing_decision.decision != MycConnectionDecision::PendingApproval 1965 || observed_at < connection.updated_at 1966 { 1967 return Err(ConnectionOperationError::Binding); 1968 } 1969 insert_permissions(transaction, connection_id, "granted", &granted_permissions).await?; 1970 let result = sqlx::query(APPROVE_CONNECTION_SQL) 1971 .bind(observed_at.sqlite_value()) 1972 .bind(authorized_until.map(MycConnectionTimeUnixMs::sqlite_value)) 1973 .bind(connection_id.as_bytes().as_slice()) 1974 .bind(policy_generation.sqlite_value()) 1975 .execute(&mut *transaction) 1976 .await 1977 .map_err(|_| ConnectionOperationError::Storage)?; 1978 require_one(result.rows_affected())?; 1979 update_approval_decision( 1980 transaction, 1981 operation_id, 1982 connection_id, 1983 policy_generation, 1984 observed_at, 1985 MycConnectionDecision::Allowed, 1986 "operator_approved", 1987 ) 1988 .await?; 1989 ( 1990 MycAuditOutcome::Succeeded, 1991 MycAuditReasonCode::OperatorApproved, 1992 ) 1993 } 1994 MycConnectionOperatorDecision::Deny => { 1995 if connection.status == MycConnectionStatus::Denied 1996 && existing_decision.decision == MycConnectionDecision::Denied 1997 { 1998 record_operator_audit( 1999 transaction, 2000 audit_correlation, 2001 connection.updated_at, 2002 MycAuditOutcome::Rejected, 2003 MycAuditReasonCode::OperatorDenied, 2004 ) 2005 .await?; 2006 return Ok(connection); 2007 } 2008 if connection.status != MycConnectionStatus::Pending 2009 || existing_decision.decision != MycConnectionDecision::PendingApproval 2010 || observed_at < connection.updated_at 2011 { 2012 return Err(ConnectionOperationError::Binding); 2013 } 2014 let result = sqlx::query(DENY_CONNECTION_SQL) 2015 .bind(observed_at.sqlite_value()) 2016 .bind(connection_id.as_bytes().as_slice()) 2017 .bind(policy_generation.sqlite_value()) 2018 .execute(&mut *transaction) 2019 .await 2020 .map_err(|_| ConnectionOperationError::Storage)?; 2021 require_one(result.rows_affected())?; 2022 update_approval_decision( 2023 transaction, 2024 operation_id, 2025 connection_id, 2026 policy_generation, 2027 observed_at, 2028 MycConnectionDecision::Denied, 2029 "operator_denied", 2030 ) 2031 .await?; 2032 ( 2033 MycAuditOutcome::Rejected, 2034 MycAuditReasonCode::OperatorDenied, 2035 ) 2036 } 2037 }; 2038 let record = read_connection(transaction, connection_id).await?; 2039 record_operator_audit( 2040 transaction, 2041 audit_correlation, 2042 observed_at, 2043 audit_outcome, 2044 audit_reason, 2045 ) 2046 .await?; 2047 Ok(record) 2048 } 2049 2050 async fn record_operator_audit( 2051 transaction: &mut ServiceSqliteTransaction<'_>, 2052 correlation: MycAuditCorrelationId, 2053 occurred_at: MycConnectionTimeUnixMs, 2054 outcome: MycAuditOutcome, 2055 reason: MycAuditReasonCode, 2056 ) -> Result<(), ConnectionOperationError> { 2057 record_audit( 2058 transaction, 2059 AuditEvidence { 2060 correlation, 2061 occurred_at, 2062 operation_id: None, 2063 }, 2064 MycAuditKind::ConnectionOperatorDecision, 2065 outcome, 2066 reason, 2067 ) 2068 .await 2069 .map_err(map_governance_error)?; 2070 Ok(()) 2071 } 2072 2073 #[allow(clippy::too_many_arguments)] 2074 async fn update_approval_decision( 2075 transaction: &mut ServiceSqliteTransaction<'_>, 2076 operation_id: MycSignerOperationId, 2077 connection_id: MycConnectionId, 2078 policy_generation: MycConnectionPolicyGeneration, 2079 observed_at: MycConnectionTimeUnixMs, 2080 decision: MycConnectionDecision, 2081 reason: &'static str, 2082 ) -> Result<(), ConnectionOperationError> { 2083 let result = sqlx::query(UPDATE_APPROVAL_DECISION_SQL) 2084 .bind(decision.as_str()) 2085 .bind(reason) 2086 .bind(observed_at.sqlite_value()) 2087 .bind(operation_id.as_bytes().as_slice()) 2088 .bind(connection_id.as_bytes().as_slice()) 2089 .bind(policy_generation.sqlite_value()) 2090 .execute(&mut *transaction) 2091 .await 2092 .map_err(|_| ConnectionOperationError::Storage)?; 2093 require_one(result.rows_affected()) 2094 } 2095 2096 async fn issue_challenge( 2097 transaction: &mut ServiceSqliteTransaction<'_>, 2098 request: &MycAuthorizationChallengeRequest, 2099 rate_policy: MycRateLimitPolicy, 2100 ) -> Result<MycAuthorizationChallengeAdmission, ConnectionOperationError> { 2101 let binding = read_request_binding(transaction, request.operation_id).await?; 2102 if binding.method == MycSignerRequestMethod::Connect { 2103 return Err(ConnectionOperationError::Binding); 2104 } 2105 let connection = read_connection(transaction, request.connection_id).await?; 2106 if connection.client_public_key != binding.client_public_key 2107 || connection.policy_generation != request.policy_generation 2108 { 2109 return Err(ConnectionOperationError::Binding); 2110 } 2111 if let Some(existing) = read_challenge(transaction, request.operation_id).await? { 2112 if existing.connection_id != request.connection_id 2113 || existing.policy_generation != request.policy_generation 2114 || existing.url != request.url 2115 || existing.issued_at != request.issued_at 2116 || existing.expires_at != request.expires_at 2117 { 2118 return Err(ConnectionOperationError::Binding); 2119 } 2120 return Ok(MycAuthorizationChallengeAdmission::ExactReplay(existing)); 2121 } 2122 if request.issued_at < binding.received_at 2123 || request.issued_at < connection.updated_at 2124 || connection.status != MycConnectionStatus::Active 2125 || connection 2126 .authorized_until 2127 .is_some_and(|until| until < request.issued_at) 2128 { 2129 return Err(ConnectionOperationError::Binding); 2130 } 2131 if read_decision(transaction, request.operation_id) 2132 .await? 2133 .is_some() 2134 { 2135 return Err(ConnectionOperationError::Binding); 2136 } 2137 let evidence = AuditEvidence { 2138 correlation: binding.correlation, 2139 occurred_at: request.issued_at, 2140 operation_id: Some(request.operation_id), 2141 }; 2142 if !govern_rate_attempt( 2143 transaction, 2144 rate_policy, 2145 MycRateLimitClass::ChallengeCreation, 2146 &[connection_subject(request.connection_id)], 2147 evidence, 2148 MycAuditKind::ChallengeCreation, 2149 ) 2150 .await 2151 .map_err(map_governance_error)? 2152 { 2153 return Ok(MycAuthorizationChallengeAdmission::RateLimited); 2154 } 2155 let challenge_id = 2156 derive_challenge_id(request.operation_id, request.connection_id, &request.nonce); 2157 let result = sqlx::query(INSERT_CHALLENGE_SQL) 2158 .bind(challenge_id.as_bytes().as_slice()) 2159 .bind(request.nonce.0.as_slice()) 2160 .bind(request.connection_id.as_bytes().as_slice()) 2161 .bind(request.operation_id.as_bytes().as_slice()) 2162 .bind(request.policy_generation.sqlite_value()) 2163 .bind(request.url.as_str()) 2164 .bind(request.issued_at.sqlite_value()) 2165 .bind(request.expires_at.sqlite_value()) 2166 .execute(&mut *transaction) 2167 .await 2168 .map_err(|_| ConnectionOperationError::Storage)?; 2169 require_one(result.rows_affected())?; 2170 insert_decision( 2171 transaction, 2172 request.operation_id, 2173 Some(request.connection_id), 2174 MycConnectionDecision::Challenged, 2175 "authorization_challenge_required", 2176 request.policy_generation, 2177 connection.requested_permissions.digest(), 2178 Some(challenge_id), 2179 request.issued_at, 2180 ) 2181 .await?; 2182 let record = read_challenge(transaction, request.operation_id) 2183 .await? 2184 .ok_or(ConnectionOperationError::Binding)?; 2185 record_audit( 2186 transaction, 2187 evidence, 2188 MycAuditKind::ChallengeCreation, 2189 MycAuditOutcome::Succeeded, 2190 MycAuditReasonCode::ChallengeRequired, 2191 ) 2192 .await 2193 .map_err(map_governance_error)?; 2194 Ok(MycAuthorizationChallengeAdmission::Created(record)) 2195 } 2196 2197 #[allow(clippy::too_many_arguments)] 2198 async fn authorize_challenge( 2199 transaction: &mut ServiceSqliteTransaction<'_>, 2200 challenge_id: MycAuthorizationChallengeId, 2201 connection_id: MycConnectionId, 2202 operation_id: MycSignerOperationId, 2203 policy_generation: MycConnectionPolicyGeneration, 2204 observed_at: MycConnectionTimeUnixMs, 2205 rate_policy: MycRateLimitPolicy, 2206 ) -> Result<MycAuthorizationChallengeAuthorization, ConnectionOperationError> { 2207 let before = read_challenge(transaction, operation_id) 2208 .await? 2209 .ok_or(ConnectionOperationError::Binding)?; 2210 if before.id != challenge_id 2211 || before.connection_id != connection_id 2212 || before.operation_id != operation_id 2213 || before.policy_generation != policy_generation 2214 { 2215 return Err(ConnectionOperationError::Binding); 2216 } 2217 let request_binding = read_request_binding(transaction, operation_id).await?; 2218 let connection = read_connection(transaction, connection_id).await?; 2219 if request_binding.method == MycSignerRequestMethod::Connect 2220 || request_binding.client_public_key != connection.client_public_key 2221 || connection.policy_generation != policy_generation 2222 || observed_at < before.issued_at 2223 { 2224 return Err(ConnectionOperationError::Binding); 2225 } 2226 if before.state != MycAuthorizationChallengeState::Pending { 2227 return Ok(MycAuthorizationChallengeAuthorization::ExactReplay(before)); 2228 } 2229 let evidence = AuditEvidence { 2230 correlation: request_binding.correlation, 2231 occurred_at: observed_at, 2232 operation_id: Some(operation_id), 2233 }; 2234 if !govern_rate_attempt( 2235 transaction, 2236 rate_policy, 2237 MycRateLimitClass::ChallengeAuthorization, 2238 &[connection_subject(connection_id)], 2239 evidence, 2240 MycAuditKind::ChallengeAuthorization, 2241 ) 2242 .await 2243 .map_err(map_governance_error)? 2244 { 2245 return Ok(MycAuthorizationChallengeAuthorization::RateLimited); 2246 } 2247 let connection_expired = connection.status != MycConnectionStatus::Active 2248 || connection 2249 .authorized_until 2250 .is_some_and(|until| until < observed_at); 2251 let expired = observed_at > before.expires_at || connection_expired; 2252 let (state, decision, reason) = if expired { 2253 ( 2254 MycAuthorizationChallengeState::Expired, 2255 MycConnectionDecision::Denied, 2256 "authorization_challenge_expired", 2257 ) 2258 } else { 2259 ( 2260 MycAuthorizationChallengeState::Authorized, 2261 MycConnectionDecision::Allowed, 2262 "authorization_challenge_authorized", 2263 ) 2264 }; 2265 let result = sqlx::query(RESOLVE_CHALLENGE_SQL) 2266 .bind(state.as_str()) 2267 .bind(observed_at.sqlite_value()) 2268 .bind(challenge_id.as_bytes().as_slice()) 2269 .bind(connection_id.as_bytes().as_slice()) 2270 .bind(operation_id.as_bytes().as_slice()) 2271 .bind(policy_generation.sqlite_value()) 2272 .execute(&mut *transaction) 2273 .await 2274 .map_err(|_| ConnectionOperationError::Storage)?; 2275 require_one(result.rows_affected())?; 2276 let result = sqlx::query(UPDATE_CHALLENGE_DECISION_SQL) 2277 .bind(decision.as_str()) 2278 .bind(reason) 2279 .bind(observed_at.sqlite_value()) 2280 .bind(operation_id.as_bytes().as_slice()) 2281 .bind(connection_id.as_bytes().as_slice()) 2282 .bind(challenge_id.as_bytes().as_slice()) 2283 .bind(policy_generation.sqlite_value()) 2284 .execute(&mut *transaction) 2285 .await 2286 .map_err(|_| ConnectionOperationError::Storage)?; 2287 require_one(result.rows_affected())?; 2288 let record = read_challenge(transaction, operation_id) 2289 .await? 2290 .ok_or(ConnectionOperationError::Binding)?; 2291 let (outcome, audit_reason) = match state { 2292 MycAuthorizationChallengeState::Authorized => ( 2293 MycAuditOutcome::Succeeded, 2294 MycAuditReasonCode::ChallengeAuthorized, 2295 ), 2296 MycAuthorizationChallengeState::Expired => ( 2297 MycAuditOutcome::Rejected, 2298 MycAuditReasonCode::ChallengeExpired, 2299 ), 2300 MycAuthorizationChallengeState::Pending => { 2301 return Err(ConnectionOperationError::Binding); 2302 } 2303 }; 2304 record_audit( 2305 transaction, 2306 evidence, 2307 MycAuditKind::ChallengeAuthorization, 2308 outcome, 2309 audit_reason, 2310 ) 2311 .await 2312 .map_err(map_governance_error)?; 2313 Ok(MycAuthorizationChallengeAuthorization::Resolved(record)) 2314 } 2315 2316 struct RequestBinding { 2317 correlation: MycAuditCorrelationId, 2318 client_public_key: MycNip46ClientPublicKey, 2319 method: MycSignerRequestMethod, 2320 received_at: MycConnectionTimeUnixMs, 2321 } 2322 2323 async fn read_request_binding( 2324 transaction: &mut ServiceSqliteTransaction<'_>, 2325 operation_id: MycSignerOperationId, 2326 ) -> Result<RequestBinding, ConnectionOperationError> { 2327 let rows = sqlx::query(READ_REQUEST_BINDING_SQL) 2328 .bind(operation_id.as_bytes().as_slice()) 2329 .fetch_all(&mut *transaction) 2330 .await 2331 .map_err(|_| ConnectionOperationError::Storage)?; 2332 if rows.len() != 1 { 2333 return Err(ConnectionOperationError::Binding); 2334 } 2335 let row = &rows[0]; 2336 let correlation = MycAuditCorrelationId::new(exact_digest(row, "correlation_id")?); 2337 let client = row 2338 .try_get::<Option<&str>, _>("client_public_key") 2339 .map_err(|_| ConnectionOperationError::Binding)? 2340 .ok_or(ConnectionOperationError::Binding) 2341 .and_then(|value| { 2342 MycNip46ClientPublicKey::new(value).map_err(|_| ConnectionOperationError::Binding) 2343 })?; 2344 let method = row 2345 .try_get::<Option<&str>, _>("method") 2346 .map_err(|_| ConnectionOperationError::Binding)? 2347 .and_then(MycSignerRequestMethod::parse) 2348 .ok_or(ConnectionOperationError::Binding)?; 2349 Ok(RequestBinding { 2350 correlation, 2351 client_public_key: client, 2352 method, 2353 received_at: time(row, "received_at_unix_ms")?, 2354 }) 2355 } 2356 2357 #[derive(Clone, Copy)] 2358 struct PersistedDecision { 2359 connection_id: Option<MycConnectionId>, 2360 decision: MycConnectionDecision, 2361 policy_generation: MycConnectionPolicyGeneration, 2362 requested_permissions_sha256: [u8; 32], 2363 challenge_id: Option<MycAuthorizationChallengeId>, 2364 decided_at: MycConnectionTimeUnixMs, 2365 admission_policy: Option<MycConnectionAdmissionPolicy>, 2366 } 2367 2368 async fn read_decision( 2369 transaction: &mut ServiceSqliteTransaction<'_>, 2370 operation_id: MycSignerOperationId, 2371 ) -> Result<Option<PersistedDecision>, ConnectionOperationError> { 2372 let rows = sqlx::query(READ_DECISION_SQL) 2373 .bind(operation_id.as_bytes().as_slice()) 2374 .fetch_all(&mut *transaction) 2375 .await 2376 .map_err(|_| ConnectionOperationError::Storage)?; 2377 if rows.len() > 1 { 2378 return Err(ConnectionOperationError::Binding); 2379 } 2380 rows.first().map(parse_decision).transpose() 2381 } 2382 2383 fn parse_decision( 2384 row: &sqlx::sqlite::SqliteRow, 2385 ) -> Result<PersistedDecision, ConnectionOperationError> { 2386 let connection_id = 2387 optional_digest(row, "connection_id", "connection_id_type")?.map(MycConnectionId); 2388 let challenge_id = 2389 optional_digest(row, "challenge_id", "challenge_id_type")?.map(MycAuthorizationChallengeId); 2390 let decision = bounded_text(row, "decision").and_then(|value| { 2391 MycConnectionDecision::parse(value).ok_or(ConnectionOperationError::Binding) 2392 })?; 2393 let reason = bounded_text(row, "reason_code")?; 2394 if !matches!( 2395 (decision, reason, connection_id, challenge_id), 2396 (MycConnectionDecision::Denied, "policy_denied", None, None) 2397 | ( 2398 MycConnectionDecision::PendingApproval, 2399 "explicit_approval_required", 2400 Some(_), 2401 None 2402 ) 2403 | ( 2404 MycConnectionDecision::Allowed, 2405 "trusted_client", 2406 Some(_), 2407 None 2408 ) 2409 | ( 2410 MycConnectionDecision::Allowed, 2411 "operator_approved", 2412 Some(_), 2413 None 2414 ) 2415 | ( 2416 MycConnectionDecision::Denied, 2417 "operator_denied", 2418 Some(_), 2419 None 2420 ) 2421 | ( 2422 MycConnectionDecision::Challenged, 2423 "authorization_challenge_required", 2424 Some(_), 2425 Some(_) 2426 ) 2427 | ( 2428 MycConnectionDecision::Allowed, 2429 "authorization_challenge_authorized", 2430 Some(_), 2431 Some(_) 2432 ) 2433 | ( 2434 MycConnectionDecision::Denied, 2435 "authorization_challenge_expired", 2436 Some(_), 2437 Some(_) 2438 ) 2439 ) { 2440 return Err(ConnectionOperationError::Binding); 2441 } 2442 let admission_policy = match reason { 2443 "trusted_client" => Some(MycConnectionAdmissionPolicy::Trusted), 2444 "explicit_approval_required" | "operator_approved" | "operator_denied" => { 2445 Some(MycConnectionAdmissionPolicy::ExplicitApproval) 2446 } 2447 "policy_denied" => Some(MycConnectionAdmissionPolicy::Denied), 2448 "authorization_challenge_required" 2449 | "authorization_challenge_authorized" 2450 | "authorization_challenge_expired" => None, 2451 _ => return Err(ConnectionOperationError::Binding), 2452 }; 2453 Ok(PersistedDecision { 2454 connection_id, 2455 decision, 2456 policy_generation: generation(row, "policy_generation")?, 2457 requested_permissions_sha256: exact_digest(row, "requested_permissions_sha256")?, 2458 challenge_id, 2459 decided_at: time(row, "decided_at_unix_ms")?, 2460 admission_policy, 2461 }) 2462 } 2463 2464 async fn decision_record( 2465 transaction: &mut ServiceSqliteTransaction<'_>, 2466 operation_id: MycSignerOperationId, 2467 decision: PersistedDecision, 2468 ) -> Result<MycConnectionDecisionRecord, ConnectionOperationError> { 2469 let connection = match decision.connection_id { 2470 Some(id) => Some(read_connection(transaction, id).await?), 2471 None => None, 2472 }; 2473 if connection.as_ref().is_some_and(|value| { 2474 value.requested_permissions.digest() != &decision.requested_permissions_sha256 2475 }) || decision.challenge_id.is_some() 2476 { 2477 return Err(ConnectionOperationError::Binding); 2478 } 2479 let state_matches = match (decision.admission_policy, decision.decision, &connection) { 2480 (Some(MycConnectionAdmissionPolicy::Denied), MycConnectionDecision::Denied, None) => true, 2481 ( 2482 Some(MycConnectionAdmissionPolicy::Trusted), 2483 MycConnectionDecision::Allowed, 2484 Some(connection), 2485 ) => matches!( 2486 connection.status, 2487 MycConnectionStatus::Active | MycConnectionStatus::Expired 2488 ), 2489 ( 2490 Some(MycConnectionAdmissionPolicy::ExplicitApproval), 2491 MycConnectionDecision::PendingApproval, 2492 Some(connection), 2493 ) => connection.status == MycConnectionStatus::Pending, 2494 ( 2495 Some(MycConnectionAdmissionPolicy::ExplicitApproval), 2496 MycConnectionDecision::Allowed, 2497 Some(connection), 2498 ) => matches!( 2499 connection.status, 2500 MycConnectionStatus::Active | MycConnectionStatus::Expired 2501 ), 2502 ( 2503 Some(MycConnectionAdmissionPolicy::ExplicitApproval), 2504 MycConnectionDecision::Denied, 2505 Some(connection), 2506 ) => connection.status == MycConnectionStatus::Denied, 2507 _ => false, 2508 }; 2509 if !state_matches { 2510 return Err(ConnectionOperationError::Binding); 2511 } 2512 Ok(MycConnectionDecisionRecord { 2513 operation_id, 2514 connection, 2515 decision: decision.decision, 2516 policy_generation: decision.policy_generation, 2517 decided_at: decision.decided_at, 2518 }) 2519 } 2520 2521 async fn read_connection( 2522 transaction: &mut ServiceSqliteTransaction<'_>, 2523 id: MycConnectionId, 2524 ) -> Result<MycConnectionRecord, ConnectionOperationError> { 2525 let rows = sqlx::query(READ_CONNECTION_SQL) 2526 .bind(id.as_bytes().as_slice()) 2527 .fetch_all(&mut *transaction) 2528 .await 2529 .map_err(|_| ConnectionOperationError::Storage)?; 2530 if rows.len() != 1 { 2531 return Err(ConnectionOperationError::Binding); 2532 } 2533 let row = &rows[0]; 2534 let actual_id = MycConnectionId(exact_digest(row, "connection_id")?); 2535 let nonce = exact_digest(row, "connection_nonce")?; 2536 let client_public_key = MycNip46ClientPublicKey::new(bounded_text(row, "client_public_key")?) 2537 .map_err(|_| ConnectionOperationError::Binding)?; 2538 let requested_permissions = read_permissions(transaction, id, "requested").await?; 2539 let granted_permissions = read_permissions(transaction, id, "granted").await?; 2540 let requested_digest = exact_digest(row, "requested_permissions_sha256")?; 2541 let policy_generation = generation(row, "policy_generation")?; 2542 let status = MycConnectionStatus::parse(bounded_text(row, "status")?) 2543 .ok_or(ConnectionOperationError::Binding)?; 2544 let created_at = time(row, "created_at_unix_ms")?; 2545 let updated_at = time(row, "updated_at_unix_ms")?; 2546 let authorized_until = optional_time(row, "authorized_until_unix_ms", "authorized_until_type")?; 2547 if actual_id != id 2548 || derive_connection_id(&client_public_key, &MycConnectionNonce(nonce)) != id 2549 || requested_permissions.digest() != &requested_digest 2550 || !granted_permissions.is_subset_of(&requested_permissions) 2551 || updated_at < created_at 2552 || (matches!( 2553 status, 2554 MycConnectionStatus::Pending | MycConnectionStatus::Denied 2555 ) && !granted_permissions.permissions.is_empty()) 2556 || (status == MycConnectionStatus::Active 2557 && authorized_until.is_some_and(|until| until <= updated_at)) 2558 || (matches!( 2559 status, 2560 MycConnectionStatus::Denied | MycConnectionStatus::Expired 2561 ) && authorized_until.is_some()) 2562 { 2563 return Err(ConnectionOperationError::Binding); 2564 } 2565 Ok(MycConnectionRecord { 2566 id, 2567 client_public_key, 2568 requested_permissions, 2569 granted_permissions, 2570 policy_generation, 2571 status, 2572 created_at, 2573 updated_at, 2574 authorized_until, 2575 }) 2576 } 2577 2578 async fn insert_permissions( 2579 transaction: &mut ServiceSqliteTransaction<'_>, 2580 connection_id: MycConnectionId, 2581 scope: &'static str, 2582 permissions: &MycConnectionPermissionSet, 2583 ) -> Result<(), ConnectionOperationError> { 2584 for permission in permissions.permissions() { 2585 let result = sqlx::query(INSERT_PERMISSION_SQL) 2586 .bind(connection_id.as_bytes().as_slice()) 2587 .bind(scope) 2588 .bind(permission.code()) 2589 .execute(&mut *transaction) 2590 .await 2591 .map_err(|_| ConnectionOperationError::Storage)?; 2592 require_one(result.rows_affected())?; 2593 } 2594 Ok(()) 2595 } 2596 2597 async fn read_permissions( 2598 transaction: &mut ServiceSqliteTransaction<'_>, 2599 connection_id: MycConnectionId, 2600 scope: &'static str, 2601 ) -> Result<MycConnectionPermissionSet, ConnectionOperationError> { 2602 let rows = sqlx::query(READ_PERMISSIONS_SQL) 2603 .bind(connection_id.as_bytes().as_slice()) 2604 .bind(scope) 2605 .fetch_all(&mut *transaction) 2606 .await 2607 .map_err(|_| ConnectionOperationError::Storage)?; 2608 if rows.len() > MYC_CONNECTION_PERMISSION_MAX_COUNT { 2609 return Err(ConnectionOperationError::Binding); 2610 } 2611 let permissions = rows 2612 .iter() 2613 .map(|row| { 2614 bounded_text(row, "permission_code").and_then(|value| { 2615 MycConnectionPermission::parse(value).ok_or(ConnectionOperationError::Binding) 2616 }) 2617 }) 2618 .collect::<Result<Vec<_>, _>>()?; 2619 MycConnectionPermissionSet::new(&permissions).map_err(|_| ConnectionOperationError::Binding) 2620 } 2621 2622 async fn read_challenge( 2623 transaction: &mut ServiceSqliteTransaction<'_>, 2624 operation_id: MycSignerOperationId, 2625 ) -> Result<Option<MycAuthorizationChallengeRecord>, ConnectionOperationError> { 2626 let rows = sqlx::query(READ_CHALLENGE_SQL) 2627 .bind(operation_id.as_bytes().as_slice()) 2628 .fetch_all(&mut *transaction) 2629 .await 2630 .map_err(|_| ConnectionOperationError::Storage)?; 2631 if rows.len() > 1 { 2632 return Err(ConnectionOperationError::Binding); 2633 } 2634 rows.first().map(parse_challenge).transpose() 2635 } 2636 2637 fn parse_challenge( 2638 row: &sqlx::sqlite::SqliteRow, 2639 ) -> Result<MycAuthorizationChallengeRecord, ConnectionOperationError> { 2640 let id = MycAuthorizationChallengeId(exact_digest(row, "challenge_id")?); 2641 let nonce = exact_digest(row, "challenge_nonce")?; 2642 let connection_id = MycConnectionId(exact_digest(row, "connection_id")?); 2643 let operation_id = MycSignerOperationId::from_persisted(exact_digest(row, "operation_id")?); 2644 let policy_generation = generation(row, "policy_generation")?; 2645 let url = MycAuthorizationChallengeUrl::new(bounded_text(row, "challenge_url")?) 2646 .map_err(|_| ConnectionOperationError::Binding)?; 2647 let state = MycAuthorizationChallengeState::parse(bounded_text(row, "state")?) 2648 .ok_or(ConnectionOperationError::Binding)?; 2649 let issued_at = time(row, "issued_at_unix_ms")?; 2650 let expires_at = time(row, "expires_at_unix_ms")?; 2651 let resolved_at = optional_time(row, "resolved_at_unix_ms", "resolved_at_type")?; 2652 if derive_challenge_id( 2653 operation_id, 2654 connection_id, 2655 &MycAuthorizationChallengeNonce(nonce), 2656 ) != id 2657 || expires_at <= issued_at 2658 || (state == MycAuthorizationChallengeState::Pending && resolved_at.is_some()) 2659 || (state != MycAuthorizationChallengeState::Pending 2660 && resolved_at.is_none_or(|resolved| resolved < issued_at)) 2661 { 2662 return Err(ConnectionOperationError::Binding); 2663 } 2664 Ok(MycAuthorizationChallengeRecord { 2665 id, 2666 operation_id, 2667 connection_id, 2668 policy_generation, 2669 url, 2670 state, 2671 issued_at, 2672 expires_at, 2673 resolved_at, 2674 }) 2675 } 2676 2677 fn permission_digest(permissions: &[MycConnectionPermission]) -> [u8; 32] { 2678 let mut hasher = Sha256::new(); 2679 hasher.update(PERMISSION_SET_DOMAIN); 2680 hasher.update( 2681 u64::try_from(permissions.len()) 2682 .expect("bounded permission count") 2683 .to_be_bytes(), 2684 ); 2685 for permission in permissions { 2686 let code = permission.code(); 2687 hasher.update( 2688 u64::try_from(code.len()) 2689 .expect("bounded permission code") 2690 .to_be_bytes(), 2691 ); 2692 hasher.update(code.as_bytes()); 2693 } 2694 hasher.finalize().into() 2695 } 2696 2697 fn derive_connection_id( 2698 client: &MycNip46ClientPublicKey, 2699 nonce: &MycConnectionNonce, 2700 ) -> MycConnectionId { 2701 let mut hasher = Sha256::new(); 2702 hasher.update(CONNECTION_ID_DOMAIN); 2703 hasher.update(client.as_hex().as_bytes()); 2704 hasher.update(nonce.0); 2705 MycConnectionId(hasher.finalize().into()) 2706 } 2707 2708 fn derive_challenge_id( 2709 operation_id: MycSignerOperationId, 2710 connection_id: MycConnectionId, 2711 nonce: &MycAuthorizationChallengeNonce, 2712 ) -> MycAuthorizationChallengeId { 2713 let mut hasher = Sha256::new(); 2714 hasher.update(CHALLENGE_ID_DOMAIN); 2715 hasher.update(operation_id.as_bytes()); 2716 hasher.update(connection_id.as_bytes()); 2717 hasher.update(nonce.0); 2718 MycAuthorizationChallengeId(hasher.finalize().into()) 2719 } 2720 2721 fn exact_digest( 2722 row: &sqlx::sqlite::SqliteRow, 2723 column: &str, 2724 ) -> Result<[u8; 32], ConnectionOperationError> { 2725 row.try_get::<Option<Vec<u8>>, _>(column) 2726 .map_err(|_| ConnectionOperationError::Binding)? 2727 .ok_or(ConnectionOperationError::Binding)? 2728 .try_into() 2729 .map_err(|_| ConnectionOperationError::Binding) 2730 } 2731 2732 fn optional_digest( 2733 row: &sqlx::sqlite::SqliteRow, 2734 column: &str, 2735 type_column: &str, 2736 ) -> Result<Option<[u8; 32]>, ConnectionOperationError> { 2737 let kind = row 2738 .try_get::<&str, _>(type_column) 2739 .map_err(|_| ConnectionOperationError::Binding)?; 2740 match kind { 2741 "null" => Ok(None), 2742 "blob" => exact_digest(row, column).map(Some), 2743 _ => Err(ConnectionOperationError::Binding), 2744 } 2745 } 2746 2747 fn bounded_text<'row>( 2748 row: &'row sqlx::sqlite::SqliteRow, 2749 column: &str, 2750 ) -> Result<&'row str, ConnectionOperationError> { 2751 row.try_get::<Option<&str>, _>(column) 2752 .map_err(|_| ConnectionOperationError::Binding)? 2753 .ok_or(ConnectionOperationError::Binding) 2754 } 2755 2756 fn generation( 2757 row: &sqlx::sqlite::SqliteRow, 2758 column: &str, 2759 ) -> Result<MycConnectionPolicyGeneration, ConnectionOperationError> { 2760 let value = row 2761 .try_get::<i64, _>(column) 2762 .map_err(|_| ConnectionOperationError::Binding)?; 2763 MycConnectionPolicyGeneration::new( 2764 u64::try_from(value).map_err(|_| ConnectionOperationError::Binding)?, 2765 ) 2766 .map_err(|_| ConnectionOperationError::Binding) 2767 } 2768 2769 fn time( 2770 row: &sqlx::sqlite::SqliteRow, 2771 column: &str, 2772 ) -> Result<MycConnectionTimeUnixMs, ConnectionOperationError> { 2773 let value = row 2774 .try_get::<i64, _>(column) 2775 .map_err(|_| ConnectionOperationError::Binding)?; 2776 MycConnectionTimeUnixMs::new( 2777 u64::try_from(value).map_err(|_| ConnectionOperationError::Binding)?, 2778 ) 2779 .map_err(|_| ConnectionOperationError::Binding) 2780 } 2781 2782 fn optional_time( 2783 row: &sqlx::sqlite::SqliteRow, 2784 column: &str, 2785 type_column: &str, 2786 ) -> Result<Option<MycConnectionTimeUnixMs>, ConnectionOperationError> { 2787 match row 2788 .try_get::<&str, _>(type_column) 2789 .map_err(|_| ConnectionOperationError::Binding)? 2790 { 2791 "null" => Ok(None), 2792 "integer" => time(row, column).map(Some), 2793 _ => Err(ConnectionOperationError::Binding), 2794 } 2795 } 2796 2797 fn require_one(rows: u64) -> Result<(), ConnectionOperationError> { 2798 (rows == 1) 2799 .then_some(()) 2800 .ok_or(ConnectionOperationError::Storage) 2801 } 2802 2803 fn map_transaction_error( 2804 error: ServiceSqliteTransactionError<ConnectionOperationError>, 2805 ) -> MycStateRepositoryError { 2806 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 2807 return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); 2808 } 2809 let kind = match error.operation_error() { 2810 Some(ConnectionOperationError::Binding) => MycStateRepositoryErrorKind::Binding, 2811 Some(ConnectionOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, 2812 }; 2813 MycStateRepositoryError::new(kind) 2814 }