state_governance.rs (53240B)
1 //! Durable bounded audit, rate-window, retention, and compaction 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 12 use crate::state_repository::{ 13 MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, 14 RepositoryOperationError, require_expected_metadata, 15 }; 16 use crate::{MycConnectionId, MycConnectionTimeUnixMs, MycSignerOperationId}; 17 18 /// Maximum governed rate-window duration. 19 pub const MYC_RATE_WINDOW_MAX_MS: u64 = 86_400_000; 20 /// Maximum attempts admitted in one governed rate window. 21 pub const MYC_RATE_MAX_ATTEMPTS: u32 = 10_000; 22 /// Maximum retained duration for rate-window evidence. 23 pub const MYC_RATE_RETENTION_MAX_MS: u64 = 2_592_000_000; 24 /// Maximum tracked non-global subjects for one rate-limit class. 25 pub const MYC_RATE_MAX_TRACKED_SUBJECTS: u32 = 65_536; 26 /// Maximum retained duration for safe audit evidence. 27 pub const MYC_AUDIT_RETENTION_MAX_MS: u64 = 31_536_000_000; 28 /// Maximum audit records returned in one page. 29 pub const MYC_AUDIT_PAGE_MAX_ITEMS: u16 = 200; 30 /// Maximum rows removed from either bounded evidence class in one compaction. 31 pub const MYC_COMPACTION_MAX_ROWS: u16 = 4_096; 32 /// Maximum UTF-8 bytes in a configured stable relay ID. 33 pub const MYC_RATE_RELAY_ID_MAX_BYTES: usize = 64; 34 35 const AUDIT_ID_DOMAIN: &[u8] = b"radroots.myc.operation_audit.v1\0"; 36 const GLOBAL_SUBJECT_DOMAIN: &[u8] = b"radroots.myc.rate_subject.global.v1\0"; 37 const RELAY_SUBJECT_DOMAIN: &[u8] = b"radroots.myc.rate_subject.relay.v1\0"; 38 const CONNECTION_SUBJECT_DOMAIN: &[u8] = b"radroots.myc.rate_subject.connection.v1\0"; 39 40 const READ_AUDIT_STATE_SQL: &str = 41 "SELECT next_sequence FROM myc_audit_state WHERE singleton = 1 LIMIT 2"; 42 const ADVANCE_AUDIT_STATE_SQL: &str = 43 "UPDATE myc_audit_state SET next_sequence = ? WHERE singleton = 1 AND next_sequence = ?"; 44 const READ_AUDIT_BY_ID_SQL: &str = r#"SELECT 45 operation_audit.audit_sequence, 46 CASE WHEN typeof(operation_audit.correlation_id) = 'blob' 47 AND length(operation_audit.correlation_id) = 32 48 THEN operation_audit.correlation_id ELSE NULL END AS correlation_id, 49 CASE WHEN typeof(operation_audit.audit_kind) = 'text' 50 AND length(CAST(operation_audit.audit_kind AS BLOB)) <= 40 51 THEN operation_audit.audit_kind ELSE NULL END AS audit_kind, 52 CASE WHEN typeof(operation_audit.outcome) = 'text' 53 AND length(CAST(operation_audit.outcome AS BLOB)) <= 16 54 THEN operation_audit.outcome ELSE NULL END AS outcome, 55 CASE WHEN typeof(operation_audit.reason_code) = 'text' 56 AND length(CAST(operation_audit.reason_code AS BLOB)) <= 48 57 THEN operation_audit.reason_code ELSE NULL END AS reason_code, 58 operation_audit.occurred_at_unix_ms, 59 CASE WHEN typeof(request_audit.operation_id) = 'blob' 60 AND length(request_audit.operation_id) = 32 61 THEN request_audit.operation_id ELSE NULL END AS operation_id, 62 CASE WHEN typeof(request_audit.audit_kind) = 'text' 63 AND length(CAST(request_audit.audit_kind AS BLOB)) <= 40 64 THEN request_audit.audit_kind ELSE NULL END AS request_audit_kind, 65 CASE WHEN typeof(request_record.correlation_id) = 'blob' 66 AND length(request_record.correlation_id) = 32 67 THEN request_record.correlation_id ELSE NULL END AS request_correlation_id 68 FROM operation_audit 69 LEFT JOIN nip46_request_audit AS request_audit 70 ON request_audit.audit_sequence = operation_audit.audit_sequence 71 LEFT JOIN nip46_requests AS request_record 72 ON request_record.operation_id = request_audit.operation_id 73 WHERE operation_audit.audit_id = ? LIMIT 2"#; 74 const INSERT_AUDIT_SQL: &str = r#"INSERT INTO operation_audit ( 75 audit_sequence, audit_id, correlation_id, audit_kind, outcome, reason_code, 76 occurred_at_unix_ms 77 ) VALUES (?, ?, ?, ?, ?, ?, ?)"#; 78 const INSERT_REQUEST_AUDIT_SQL: &str = r#"INSERT INTO nip46_request_audit ( 79 operation_id, audit_kind, audit_sequence 80 ) VALUES (?, ?, ?)"#; 81 82 const READ_RATE_SQL: &str = r#"SELECT 83 window_started_at_unix_ms, window_ends_at_unix_ms, 84 accepted_count, rejected_count, lifetime_accepted_count, 85 lifetime_rejected_count, last_observed_at_unix_ms, retention_expires_at_unix_ms 86 FROM connection_rate_windows 87 WHERE rate_kind = ? AND subject_scope = ? AND subject_sha256 = ? 88 LIMIT 2"#; 89 const COUNT_RATE_SUBJECTS_SQL: &str = r#"SELECT COUNT(*) 90 FROM connection_rate_windows 91 WHERE rate_kind = ? AND subject_scope = ?"#; 92 const INSERT_RATE_SQL: &str = r#"INSERT INTO connection_rate_windows ( 93 rate_kind, subject_scope, subject_sha256, window_started_at_unix_ms, 94 window_ends_at_unix_ms, accepted_count, rejected_count, 95 lifetime_accepted_count, lifetime_rejected_count, 96 last_observed_at_unix_ms, retention_expires_at_unix_ms 97 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 98 const UPDATE_RATE_SQL: &str = r#"UPDATE connection_rate_windows SET 99 window_started_at_unix_ms = ?, window_ends_at_unix_ms = ?, 100 accepted_count = ?, rejected_count = ?, lifetime_accepted_count = ?, 101 lifetime_rejected_count = ?, last_observed_at_unix_ms = ?, 102 retention_expires_at_unix_ms = ? 103 WHERE rate_kind = ? AND subject_scope = ? AND subject_sha256 = ? 104 AND last_observed_at_unix_ms = ?"#; 105 106 /// Stable construction and validation failure classes. 107 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 108 pub enum MycGovernanceStateErrorKind { 109 InvalidRatePolicy, 110 InvalidRelayId, 111 InvalidPageLimit, 112 InvalidCompactionPolicy, 113 } 114 115 impl MycGovernanceStateErrorKind { 116 /// Returns the stable machine-readable error code. 117 #[must_use] 118 pub const fn code(self) -> &'static str { 119 match self { 120 Self::InvalidRatePolicy => "governance_rate_policy_invalid", 121 Self::InvalidRelayId => "governance_relay_id_invalid", 122 Self::InvalidPageLimit => "governance_audit_page_limit_invalid", 123 Self::InvalidCompactionPolicy => "governance_compaction_policy_invalid", 124 } 125 } 126 } 127 128 /// Source-free governance-state validation failure. 129 #[derive(Clone, Copy, PartialEq, Eq)] 130 pub struct MycGovernanceStateError { 131 kind: MycGovernanceStateErrorKind, 132 } 133 134 impl MycGovernanceStateError { 135 const fn new(kind: MycGovernanceStateErrorKind) -> Self { 136 Self { kind } 137 } 138 139 /// Returns the stable failure class. 140 #[must_use] 141 pub const fn kind(self) -> MycGovernanceStateErrorKind { 142 self.kind 143 } 144 145 /// Returns the stable machine-readable error code. 146 #[must_use] 147 pub const fn code(self) -> &'static str { 148 self.kind.code() 149 } 150 } 151 152 impl fmt::Display for MycGovernanceStateError { 153 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 154 formatter.write_str(match self.kind { 155 MycGovernanceStateErrorKind::InvalidRatePolicy => "rate policy is invalid", 156 MycGovernanceStateErrorKind::InvalidRelayId => "rate relay ID is invalid", 157 MycGovernanceStateErrorKind::InvalidPageLimit => "audit page limit is invalid", 158 MycGovernanceStateErrorKind::InvalidCompactionPolicy => "compaction policy is invalid", 159 }) 160 } 161 } 162 163 impl fmt::Debug for MycGovernanceStateError { 164 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 165 formatter 166 .debug_struct("MycGovernanceStateError") 167 .field("kind", &self.kind) 168 .finish() 169 } 170 } 171 172 impl Error for MycGovernanceStateError {} 173 174 /// Closed rate-window classes. Their subject scope is fixed by contract. 175 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 176 pub enum MycRateLimitClass { 177 ConnectionAdmission, 178 ChallengeCreation, 179 ChallengeAuthorization, 180 } 181 182 impl MycRateLimitClass { 183 pub(crate) const fn as_str(self) -> &'static str { 184 match self { 185 Self::ConnectionAdmission => "connection_admission", 186 Self::ChallengeCreation => "challenge_creation", 187 Self::ChallengeAuthorization => "challenge_authorization", 188 } 189 } 190 } 191 192 /// Validated bounded rate-window policy. 193 #[derive(Clone, Copy, PartialEq, Eq)] 194 pub struct MycRateLimitPolicy { 195 class: MycRateLimitClass, 196 window_ms: u64, 197 max_attempts: u32, 198 retention_ms: u64, 199 maximum_tracked_subjects: u32, 200 } 201 202 impl MycRateLimitPolicy { 203 /// Constructs one exact policy from fully normalized configuration values. 204 pub fn new( 205 class: MycRateLimitClass, 206 window_ms: u64, 207 max_attempts: u32, 208 retention_ms: u64, 209 maximum_tracked_subjects: u32, 210 ) -> Result<Self, MycGovernanceStateError> { 211 if window_ms == 0 212 || window_ms > MYC_RATE_WINDOW_MAX_MS 213 || max_attempts == 0 214 || max_attempts > MYC_RATE_MAX_ATTEMPTS 215 || retention_ms < window_ms 216 || retention_ms > MYC_RATE_RETENTION_MAX_MS 217 || maximum_tracked_subjects == 0 218 || maximum_tracked_subjects > MYC_RATE_MAX_TRACKED_SUBJECTS 219 { 220 return Err(MycGovernanceStateError::new( 221 MycGovernanceStateErrorKind::InvalidRatePolicy, 222 )); 223 } 224 Ok(Self { 225 class, 226 window_ms, 227 max_attempts, 228 retention_ms, 229 maximum_tracked_subjects, 230 }) 231 } 232 233 /// Returns the closed policy class. 234 #[must_use] 235 pub const fn class(self) -> MycRateLimitClass { 236 self.class 237 } 238 } 239 240 impl fmt::Debug for MycRateLimitPolicy { 241 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 242 formatter 243 .debug_struct("MycRateLimitPolicy") 244 .field("class", &self.class) 245 .field("values", &"[bounded]") 246 .finish() 247 } 248 } 249 250 #[derive(Clone, PartialEq, Eq)] 251 pub(crate) struct MycGovernancePolicies { 252 connection_admission: MycRateLimitPolicy, 253 challenge_creation: MycRateLimitPolicy, 254 challenge_authorization: MycRateLimitPolicy, 255 audit_retention_ms: u64, 256 relays: Box<[MycRateRelayId]>, 257 } 258 259 impl MycGovernancePolicies { 260 pub(crate) fn new( 261 connection_admission: MycRateLimitPolicy, 262 challenge_creation: MycRateLimitPolicy, 263 challenge_authorization: MycRateLimitPolicy, 264 audit_retention_ms: u64, 265 relays: Box<[MycRateRelayId]>, 266 ) -> Result<Self, MycGovernanceStateError> { 267 if connection_admission.class != MycRateLimitClass::ConnectionAdmission 268 || challenge_creation.class != MycRateLimitClass::ChallengeCreation 269 || challenge_authorization.class != MycRateLimitClass::ChallengeAuthorization 270 || audit_retention_ms == 0 271 || audit_retention_ms > MYC_AUDIT_RETENTION_MAX_MS 272 || relays.is_empty() 273 { 274 return Err(MycGovernanceStateError::new( 275 MycGovernanceStateErrorKind::InvalidRatePolicy, 276 )); 277 } 278 Ok(Self { 279 connection_admission, 280 challenge_creation, 281 challenge_authorization, 282 audit_retention_ms, 283 relays, 284 }) 285 } 286 287 pub(crate) const fn rate_policy(&self, class: MycRateLimitClass) -> MycRateLimitPolicy { 288 match class { 289 MycRateLimitClass::ConnectionAdmission => self.connection_admission, 290 MycRateLimitClass::ChallengeCreation => self.challenge_creation, 291 MycRateLimitClass::ChallengeAuthorization => self.challenge_authorization, 292 } 293 } 294 295 pub(crate) fn admits_relay(&self, relay: &MycRateRelayId) -> bool { 296 self.relays.iter().any(|configured| configured == relay) 297 } 298 299 pub(crate) const fn audit_retention_ms(&self) -> u64 { 300 self.audit_retention_ms 301 } 302 } 303 304 /// Validated configured relay identity used only to derive bounded rate scope. 305 #[derive(Clone, PartialEq, Eq)] 306 pub struct MycRateRelayId(Box<str>); 307 308 impl MycRateRelayId { 309 /// Validates one stable lower-snake relay identity before allocation. 310 pub fn new(value: &str) -> Result<Self, MycGovernanceStateError> { 311 let bytes = value.as_bytes(); 312 if bytes.is_empty() 313 || bytes.len() > MYC_RATE_RELAY_ID_MAX_BYTES 314 || !bytes[0].is_ascii_lowercase() 315 || !bytes[bytes.len() - 1].is_ascii_alphanumeric() 316 || bytes.windows(2).any(|pair| pair == b"__") 317 || !bytes 318 .iter() 319 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_') 320 { 321 return Err(MycGovernanceStateError::new( 322 MycGovernanceStateErrorKind::InvalidRelayId, 323 )); 324 } 325 Ok(Self(value.into())) 326 } 327 } 328 329 impl fmt::Debug for MycRateRelayId { 330 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 331 formatter.write_str("MycRateRelayId([redacted])") 332 } 333 } 334 335 /// Injected stable correlation identity for non-NIP-46 operator work. 336 #[derive(Clone, Copy, PartialEq, Eq)] 337 pub struct MycAuditCorrelationId([u8; 32]); 338 339 impl MycAuditCorrelationId { 340 /// Constructs an opaque caller-injected correlation identity. 341 #[must_use] 342 pub const fn new(bytes: [u8; 32]) -> Self { 343 Self(bytes) 344 } 345 346 pub(crate) const fn as_bytes(&self) -> &[u8; 32] { 347 &self.0 348 } 349 } 350 351 impl fmt::Debug for MycAuditCorrelationId { 352 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 353 formatter.write_str("MycAuditCorrelationId([redacted])") 354 } 355 } 356 357 /// Closed audit event classes. 358 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 359 pub enum MycAuditKind { 360 ConnectionAdmission, 361 ConnectionOperatorDecision, 362 ConnectionExpiry, 363 ChallengeCreation, 364 ChallengeAuthorization, 365 GovernanceCompaction, 366 } 367 368 impl MycAuditKind { 369 pub(crate) const fn as_str(self) -> &'static str { 370 match self { 371 Self::ConnectionAdmission => "connection_admission", 372 Self::ConnectionOperatorDecision => "connection_operator_decision", 373 Self::ConnectionExpiry => "connection_expiry", 374 Self::ChallengeCreation => "challenge_creation", 375 Self::ChallengeAuthorization => "challenge_authorization", 376 Self::GovernanceCompaction => "governance_compaction", 377 } 378 } 379 380 pub(crate) fn parse(value: &str) -> Option<Self> { 381 match value { 382 "connection_admission" => Some(Self::ConnectionAdmission), 383 "connection_operator_decision" => Some(Self::ConnectionOperatorDecision), 384 "connection_expiry" => Some(Self::ConnectionExpiry), 385 "challenge_creation" => Some(Self::ChallengeCreation), 386 "challenge_authorization" => Some(Self::ChallengeAuthorization), 387 "governance_compaction" => Some(Self::GovernanceCompaction), 388 _ => None, 389 } 390 } 391 } 392 393 /// Closed audit outcomes. 394 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 395 pub enum MycAuditOutcome { 396 Succeeded, 397 Rejected, 398 Failed, 399 } 400 401 impl MycAuditOutcome { 402 pub(crate) const fn as_str(self) -> &'static str { 403 match self { 404 Self::Succeeded => "succeeded", 405 Self::Rejected => "rejected", 406 Self::Failed => "failed", 407 } 408 } 409 410 pub(crate) fn parse(value: &str) -> Option<Self> { 411 match value { 412 "succeeded" => Some(Self::Succeeded), 413 "rejected" => Some(Self::Rejected), 414 "failed" => Some(Self::Failed), 415 _ => None, 416 } 417 } 418 } 419 420 /// Closed safe audit reasons. 421 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 422 pub enum MycAuditReasonCode { 423 Trusted, 424 ApprovalRequired, 425 PolicyDenied, 426 OperatorApproved, 427 OperatorDenied, 428 ConnectionExpired, 429 ChallengeRequired, 430 ChallengeAuthorized, 431 ChallengeExpired, 432 RateLimited, 433 Compacted, 434 } 435 436 impl MycAuditReasonCode { 437 pub(crate) const fn as_str(self) -> &'static str { 438 match self { 439 Self::Trusted => "trusted", 440 Self::ApprovalRequired => "approval_required", 441 Self::PolicyDenied => "policy_denied", 442 Self::OperatorApproved => "operator_approved", 443 Self::OperatorDenied => "operator_denied", 444 Self::ConnectionExpired => "connection_expired", 445 Self::ChallengeRequired => "challenge_required", 446 Self::ChallengeAuthorized => "challenge_authorized", 447 Self::ChallengeExpired => "challenge_expired", 448 Self::RateLimited => "rate_limited", 449 Self::Compacted => "compacted", 450 } 451 } 452 453 fn parse(value: &str) -> Option<Self> { 454 match value { 455 "trusted" => Some(Self::Trusted), 456 "approval_required" => Some(Self::ApprovalRequired), 457 "policy_denied" => Some(Self::PolicyDenied), 458 "operator_approved" => Some(Self::OperatorApproved), 459 "operator_denied" => Some(Self::OperatorDenied), 460 "connection_expired" => Some(Self::ConnectionExpired), 461 "challenge_required" => Some(Self::ChallengeRequired), 462 "challenge_authorized" => Some(Self::ChallengeAuthorized), 463 "challenge_expired" => Some(Self::ChallengeExpired), 464 "rate_limited" => Some(Self::RateLimited), 465 "compacted" => Some(Self::Compacted), 466 _ => None, 467 } 468 } 469 } 470 471 /// One validated safe audit record. 472 #[derive(Clone, Copy, PartialEq, Eq)] 473 pub struct MycAuditRecord { 474 sequence: u64, 475 correlation_id: MycAuditCorrelationId, 476 operation_id: Option<MycSignerOperationId>, 477 kind: MycAuditKind, 478 outcome: MycAuditOutcome, 479 reason: MycAuditReasonCode, 480 occurred_at: MycConnectionTimeUnixMs, 481 } 482 483 impl MycAuditRecord { 484 #[must_use] 485 pub const fn sequence(&self) -> u64 { 486 self.sequence 487 } 488 489 /// Returns the opaque correlation identity bound to this record. 490 #[must_use] 491 pub const fn correlation_id(&self) -> MycAuditCorrelationId { 492 self.correlation_id 493 } 494 495 /// Returns the typed signer-operation identity when this is request audit. 496 #[must_use] 497 pub const fn operation_id(&self) -> Option<MycSignerOperationId> { 498 self.operation_id 499 } 500 501 #[must_use] 502 pub const fn kind(&self) -> MycAuditKind { 503 self.kind 504 } 505 506 #[must_use] 507 pub const fn outcome(&self) -> MycAuditOutcome { 508 self.outcome 509 } 510 511 #[must_use] 512 pub const fn reason(&self) -> MycAuditReasonCode { 513 self.reason 514 } 515 516 #[must_use] 517 pub const fn occurred_at(&self) -> MycConnectionTimeUnixMs { 518 self.occurred_at 519 } 520 } 521 522 impl fmt::Debug for MycAuditRecord { 523 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 524 formatter 525 .debug_struct("MycAuditRecord") 526 .field("sequence", &self.sequence) 527 .field("kind", &self.kind) 528 .field("outcome", &self.outcome) 529 .field("reason", &self.reason) 530 .field("correlation", &"[redacted]") 531 .finish() 532 } 533 } 534 535 /// Validated audit page limit. 536 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 537 pub struct MycAuditPageLimit(u16); 538 539 impl MycAuditPageLimit { 540 pub fn new(value: u16) -> Result<Self, MycGovernanceStateError> { 541 (value != 0 && value <= MYC_AUDIT_PAGE_MAX_ITEMS) 542 .then_some(Self(value)) 543 .ok_or_else(|| { 544 MycGovernanceStateError::new(MycGovernanceStateErrorKind::InvalidPageLimit) 545 }) 546 } 547 } 548 549 /// Immutable bounded audit page. 550 pub struct MycAuditPage { 551 snapshot_sequence: u64, 552 items: Box<[MycAuditRecord]>, 553 next_before_sequence: Option<u64>, 554 } 555 556 impl MycAuditPage { 557 #[must_use] 558 pub const fn snapshot_sequence(&self) -> u64 { 559 self.snapshot_sequence 560 } 561 562 #[must_use] 563 pub fn items(&self) -> &[MycAuditRecord] { 564 &self.items 565 } 566 567 #[must_use] 568 pub const fn next_before_sequence(&self) -> Option<u64> { 569 self.next_before_sequence 570 } 571 } 572 573 impl fmt::Debug for MycAuditPage { 574 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 575 formatter 576 .debug_struct("MycAuditPage") 577 .field("snapshot_sequence", &self.snapshot_sequence) 578 .field("item_count", &self.items.len()) 579 .finish() 580 } 581 } 582 583 #[derive(Clone, Copy)] 584 pub(crate) struct MycAdminAuditQuery { 585 limit: MycAuditPageLimit, 586 snapshot_sequence: Option<u64>, 587 before_sequence: Option<u64>, 588 from_unix_ms: Option<u64>, 589 to_unix_ms: Option<u64>, 590 kind: Option<MycAuditKind>, 591 outcome: Option<MycAuditOutcome>, 592 } 593 594 impl MycAdminAuditQuery { 595 pub(crate) fn new( 596 limit: MycAuditPageLimit, 597 snapshot_sequence: Option<u64>, 598 before_sequence: Option<u64>, 599 from_unix_ms: Option<u64>, 600 to_unix_ms: Option<u64>, 601 kind: Option<MycAuditKind>, 602 outcome: Option<MycAuditOutcome>, 603 ) -> Result<Self, MycStateRepositoryError> { 604 if from_unix_ms.is_some_and(|from| to_unix_ms.is_some_and(|to| from > to)) 605 || from_unix_ms.is_some_and(|value| i64::try_from(value).is_err()) 606 || to_unix_ms.is_some_and(|value| i64::try_from(value).is_err()) 607 { 608 return Err(MycStateRepositoryError::new( 609 MycStateRepositoryErrorKind::Binding, 610 )); 611 } 612 Ok(Self { 613 limit, 614 snapshot_sequence, 615 before_sequence, 616 from_unix_ms, 617 to_unix_ms, 618 kind, 619 outcome, 620 }) 621 } 622 } 623 624 /// Explicit bounded retention and compaction policy. 625 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 626 pub struct MycGovernanceCompactionPolicy { 627 audit_retention_ms: u64, 628 maximum_rows_per_class: u16, 629 } 630 631 impl MycGovernanceCompactionPolicy { 632 pub fn new( 633 audit_retention_ms: u64, 634 maximum_rows_per_class: u16, 635 ) -> Result<Self, MycGovernanceStateError> { 636 if audit_retention_ms == 0 637 || audit_retention_ms > MYC_AUDIT_RETENTION_MAX_MS 638 || maximum_rows_per_class == 0 639 || maximum_rows_per_class > MYC_COMPACTION_MAX_ROWS 640 { 641 return Err(MycGovernanceStateError::new( 642 MycGovernanceStateErrorKind::InvalidCompactionPolicy, 643 )); 644 } 645 Ok(Self { 646 audit_retention_ms, 647 maximum_rows_per_class, 648 }) 649 } 650 } 651 652 /// Bounded compaction result. Authoritative domain state is never included. 653 /// Counts describe this invocation; an exact correlation replay is a zero-work success. 654 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 655 pub struct MycGovernanceCompactionOutcome { 656 removed_audit_records: u16, 657 removed_rate_subjects: u16, 658 } 659 660 impl MycGovernanceCompactionOutcome { 661 #[must_use] 662 pub const fn removed_audit_records(self) -> u16 { 663 self.removed_audit_records 664 } 665 666 #[must_use] 667 pub const fn removed_rate_subjects(self) -> u16 { 668 self.removed_rate_subjects 669 } 670 } 671 672 #[derive(Clone, Copy)] 673 pub(crate) struct AuditEvidence { 674 pub(crate) correlation: MycAuditCorrelationId, 675 pub(crate) occurred_at: MycConnectionTimeUnixMs, 676 pub(crate) operation_id: Option<MycSignerOperationId>, 677 } 678 679 #[derive(Clone, Copy)] 680 pub(crate) enum RateSubject { 681 Global, 682 Relay([u8; 32]), 683 Connection(MycConnectionId), 684 } 685 686 pub(crate) fn relay_subject(relay: &MycRateRelayId) -> RateSubject { 687 RateSubject::Relay(hash_framed(RELAY_SUBJECT_DOMAIN, relay.0.as_bytes())) 688 } 689 690 pub(crate) fn global_subject() -> RateSubject { 691 RateSubject::Global 692 } 693 694 pub(crate) fn connection_subject(connection: MycConnectionId) -> RateSubject { 695 RateSubject::Connection(connection) 696 } 697 698 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 699 pub(crate) enum GovernanceOperationError { 700 Binding, 701 Storage, 702 } 703 704 pub(crate) async fn govern_rate_attempt( 705 transaction: &mut ServiceSqliteTransaction<'_>, 706 policy: MycRateLimitPolicy, 707 expected_class: MycRateLimitClass, 708 subjects: &[RateSubject], 709 evidence: AuditEvidence, 710 audit_kind: MycAuditKind, 711 ) -> Result<bool, GovernanceOperationError> { 712 if policy.class != expected_class || subjects.is_empty() || subjects.len() > 2 { 713 return Err(GovernanceOperationError::Binding); 714 } 715 if let Some(existing) = read_audit_by_id( 716 transaction, 717 derive_audit_id(evidence.correlation, audit_kind), 718 ) 719 .await? 720 { 721 return (existing.correlation_id == evidence.correlation 722 && existing.operation_id == evidence.operation_id 723 && existing.kind == audit_kind 724 && existing.outcome == MycAuditOutcome::Rejected 725 && existing.reason == MycAuditReasonCode::RateLimited) 726 .then_some(false) 727 .ok_or(GovernanceOperationError::Binding); 728 } 729 let mut states = Vec::with_capacity(subjects.len()); 730 for subject in subjects { 731 states.push(load_rate_state(transaction, policy, *subject, evidence.occurred_at).await?); 732 } 733 let admitted = states.iter().all(|state| state.can_accept(policy)); 734 for state in &states { 735 persist_rate_state(transaction, policy, *state, evidence.occurred_at, admitted).await?; 736 } 737 if !admitted { 738 record_audit( 739 transaction, 740 evidence, 741 audit_kind, 742 MycAuditOutcome::Rejected, 743 MycAuditReasonCode::RateLimited, 744 ) 745 .await?; 746 } 747 Ok(admitted) 748 } 749 750 pub(crate) async fn record_audit( 751 transaction: &mut ServiceSqliteTransaction<'_>, 752 evidence: AuditEvidence, 753 kind: MycAuditKind, 754 outcome: MycAuditOutcome, 755 reason: MycAuditReasonCode, 756 ) -> Result<MycAuditRecord, GovernanceOperationError> { 757 let audit_id = derive_audit_id(evidence.correlation, kind); 758 if let Some(existing) = read_audit_by_id(transaction, audit_id).await? { 759 return (existing.correlation_id == evidence.correlation 760 && existing.operation_id == evidence.operation_id 761 && existing.kind == kind 762 && existing.outcome == outcome 763 && existing.reason == reason 764 && existing.occurred_at == evidence.occurred_at) 765 .then_some(existing) 766 .ok_or(GovernanceOperationError::Binding); 767 } 768 let current = read_audit_state(transaction).await?; 769 let sequence = current 770 .checked_add(1) 771 .filter(|value| *value <= i64::MAX as u64) 772 .ok_or(GovernanceOperationError::Binding)?; 773 let result = sqlx::query(ADVANCE_AUDIT_STATE_SQL) 774 .bind(i64::try_from(sequence).map_err(|_| GovernanceOperationError::Binding)?) 775 .bind(i64::try_from(current).map_err(|_| GovernanceOperationError::Binding)?) 776 .execute(&mut *transaction) 777 .await 778 .map_err(|_| GovernanceOperationError::Storage)?; 779 require_one(result.rows_affected())?; 780 let result = sqlx::query(INSERT_AUDIT_SQL) 781 .bind(i64::try_from(sequence).map_err(|_| GovernanceOperationError::Binding)?) 782 .bind(audit_id.as_slice()) 783 .bind(evidence.correlation.0.as_slice()) 784 .bind(kind.as_str()) 785 .bind(outcome.as_str()) 786 .bind(reason.as_str()) 787 .bind(to_i64(evidence.occurred_at.get())?) 788 .execute(&mut *transaction) 789 .await 790 .map_err(|_| GovernanceOperationError::Storage)?; 791 require_one(result.rows_affected())?; 792 if let Some(operation_id) = evidence.operation_id { 793 let result = sqlx::query(INSERT_REQUEST_AUDIT_SQL) 794 .bind(operation_id.as_bytes().as_slice()) 795 .bind(kind.as_str()) 796 .bind(i64::try_from(sequence).map_err(|_| GovernanceOperationError::Binding)?) 797 .execute(&mut *transaction) 798 .await 799 .map_err(|_| GovernanceOperationError::Storage)?; 800 require_one(result.rows_affected())?; 801 } 802 Ok(MycAuditRecord { 803 sequence, 804 correlation_id: evidence.correlation, 805 operation_id: evidence.operation_id, 806 kind, 807 outcome, 808 reason, 809 occurred_at: evidence.occurred_at, 810 }) 811 } 812 813 impl MycStateRepository<'_> { 814 pub(crate) async fn read_admin_audit_page( 815 &self, 816 query: MycAdminAuditQuery, 817 ) -> Result<MycAuditPage, MycStateRepositoryError> { 818 let expected = PersistedMetadata::from(self.expected()); 819 self.host() 820 .transaction(move |transaction| { 821 Box::pin(async move { 822 verify_metadata(transaction, &expected).await?; 823 read_admin_audit_page(transaction, query).await 824 }) 825 }) 826 .await 827 .map_err(map_transaction_error) 828 } 829 830 pub(crate) async fn read_admin_audit_summary( 831 &self, 832 from_unix_ms: u64, 833 to_unix_ms: u64, 834 ) -> Result<Box<[(String, u64)]>, MycStateRepositoryError> { 835 if from_unix_ms > to_unix_ms 836 || i64::try_from(from_unix_ms).is_err() 837 || i64::try_from(to_unix_ms).is_err() 838 { 839 return Err(MycStateRepositoryError::new( 840 MycStateRepositoryErrorKind::Binding, 841 )); 842 } 843 let expected = PersistedMetadata::from(self.expected()); 844 self.host() 845 .transaction(move |transaction| { 846 Box::pin(async move { 847 verify_metadata(transaction, &expected).await?; 848 let rows = sqlx::query( 849 r#"SELECT audit_kind, outcome, COUNT(*) AS item_count 850 FROM operation_audit 851 WHERE occurred_at_unix_ms BETWEEN ? AND ? 852 GROUP BY audit_kind, outcome 853 ORDER BY audit_kind ASC, outcome ASC 854 LIMIT 19"#, 855 ) 856 .bind( 857 i64::try_from(from_unix_ms) 858 .map_err(|_| GovernanceOperationError::Binding)?, 859 ) 860 .bind(i64::try_from(to_unix_ms).map_err(|_| GovernanceOperationError::Binding)?) 861 .fetch_all(&mut *transaction) 862 .await 863 .map_err(|_| GovernanceOperationError::Storage)?; 864 let mut counts = Vec::with_capacity(rows.len()); 865 for row in rows { 866 let kind = row 867 .try_get::<&str, _>("audit_kind") 868 .ok() 869 .and_then(MycAuditKind::parse) 870 .ok_or(GovernanceOperationError::Binding)?; 871 let outcome = row 872 .try_get::<&str, _>("outcome") 873 .ok() 874 .and_then(MycAuditOutcome::parse) 875 .ok_or(GovernanceOperationError::Binding)?; 876 let count = row 877 .try_get::<i64, _>("item_count") 878 .ok() 879 .and_then(|value| u64::try_from(value).ok()) 880 .ok_or(GovernanceOperationError::Binding)?; 881 counts.push((format!("{}.{}", kind.as_str(), outcome.as_str()), count)); 882 } 883 Ok(counts.into_boxed_slice()) 884 }) 885 }) 886 .await 887 .map_err(map_transaction_error) 888 } 889 890 /// Reads one deterministic, snapshot-bounded page of safe audit evidence. 891 pub async fn read_audit_page( 892 &self, 893 limit: MycAuditPageLimit, 894 snapshot_sequence: Option<u64>, 895 before_sequence: Option<u64>, 896 ) -> Result<MycAuditPage, MycStateRepositoryError> { 897 let expected = PersistedMetadata::from(self.expected()); 898 self.host() 899 .transaction(move |transaction| { 900 Box::pin(async move { 901 verify_metadata(transaction, &expected).await?; 902 read_audit_page(transaction, limit, snapshot_sequence, before_sequence).await 903 }) 904 }) 905 .await 906 .map_err(map_transaction_error) 907 } 908 909 /// Removes only expired safe audit and rate-window evidence in bounded batches. 910 pub async fn compact_governance_evidence( 911 &self, 912 observed_at: MycConnectionTimeUnixMs, 913 policy: MycGovernanceCompactionPolicy, 914 correlation: MycAuditCorrelationId, 915 ) -> Result<MycGovernanceCompactionOutcome, MycStateRepositoryError> { 916 if policy.audit_retention_ms != self.expected().governance_audit_retention_ms() { 917 return Err(MycStateRepositoryError::new( 918 MycStateRepositoryErrorKind::Binding, 919 )); 920 } 921 let expected = PersistedMetadata::from(self.expected()); 922 self.host() 923 .transaction(move |transaction| { 924 Box::pin(async move { 925 verify_metadata(transaction, &expected).await?; 926 compact_governance(transaction, observed_at, policy, correlation).await 927 }) 928 }) 929 .await 930 .map_err(map_transaction_error) 931 } 932 } 933 934 #[derive(Clone, Copy)] 935 struct RateState { 936 subject: RateSubject, 937 existing_last: Option<u64>, 938 window_started: u64, 939 window_ends: u64, 940 accepted: u64, 941 rejected: u64, 942 lifetime_accepted: u64, 943 lifetime_rejected: u64, 944 } 945 946 impl RateState { 947 fn can_accept(self, policy: MycRateLimitPolicy) -> bool { 948 self.accepted < u64::from(policy.max_attempts) 949 } 950 } 951 952 async fn load_rate_state( 953 transaction: &mut ServiceSqliteTransaction<'_>, 954 policy: MycRateLimitPolicy, 955 subject: RateSubject, 956 observed_at: MycConnectionTimeUnixMs, 957 ) -> Result<RateState, GovernanceOperationError> { 958 let (scope, digest) = subject_parts(subject); 959 let rows = sqlx::query(READ_RATE_SQL) 960 .bind(policy.class.as_str()) 961 .bind(scope) 962 .bind(digest.as_slice()) 963 .fetch_all(&mut *transaction) 964 .await 965 .map_err(|_| GovernanceOperationError::Storage)?; 966 if rows.len() > 1 { 967 return Err(GovernanceOperationError::Binding); 968 } 969 let now = observed_at.get(); 970 let Some(row) = rows.first() else { 971 if scope != "global" { 972 let count = sqlx::query_scalar::<_, i64>(COUNT_RATE_SUBJECTS_SQL) 973 .bind(policy.class.as_str()) 974 .bind(scope) 975 .fetch_one(&mut *transaction) 976 .await 977 .map_err(|_| GovernanceOperationError::Storage)?; 978 let count = u64::try_from(count).map_err(|_| GovernanceOperationError::Binding)?; 979 if count >= u64::from(policy.maximum_tracked_subjects) { 980 return Err(GovernanceOperationError::Binding); 981 } 982 } 983 return Ok(RateState { 984 subject, 985 existing_last: None, 986 window_started: now, 987 window_ends: checked_time_add(now, policy.window_ms)?, 988 accepted: 0, 989 rejected: 0, 990 lifetime_accepted: 0, 991 lifetime_rejected: 0, 992 }); 993 }; 994 let mut state = RateState { 995 subject, 996 existing_last: Some(bounded_i64(row, "last_observed_at_unix_ms")?), 997 window_started: bounded_i64(row, "window_started_at_unix_ms")?, 998 window_ends: bounded_i64(row, "window_ends_at_unix_ms")?, 999 accepted: bounded_nonnegative(row, "accepted_count")?, 1000 rejected: bounded_nonnegative(row, "rejected_count")?, 1001 lifetime_accepted: bounded_nonnegative(row, "lifetime_accepted_count")?, 1002 lifetime_rejected: bounded_nonnegative(row, "lifetime_rejected_count")?, 1003 }; 1004 let retention_expires = bounded_i64(row, "retention_expires_at_unix_ms")?; 1005 let last = state 1006 .existing_last 1007 .ok_or(GovernanceOperationError::Binding)?; 1008 if now < last 1009 || state.window_started > last 1010 || state.window_ends < last 1011 || retention_expires < last 1012 || state.accepted > u64::from(policy.max_attempts) 1013 { 1014 return Err(GovernanceOperationError::Binding); 1015 } 1016 if now > state.window_ends { 1017 state.window_started = now; 1018 state.window_ends = checked_time_add(now, policy.window_ms)?; 1019 state.accepted = 0; 1020 state.rejected = 0; 1021 } 1022 Ok(state) 1023 } 1024 1025 async fn persist_rate_state( 1026 transaction: &mut ServiceSqliteTransaction<'_>, 1027 policy: MycRateLimitPolicy, 1028 state: RateState, 1029 observed_at: MycConnectionTimeUnixMs, 1030 admitted: bool, 1031 ) -> Result<(), GovernanceOperationError> { 1032 let (scope, digest) = subject_parts(state.subject); 1033 let accepted = state.accepted + u64::from(admitted); 1034 let rejected = state.rejected + u64::from(!admitted); 1035 let lifetime_accepted = state.lifetime_accepted + u64::from(admitted); 1036 let lifetime_rejected = state.lifetime_rejected + u64::from(!admitted); 1037 for value in [accepted, rejected, lifetime_accepted, lifetime_rejected] { 1038 if value > i64::MAX as u64 { 1039 return Err(GovernanceOperationError::Binding); 1040 } 1041 } 1042 let retention_expires = checked_time_add(observed_at.get(), policy.retention_ms)?; 1043 let mut query = if state.existing_last.is_some() { 1044 sqlx::query(UPDATE_RATE_SQL) 1045 } else { 1046 sqlx::query(INSERT_RATE_SQL) 1047 }; 1048 if state.existing_last.is_some() { 1049 query = query 1050 .bind(to_i64(state.window_started)?) 1051 .bind(to_i64(state.window_ends)?) 1052 .bind(to_i64(accepted)?) 1053 .bind(to_i64(rejected)?) 1054 .bind(to_i64(lifetime_accepted)?) 1055 .bind(to_i64(lifetime_rejected)?) 1056 .bind(to_i64(observed_at.get())?) 1057 .bind(to_i64(retention_expires)?) 1058 .bind(policy.class.as_str()) 1059 .bind(scope) 1060 .bind(digest.as_slice()) 1061 .bind(to_i64( 1062 state 1063 .existing_last 1064 .ok_or(GovernanceOperationError::Binding)?, 1065 )?); 1066 } else { 1067 query = query 1068 .bind(policy.class.as_str()) 1069 .bind(scope) 1070 .bind(digest.as_slice()) 1071 .bind(to_i64(state.window_started)?) 1072 .bind(to_i64(state.window_ends)?) 1073 .bind(to_i64(accepted)?) 1074 .bind(to_i64(rejected)?) 1075 .bind(to_i64(lifetime_accepted)?) 1076 .bind(to_i64(lifetime_rejected)?) 1077 .bind(to_i64(observed_at.get())?) 1078 .bind(to_i64(retention_expires)?); 1079 } 1080 let result = query 1081 .execute(&mut *transaction) 1082 .await 1083 .map_err(|_| GovernanceOperationError::Storage)?; 1084 require_one(result.rows_affected()) 1085 } 1086 1087 async fn read_audit_page( 1088 transaction: &mut ServiceSqliteTransaction<'_>, 1089 limit: MycAuditPageLimit, 1090 requested_snapshot: Option<u64>, 1091 before: Option<u64>, 1092 ) -> Result<MycAuditPage, GovernanceOperationError> { 1093 let high_water = read_audit_state(transaction).await?; 1094 let snapshot = requested_snapshot.unwrap_or(high_water); 1095 if snapshot > high_water 1096 || before.is_some_and(|value| value == 0 || value > snapshot.saturating_add(1)) 1097 { 1098 return Err(GovernanceOperationError::Binding); 1099 } 1100 let before = before.unwrap_or_else(|| snapshot.saturating_add(1)); 1101 let fetch_limit = u32::from(limit.0) + 1; 1102 let rows = sqlx::query( 1103 r#"SELECT audit_id FROM operation_audit 1104 WHERE audit_sequence <= ? AND audit_sequence < ? 1105 ORDER BY audit_sequence DESC LIMIT ?"#, 1106 ) 1107 .bind(to_i64(snapshot)?) 1108 .bind(to_i64(before)?) 1109 .bind(i64::from(fetch_limit)) 1110 .fetch_all(&mut *transaction) 1111 .await 1112 .map_err(|_| GovernanceOperationError::Storage)?; 1113 if rows.len() > usize::try_from(fetch_limit).map_err(|_| GovernanceOperationError::Binding)? { 1114 return Err(GovernanceOperationError::Binding); 1115 } 1116 let has_more = rows.len() > usize::from(limit.0); 1117 let mut items = Vec::with_capacity(rows.len().min(usize::from(limit.0))); 1118 for row in rows.iter().take(usize::from(limit.0)) { 1119 let id = bounded_digest(row, "audit_id")?; 1120 items.push( 1121 read_audit_by_id(transaction, id) 1122 .await? 1123 .ok_or(GovernanceOperationError::Binding)?, 1124 ); 1125 } 1126 let next = has_more 1127 .then(|| items.last().map(|item| item.sequence)) 1128 .flatten(); 1129 Ok(MycAuditPage { 1130 snapshot_sequence: snapshot, 1131 items: items.into_boxed_slice(), 1132 next_before_sequence: next, 1133 }) 1134 } 1135 1136 #[allow(clippy::too_many_arguments)] 1137 async fn read_admin_audit_page( 1138 transaction: &mut ServiceSqliteTransaction<'_>, 1139 query: MycAdminAuditQuery, 1140 ) -> Result<MycAuditPage, GovernanceOperationError> { 1141 let high_water = read_audit_state(transaction).await?; 1142 let snapshot = query.snapshot_sequence.unwrap_or(high_water); 1143 if snapshot > high_water 1144 || query 1145 .before_sequence 1146 .is_some_and(|value| value == 0 || value > snapshot.saturating_add(1)) 1147 { 1148 return Err(GovernanceOperationError::Binding); 1149 } 1150 let before = query 1151 .before_sequence 1152 .unwrap_or_else(|| snapshot.saturating_add(1)); 1153 let fetch_limit = u32::from(query.limit.0) + 1; 1154 let from = query 1155 .from_unix_ms 1156 .map(i64::try_from) 1157 .transpose() 1158 .map_err(|_| GovernanceOperationError::Binding)?; 1159 let to = query 1160 .to_unix_ms 1161 .map(i64::try_from) 1162 .transpose() 1163 .map_err(|_| GovernanceOperationError::Binding)?; 1164 let kind = query.kind.map(MycAuditKind::as_str); 1165 let outcome = query.outcome.map(MycAuditOutcome::as_str); 1166 let rows = sqlx::query( 1167 r#"SELECT audit_id FROM operation_audit 1168 WHERE audit_sequence <= ? AND audit_sequence < ? 1169 AND (? IS NULL OR occurred_at_unix_ms >= ?) 1170 AND (? IS NULL OR occurred_at_unix_ms <= ?) 1171 AND (? IS NULL OR audit_kind = ?) 1172 AND (? IS NULL OR outcome = ?) 1173 ORDER BY audit_sequence DESC LIMIT ?"#, 1174 ) 1175 .bind(to_i64(snapshot)?) 1176 .bind(to_i64(before)?) 1177 .bind(from) 1178 .bind(from) 1179 .bind(to) 1180 .bind(to) 1181 .bind(kind) 1182 .bind(kind) 1183 .bind(outcome) 1184 .bind(outcome) 1185 .bind(i64::from(fetch_limit)) 1186 .fetch_all(&mut *transaction) 1187 .await 1188 .map_err(|_| GovernanceOperationError::Storage)?; 1189 if rows.len() > usize::try_from(fetch_limit).map_err(|_| GovernanceOperationError::Binding)? { 1190 return Err(GovernanceOperationError::Binding); 1191 } 1192 let has_more = rows.len() > usize::from(query.limit.0); 1193 let mut items = Vec::with_capacity(rows.len().min(usize::from(query.limit.0))); 1194 for row in rows.iter().take(usize::from(query.limit.0)) { 1195 let id = bounded_digest(row, "audit_id")?; 1196 items.push( 1197 read_audit_by_id(transaction, id) 1198 .await? 1199 .ok_or(GovernanceOperationError::Binding)?, 1200 ); 1201 } 1202 let next = has_more 1203 .then(|| items.last().map(|item| item.sequence)) 1204 .flatten(); 1205 Ok(MycAuditPage { 1206 snapshot_sequence: snapshot, 1207 items: items.into_boxed_slice(), 1208 next_before_sequence: next, 1209 }) 1210 } 1211 1212 async fn compact_governance( 1213 transaction: &mut ServiceSqliteTransaction<'_>, 1214 observed_at: MycConnectionTimeUnixMs, 1215 policy: MycGovernanceCompactionPolicy, 1216 correlation: MycAuditCorrelationId, 1217 ) -> Result<MycGovernanceCompactionOutcome, GovernanceOperationError> { 1218 if let Some(existing) = read_audit_by_id( 1219 transaction, 1220 derive_audit_id(correlation, MycAuditKind::GovernanceCompaction), 1221 ) 1222 .await? 1223 { 1224 return (existing.correlation_id == correlation 1225 && existing.operation_id.is_none() 1226 && existing.kind == MycAuditKind::GovernanceCompaction 1227 && existing.outcome == MycAuditOutcome::Succeeded 1228 && existing.reason == MycAuditReasonCode::Compacted 1229 && existing.occurred_at == observed_at) 1230 .then_some(MycGovernanceCompactionOutcome { 1231 removed_audit_records: 0, 1232 removed_rate_subjects: 0, 1233 }) 1234 .ok_or(GovernanceOperationError::Binding); 1235 } 1236 let cutoff = observed_at.get().saturating_sub(policy.audit_retention_ms); 1237 let limit = i64::from(policy.maximum_rows_per_class); 1238 let request_links = sqlx::query( 1239 r#"DELETE FROM nip46_request_audit WHERE audit_sequence IN ( 1240 SELECT audit_sequence FROM operation_audit 1241 WHERE occurred_at_unix_ms < ? ORDER BY audit_sequence ASC LIMIT ? 1242 )"#, 1243 ) 1244 .bind(to_i64(cutoff)?) 1245 .bind(limit) 1246 .execute(&mut *transaction) 1247 .await 1248 .map_err(|_| GovernanceOperationError::Storage)?; 1249 let _ = request_links; 1250 let audit = sqlx::query( 1251 r#"DELETE FROM operation_audit WHERE audit_sequence IN ( 1252 SELECT audit_sequence FROM operation_audit 1253 WHERE occurred_at_unix_ms < ? ORDER BY audit_sequence ASC LIMIT ? 1254 )"#, 1255 ) 1256 .bind(to_i64(cutoff)?) 1257 .bind(limit) 1258 .execute(&mut *transaction) 1259 .await 1260 .map_err(|_| GovernanceOperationError::Storage)?; 1261 let rates = sqlx::query( 1262 r#"DELETE FROM connection_rate_windows WHERE rowid IN ( 1263 SELECT rowid FROM connection_rate_windows 1264 WHERE retention_expires_at_unix_ms < ? 1265 ORDER BY retention_expires_at_unix_ms ASC, rate_kind ASC, 1266 subject_scope ASC, subject_sha256 ASC LIMIT ? 1267 )"#, 1268 ) 1269 .bind(to_i64(observed_at.get())?) 1270 .bind(limit) 1271 .execute(&mut *transaction) 1272 .await 1273 .map_err(|_| GovernanceOperationError::Storage)?; 1274 let removed_audit_records = 1275 u16::try_from(audit.rows_affected()).map_err(|_| GovernanceOperationError::Binding)?; 1276 let removed_rate_subjects = 1277 u16::try_from(rates.rows_affected()).map_err(|_| GovernanceOperationError::Binding)?; 1278 record_audit( 1279 transaction, 1280 AuditEvidence { 1281 correlation, 1282 occurred_at: observed_at, 1283 operation_id: None, 1284 }, 1285 MycAuditKind::GovernanceCompaction, 1286 MycAuditOutcome::Succeeded, 1287 MycAuditReasonCode::Compacted, 1288 ) 1289 .await?; 1290 Ok(MycGovernanceCompactionOutcome { 1291 removed_audit_records, 1292 removed_rate_subjects, 1293 }) 1294 } 1295 1296 async fn read_audit_state( 1297 transaction: &mut ServiceSqliteTransaction<'_>, 1298 ) -> Result<u64, GovernanceOperationError> { 1299 let rows = sqlx::query(READ_AUDIT_STATE_SQL) 1300 .fetch_all(&mut *transaction) 1301 .await 1302 .map_err(|_| GovernanceOperationError::Storage)?; 1303 if rows.len() != 1 { 1304 return Err(GovernanceOperationError::Binding); 1305 } 1306 bounded_nonnegative(&rows[0], "next_sequence") 1307 } 1308 1309 async fn read_audit_by_id( 1310 transaction: &mut ServiceSqliteTransaction<'_>, 1311 audit_id: [u8; 32], 1312 ) -> Result<Option<MycAuditRecord>, GovernanceOperationError> { 1313 let rows = sqlx::query(READ_AUDIT_BY_ID_SQL) 1314 .bind(audit_id.as_slice()) 1315 .fetch_all(&mut *transaction) 1316 .await 1317 .map_err(|_| GovernanceOperationError::Storage)?; 1318 if rows.len() > 1 { 1319 return Err(GovernanceOperationError::Binding); 1320 } 1321 rows.first().map(decode_audit).transpose() 1322 } 1323 1324 fn decode_audit(row: &sqlx::sqlite::SqliteRow) -> Result<MycAuditRecord, GovernanceOperationError> { 1325 let correlation_id = MycAuditCorrelationId(bounded_digest(row, "correlation_id")?); 1326 let kind = bounded_text(row, "audit_kind") 1327 .and_then(|value| MycAuditKind::parse(&value).ok_or(GovernanceOperationError::Binding))?; 1328 let outcome = bounded_text(row, "outcome").and_then(|value| { 1329 MycAuditOutcome::parse(&value).ok_or(GovernanceOperationError::Binding) 1330 })?; 1331 let reason = bounded_text(row, "reason_code").and_then(|value| { 1332 MycAuditReasonCode::parse(&value).ok_or(GovernanceOperationError::Binding) 1333 })?; 1334 let operation_id = 1335 optional_digest(row, "operation_id")?.map(MycSignerOperationId::from_persisted); 1336 let request_audit_kind = optional_text(row, "request_audit_kind")?; 1337 let request_correlation_id = optional_digest(row, "request_correlation_id")?; 1338 let request_bound = matches!( 1339 kind, 1340 MycAuditKind::ConnectionAdmission 1341 | MycAuditKind::ChallengeCreation 1342 | MycAuditKind::ChallengeAuthorization 1343 ); 1344 if request_bound 1345 != (operation_id.is_some() 1346 && request_audit_kind.as_deref() == Some(kind.as_str()) 1347 && request_correlation_id == Some(correlation_id.0)) 1348 || (!request_bound 1349 && (operation_id.is_some() 1350 || request_audit_kind.is_some() 1351 || request_correlation_id.is_some())) 1352 { 1353 return Err(GovernanceOperationError::Binding); 1354 } 1355 let sequence = bounded_i64(row, "audit_sequence")?; 1356 let occurred_at = MycConnectionTimeUnixMs::new(bounded_i64(row, "occurred_at_unix_ms")?) 1357 .map_err(|_| GovernanceOperationError::Binding)?; 1358 Ok(MycAuditRecord { 1359 sequence, 1360 correlation_id, 1361 operation_id, 1362 kind, 1363 outcome, 1364 reason, 1365 occurred_at, 1366 }) 1367 } 1368 1369 async fn verify_metadata( 1370 transaction: &mut ServiceSqliteTransaction<'_>, 1371 expected: &PersistedMetadata, 1372 ) -> Result<(), GovernanceOperationError> { 1373 require_expected_metadata(transaction, expected) 1374 .await 1375 .map_err(|error| match error { 1376 RepositoryOperationError::Binding => GovernanceOperationError::Binding, 1377 RepositoryOperationError::Storage => GovernanceOperationError::Storage, 1378 }) 1379 } 1380 1381 fn map_transaction_error( 1382 error: ServiceSqliteTransactionError<GovernanceOperationError>, 1383 ) -> MycStateRepositoryError { 1384 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 1385 return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); 1386 } 1387 let kind = match error.operation_error() { 1388 Some(GovernanceOperationError::Binding) => MycStateRepositoryErrorKind::Binding, 1389 Some(GovernanceOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, 1390 }; 1391 MycStateRepositoryError::new(kind) 1392 } 1393 1394 fn subject_parts(subject: RateSubject) -> (&'static str, [u8; 32]) { 1395 match subject { 1396 RateSubject::Global => ("global", hash_framed(GLOBAL_SUBJECT_DOMAIN, b"global")), 1397 RateSubject::Relay(digest) => ("relay", digest), 1398 RateSubject::Connection(connection) => ( 1399 "connection", 1400 hash_framed(CONNECTION_SUBJECT_DOMAIN, connection.as_bytes()), 1401 ), 1402 } 1403 } 1404 1405 fn derive_audit_id(correlation: MycAuditCorrelationId, kind: MycAuditKind) -> [u8; 32] { 1406 let mut hasher = Sha256::new(); 1407 hasher.update(AUDIT_ID_DOMAIN); 1408 hasher.update(correlation.0); 1409 hasher.update((kind.as_str().len() as u64).to_be_bytes()); 1410 hasher.update(kind.as_str().as_bytes()); 1411 hasher.finalize().into() 1412 } 1413 1414 fn hash_framed(domain: &[u8], value: &[u8]) -> [u8; 32] { 1415 let mut hasher = Sha256::new(); 1416 hasher.update(domain); 1417 hasher.update((value.len() as u64).to_be_bytes()); 1418 hasher.update(value); 1419 hasher.finalize().into() 1420 } 1421 1422 fn checked_time_add(left: u64, right: u64) -> Result<u64, GovernanceOperationError> { 1423 left.checked_add(right) 1424 .filter(|value| *value <= i64::MAX as u64) 1425 .ok_or(GovernanceOperationError::Binding) 1426 } 1427 1428 fn to_i64(value: u64) -> Result<i64, GovernanceOperationError> { 1429 i64::try_from(value).map_err(|_| GovernanceOperationError::Binding) 1430 } 1431 1432 fn bounded_i64( 1433 row: &sqlx::sqlite::SqliteRow, 1434 column: &str, 1435 ) -> Result<u64, GovernanceOperationError> { 1436 let value = row 1437 .try_get::<i64, _>(column) 1438 .map_err(|_| GovernanceOperationError::Binding)?; 1439 u64::try_from(value).map_err(|_| GovernanceOperationError::Binding) 1440 } 1441 1442 fn bounded_nonnegative( 1443 row: &sqlx::sqlite::SqliteRow, 1444 column: &str, 1445 ) -> Result<u64, GovernanceOperationError> { 1446 bounded_i64(row, column) 1447 } 1448 1449 fn bounded_digest( 1450 row: &sqlx::sqlite::SqliteRow, 1451 column: &str, 1452 ) -> Result<[u8; 32], GovernanceOperationError> { 1453 row.try_get::<Option<Vec<u8>>, _>(column) 1454 .map_err(|_| GovernanceOperationError::Binding)? 1455 .ok_or(GovernanceOperationError::Binding)? 1456 .try_into() 1457 .map_err(|_| GovernanceOperationError::Binding) 1458 } 1459 1460 fn bounded_text( 1461 row: &sqlx::sqlite::SqliteRow, 1462 column: &str, 1463 ) -> Result<Box<str>, GovernanceOperationError> { 1464 row.try_get::<Option<String>, _>(column) 1465 .map_err(|_| GovernanceOperationError::Binding)? 1466 .ok_or(GovernanceOperationError::Binding) 1467 .map(String::into_boxed_str) 1468 } 1469 1470 fn optional_digest( 1471 row: &sqlx::sqlite::SqliteRow, 1472 column: &str, 1473 ) -> Result<Option<[u8; 32]>, GovernanceOperationError> { 1474 row.try_get::<Option<Vec<u8>>, _>(column) 1475 .map_err(|_| GovernanceOperationError::Binding)? 1476 .map(|value| { 1477 value 1478 .try_into() 1479 .map_err(|_| GovernanceOperationError::Binding) 1480 }) 1481 .transpose() 1482 } 1483 1484 fn optional_text( 1485 row: &sqlx::sqlite::SqliteRow, 1486 column: &str, 1487 ) -> Result<Option<Box<str>>, GovernanceOperationError> { 1488 row.try_get::<Option<String>, _>(column) 1489 .map_err(|_| GovernanceOperationError::Binding) 1490 .map(|value| value.map(String::into_boxed_str)) 1491 } 1492 1493 fn require_one(rows: u64) -> Result<(), GovernanceOperationError> { 1494 (rows == 1) 1495 .then_some(()) 1496 .ok_or(GovernanceOperationError::Binding) 1497 }