myc

Self-custodial remote signer for Radroots apps
git clone https://radroots.dev/git/myc.git
Log | Files | Refs | README | LICENSE

commit fdde56bc7f851660efc46d0e995044eb43427af7
parent 9368d381f9b470e855396a586e078e8ff33d497d
Author: triesap <tyson@radroots.org>
Date:   Fri, 21 Aug 2026 17:33:07 +0000

state: govern audit and rate evidence

- bind rate policy to normalized configuration and persist exact schema-v5 catalogs
- add atomic replay-safe audit, bounded rate windows, pagination, and compaction
- preserve authoritative connection and challenge state across retention work

Diffstat:
MREADME | 21+++++++++++++++++----
Msrc/lib.rs | 18++++++++++++++----
Msrc/state_catalog.rs | 382++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Msrc/state_connection.rs | 313++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Asrc/state_governance.rs | 1300+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/state_host.rs | 8++++++--
Msrc/state_metadata.rs | 78++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_connection_state.rs | 771+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Mtests/services_hardening_signer_request_state.rs | 13+++++++++++--
Mtests/services_hardening_state_catalog.rs | 62+++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mtests/services_hardening_state_repository.rs | 7++++++-
11 files changed, 2881 insertions(+), 92 deletions(-)

diff --git a/README b/README @@ -47,10 +47,10 @@ without changing the canonical common artifact inventory. Create-new initialization reserves the shared schema-v1 metadata and migration ledger, retains exclusive writer authority, applies the exact Myc schema-v2 -through schema-v4 migrations, binds the normalized configuration, expected +through schema-v5 migrations, binds the normalized configuration, expected identity roles, and policy versions through a sealed typed repository, and explicitly closes the host before reporting success. Existing writable open can -resume the exact v1, v2, or v3 prefix; read-only inspection requires the current +resume any exact v1 through v4 prefix; read-only inspection requires the current catalog and exact immutable Myc binding. The public Myc repository exposes no raw pool, connection, transaction-control @@ -59,8 +59,9 @@ inside the shared `ServiceSqliteTransaction` runner. Provider and relay work cannot occur inside that transaction boundary. The v2 binding table, v3 request/dedup tables, and v4 connection, permission, request-decision, and authorization-challenge tables and guards are checksum-pinned service-owned -schema objects; later workflow tables remain owned by their ordered repository -steps. +schema objects. Schema v5 adds checksum-pinned audit-sequence, safe operation +audit, request-audit binding, and bounded rate-window objects; later workflow +tables remain owned by their ordered repository steps. Signer-request admission validates bounded client, request, event, method, canonical request, injected operation entropy, and injected time evidence @@ -85,6 +86,18 @@ injected entropy and time, and transition once to authorized or expired. Exact retries return the original durable identity or terminal state, including after explicit close and reopen. +Connection admission uses one bounded global window plus one configured relay +window, so arbitrary client keys cannot create persistent subjects before a +connection exists. Challenge creation and authorization use distinct stable +connection-scoped windows. Exact saturation creates no connection, decision, +challenge, or authorization transition beyond the accepted bound, and a +retained rate rejection is replay-stable. Safe audit records preserve exact +typed operation/correlation bindings and use closed category, outcome, and +reason vocabularies; snapshot pagination is bounded and Debug output redacts +identity. Explicit bounded compaction is correlation-idempotent and retires +only expired rate-window and old safe-audit evidence. It never removes schema +history, metadata, connections, decisions, permissions, or challenges. + Writable hosts expose Myc-bound online-backup and active-integrity operations. Backup verification retains the exact admitted member inode, and offline staging derives the same runtime paths, database identity, migration diff --git a/src/lib.rs b/src/lib.rs @@ -27,6 +27,7 @@ mod signing_adapter; pub mod sql; mod state_catalog; mod state_connection; +mod state_governance; mod state_host; mod state_maintenance; mod state_metadata; @@ -119,13 +120,14 @@ pub use state_catalog::{ MYC_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_3_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_3_SHA256, MYC_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_4_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_4_SHA256, - MycStateCatalogError, MycStateCatalogErrorKind, myc_migration_catalog, myc_schema_catalog, - validate_myc_state_catalogs, + MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT, + MYC_STATE_SCHEMA_VERSION_5_SHA256, MycStateCatalogError, MycStateCatalogErrorKind, + myc_migration_catalog, myc_schema_catalog, validate_myc_state_catalogs, }; pub use state_connection::{ MYC_AUTHORIZATION_CHALLENGE_URL_MAX_BYTES, MYC_CONNECTION_PERMISSION_MAX_COUNT, - MycAuthorizationChallengeAdmission, MycAuthorizationChallengeId, - MycAuthorizationChallengeNonce, MycAuthorizationChallengeRecord, + MycAuthorizationChallengeAdmission, MycAuthorizationChallengeAuthorization, + MycAuthorizationChallengeId, MycAuthorizationChallengeNonce, MycAuthorizationChallengeRecord, MycAuthorizationChallengeRequest, MycAuthorizationChallengeState, MycAuthorizationChallengeUrl, MycConnectionAdmission, MycConnectionAdmissionPolicy, MycConnectionAdmissionRequest, MycConnectionDecision, MycConnectionDecisionRecord, MycConnectionId, MycConnectionNonce, @@ -133,6 +135,14 @@ pub use state_connection::{ MycConnectionPolicyGeneration, MycConnectionRecord, MycConnectionStateError, MycConnectionStateErrorKind, MycConnectionStatus, MycConnectionTimeUnixMs, }; +pub use state_governance::{ + MYC_AUDIT_PAGE_MAX_ITEMS, MYC_AUDIT_RETENTION_MAX_MS, MYC_COMPACTION_MAX_ROWS, + MYC_RATE_MAX_ATTEMPTS, MYC_RATE_MAX_TRACKED_SUBJECTS, MYC_RATE_RELAY_ID_MAX_BYTES, + MYC_RATE_RETENTION_MAX_MS, MYC_RATE_WINDOW_MAX_MS, MycAuditCorrelationId, MycAuditKind, + MycAuditOutcome, MycAuditPage, MycAuditPageLimit, MycAuditReasonCode, MycAuditRecord, + MycGovernanceCompactionOutcome, MycGovernanceCompactionPolicy, MycGovernanceStateError, + MycGovernanceStateErrorKind, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, +}; pub use state_host::{ MycStateHost, MycStateHostError, MycStateHostErrorKind, MycStateHostMode, initialize_myc_state, open_myc_state_inspection, open_myc_state_read_write, diff --git a/src/state_catalog.rs b/src/state_catalog.rs @@ -12,7 +12,7 @@ use radroots_service_sqlite::{ pub const MYC_STATE_BASE_SCHEMA_VERSION: u32 = 1; /// The newest governed Myc state schema understood by this binary. -pub const MYC_STATE_SCHEMA_VERSION: u32 = 4; +pub const MYC_STATE_SCHEMA_VERSION: u32 = 5; /// The shared metadata and migration-ledger objects present at schema v1. pub const MYC_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6; @@ -26,6 +26,9 @@ pub const MYC_STATE_SCHEMA_VERSION_3_OBJECT_COUNT: u32 = 13; /// The shared objects plus Myc metadata, request, and connection objects at schema v4. pub const MYC_STATE_SCHEMA_VERSION_4_OBJECT_COUNT: u32 = 25; +/// The shared objects plus all Myc metadata, request, connection, and governance objects at v5. +pub const MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT: u32 = 34; + /// SHA-256 identity of the exact schema-v1 object snapshot. pub const MYC_STATE_SCHEMA_VERSION_1_SHA256: [u8; 32] = [ 0x94, 0xdc, 0x66, 0xfb, 0xca, 0x60, 0x16, 0x79, 0x61, 0x5c, 0x05, 0x52, 0x29, 0xdc, 0x0d, 0xb6, @@ -40,14 +43,14 @@ pub const MYC_STATE_SCHEMA_VERSION_2_SHA256: [u8; 32] = [ /// SHA-256 identity of the ordered Myc migration catalog. pub const MYC_MIGRATION_CATALOG_SHA256: [u8; 32] = [ - 0x45, 0x3d, 0x99, 0xf4, 0xc1, 0x9c, 0x09, 0x4c, 0x59, 0x2f, 0x1a, 0x3f, 0xe7, 0xe2, 0x8e, 0x82, - 0xc1, 0xde, 0xa6, 0x74, 0x51, 0x4a, 0x9e, 0x7d, 0x06, 0x3b, 0xc4, 0x66, 0xc2, 0x0b, 0xd4, 0xbd, + 0x1d, 0x06, 0x99, 0x96, 0x43, 0x55, 0x60, 0x21, 0x7d, 0xbd, 0x36, 0x92, 0xd5, 0x43, 0x06, 0x12, + 0xad, 0x3a, 0x56, 0x67, 0xee, 0x0b, 0x12, 0xb1, 0x91, 0x9f, 0x0b, 0x5e, 0xdf, 0x2c, 0xbe, 0x7d, ]; /// SHA-256 identity of the schema catalog bound to the migration catalog. pub const MYC_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [ - 0xa4, 0x7d, 0x9f, 0x0a, 0xf8, 0x04, 0x20, 0x2b, 0xfa, 0x07, 0x26, 0x8e, 0x08, 0xfb, 0x48, 0x4d, - 0xed, 0x59, 0x0b, 0x78, 0x5c, 0xc4, 0x2a, 0x01, 0x09, 0x4d, 0x1a, 0x8b, 0x9f, 0x37, 0xee, 0xca, + 0x21, 0x19, 0xef, 0xef, 0xcf, 0xd4, 0xac, 0x54, 0x77, 0x60, 0x96, 0x55, 0x34, 0x1a, 0xa4, 0xc5, + 0xb9, 0xb9, 0xa1, 0xbe, 0xfb, 0x88, 0xcc, 0x16, 0x8a, 0x2b, 0x93, 0x73, 0xe9, 0x9e, 0xbb, 0xd5, ]; /// SHA-256 identity of the schema-v2 migration content. @@ -80,6 +83,18 @@ pub const MYC_STATE_SCHEMA_VERSION_4_SHA256: [u8; 32] = [ 0x7c, 0xdd, 0x3b, 0x24, 0xa9, 0x3a, 0xb4, 0x82, 0xd7, 0x74, 0xbc, 0xec, 0x02, 0x18, 0xc1, 0x74, ]; +/// SHA-256 identity of the schema-v5 migration content. +pub const MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256: [u8; 32] = [ + 0x0e, 0x00, 0x4f, 0xcb, 0x5d, 0x0b, 0xc7, 0xc9, 0x51, 0xb1, 0x6f, 0x43, 0x34, 0x53, 0x3d, 0x10, + 0xef, 0xcb, 0x18, 0xa8, 0x26, 0xf6, 0x24, 0xde, 0x09, 0x22, 0xd8, 0x6d, 0xc8, 0x50, 0x4b, 0x33, +]; + +/// SHA-256 identity of the schema-v5 object snapshot. +pub const MYC_STATE_SCHEMA_VERSION_5_SHA256: [u8; 32] = [ + 0xfe, 0x89, 0xd4, 0xaf, 0x7d, 0xe4, 0xed, 0xed, 0x78, 0xa3, 0xf0, 0x62, 0xdc, 0x9d, 0x30, 0x0f, + 0x06, 0xcf, 0x0f, 0x4c, 0xfc, 0xa5, 0xfa, 0x5c, 0xff, 0xc9, 0x3a, 0xf1, 0x14, 0x6f, 0x21, 0xc5, +]; + /// SHA-256 identity of the Myc metadata table definition. const MYC_STATE_METADATA_TABLE_SHA256: [u8; 32] = [ 0x16, 0x17, 0x46, 0xa2, 0x26, 0x42, 0x46, 0x2f, 0x2b, 0xdb, 0x08, 0x5b, 0xae, 0xde, 0xb2, 0x3b, @@ -567,6 +582,218 @@ const CREATE_CONNECTION_STATE_MIGRATION_SQL: &str = concat!( connection_auth_challenges_no_delete_sql!(), ); +macro_rules! myc_audit_state_table_sql { + () => { + r#"CREATE TABLE myc_audit_state ( + singleton INTEGER NOT NULL PRIMARY KEY CHECK (singleton = 1), + next_sequence INTEGER NOT NULL CHECK (next_sequence BETWEEN 0 AND 9223372036854775807) +) STRICT"# + }; +} + +macro_rules! operation_audit_table_sql { + () => { + r#"CREATE TABLE operation_audit ( + audit_sequence INTEGER NOT NULL PRIMARY KEY + CHECK (audit_sequence BETWEEN 1 AND 9223372036854775807), + audit_id BLOB NOT NULL UNIQUE CHECK (length(audit_id) = 32), + correlation_id BLOB NOT NULL CHECK (length(correlation_id) = 32), + audit_kind TEXT NOT NULL CHECK (audit_kind IN ( + 'connection_admission', + 'connection_operator_decision', + 'connection_expiry', + 'challenge_creation', + 'challenge_authorization', + 'governance_compaction' + )), + outcome TEXT NOT NULL CHECK (outcome IN ('succeeded', 'rejected', 'failed')), + reason_code TEXT NOT NULL CHECK (reason_code IN ( + 'trusted', + 'approval_required', + 'policy_denied', + 'operator_approved', + 'operator_denied', + 'connection_expired', + 'challenge_required', + 'challenge_authorized', + 'challenge_expired', + 'rate_limited', + 'compacted' + )), + occurred_at_unix_ms INTEGER NOT NULL + CHECK (occurred_at_unix_ms BETWEEN 1 AND 9223372036854775807), + UNIQUE (correlation_id, audit_kind) +) STRICT"# + }; +} + +macro_rules! nip46_request_audit_table_sql { + () => { + r#"CREATE TABLE nip46_request_audit ( + operation_id BLOB NOT NULL CHECK (length(operation_id) = 32) + REFERENCES nip46_requests(operation_id), + audit_kind TEXT NOT NULL CHECK (audit_kind IN ( + 'connection_admission', 'challenge_creation', 'challenge_authorization' + )), + audit_sequence INTEGER NOT NULL UNIQUE + CHECK (audit_sequence BETWEEN 1 AND 9223372036854775807) + REFERENCES operation_audit(audit_sequence), + PRIMARY KEY (operation_id, audit_kind) +) STRICT"# + }; +} + +macro_rules! connection_rate_windows_table_sql { + () => { + r#"CREATE TABLE connection_rate_windows ( + rate_kind TEXT NOT NULL CHECK (rate_kind IN ( + 'connection_admission', 'challenge_creation', 'challenge_authorization' + )), + subject_scope TEXT NOT NULL CHECK (subject_scope IN ('global', 'relay', 'connection')), + subject_sha256 BLOB NOT NULL CHECK (length(subject_sha256) = 32), + window_started_at_unix_ms INTEGER NOT NULL + CHECK (window_started_at_unix_ms BETWEEN 1 AND 9223372036854775807), + window_ends_at_unix_ms INTEGER NOT NULL + CHECK (window_ends_at_unix_ms BETWEEN window_started_at_unix_ms AND 9223372036854775807), + accepted_count INTEGER NOT NULL CHECK (accepted_count BETWEEN 0 AND 10000), + rejected_count INTEGER NOT NULL CHECK (rejected_count BETWEEN 0 AND 9223372036854775807), + lifetime_accepted_count INTEGER NOT NULL + CHECK (lifetime_accepted_count BETWEEN accepted_count AND 9223372036854775807), + lifetime_rejected_count INTEGER NOT NULL + CHECK (lifetime_rejected_count BETWEEN rejected_count AND 9223372036854775807), + last_observed_at_unix_ms INTEGER NOT NULL + CHECK (last_observed_at_unix_ms BETWEEN window_started_at_unix_ms AND 9223372036854775807), + retention_expires_at_unix_ms INTEGER NOT NULL + CHECK (retention_expires_at_unix_ms BETWEEN last_observed_at_unix_ms AND 9223372036854775807), + CHECK ( + (rate_kind = 'connection_admission' AND subject_scope IN ('global', 'relay')) + OR (rate_kind IN ('challenge_creation', 'challenge_authorization') + AND subject_scope = 'connection') + ), + PRIMARY KEY (rate_kind, subject_scope, subject_sha256) +) STRICT"# + }; +} + +macro_rules! myc_audit_state_guard_update_sql { + () => { + r#"CREATE TRIGGER myc_audit_state_guard_update +BEFORE UPDATE ON myc_audit_state +WHEN NEW.singleton != OLD.singleton + OR OLD.next_sequence = 9223372036854775807 + OR NEW.next_sequence != OLD.next_sequence + 1 +BEGIN + SELECT RAISE(ABORT, 'audit sequence transition is invalid'); +END"# + }; +} + +macro_rules! myc_audit_state_no_delete_sql { + () => { + r#"CREATE TRIGGER myc_audit_state_no_delete +BEFORE DELETE ON myc_audit_state +BEGIN + SELECT RAISE(ABORT, 'audit sequence authority is retained'); +END"# + }; +} + +macro_rules! operation_audit_no_update_sql { + () => { + r#"CREATE TRIGGER operation_audit_no_update +BEFORE UPDATE ON operation_audit +BEGIN + SELECT RAISE(ABORT, 'operation audit is immutable'); +END"# + }; +} + +macro_rules! nip46_request_audit_no_update_sql { + () => { + r#"CREATE TRIGGER nip46_request_audit_no_update +BEFORE UPDATE ON nip46_request_audit +BEGIN + SELECT RAISE(ABORT, 'request audit binding is immutable'); +END"# + }; +} + +macro_rules! connection_rate_windows_guard_update_sql { + () => { + r#"CREATE TRIGGER connection_rate_windows_guard_update +BEFORE UPDATE ON connection_rate_windows +WHEN NEW.rate_kind != OLD.rate_kind + OR NEW.subject_scope != OLD.subject_scope + OR NEW.subject_sha256 != OLD.subject_sha256 + OR NEW.window_started_at_unix_ms < OLD.window_started_at_unix_ms + OR NEW.window_ends_at_unix_ms < NEW.window_started_at_unix_ms + OR NEW.lifetime_accepted_count < OLD.lifetime_accepted_count + OR NEW.lifetime_rejected_count < OLD.lifetime_rejected_count + OR NEW.last_observed_at_unix_ms < OLD.last_observed_at_unix_ms + OR NEW.retention_expires_at_unix_ms < NEW.last_observed_at_unix_ms + OR NOT ( + (NEW.window_started_at_unix_ms = OLD.window_started_at_unix_ms + AND NEW.window_ends_at_unix_ms = OLD.window_ends_at_unix_ms + AND ( + (NEW.accepted_count = OLD.accepted_count + 1 + AND NEW.rejected_count = OLD.rejected_count + AND NEW.lifetime_accepted_count = OLD.lifetime_accepted_count + 1 + AND NEW.lifetime_rejected_count = OLD.lifetime_rejected_count) + OR (NEW.accepted_count = OLD.accepted_count + AND NEW.rejected_count = OLD.rejected_count + 1 + AND NEW.lifetime_accepted_count = OLD.lifetime_accepted_count + AND NEW.lifetime_rejected_count = OLD.lifetime_rejected_count + 1) + )) + OR (NEW.window_started_at_unix_ms > OLD.window_ends_at_unix_ms + AND NEW.window_ends_at_unix_ms > NEW.window_started_at_unix_ms + AND ((NEW.accepted_count = 1 AND NEW.rejected_count = 0) + OR (NEW.accepted_count = 0 AND NEW.rejected_count = 1)) + AND NEW.lifetime_accepted_count = OLD.lifetime_accepted_count + NEW.accepted_count + AND NEW.lifetime_rejected_count = OLD.lifetime_rejected_count + NEW.rejected_count) + ) +BEGIN + SELECT RAISE(ABORT, 'rate-window transition is invalid'); +END"# + }; +} + +const CREATE_MYC_AUDIT_STATE_TABLE_SQL: &str = myc_audit_state_table_sql!(); +const CREATE_OPERATION_AUDIT_TABLE_SQL: &str = operation_audit_table_sql!(); +const CREATE_NIP46_REQUEST_AUDIT_TABLE_SQL: &str = nip46_request_audit_table_sql!(); +const CREATE_CONNECTION_RATE_WINDOWS_TABLE_SQL: &str = connection_rate_windows_table_sql!(); +const CREATE_MYC_AUDIT_STATE_GUARD_UPDATE_SQL: &str = myc_audit_state_guard_update_sql!(); +const CREATE_MYC_AUDIT_STATE_NO_DELETE_SQL: &str = myc_audit_state_no_delete_sql!(); +const CREATE_OPERATION_AUDIT_NO_UPDATE_SQL: &str = operation_audit_no_update_sql!(); +const CREATE_NIP46_REQUEST_AUDIT_NO_UPDATE_SQL: &str = nip46_request_audit_no_update_sql!(); +const CREATE_CONNECTION_RATE_WINDOWS_GUARD_UPDATE_SQL: &str = + connection_rate_windows_guard_update_sql!(); + +const CREATE_GOVERNANCE_STATE_MIGRATION_SQL: &str = concat!( + "DROP TRIGGER myc_state_metadata_no_update;\n", + "UPDATE myc_state_metadata SET state_contract_version = CASE ", + "WHEN state_contract_version = 4 THEN 5 ELSE 0 END WHERE singleton = 1;\n", + myc_state_metadata_no_update_sql!(), + ";\n", + myc_audit_state_table_sql!(), + ";\n", + "INSERT INTO myc_audit_state (singleton, next_sequence) VALUES (1, 0);\n", + operation_audit_table_sql!(), + ";\n", + nip46_request_audit_table_sql!(), + ";\n", + connection_rate_windows_table_sql!(), + ";\n", + myc_audit_state_guard_update_sql!(), + ";\n", + myc_audit_state_no_delete_sql!(), + ";\n", + operation_audit_no_update_sql!(), + ";\n", + nip46_request_audit_no_update_sql!(), + ";\n", + connection_rate_windows_guard_update_sql!(), +); + const CONNECTIONS_TABLE_SHA256: [u8; 32] = [ 0x72, 0xd5, 0xd8, 0xba, 0x24, 0x68, 0x9c, 0x93, 0x34, 0xb3, 0x8f, 0xbf, 0x64, 0x21, 0xe1, 0x65, 0xfd, 0xc3, 0x80, 0x46, 0xf1, 0x3f, 0x56, 0x49, 0x3a, 0xef, 0xd7, 0x42, 0xc4, 0xe6, 0x49, 0x85, @@ -615,6 +842,42 @@ const CONNECTION_AUTH_CHALLENGES_NO_DELETE_SHA256: [u8; 32] = [ 0xf3, 0xf8, 0xd1, 0x48, 0xbc, 0xde, 0x89, 0xd3, 0x34, 0xcd, 0xde, 0x51, 0x4b, 0x83, 0xa2, 0x19, 0x16, 0xe5, 0xd6, 0x72, 0xb7, 0xc3, 0x1e, 0x59, 0xcb, 0xdf, 0x3a, 0x3c, 0x33, 0x80, 0x21, 0x4f, ]; +const MYC_AUDIT_STATE_TABLE_SHA256: [u8; 32] = [ + 0xc8, 0x6b, 0x49, 0xf6, 0x55, 0xac, 0x2c, 0xa4, 0xcb, 0x13, 0xed, 0x31, 0xff, 0x9e, 0xce, 0xe3, + 0x2f, 0xd4, 0x2a, 0x3e, 0xb0, 0xf1, 0xc8, 0x52, 0x9d, 0x92, 0xae, 0x5b, 0x48, 0x63, 0x38, 0x3b, +]; +const OPERATION_AUDIT_TABLE_SHA256: [u8; 32] = [ + 0xc1, 0x8d, 0x0b, 0x72, 0x34, 0xff, 0x8b, 0x21, 0x9f, 0x54, 0x20, 0x7f, 0x6c, 0x0b, 0x64, 0xae, + 0xd6, 0x8d, 0xd6, 0x1b, 0x48, 0xb3, 0x5b, 0xbe, 0x13, 0x2c, 0x0d, 0xb0, 0x9b, 0xea, 0x16, 0x2d, +]; +const NIP46_REQUEST_AUDIT_TABLE_SHA256: [u8; 32] = [ + 0xe2, 0x42, 0x82, 0x0b, 0xbb, 0xb2, 0x31, 0xa3, 0x8e, 0x9e, 0x7d, 0xf3, 0xf0, 0xe4, 0xd3, 0xc6, + 0x85, 0x93, 0x19, 0xe1, 0x2e, 0x47, 0x84, 0x48, 0x98, 0x8b, 0xc0, 0xdf, 0xdf, 0x96, 0xc6, 0xe3, +]; +const CONNECTION_RATE_WINDOWS_TABLE_SHA256: [u8; 32] = [ + 0x57, 0x53, 0xdf, 0x7b, 0x74, 0x44, 0x96, 0x9d, 0x88, 0x56, 0xe6, 0x1f, 0x15, 0x36, 0xdc, 0xab, + 0xa0, 0x02, 0x0a, 0x78, 0x77, 0x8e, 0x48, 0xa2, 0x80, 0x97, 0xc6, 0x37, 0xa6, 0x31, 0x17, 0xfc, +]; +const MYC_AUDIT_STATE_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0xda, 0xba, 0xc9, 0x86, 0x0f, 0x2a, 0xd8, 0xa0, 0x60, 0x75, 0x4b, 0x87, 0x78, 0xc8, 0xd8, 0xce, + 0x78, 0xa3, 0x57, 0x6e, 0xea, 0x08, 0x2b, 0x0c, 0x9c, 0x56, 0x5d, 0xa2, 0xb5, 0x6c, 0x68, 0xad, +]; +const MYC_AUDIT_STATE_NO_DELETE_SHA256: [u8; 32] = [ + 0x0a, 0x88, 0xa1, 0xbe, 0xe8, 0x23, 0x1f, 0xf0, 0xaf, 0x41, 0x91, 0xd7, 0x38, 0x64, 0x67, 0xb6, + 0xa8, 0xac, 0xda, 0xf0, 0x38, 0x8e, 0xd4, 0xb0, 0xac, 0x2c, 0xb6, 0xf2, 0x0a, 0xa1, 0xf5, 0x5c, +]; +const OPERATION_AUDIT_NO_UPDATE_SHA256: [u8; 32] = [ + 0xd8, 0xe4, 0x39, 0x63, 0x67, 0x74, 0x97, 0x4c, 0xa7, 0x97, 0x54, 0x6c, 0xed, 0x39, 0x9a, 0x7b, + 0xbc, 0x6c, 0x36, 0xc5, 0xd7, 0x8e, 0xf4, 0x08, 0xbd, 0xfd, 0xb7, 0x9b, 0xd0, 0x36, 0x74, 0x40, +]; +const NIP46_REQUEST_AUDIT_NO_UPDATE_SHA256: [u8; 32] = [ + 0x8a, 0x4a, 0x99, 0x4c, 0x13, 0x5f, 0x8d, 0x4d, 0x85, 0xd1, 0x25, 0x1d, 0x65, 0x47, 0x4d, 0x23, + 0x37, 0x62, 0xb7, 0x27, 0xb6, 0x7f, 0x31, 0xd4, 0x9c, 0x93, 0xe5, 0xce, 0x18, 0x43, 0xba, 0x81, +]; +const CONNECTION_RATE_WINDOWS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x59, 0x3c, 0xfb, 0xff, 0x20, 0x95, 0x32, 0x4e, 0x61, 0xdc, 0xd6, 0x09, 0xea, 0x7b, 0x1b, 0x1d, + 0xc4, 0xb9, 0xb7, 0x38, 0xcf, 0x80, 0x56, 0x00, 0xa5, 0x3b, 0xb4, 0xcd, 0x02, 0xcd, 0xe1, 0xef, +]; /// Stable classes for invalid embedded Myc catalog definitions. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -708,10 +971,17 @@ pub fn myc_migration_catalog() -> Result<MigrationCatalog, MycStateCatalogError> MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256), ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; - let catalog = MigrationCatalog::new([metadata, requests, connections]) + let governance = MigrationDescriptor::sql( + 5, + "create_bounded_governance_state", + CREATE_GOVERNANCE_STATE_MIGRATION_SQL, + MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256), + ) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; + let catalog = MigrationCatalog::new([metadata, requests, connections, governance]) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; if catalog.current_version() != MYC_STATE_SCHEMA_VERSION - || catalog.descriptors().len() != 3 + || catalog.descriptors().len() != 4 || catalog.digest().as_bytes() != &MYC_MIGRATION_CATALOG_SHA256 { return Err(MycStateCatalogError::new( @@ -748,9 +1018,21 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> { SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_4_SHA256), ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; + let version_five = SchemaVersionCatalog::new( + 5, + myc_state_governance_objects()?, + SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_5_SHA256), + ) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; let catalog = SchemaCatalog::new( &migrations, - [version_one, version_two, version_three, version_four], + [ + version_one, + version_two, + version_three, + version_four, + version_five, + ], ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; validate_myc_state_catalogs(&migrations, &catalog)?; @@ -935,6 +1217,80 @@ fn myc_state_connection_objects() -> Result<[SchemaObject; 19], MycStateCatalogE ]) } +fn myc_state_governance_objects() -> Result<Vec<SchemaObject>, MycStateCatalogError> { + let mut objects = Vec::from(myc_state_connection_objects()?); + let object = |kind, name, table, sql, digest| { + SchemaObject::new(kind, name, table, sql, SchemaDigest::from_bytes(digest)) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog)) + }; + objects.extend([ + object( + SchemaObjectKind::Table, + "myc_audit_state", + "myc_audit_state", + CREATE_MYC_AUDIT_STATE_TABLE_SQL, + MYC_AUDIT_STATE_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "operation_audit", + "operation_audit", + CREATE_OPERATION_AUDIT_TABLE_SQL, + OPERATION_AUDIT_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "nip46_request_audit", + "nip46_request_audit", + CREATE_NIP46_REQUEST_AUDIT_TABLE_SQL, + NIP46_REQUEST_AUDIT_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "connection_rate_windows", + "connection_rate_windows", + CREATE_CONNECTION_RATE_WINDOWS_TABLE_SQL, + CONNECTION_RATE_WINDOWS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "myc_audit_state_guard_update", + "myc_audit_state", + CREATE_MYC_AUDIT_STATE_GUARD_UPDATE_SQL, + MYC_AUDIT_STATE_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "myc_audit_state_no_delete", + "myc_audit_state", + CREATE_MYC_AUDIT_STATE_NO_DELETE_SQL, + MYC_AUDIT_STATE_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "operation_audit_no_update", + "operation_audit", + CREATE_OPERATION_AUDIT_NO_UPDATE_SQL, + OPERATION_AUDIT_NO_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "nip46_request_audit_no_update", + "nip46_request_audit", + CREATE_NIP46_REQUEST_AUDIT_NO_UPDATE_SQL, + NIP46_REQUEST_AUDIT_NO_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "connection_rate_windows_guard_update", + "connection_rate_windows", + CREATE_CONNECTION_RATE_WINDOWS_GUARD_UPDATE_SQL, + CONNECTION_RATE_WINDOWS_GUARD_UPDATE_SHA256, + )?, + ]); + Ok(objects) +} + /// Independently validates exact catalog versions, counts, and digests. pub fn validate_myc_state_catalogs( migrations: &MigrationCatalog, @@ -943,7 +1299,7 @@ pub fn validate_myc_state_catalogs( let versions = schema.versions(); let descriptors = migrations.descriptors(); let valid = migrations.current_version() == MYC_STATE_SCHEMA_VERSION - && descriptors.len() == 3 + && descriptors.len() == 4 && descriptors[0].target_version() == 2 && descriptors[0].name().as_str() == "create_myc_state_metadata" && descriptors[0].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256 @@ -953,9 +1309,12 @@ pub fn validate_myc_state_catalogs( && descriptors[2].target_version() == 4 && descriptors[2].name().as_str() == "create_connection_authorization_state" && descriptors[2].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256 + && descriptors[3].target_version() == 5 + && descriptors[3].name().as_str() == "create_bounded_governance_state" + && descriptors[3].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256 && migrations.digest().as_bytes() == &MYC_MIGRATION_CATALOG_SHA256 && schema.migration_catalog_digest() == migrations.digest() - && versions.len() == 4 + && versions.len() == 5 && versions[0].version() == MYC_STATE_BASE_SCHEMA_VERSION && versions[0].object_count() == MYC_STATE_SCHEMA_VERSION_1_OBJECT_COUNT && versions[0].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_1_SHA256 @@ -968,6 +1327,9 @@ pub fn validate_myc_state_catalogs( && versions[3].version() == 4 && versions[3].object_count() == MYC_STATE_SCHEMA_VERSION_4_OBJECT_COUNT && versions[3].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_4_SHA256 + && versions[4].version() == 5 + && versions[4].object_count() == MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT + && versions[4].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_5_SHA256 && schema.digest().as_bytes() == &MYC_STATE_SCHEMA_CATALOG_SHA256; if valid { Ok(()) diff --git a/src/state_connection.rs b/src/state_connection.rs @@ -10,6 +10,11 @@ use sha2::{Digest, Sha256}; use sqlx::Row; use url::{Host, Url}; +use crate::state_governance::{ + AuditEvidence, GovernanceOperationError, MycAuditCorrelationId, MycAuditKind, MycAuditOutcome, + MycAuditReasonCode, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, connection_subject, + global_subject, govern_rate_attempt, record_audit, relay_subject, +}; use crate::state_repository::{ MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, RepositoryOperationError, require_expected_metadata, @@ -26,6 +31,8 @@ const CHALLENGE_ID_DOMAIN: &[u8] = b"radroots.myc.authorization_challenge.v1\0"; const PERMISSION_SET_DOMAIN: &[u8] = b"radroots.myc.connection_permissions.v1\0"; const READ_REQUEST_BINDING_SQL: &str = r#"SELECT + CASE WHEN typeof(correlation_id) = 'blob' AND length(correlation_id) = 32 + THEN correlation_id ELSE NULL END AS correlation_id, CASE WHEN typeof(client_public_key) = 'text' AND length(CAST(client_public_key AS BLOB)) = 64 THEN client_public_key ELSE NULL END AS client_public_key, @@ -521,6 +528,7 @@ pub struct MycConnectionAdmissionRequest { observed_at: MycConnectionTimeUnixMs, authorized_until: Option<MycConnectionTimeUnixMs>, policy: MycConnectionAdmissionPolicy, + relay_id: MycRateRelayId, } impl MycConnectionAdmissionRequest { @@ -535,6 +543,7 @@ impl MycConnectionAdmissionRequest { observed_at: MycConnectionTimeUnixMs, authorized_until: Option<MycConnectionTimeUnixMs>, policy: MycConnectionAdmissionPolicy, + relay_id: MycRateRelayId, ) -> Result<Self, MycConnectionStateError> { if (matches!(policy, MycConnectionAdmissionPolicy::Trusted) && authorized_until.is_some_and(|until| until <= observed_at)) @@ -554,6 +563,7 @@ impl MycConnectionAdmissionRequest { observed_at, authorized_until, policy, + relay_id, }) } @@ -567,6 +577,7 @@ impl MycConnectionAdmissionRequest { observed_at: self.observed_at, authorized_until: self.authorized_until, policy: self.policy, + relay_id: self.relay_id.clone(), } } } @@ -720,13 +731,15 @@ impl fmt::Debug for MycConnectionDecisionRecord { pub enum MycConnectionAdmission { Admitted(MycConnectionDecisionRecord), ExactReplay(MycConnectionDecisionRecord), + RateLimited, } impl MycConnectionAdmission { #[must_use] - pub const fn record(&self) -> &MycConnectionDecisionRecord { + pub const fn record(&self) -> Option<&MycConnectionDecisionRecord> { match self { - Self::Admitted(record) | Self::ExactReplay(record) => record, + Self::Admitted(record) | Self::ExactReplay(record) => Some(record), + Self::RateLimited => None, } } } @@ -736,6 +749,7 @@ impl fmt::Debug for MycConnectionAdmission { formatter.write_str(match self { Self::Admitted(_) => "MycConnectionAdmission::Admitted([redacted])", Self::ExactReplay(_) => "MycConnectionAdmission::ExactReplay([redacted])", + Self::RateLimited => "MycConnectionAdmission::RateLimited", }) } } @@ -956,13 +970,15 @@ impl fmt::Debug for MycAuthorizationChallengeRecord { pub enum MycAuthorizationChallengeAdmission { Created(MycAuthorizationChallengeRecord), ExactReplay(MycAuthorizationChallengeRecord), + RateLimited, } impl MycAuthorizationChallengeAdmission { #[must_use] - pub const fn record(&self) -> &MycAuthorizationChallengeRecord { + pub const fn record(&self) -> Option<&MycAuthorizationChallengeRecord> { match self { - Self::Created(record) | Self::ExactReplay(record) => record, + Self::Created(record) | Self::ExactReplay(record) => Some(record), + Self::RateLimited => None, } } } @@ -972,6 +988,37 @@ impl fmt::Debug for MycAuthorizationChallengeAdmission { formatter.write_str(match self { Self::Created(_) => "MycAuthorizationChallengeAdmission::Created([redacted])", Self::ExactReplay(_) => "MycAuthorizationChallengeAdmission::ExactReplay([redacted])", + Self::RateLimited => "MycAuthorizationChallengeAdmission::RateLimited", + }) + } +} + +/// New, replayed, or rate-limited authorization result. +#[derive(Clone, PartialEq, Eq)] +pub enum MycAuthorizationChallengeAuthorization { + Resolved(MycAuthorizationChallengeRecord), + ExactReplay(MycAuthorizationChallengeRecord), + RateLimited, +} + +impl MycAuthorizationChallengeAuthorization { + #[must_use] + pub const fn record(&self) -> Option<&MycAuthorizationChallengeRecord> { + match self { + Self::Resolved(record) | Self::ExactReplay(record) => Some(record), + Self::RateLimited => None, + } + } +} + +impl fmt::Debug for MycAuthorizationChallengeAuthorization { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::Resolved(_) => "MycAuthorizationChallengeAuthorization::Resolved([redacted])", + Self::ExactReplay(_) => { + "MycAuthorizationChallengeAuthorization::ExactReplay([redacted])" + } + Self::RateLimited => "MycAuthorizationChallengeAuthorization::RateLimited", }) } } @@ -982,13 +1029,21 @@ impl MycStateRepository<'_> { &self, request: &MycConnectionAdmissionRequest, ) -> Result<MycConnectionAdmission, MycStateRepositoryError> { + if !self.expected().admits_rate_relay(&request.relay_id) { + return Err(MycStateRepositoryError::new( + MycStateRepositoryErrorKind::Binding, + )); + } let request = request.owned(); + let rate_policy = self + .expected() + .governance_rate_policy(MycRateLimitClass::ConnectionAdmission); let expected = PersistedMetadata::from(self.expected()); self.host() .transaction(move |transaction| { Box::pin(async move { verify_metadata(transaction, &expected).await?; - admit_connection(transaction, &request).await + admit_connection(transaction, &request, rate_policy).await }) }) .await @@ -1002,6 +1057,7 @@ impl MycStateRepository<'_> { connection_id: MycConnectionId, policy_generation: MycConnectionPolicyGeneration, observed_at: MycConnectionTimeUnixMs, + audit_correlation: MycAuditCorrelationId, decision: MycConnectionOperatorDecision, ) -> Result<MycConnectionRecord, MycStateRepositoryError> { let expected = PersistedMetadata::from(self.expected()); @@ -1015,6 +1071,7 @@ impl MycStateRepository<'_> { connection_id, policy_generation, observed_at, + audit_correlation, decision, ) .await @@ -1030,6 +1087,7 @@ impl MycStateRepository<'_> { connection_id: MycConnectionId, policy_generation: MycConnectionPolicyGeneration, observed_at: MycConnectionTimeUnixMs, + audit_correlation: MycAuditCorrelationId, ) -> Result<MycConnectionRecord, MycStateRepositoryError> { let expected = PersistedMetadata::from(self.expected()); self.host() @@ -1041,6 +1099,12 @@ impl MycStateRepository<'_> { return Err(ConnectionOperationError::Binding); } if before.status == MycConnectionStatus::Expired { + record_connection_expiry_audit( + transaction, + audit_correlation, + before.updated_at, + ) + .await?; return Ok(before); } if before.status != MycConnectionStatus::Active @@ -1059,7 +1123,10 @@ impl MycStateRepository<'_> { .await .map_err(|_| ConnectionOperationError::Storage)?; require_one(result.rows_affected())?; - read_connection(transaction, connection_id).await + let record = read_connection(transaction, connection_id).await?; + record_connection_expiry_audit(transaction, audit_correlation, observed_at) + .await?; + Ok(record) }) }) .await @@ -1072,12 +1139,15 @@ impl MycStateRepository<'_> { request: &MycAuthorizationChallengeRequest, ) -> Result<MycAuthorizationChallengeAdmission, MycStateRepositoryError> { let request = request.owned(); + let rate_policy = self + .expected() + .governance_rate_policy(MycRateLimitClass::ChallengeCreation); let expected = PersistedMetadata::from(self.expected()); self.host() .transaction(move |transaction| { Box::pin(async move { verify_metadata(transaction, &expected).await?; - issue_challenge(transaction, &request).await + issue_challenge(transaction, &request, rate_policy).await }) }) .await @@ -1092,7 +1162,10 @@ impl MycStateRepository<'_> { operation_id: MycSignerOperationId, policy_generation: MycConnectionPolicyGeneration, observed_at: MycConnectionTimeUnixMs, - ) -> Result<MycAuthorizationChallengeRecord, MycStateRepositoryError> { + ) -> Result<MycAuthorizationChallengeAuthorization, MycStateRepositoryError> { + let rate_policy = self + .expected() + .governance_rate_policy(MycRateLimitClass::ChallengeAuthorization); let expected = PersistedMetadata::from(self.expected()); self.host() .transaction(move |transaction| { @@ -1105,6 +1178,7 @@ impl MycStateRepository<'_> { operation_id, policy_generation, observed_at, + rate_policy, ) .await }) @@ -1114,12 +1188,40 @@ impl MycStateRepository<'_> { } } +async fn record_connection_expiry_audit( + transaction: &mut ServiceSqliteTransaction<'_>, + correlation: MycAuditCorrelationId, + occurred_at: MycConnectionTimeUnixMs, +) -> Result<(), ConnectionOperationError> { + record_audit( + transaction, + AuditEvidence { + correlation, + occurred_at, + operation_id: None, + }, + MycAuditKind::ConnectionExpiry, + MycAuditOutcome::Succeeded, + MycAuditReasonCode::ConnectionExpired, + ) + .await + .map_err(map_governance_error)?; + Ok(()) +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum ConnectionOperationError { Binding, Storage, } +const fn map_governance_error(error: GovernanceOperationError) -> ConnectionOperationError { + match error { + GovernanceOperationError::Binding => ConnectionOperationError::Binding, + GovernanceOperationError::Storage => ConnectionOperationError::Storage, + } +} + async fn verify_metadata( transaction: &mut ServiceSqliteTransaction<'_>, expected: &PersistedMetadata, @@ -1135,6 +1237,7 @@ async fn verify_metadata( async fn admit_connection( transaction: &mut ServiceSqliteTransaction<'_>, request: &MycConnectionAdmissionRequest, + rate_policy: MycRateLimitPolicy, ) -> Result<MycConnectionAdmission, ConnectionOperationError> { let binding = read_request_binding(transaction, request.operation_id).await?; if binding.client_public_key != request.client_public_key @@ -1160,6 +1263,25 @@ async fn admit_connection( return Ok(MycConnectionAdmission::ExactReplay(record)); } + let evidence = AuditEvidence { + correlation: binding.correlation, + occurred_at: request.observed_at, + operation_id: Some(request.operation_id), + }; + if !govern_rate_attempt( + transaction, + rate_policy, + MycRateLimitClass::ConnectionAdmission, + &[global_subject(), relay_subject(&request.relay_id)], + evidence, + MycAuditKind::ConnectionAdmission, + ) + .await + .map_err(map_governance_error)? + { + return Ok(MycConnectionAdmission::RateLimited); + } + let decision = request.policy.decision(); let connection_id = if request.policy == MycConnectionAdmissionPolicy::Denied { None @@ -1183,9 +1305,29 @@ async fn admit_connection( let persisted = read_decision(transaction, request.operation_id) .await? .ok_or(ConnectionOperationError::Binding)?; - decision_record(transaction, request.operation_id, persisted) - .await - .map(MycConnectionAdmission::Admitted) + let record = decision_record(transaction, request.operation_id, persisted).await?; + let (outcome, reason) = match request.policy { + MycConnectionAdmissionPolicy::Trusted => { + (MycAuditOutcome::Succeeded, MycAuditReasonCode::Trusted) + } + MycConnectionAdmissionPolicy::ExplicitApproval => ( + MycAuditOutcome::Succeeded, + MycAuditReasonCode::ApprovalRequired, + ), + MycConnectionAdmissionPolicy::Denied => { + (MycAuditOutcome::Rejected, MycAuditReasonCode::PolicyDenied) + } + }; + record_audit( + transaction, + evidence, + MycAuditKind::ConnectionAdmission, + outcome, + reason, + ) + .await + .map_err(map_governance_error)?; + Ok(MycConnectionAdmission::Admitted(record)) } async fn insert_connection( @@ -1256,6 +1398,7 @@ async fn decide_connection( connection_id: MycConnectionId, policy_generation: MycConnectionPolicyGeneration, observed_at: MycConnectionTimeUnixMs, + audit_correlation: MycAuditCorrelationId, decision: MycConnectionOperatorDecision, ) -> Result<MycConnectionRecord, ConnectionOperationError> { let existing_decision = read_decision(transaction, operation_id) @@ -1271,7 +1414,7 @@ async fn decide_connection( return Err(ConnectionOperationError::Binding); } - match decision { + let (audit_outcome, audit_reason) = match decision { MycConnectionOperatorDecision::Approve { granted_permissions, authorized_until, @@ -1284,10 +1427,19 @@ async fn decide_connection( if connection.status == MycConnectionStatus::Active && existing_decision.decision == MycConnectionDecision::Allowed { - return (connection.granted_permissions == granted_permissions + let replay = (connection.granted_permissions == granted_permissions && connection.authorized_until == authorized_until) .then_some(connection) - .ok_or(ConnectionOperationError::Binding); + .ok_or(ConnectionOperationError::Binding)?; + record_operator_audit( + transaction, + audit_correlation, + replay.updated_at, + MycAuditOutcome::Succeeded, + MycAuditReasonCode::OperatorApproved, + ) + .await?; + return Ok(replay); } if connection.status != MycConnectionStatus::Pending || existing_decision.decision != MycConnectionDecision::PendingApproval @@ -1315,11 +1467,23 @@ async fn decide_connection( "operator_approved", ) .await?; + ( + MycAuditOutcome::Succeeded, + MycAuditReasonCode::OperatorApproved, + ) } MycConnectionOperatorDecision::Deny => { if connection.status == MycConnectionStatus::Denied && existing_decision.decision == MycConnectionDecision::Denied { + record_operator_audit( + transaction, + audit_correlation, + connection.updated_at, + MycAuditOutcome::Rejected, + MycAuditReasonCode::OperatorDenied, + ) + .await?; return Ok(connection); } if connection.status != MycConnectionStatus::Pending @@ -1346,9 +1510,45 @@ async fn decide_connection( "operator_denied", ) .await?; + ( + MycAuditOutcome::Rejected, + MycAuditReasonCode::OperatorDenied, + ) } - } - read_connection(transaction, connection_id).await + }; + let record = read_connection(transaction, connection_id).await?; + record_operator_audit( + transaction, + audit_correlation, + observed_at, + audit_outcome, + audit_reason, + ) + .await?; + Ok(record) +} + +async fn record_operator_audit( + transaction: &mut ServiceSqliteTransaction<'_>, + correlation: MycAuditCorrelationId, + occurred_at: MycConnectionTimeUnixMs, + outcome: MycAuditOutcome, + reason: MycAuditReasonCode, +) -> Result<(), ConnectionOperationError> { + record_audit( + transaction, + AuditEvidence { + correlation, + occurred_at, + operation_id: None, + }, + MycAuditKind::ConnectionOperatorDecision, + outcome, + reason, + ) + .await + .map_err(map_governance_error)?; + Ok(()) } #[allow(clippy::too_many_arguments)] @@ -1377,6 +1577,7 @@ async fn update_approval_decision( async fn issue_challenge( transaction: &mut ServiceSqliteTransaction<'_>, request: &MycAuthorizationChallengeRequest, + rate_policy: MycRateLimitPolicy, ) -> Result<MycAuthorizationChallengeAdmission, ConnectionOperationError> { let binding = read_request_binding(transaction, request.operation_id).await?; if binding.method == MycSignerRequestMethod::Connect { @@ -1414,6 +1615,24 @@ async fn issue_challenge( { return Err(ConnectionOperationError::Binding); } + let evidence = AuditEvidence { + correlation: binding.correlation, + occurred_at: request.issued_at, + operation_id: Some(request.operation_id), + }; + if !govern_rate_attempt( + transaction, + rate_policy, + MycRateLimitClass::ChallengeCreation, + &[connection_subject(request.connection_id)], + evidence, + MycAuditKind::ChallengeCreation, + ) + .await + .map_err(map_governance_error)? + { + return Ok(MycAuthorizationChallengeAdmission::RateLimited); + } let challenge_id = derive_challenge_id(request.operation_id, request.connection_id, &request.nonce); let result = sqlx::query(INSERT_CHALLENGE_SQL) @@ -1444,6 +1663,15 @@ async fn issue_challenge( let record = read_challenge(transaction, request.operation_id) .await? .ok_or(ConnectionOperationError::Binding)?; + record_audit( + transaction, + evidence, + MycAuditKind::ChallengeCreation, + MycAuditOutcome::Succeeded, + MycAuditReasonCode::ChallengeRequired, + ) + .await + .map_err(map_governance_error)?; Ok(MycAuthorizationChallengeAdmission::Created(record)) } @@ -1455,7 +1683,8 @@ async fn authorize_challenge( operation_id: MycSignerOperationId, policy_generation: MycConnectionPolicyGeneration, observed_at: MycConnectionTimeUnixMs, -) -> Result<MycAuthorizationChallengeRecord, ConnectionOperationError> { + rate_policy: MycRateLimitPolicy, +) -> Result<MycAuthorizationChallengeAuthorization, ConnectionOperationError> { let before = read_challenge(transaction, operation_id) .await? .ok_or(ConnectionOperationError::Binding)?; @@ -1476,7 +1705,25 @@ async fn authorize_challenge( return Err(ConnectionOperationError::Binding); } if before.state != MycAuthorizationChallengeState::Pending { - return Ok(before); + return Ok(MycAuthorizationChallengeAuthorization::ExactReplay(before)); + } + let evidence = AuditEvidence { + correlation: request_binding.correlation, + occurred_at: observed_at, + operation_id: Some(operation_id), + }; + if !govern_rate_attempt( + transaction, + rate_policy, + MycRateLimitClass::ChallengeAuthorization, + &[connection_subject(connection_id)], + evidence, + MycAuditKind::ChallengeAuthorization, + ) + .await + .map_err(map_governance_error)? + { + return Ok(MycAuthorizationChallengeAuthorization::RateLimited); } let connection_expired = connection.status != MycConnectionStatus::Active || connection @@ -1518,12 +1765,36 @@ async fn authorize_challenge( .await .map_err(|_| ConnectionOperationError::Storage)?; require_one(result.rows_affected())?; - read_challenge(transaction, operation_id) + let record = read_challenge(transaction, operation_id) .await? - .ok_or(ConnectionOperationError::Binding) + .ok_or(ConnectionOperationError::Binding)?; + let (outcome, audit_reason) = match state { + MycAuthorizationChallengeState::Authorized => ( + MycAuditOutcome::Succeeded, + MycAuditReasonCode::ChallengeAuthorized, + ), + MycAuthorizationChallengeState::Expired => ( + MycAuditOutcome::Rejected, + MycAuditReasonCode::ChallengeExpired, + ), + MycAuthorizationChallengeState::Pending => { + return Err(ConnectionOperationError::Binding); + } + }; + record_audit( + transaction, + evidence, + MycAuditKind::ChallengeAuthorization, + outcome, + audit_reason, + ) + .await + .map_err(map_governance_error)?; + Ok(MycAuthorizationChallengeAuthorization::Resolved(record)) } struct RequestBinding { + correlation: MycAuditCorrelationId, client_public_key: MycNip46ClientPublicKey, method: MycSignerRequestMethod, received_at: MycConnectionTimeUnixMs, @@ -1542,6 +1813,7 @@ async fn read_request_binding( return Err(ConnectionOperationError::Binding); } let row = &rows[0]; + let correlation = MycAuditCorrelationId::new(exact_digest(row, "correlation_id")?); let client = row .try_get::<Option<&str>, _>("client_public_key") .map_err(|_| ConnectionOperationError::Binding)? @@ -1555,6 +1827,7 @@ async fn read_request_binding( .and_then(MycSignerRequestMethod::parse) .ok_or(ConnectionOperationError::Binding)?; Ok(RequestBinding { + correlation, client_public_key: client, method, received_at: time(row, "received_at_unix_ms")?, diff --git a/src/state_governance.rs b/src/state_governance.rs @@ -0,0 +1,1300 @@ +//! Durable bounded audit, rate-window, retention, and compaction state. + +use core::fmt; +use std::error::Error; + +use radroots_service_sqlite::{ + ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, +}; +use sha2::{Digest, Sha256}; +use sqlx::Row; + +use crate::state_repository::{ + MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, + RepositoryOperationError, require_expected_metadata, +}; +use crate::{MycConnectionId, MycConnectionTimeUnixMs, MycSignerOperationId}; + +/// Maximum governed rate-window duration. +pub const MYC_RATE_WINDOW_MAX_MS: u64 = 86_400_000; +/// Maximum attempts admitted in one governed rate window. +pub const MYC_RATE_MAX_ATTEMPTS: u32 = 10_000; +/// Maximum retained duration for rate-window evidence. +pub const MYC_RATE_RETENTION_MAX_MS: u64 = 2_592_000_000; +/// Maximum tracked non-global subjects for one rate-limit class. +pub const MYC_RATE_MAX_TRACKED_SUBJECTS: u32 = 65_536; +/// Maximum retained duration for safe audit evidence. +pub const MYC_AUDIT_RETENTION_MAX_MS: u64 = 31_536_000_000; +/// Maximum audit records returned in one page. +pub const MYC_AUDIT_PAGE_MAX_ITEMS: u16 = 200; +/// Maximum rows removed from either bounded evidence class in one compaction. +pub const MYC_COMPACTION_MAX_ROWS: u16 = 4_096; +/// Maximum UTF-8 bytes in a configured stable relay ID. +pub const MYC_RATE_RELAY_ID_MAX_BYTES: usize = 64; + +const AUDIT_ID_DOMAIN: &[u8] = b"radroots.myc.operation_audit.v1\0"; +const GLOBAL_SUBJECT_DOMAIN: &[u8] = b"radroots.myc.rate_subject.global.v1\0"; +const RELAY_SUBJECT_DOMAIN: &[u8] = b"radroots.myc.rate_subject.relay.v1\0"; +const CONNECTION_SUBJECT_DOMAIN: &[u8] = b"radroots.myc.rate_subject.connection.v1\0"; + +const READ_AUDIT_STATE_SQL: &str = + "SELECT next_sequence FROM myc_audit_state WHERE singleton = 1 LIMIT 2"; +const ADVANCE_AUDIT_STATE_SQL: &str = + "UPDATE myc_audit_state SET next_sequence = ? WHERE singleton = 1 AND next_sequence = ?"; +const READ_AUDIT_BY_ID_SQL: &str = r#"SELECT + operation_audit.audit_sequence, + CASE WHEN typeof(operation_audit.correlation_id) = 'blob' + AND length(operation_audit.correlation_id) = 32 + THEN operation_audit.correlation_id ELSE NULL END AS correlation_id, + CASE WHEN typeof(operation_audit.audit_kind) = 'text' + AND length(CAST(operation_audit.audit_kind AS BLOB)) <= 40 + THEN operation_audit.audit_kind ELSE NULL END AS audit_kind, + CASE WHEN typeof(operation_audit.outcome) = 'text' + AND length(CAST(operation_audit.outcome AS BLOB)) <= 16 + THEN operation_audit.outcome ELSE NULL END AS outcome, + CASE WHEN typeof(operation_audit.reason_code) = 'text' + AND length(CAST(operation_audit.reason_code AS BLOB)) <= 48 + THEN operation_audit.reason_code ELSE NULL END AS reason_code, + operation_audit.occurred_at_unix_ms, + CASE WHEN typeof(request_audit.operation_id) = 'blob' + AND length(request_audit.operation_id) = 32 + THEN request_audit.operation_id ELSE NULL END AS operation_id, + CASE WHEN typeof(request_audit.audit_kind) = 'text' + AND length(CAST(request_audit.audit_kind AS BLOB)) <= 40 + THEN request_audit.audit_kind ELSE NULL END AS request_audit_kind, + CASE WHEN typeof(request_record.correlation_id) = 'blob' + AND length(request_record.correlation_id) = 32 + THEN request_record.correlation_id ELSE NULL END AS request_correlation_id +FROM operation_audit +LEFT JOIN nip46_request_audit AS request_audit + ON request_audit.audit_sequence = operation_audit.audit_sequence +LEFT JOIN nip46_requests AS request_record + ON request_record.operation_id = request_audit.operation_id +WHERE operation_audit.audit_id = ? LIMIT 2"#; +const INSERT_AUDIT_SQL: &str = r#"INSERT INTO operation_audit ( + audit_sequence, audit_id, correlation_id, audit_kind, outcome, reason_code, + occurred_at_unix_ms +) VALUES (?, ?, ?, ?, ?, ?, ?)"#; +const INSERT_REQUEST_AUDIT_SQL: &str = r#"INSERT INTO nip46_request_audit ( + operation_id, audit_kind, audit_sequence +) VALUES (?, ?, ?)"#; + +const READ_RATE_SQL: &str = r#"SELECT + window_started_at_unix_ms, window_ends_at_unix_ms, + accepted_count, rejected_count, lifetime_accepted_count, + lifetime_rejected_count, last_observed_at_unix_ms, retention_expires_at_unix_ms +FROM connection_rate_windows +WHERE rate_kind = ? AND subject_scope = ? AND subject_sha256 = ? +LIMIT 2"#; +const COUNT_RATE_SUBJECTS_SQL: &str = r#"SELECT COUNT(*) +FROM connection_rate_windows +WHERE rate_kind = ? AND subject_scope = ?"#; +const INSERT_RATE_SQL: &str = r#"INSERT INTO connection_rate_windows ( + rate_kind, subject_scope, subject_sha256, window_started_at_unix_ms, + window_ends_at_unix_ms, accepted_count, rejected_count, + lifetime_accepted_count, lifetime_rejected_count, + last_observed_at_unix_ms, retention_expires_at_unix_ms +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; +const UPDATE_RATE_SQL: &str = r#"UPDATE connection_rate_windows SET + window_started_at_unix_ms = ?, window_ends_at_unix_ms = ?, + accepted_count = ?, rejected_count = ?, lifetime_accepted_count = ?, + lifetime_rejected_count = ?, last_observed_at_unix_ms = ?, + retention_expires_at_unix_ms = ? +WHERE rate_kind = ? AND subject_scope = ? AND subject_sha256 = ? + AND last_observed_at_unix_ms = ?"#; + +/// Stable construction and validation failure classes. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycGovernanceStateErrorKind { + InvalidRatePolicy, + InvalidRelayId, + InvalidPageLimit, + InvalidCompactionPolicy, +} + +impl MycGovernanceStateErrorKind { + /// Returns the stable machine-readable error code. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::InvalidRatePolicy => "governance_rate_policy_invalid", + Self::InvalidRelayId => "governance_relay_id_invalid", + Self::InvalidPageLimit => "governance_audit_page_limit_invalid", + Self::InvalidCompactionPolicy => "governance_compaction_policy_invalid", + } + } +} + +/// Source-free governance-state validation failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct MycGovernanceStateError { + kind: MycGovernanceStateErrorKind, +} + +impl MycGovernanceStateError { + const fn new(kind: MycGovernanceStateErrorKind) -> Self { + Self { kind } + } + + /// Returns the stable failure class. + #[must_use] + pub const fn kind(self) -> MycGovernanceStateErrorKind { + self.kind + } + + /// Returns the stable machine-readable error code. + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Display for MycGovernanceStateError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + MycGovernanceStateErrorKind::InvalidRatePolicy => "rate policy is invalid", + MycGovernanceStateErrorKind::InvalidRelayId => "rate relay ID is invalid", + MycGovernanceStateErrorKind::InvalidPageLimit => "audit page limit is invalid", + MycGovernanceStateErrorKind::InvalidCompactionPolicy => "compaction policy is invalid", + }) + } +} + +impl fmt::Debug for MycGovernanceStateError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycGovernanceStateError") + .field("kind", &self.kind) + .finish() + } +} + +impl Error for MycGovernanceStateError {} + +/// Closed rate-window classes. Their subject scope is fixed by contract. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycRateLimitClass { + ConnectionAdmission, + ChallengeCreation, + ChallengeAuthorization, +} + +impl MycRateLimitClass { + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::ConnectionAdmission => "connection_admission", + Self::ChallengeCreation => "challenge_creation", + Self::ChallengeAuthorization => "challenge_authorization", + } + } +} + +/// Validated bounded rate-window policy. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct MycRateLimitPolicy { + class: MycRateLimitClass, + window_ms: u64, + max_attempts: u32, + retention_ms: u64, + maximum_tracked_subjects: u32, +} + +impl MycRateLimitPolicy { + /// Constructs one exact policy from fully normalized configuration values. + pub fn new( + class: MycRateLimitClass, + window_ms: u64, + max_attempts: u32, + retention_ms: u64, + maximum_tracked_subjects: u32, + ) -> Result<Self, MycGovernanceStateError> { + if window_ms == 0 + || window_ms > MYC_RATE_WINDOW_MAX_MS + || max_attempts == 0 + || max_attempts > MYC_RATE_MAX_ATTEMPTS + || retention_ms < window_ms + || retention_ms > MYC_RATE_RETENTION_MAX_MS + || maximum_tracked_subjects == 0 + || maximum_tracked_subjects > MYC_RATE_MAX_TRACKED_SUBJECTS + { + return Err(MycGovernanceStateError::new( + MycGovernanceStateErrorKind::InvalidRatePolicy, + )); + } + Ok(Self { + class, + window_ms, + max_attempts, + retention_ms, + maximum_tracked_subjects, + }) + } + + /// Returns the closed policy class. + #[must_use] + pub const fn class(self) -> MycRateLimitClass { + self.class + } +} + +impl fmt::Debug for MycRateLimitPolicy { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycRateLimitPolicy") + .field("class", &self.class) + .field("values", &"[bounded]") + .finish() + } +} + +#[derive(Clone, PartialEq, Eq)] +pub(crate) struct MycGovernancePolicies { + connection_admission: MycRateLimitPolicy, + challenge_creation: MycRateLimitPolicy, + challenge_authorization: MycRateLimitPolicy, + audit_retention_ms: u64, + relays: Box<[MycRateRelayId]>, +} + +impl MycGovernancePolicies { + pub(crate) fn new( + connection_admission: MycRateLimitPolicy, + challenge_creation: MycRateLimitPolicy, + challenge_authorization: MycRateLimitPolicy, + audit_retention_ms: u64, + relays: Box<[MycRateRelayId]>, + ) -> Result<Self, MycGovernanceStateError> { + if connection_admission.class != MycRateLimitClass::ConnectionAdmission + || challenge_creation.class != MycRateLimitClass::ChallengeCreation + || challenge_authorization.class != MycRateLimitClass::ChallengeAuthorization + || audit_retention_ms == 0 + || audit_retention_ms > MYC_AUDIT_RETENTION_MAX_MS + || relays.is_empty() + { + return Err(MycGovernanceStateError::new( + MycGovernanceStateErrorKind::InvalidRatePolicy, + )); + } + Ok(Self { + connection_admission, + challenge_creation, + challenge_authorization, + audit_retention_ms, + relays, + }) + } + + pub(crate) const fn rate_policy(&self, class: MycRateLimitClass) -> MycRateLimitPolicy { + match class { + MycRateLimitClass::ConnectionAdmission => self.connection_admission, + MycRateLimitClass::ChallengeCreation => self.challenge_creation, + MycRateLimitClass::ChallengeAuthorization => self.challenge_authorization, + } + } + + pub(crate) fn admits_relay(&self, relay: &MycRateRelayId) -> bool { + self.relays.iter().any(|configured| configured == relay) + } + + pub(crate) const fn audit_retention_ms(&self) -> u64 { + self.audit_retention_ms + } +} + +/// Validated configured relay identity used only to derive bounded rate scope. +#[derive(Clone, PartialEq, Eq)] +pub struct MycRateRelayId(Box<str>); + +impl MycRateRelayId { + /// Validates one stable lower-snake relay identity before allocation. + pub fn new(value: &str) -> Result<Self, MycGovernanceStateError> { + let bytes = value.as_bytes(); + if bytes.is_empty() + || bytes.len() > MYC_RATE_RELAY_ID_MAX_BYTES + || !bytes[0].is_ascii_lowercase() + || !bytes[bytes.len() - 1].is_ascii_alphanumeric() + || bytes.windows(2).any(|pair| pair == b"__") + || !bytes + .iter() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_') + { + return Err(MycGovernanceStateError::new( + MycGovernanceStateErrorKind::InvalidRelayId, + )); + } + Ok(Self(value.into())) + } +} + +impl fmt::Debug for MycRateRelayId { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycRateRelayId([redacted])") + } +} + +/// Injected stable correlation identity for non-NIP-46 operator work. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct MycAuditCorrelationId([u8; 32]); + +impl MycAuditCorrelationId { + /// Constructs an opaque caller-injected correlation identity. + #[must_use] + pub const fn new(bytes: [u8; 32]) -> Self { + Self(bytes) + } +} + +impl fmt::Debug for MycAuditCorrelationId { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycAuditCorrelationId([redacted])") + } +} + +/// Closed audit event classes. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycAuditKind { + ConnectionAdmission, + ConnectionOperatorDecision, + ConnectionExpiry, + ChallengeCreation, + ChallengeAuthorization, + GovernanceCompaction, +} + +impl MycAuditKind { + const fn as_str(self) -> &'static str { + match self { + Self::ConnectionAdmission => "connection_admission", + Self::ConnectionOperatorDecision => "connection_operator_decision", + Self::ConnectionExpiry => "connection_expiry", + Self::ChallengeCreation => "challenge_creation", + Self::ChallengeAuthorization => "challenge_authorization", + Self::GovernanceCompaction => "governance_compaction", + } + } + + fn parse(value: &str) -> Option<Self> { + match value { + "connection_admission" => Some(Self::ConnectionAdmission), + "connection_operator_decision" => Some(Self::ConnectionOperatorDecision), + "connection_expiry" => Some(Self::ConnectionExpiry), + "challenge_creation" => Some(Self::ChallengeCreation), + "challenge_authorization" => Some(Self::ChallengeAuthorization), + "governance_compaction" => Some(Self::GovernanceCompaction), + _ => None, + } + } +} + +/// Closed audit outcomes. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycAuditOutcome { + Succeeded, + Rejected, + Failed, +} + +impl MycAuditOutcome { + const fn as_str(self) -> &'static str { + match self { + Self::Succeeded => "succeeded", + Self::Rejected => "rejected", + Self::Failed => "failed", + } + } + + fn parse(value: &str) -> Option<Self> { + match value { + "succeeded" => Some(Self::Succeeded), + "rejected" => Some(Self::Rejected), + "failed" => Some(Self::Failed), + _ => None, + } + } +} + +/// Closed safe audit reasons. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycAuditReasonCode { + Trusted, + ApprovalRequired, + PolicyDenied, + OperatorApproved, + OperatorDenied, + ConnectionExpired, + ChallengeRequired, + ChallengeAuthorized, + ChallengeExpired, + RateLimited, + Compacted, +} + +impl MycAuditReasonCode { + const fn as_str(self) -> &'static str { + match self { + Self::Trusted => "trusted", + Self::ApprovalRequired => "approval_required", + Self::PolicyDenied => "policy_denied", + Self::OperatorApproved => "operator_approved", + Self::OperatorDenied => "operator_denied", + Self::ConnectionExpired => "connection_expired", + Self::ChallengeRequired => "challenge_required", + Self::ChallengeAuthorized => "challenge_authorized", + Self::ChallengeExpired => "challenge_expired", + Self::RateLimited => "rate_limited", + Self::Compacted => "compacted", + } + } + + fn parse(value: &str) -> Option<Self> { + match value { + "trusted" => Some(Self::Trusted), + "approval_required" => Some(Self::ApprovalRequired), + "policy_denied" => Some(Self::PolicyDenied), + "operator_approved" => Some(Self::OperatorApproved), + "operator_denied" => Some(Self::OperatorDenied), + "connection_expired" => Some(Self::ConnectionExpired), + "challenge_required" => Some(Self::ChallengeRequired), + "challenge_authorized" => Some(Self::ChallengeAuthorized), + "challenge_expired" => Some(Self::ChallengeExpired), + "rate_limited" => Some(Self::RateLimited), + "compacted" => Some(Self::Compacted), + _ => None, + } + } +} + +/// One validated safe audit record. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct MycAuditRecord { + sequence: u64, + correlation_id: MycAuditCorrelationId, + operation_id: Option<MycSignerOperationId>, + kind: MycAuditKind, + outcome: MycAuditOutcome, + reason: MycAuditReasonCode, + occurred_at: MycConnectionTimeUnixMs, +} + +impl MycAuditRecord { + #[must_use] + pub const fn sequence(&self) -> u64 { + self.sequence + } + + /// Returns the opaque correlation identity bound to this record. + #[must_use] + pub const fn correlation_id(&self) -> MycAuditCorrelationId { + self.correlation_id + } + + /// Returns the typed signer-operation identity when this is request audit. + #[must_use] + pub const fn operation_id(&self) -> Option<MycSignerOperationId> { + self.operation_id + } + + #[must_use] + pub const fn kind(&self) -> MycAuditKind { + self.kind + } + + #[must_use] + pub const fn outcome(&self) -> MycAuditOutcome { + self.outcome + } + + #[must_use] + pub const fn reason(&self) -> MycAuditReasonCode { + self.reason + } + + #[must_use] + pub const fn occurred_at(&self) -> MycConnectionTimeUnixMs { + self.occurred_at + } +} + +impl fmt::Debug for MycAuditRecord { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycAuditRecord") + .field("sequence", &self.sequence) + .field("kind", &self.kind) + .field("outcome", &self.outcome) + .field("reason", &self.reason) + .field("correlation", &"[redacted]") + .finish() + } +} + +/// Validated audit page limit. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct MycAuditPageLimit(u16); + +impl MycAuditPageLimit { + pub fn new(value: u16) -> Result<Self, MycGovernanceStateError> { + (value != 0 && value <= MYC_AUDIT_PAGE_MAX_ITEMS) + .then_some(Self(value)) + .ok_or_else(|| { + MycGovernanceStateError::new(MycGovernanceStateErrorKind::InvalidPageLimit) + }) + } +} + +/// Immutable bounded audit page. +pub struct MycAuditPage { + snapshot_sequence: u64, + items: Box<[MycAuditRecord]>, + next_before_sequence: Option<u64>, +} + +impl MycAuditPage { + #[must_use] + pub const fn snapshot_sequence(&self) -> u64 { + self.snapshot_sequence + } + + #[must_use] + pub fn items(&self) -> &[MycAuditRecord] { + &self.items + } + + #[must_use] + pub const fn next_before_sequence(&self) -> Option<u64> { + self.next_before_sequence + } +} + +impl fmt::Debug for MycAuditPage { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycAuditPage") + .field("snapshot_sequence", &self.snapshot_sequence) + .field("item_count", &self.items.len()) + .finish() + } +} + +/// Explicit bounded retention and compaction policy. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct MycGovernanceCompactionPolicy { + audit_retention_ms: u64, + maximum_rows_per_class: u16, +} + +impl MycGovernanceCompactionPolicy { + pub fn new( + audit_retention_ms: u64, + maximum_rows_per_class: u16, + ) -> Result<Self, MycGovernanceStateError> { + if audit_retention_ms == 0 + || audit_retention_ms > MYC_AUDIT_RETENTION_MAX_MS + || maximum_rows_per_class == 0 + || maximum_rows_per_class > MYC_COMPACTION_MAX_ROWS + { + return Err(MycGovernanceStateError::new( + MycGovernanceStateErrorKind::InvalidCompactionPolicy, + )); + } + Ok(Self { + audit_retention_ms, + maximum_rows_per_class, + }) + } +} + +/// Bounded compaction result. Authoritative domain state is never included. +/// Counts describe this invocation; an exact correlation replay is a zero-work success. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct MycGovernanceCompactionOutcome { + removed_audit_records: u16, + removed_rate_subjects: u16, +} + +impl MycGovernanceCompactionOutcome { + #[must_use] + pub const fn removed_audit_records(self) -> u16 { + self.removed_audit_records + } + + #[must_use] + pub const fn removed_rate_subjects(self) -> u16 { + self.removed_rate_subjects + } +} + +#[derive(Clone, Copy)] +pub(crate) struct AuditEvidence { + pub(crate) correlation: MycAuditCorrelationId, + pub(crate) occurred_at: MycConnectionTimeUnixMs, + pub(crate) operation_id: Option<MycSignerOperationId>, +} + +#[derive(Clone, Copy)] +pub(crate) enum RateSubject { + Global, + Relay([u8; 32]), + Connection(MycConnectionId), +} + +pub(crate) fn relay_subject(relay: &MycRateRelayId) -> RateSubject { + RateSubject::Relay(hash_framed(RELAY_SUBJECT_DOMAIN, relay.0.as_bytes())) +} + +pub(crate) fn global_subject() -> RateSubject { + RateSubject::Global +} + +pub(crate) fn connection_subject(connection: MycConnectionId) -> RateSubject { + RateSubject::Connection(connection) +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum GovernanceOperationError { + Binding, + Storage, +} + +pub(crate) async fn govern_rate_attempt( + transaction: &mut ServiceSqliteTransaction<'_>, + policy: MycRateLimitPolicy, + expected_class: MycRateLimitClass, + subjects: &[RateSubject], + evidence: AuditEvidence, + audit_kind: MycAuditKind, +) -> Result<bool, GovernanceOperationError> { + if policy.class != expected_class || subjects.is_empty() || subjects.len() > 2 { + return Err(GovernanceOperationError::Binding); + } + if let Some(existing) = read_audit_by_id( + transaction, + derive_audit_id(evidence.correlation, audit_kind), + ) + .await? + { + return (existing.correlation_id == evidence.correlation + && existing.operation_id == evidence.operation_id + && existing.kind == audit_kind + && existing.outcome == MycAuditOutcome::Rejected + && existing.reason == MycAuditReasonCode::RateLimited) + .then_some(false) + .ok_or(GovernanceOperationError::Binding); + } + let mut states = Vec::with_capacity(subjects.len()); + for subject in subjects { + states.push(load_rate_state(transaction, policy, *subject, evidence.occurred_at).await?); + } + let admitted = states.iter().all(|state| state.can_accept(policy)); + for state in &states { + persist_rate_state(transaction, policy, *state, evidence.occurred_at, admitted).await?; + } + if !admitted { + record_audit( + transaction, + evidence, + audit_kind, + MycAuditOutcome::Rejected, + MycAuditReasonCode::RateLimited, + ) + .await?; + } + Ok(admitted) +} + +pub(crate) async fn record_audit( + transaction: &mut ServiceSqliteTransaction<'_>, + evidence: AuditEvidence, + kind: MycAuditKind, + outcome: MycAuditOutcome, + reason: MycAuditReasonCode, +) -> Result<MycAuditRecord, GovernanceOperationError> { + let audit_id = derive_audit_id(evidence.correlation, kind); + if let Some(existing) = read_audit_by_id(transaction, audit_id).await? { + return (existing.correlation_id == evidence.correlation + && existing.operation_id == evidence.operation_id + && existing.kind == kind + && existing.outcome == outcome + && existing.reason == reason + && existing.occurred_at == evidence.occurred_at) + .then_some(existing) + .ok_or(GovernanceOperationError::Binding); + } + let current = read_audit_state(transaction).await?; + let sequence = current + .checked_add(1) + .filter(|value| *value <= i64::MAX as u64) + .ok_or(GovernanceOperationError::Binding)?; + let result = sqlx::query(ADVANCE_AUDIT_STATE_SQL) + .bind(i64::try_from(sequence).map_err(|_| GovernanceOperationError::Binding)?) + .bind(i64::try_from(current).map_err(|_| GovernanceOperationError::Binding)?) + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + require_one(result.rows_affected())?; + let result = sqlx::query(INSERT_AUDIT_SQL) + .bind(i64::try_from(sequence).map_err(|_| GovernanceOperationError::Binding)?) + .bind(audit_id.as_slice()) + .bind(evidence.correlation.0.as_slice()) + .bind(kind.as_str()) + .bind(outcome.as_str()) + .bind(reason.as_str()) + .bind(to_i64(evidence.occurred_at.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + require_one(result.rows_affected())?; + if let Some(operation_id) = evidence.operation_id { + let result = sqlx::query(INSERT_REQUEST_AUDIT_SQL) + .bind(operation_id.as_bytes().as_slice()) + .bind(kind.as_str()) + .bind(i64::try_from(sequence).map_err(|_| GovernanceOperationError::Binding)?) + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + require_one(result.rows_affected())?; + } + Ok(MycAuditRecord { + sequence, + correlation_id: evidence.correlation, + operation_id: evidence.operation_id, + kind, + outcome, + reason, + occurred_at: evidence.occurred_at, + }) +} + +impl MycStateRepository<'_> { + /// Reads one deterministic, snapshot-bounded page of safe audit evidence. + pub async fn read_audit_page( + &self, + limit: MycAuditPageLimit, + snapshot_sequence: Option<u64>, + before_sequence: Option<u64>, + ) -> Result<MycAuditPage, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + read_audit_page(transaction, limit, snapshot_sequence, before_sequence).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Removes only expired safe audit and rate-window evidence in bounded batches. + pub async fn compact_governance_evidence( + &self, + observed_at: MycConnectionTimeUnixMs, + policy: MycGovernanceCompactionPolicy, + correlation: MycAuditCorrelationId, + ) -> Result<MycGovernanceCompactionOutcome, MycStateRepositoryError> { + if policy.audit_retention_ms != self.expected().governance_audit_retention_ms() { + return Err(MycStateRepositoryError::new( + MycStateRepositoryErrorKind::Binding, + )); + } + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + compact_governance(transaction, observed_at, policy, correlation).await + }) + }) + .await + .map_err(map_transaction_error) + } +} + +#[derive(Clone, Copy)] +struct RateState { + subject: RateSubject, + existing_last: Option<u64>, + window_started: u64, + window_ends: u64, + accepted: u64, + rejected: u64, + lifetime_accepted: u64, + lifetime_rejected: u64, +} + +impl RateState { + fn can_accept(self, policy: MycRateLimitPolicy) -> bool { + self.accepted < u64::from(policy.max_attempts) + } +} + +async fn load_rate_state( + transaction: &mut ServiceSqliteTransaction<'_>, + policy: MycRateLimitPolicy, + subject: RateSubject, + observed_at: MycConnectionTimeUnixMs, +) -> Result<RateState, GovernanceOperationError> { + let (scope, digest) = subject_parts(subject); + let rows = sqlx::query(READ_RATE_SQL) + .bind(policy.class.as_str()) + .bind(scope) + .bind(digest.as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + if rows.len() > 1 { + return Err(GovernanceOperationError::Binding); + } + let now = observed_at.get(); + let Some(row) = rows.first() else { + if scope != "global" { + let count = sqlx::query_scalar::<_, i64>(COUNT_RATE_SUBJECTS_SQL) + .bind(policy.class.as_str()) + .bind(scope) + .fetch_one(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + let count = u64::try_from(count).map_err(|_| GovernanceOperationError::Binding)?; + if count >= u64::from(policy.maximum_tracked_subjects) { + return Err(GovernanceOperationError::Binding); + } + } + return Ok(RateState { + subject, + existing_last: None, + window_started: now, + window_ends: checked_time_add(now, policy.window_ms)?, + accepted: 0, + rejected: 0, + lifetime_accepted: 0, + lifetime_rejected: 0, + }); + }; + let mut state = RateState { + subject, + existing_last: Some(bounded_i64(row, "last_observed_at_unix_ms")?), + window_started: bounded_i64(row, "window_started_at_unix_ms")?, + window_ends: bounded_i64(row, "window_ends_at_unix_ms")?, + accepted: bounded_nonnegative(row, "accepted_count")?, + rejected: bounded_nonnegative(row, "rejected_count")?, + lifetime_accepted: bounded_nonnegative(row, "lifetime_accepted_count")?, + lifetime_rejected: bounded_nonnegative(row, "lifetime_rejected_count")?, + }; + let retention_expires = bounded_i64(row, "retention_expires_at_unix_ms")?; + let last = state + .existing_last + .ok_or(GovernanceOperationError::Binding)?; + if now < last + || state.window_started > last + || state.window_ends < last + || retention_expires < last + || state.accepted > u64::from(policy.max_attempts) + { + return Err(GovernanceOperationError::Binding); + } + if now > state.window_ends { + state.window_started = now; + state.window_ends = checked_time_add(now, policy.window_ms)?; + state.accepted = 0; + state.rejected = 0; + } + Ok(state) +} + +async fn persist_rate_state( + transaction: &mut ServiceSqliteTransaction<'_>, + policy: MycRateLimitPolicy, + state: RateState, + observed_at: MycConnectionTimeUnixMs, + admitted: bool, +) -> Result<(), GovernanceOperationError> { + let (scope, digest) = subject_parts(state.subject); + let accepted = state.accepted + u64::from(admitted); + let rejected = state.rejected + u64::from(!admitted); + let lifetime_accepted = state.lifetime_accepted + u64::from(admitted); + let lifetime_rejected = state.lifetime_rejected + u64::from(!admitted); + for value in [accepted, rejected, lifetime_accepted, lifetime_rejected] { + if value > i64::MAX as u64 { + return Err(GovernanceOperationError::Binding); + } + } + let retention_expires = checked_time_add(observed_at.get(), policy.retention_ms)?; + let mut query = if state.existing_last.is_some() { + sqlx::query(UPDATE_RATE_SQL) + } else { + sqlx::query(INSERT_RATE_SQL) + }; + if state.existing_last.is_some() { + query = query + .bind(to_i64(state.window_started)?) + .bind(to_i64(state.window_ends)?) + .bind(to_i64(accepted)?) + .bind(to_i64(rejected)?) + .bind(to_i64(lifetime_accepted)?) + .bind(to_i64(lifetime_rejected)?) + .bind(to_i64(observed_at.get())?) + .bind(to_i64(retention_expires)?) + .bind(policy.class.as_str()) + .bind(scope) + .bind(digest.as_slice()) + .bind(to_i64( + state + .existing_last + .ok_or(GovernanceOperationError::Binding)?, + )?); + } else { + query = query + .bind(policy.class.as_str()) + .bind(scope) + .bind(digest.as_slice()) + .bind(to_i64(state.window_started)?) + .bind(to_i64(state.window_ends)?) + .bind(to_i64(accepted)?) + .bind(to_i64(rejected)?) + .bind(to_i64(lifetime_accepted)?) + .bind(to_i64(lifetime_rejected)?) + .bind(to_i64(observed_at.get())?) + .bind(to_i64(retention_expires)?); + } + let result = query + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + require_one(result.rows_affected()) +} + +async fn read_audit_page( + transaction: &mut ServiceSqliteTransaction<'_>, + limit: MycAuditPageLimit, + requested_snapshot: Option<u64>, + before: Option<u64>, +) -> Result<MycAuditPage, GovernanceOperationError> { + let high_water = read_audit_state(transaction).await?; + let snapshot = requested_snapshot.unwrap_or(high_water); + if snapshot > high_water + || before.is_some_and(|value| value == 0 || value > snapshot.saturating_add(1)) + { + return Err(GovernanceOperationError::Binding); + } + let before = before.unwrap_or_else(|| snapshot.saturating_add(1)); + let fetch_limit = u32::from(limit.0) + 1; + let rows = sqlx::query( + r#"SELECT audit_id FROM operation_audit + WHERE audit_sequence <= ? AND audit_sequence < ? + ORDER BY audit_sequence DESC LIMIT ?"#, + ) + .bind(to_i64(snapshot)?) + .bind(to_i64(before)?) + .bind(i64::from(fetch_limit)) + .fetch_all(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + if rows.len() > usize::try_from(fetch_limit).map_err(|_| GovernanceOperationError::Binding)? { + return Err(GovernanceOperationError::Binding); + } + let has_more = rows.len() > usize::from(limit.0); + let mut items = Vec::with_capacity(rows.len().min(usize::from(limit.0))); + for row in rows.iter().take(usize::from(limit.0)) { + let id = bounded_digest(row, "audit_id")?; + items.push( + read_audit_by_id(transaction, id) + .await? + .ok_or(GovernanceOperationError::Binding)?, + ); + } + let next = has_more + .then(|| items.last().map(|item| item.sequence)) + .flatten(); + Ok(MycAuditPage { + snapshot_sequence: snapshot, + items: items.into_boxed_slice(), + next_before_sequence: next, + }) +} + +async fn compact_governance( + transaction: &mut ServiceSqliteTransaction<'_>, + observed_at: MycConnectionTimeUnixMs, + policy: MycGovernanceCompactionPolicy, + correlation: MycAuditCorrelationId, +) -> Result<MycGovernanceCompactionOutcome, GovernanceOperationError> { + if let Some(existing) = read_audit_by_id( + transaction, + derive_audit_id(correlation, MycAuditKind::GovernanceCompaction), + ) + .await? + { + return (existing.correlation_id == correlation + && existing.operation_id.is_none() + && existing.kind == MycAuditKind::GovernanceCompaction + && existing.outcome == MycAuditOutcome::Succeeded + && existing.reason == MycAuditReasonCode::Compacted + && existing.occurred_at == observed_at) + .then_some(MycGovernanceCompactionOutcome { + removed_audit_records: 0, + removed_rate_subjects: 0, + }) + .ok_or(GovernanceOperationError::Binding); + } + let cutoff = observed_at.get().saturating_sub(policy.audit_retention_ms); + let limit = i64::from(policy.maximum_rows_per_class); + let request_links = sqlx::query( + r#"DELETE FROM nip46_request_audit WHERE audit_sequence IN ( + SELECT audit_sequence FROM operation_audit + WHERE occurred_at_unix_ms < ? ORDER BY audit_sequence ASC LIMIT ? + )"#, + ) + .bind(to_i64(cutoff)?) + .bind(limit) + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + let _ = request_links; + let audit = sqlx::query( + r#"DELETE FROM operation_audit WHERE audit_sequence IN ( + SELECT audit_sequence FROM operation_audit + WHERE occurred_at_unix_ms < ? ORDER BY audit_sequence ASC LIMIT ? + )"#, + ) + .bind(to_i64(cutoff)?) + .bind(limit) + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + let rates = sqlx::query( + r#"DELETE FROM connection_rate_windows WHERE rowid IN ( + SELECT rowid FROM connection_rate_windows + WHERE retention_expires_at_unix_ms < ? + ORDER BY retention_expires_at_unix_ms ASC, rate_kind ASC, + subject_scope ASC, subject_sha256 ASC LIMIT ? + )"#, + ) + .bind(to_i64(observed_at.get())?) + .bind(limit) + .execute(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + let removed_audit_records = + u16::try_from(audit.rows_affected()).map_err(|_| GovernanceOperationError::Binding)?; + let removed_rate_subjects = + u16::try_from(rates.rows_affected()).map_err(|_| GovernanceOperationError::Binding)?; + record_audit( + transaction, + AuditEvidence { + correlation, + occurred_at: observed_at, + operation_id: None, + }, + MycAuditKind::GovernanceCompaction, + MycAuditOutcome::Succeeded, + MycAuditReasonCode::Compacted, + ) + .await?; + Ok(MycGovernanceCompactionOutcome { + removed_audit_records, + removed_rate_subjects, + }) +} + +async fn read_audit_state( + transaction: &mut ServiceSqliteTransaction<'_>, +) -> Result<u64, GovernanceOperationError> { + let rows = sqlx::query(READ_AUDIT_STATE_SQL) + .fetch_all(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + if rows.len() != 1 { + return Err(GovernanceOperationError::Binding); + } + bounded_nonnegative(&rows[0], "next_sequence") +} + +async fn read_audit_by_id( + transaction: &mut ServiceSqliteTransaction<'_>, + audit_id: [u8; 32], +) -> Result<Option<MycAuditRecord>, GovernanceOperationError> { + let rows = sqlx::query(READ_AUDIT_BY_ID_SQL) + .bind(audit_id.as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + if rows.len() > 1 { + return Err(GovernanceOperationError::Binding); + } + rows.first().map(decode_audit).transpose() +} + +fn decode_audit(row: &sqlx::sqlite::SqliteRow) -> Result<MycAuditRecord, GovernanceOperationError> { + let correlation_id = MycAuditCorrelationId(bounded_digest(row, "correlation_id")?); + let kind = bounded_text(row, "audit_kind") + .and_then(|value| MycAuditKind::parse(&value).ok_or(GovernanceOperationError::Binding))?; + let outcome = bounded_text(row, "outcome").and_then(|value| { + MycAuditOutcome::parse(&value).ok_or(GovernanceOperationError::Binding) + })?; + let reason = bounded_text(row, "reason_code").and_then(|value| { + MycAuditReasonCode::parse(&value).ok_or(GovernanceOperationError::Binding) + })?; + let operation_id = + optional_digest(row, "operation_id")?.map(MycSignerOperationId::from_persisted); + let request_audit_kind = optional_text(row, "request_audit_kind")?; + let request_correlation_id = optional_digest(row, "request_correlation_id")?; + let request_bound = matches!( + kind, + MycAuditKind::ConnectionAdmission + | MycAuditKind::ChallengeCreation + | MycAuditKind::ChallengeAuthorization + ); + if request_bound + != (operation_id.is_some() + && request_audit_kind.as_deref() == Some(kind.as_str()) + && request_correlation_id == Some(correlation_id.0)) + || (!request_bound + && (operation_id.is_some() + || request_audit_kind.is_some() + || request_correlation_id.is_some())) + { + return Err(GovernanceOperationError::Binding); + } + let sequence = bounded_i64(row, "audit_sequence")?; + let occurred_at = MycConnectionTimeUnixMs::new(bounded_i64(row, "occurred_at_unix_ms")?) + .map_err(|_| GovernanceOperationError::Binding)?; + Ok(MycAuditRecord { + sequence, + correlation_id, + operation_id, + kind, + outcome, + reason, + occurred_at, + }) +} + +async fn verify_metadata( + transaction: &mut ServiceSqliteTransaction<'_>, + expected: &PersistedMetadata, +) -> Result<(), GovernanceOperationError> { + require_expected_metadata(transaction, expected) + .await + .map_err(|error| match error { + RepositoryOperationError::Binding => GovernanceOperationError::Binding, + RepositoryOperationError::Storage => GovernanceOperationError::Storage, + }) +} + +fn map_transaction_error( + error: ServiceSqliteTransactionError<GovernanceOperationError>, +) -> MycStateRepositoryError { + if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { + return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); + } + let kind = match error.operation_error() { + Some(GovernanceOperationError::Binding) => MycStateRepositoryErrorKind::Binding, + Some(GovernanceOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, + }; + MycStateRepositoryError::new(kind) +} + +fn subject_parts(subject: RateSubject) -> (&'static str, [u8; 32]) { + match subject { + RateSubject::Global => ("global", hash_framed(GLOBAL_SUBJECT_DOMAIN, b"global")), + RateSubject::Relay(digest) => ("relay", digest), + RateSubject::Connection(connection) => ( + "connection", + hash_framed(CONNECTION_SUBJECT_DOMAIN, connection.as_bytes()), + ), + } +} + +fn derive_audit_id(correlation: MycAuditCorrelationId, kind: MycAuditKind) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(AUDIT_ID_DOMAIN); + hasher.update(correlation.0); + hasher.update((kind.as_str().len() as u64).to_be_bytes()); + hasher.update(kind.as_str().as_bytes()); + hasher.finalize().into() +} + +fn hash_framed(domain: &[u8], value: &[u8]) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(domain); + hasher.update((value.len() as u64).to_be_bytes()); + hasher.update(value); + hasher.finalize().into() +} + +fn checked_time_add(left: u64, right: u64) -> Result<u64, GovernanceOperationError> { + left.checked_add(right) + .filter(|value| *value <= i64::MAX as u64) + .ok_or(GovernanceOperationError::Binding) +} + +fn to_i64(value: u64) -> Result<i64, GovernanceOperationError> { + i64::try_from(value).map_err(|_| GovernanceOperationError::Binding) +} + +fn bounded_i64( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u64, GovernanceOperationError> { + let value = row + .try_get::<i64, _>(column) + .map_err(|_| GovernanceOperationError::Binding)?; + u64::try_from(value).map_err(|_| GovernanceOperationError::Binding) +} + +fn bounded_nonnegative( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u64, GovernanceOperationError> { + bounded_i64(row, column) +} + +fn bounded_digest( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<[u8; 32], GovernanceOperationError> { + row.try_get::<Option<Vec<u8>>, _>(column) + .map_err(|_| GovernanceOperationError::Binding)? + .ok_or(GovernanceOperationError::Binding)? + .try_into() + .map_err(|_| GovernanceOperationError::Binding) +} + +fn bounded_text( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<Box<str>, GovernanceOperationError> { + row.try_get::<Option<String>, _>(column) + .map_err(|_| GovernanceOperationError::Binding)? + .ok_or(GovernanceOperationError::Binding) + .map(String::into_boxed_str) +} + +fn optional_digest( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<Option<[u8; 32]>, GovernanceOperationError> { + row.try_get::<Option<Vec<u8>>, _>(column) + .map_err(|_| GovernanceOperationError::Binding)? + .map(|value| { + value + .try_into() + .map_err(|_| GovernanceOperationError::Binding) + }) + .transpose() +} + +fn optional_text( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<Option<Box<str>>, GovernanceOperationError> { + row.try_get::<Option<String>, _>(column) + .map_err(|_| GovernanceOperationError::Binding) + .map(|value| value.map(String::into_boxed_str)) +} + +fn require_one(rows: u64) -> Result<(), GovernanceOperationError> { + (rows == 1) + .then_some(()) + .ok_or(GovernanceOperationError::Binding) +} diff --git a/src/state_host.rs b/src/state_host.rs @@ -377,14 +377,18 @@ fn require_migration_build( fn exact_initialization_outcome(outcome: MigrationApplicationOutcome) -> bool { outcome.initial_version() == MYC_STATE_BASE_SCHEMA_VERSION && outcome.final_version() == MYC_STATE_SCHEMA_VERSION - && outcome.applied_count() == 3 + && outcome.applied_count() == 4 } fn exact_existing_outcome(outcome: MigrationApplicationOutcome) -> bool { outcome.final_version() == MYC_STATE_SCHEMA_VERSION && matches!( (outcome.initial_version(), outcome.applied_count()), - (MYC_STATE_BASE_SCHEMA_VERSION, 3) | (2, 2) | (3, 1) | (MYC_STATE_SCHEMA_VERSION, 0) + (MYC_STATE_BASE_SCHEMA_VERSION, 4) + | (2, 3) + | (3, 2) + | (4, 1) + | (MYC_STATE_SCHEMA_VERSION, 0) ) } diff --git a/src/state_metadata.rs b/src/state_metadata.rs @@ -11,6 +11,9 @@ use radroots_service_sqlite::{ use radroots_storage::event::SourceGeneration; use sha2::{Digest, Sha256}; +use crate::state_governance::{ + MycGovernancePolicies, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, +}; use crate::{ MYC_CONFIG_SCHEMA_VERSION, MYC_SIGNER_STATUS_CONTRACT_VERSION, MYC_STATE_BASE_SCHEMA_VERSION, MYC_STATE_SCHEMA_VERSION, MycBootstrapProfileV1, MycConfigDocumentV1, MycConfigProfile, @@ -167,6 +170,7 @@ pub struct MycStateMetadata { database_identity: ServiceDatabaseIdentity, configuration: MycNormalizedConfigDigest, identities: MycExpectedIdentities, + governance: MycGovernancePolicies, policy_versions: MycStatePolicyVersions, } @@ -204,6 +208,7 @@ impl MycStateMetadata { application_id, ); let normalized = configuration.normalized(); + let governance = governance_policies(normalized)?; let configuration = normalized_config_digest(configuration.profile(), normalized)?; let identities = expected_identities(normalized)?; let policy_versions = MycStatePolicyVersions::governed(); @@ -225,6 +230,7 @@ impl MycStateMetadata { database_identity, configuration, identities, + governance, policy_versions, }) } @@ -259,6 +265,21 @@ impl MycStateMetadata { self.policy_versions } + pub(crate) const fn governance_rate_policy( + &self, + class: MycRateLimitClass, + ) -> MycRateLimitPolicy { + self.governance.rate_policy(class) + } + + pub(crate) fn admits_rate_relay(&self, relay: &MycRateRelayId) -> bool { + self.governance.admits_relay(relay) + } + + pub(crate) const fn governance_audit_retention_ms(&self) -> u64 { + self.governance.audit_retention_ms() + } + pub(crate) fn matches_runtime(&self, runtime: &MycRuntimeContext) -> bool { ServiceSqlitePaths::from_runtime_context(runtime.context()) .is_ok_and(|paths| paths == self.paths) @@ -273,6 +294,7 @@ impl fmt::Debug for MycStateMetadata { .field("database_identity", &self.database_identity) .field("configuration", &self.configuration) .field("identities", &self.identities) + .field("governance", &"[redacted]") .field("policy_versions", &self.policy_versions) .field("paths", &"[redacted]") .finish() @@ -367,6 +389,62 @@ fn normalized_config_digest( Ok(MycNormalizedConfigDigest(hasher.finalize().into())) } +fn governance_policies( + normalized: &serde_json::Value, +) -> Result<MycGovernancePolicies, MycStateMetadataError> { + let integer = |pointer: &str| { + normalized + .pointer(pointer) + .and_then(serde_json::Value::as_u64) + .ok_or_else(|| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant)) + }; + let policy = |class, name: &str| { + let prefix = format!("/rate_limits/{name}"); + MycRateLimitPolicy::new( + class, + integer(&format!("{prefix}/window_ms"))?, + u32::try_from(integer(&format!("{prefix}/max_attempts"))?) + .map_err(|_| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant))?, + integer(&format!("{prefix}/retention_ms"))?, + u32::try_from(integer(&format!("{prefix}/maximum_tracked_subjects"))?) + .map_err(|_| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant))?, + ) + .map_err(|_| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant)) + }; + let relays = normalized + .pointer("/relays") + .and_then(serde_json::Value::as_array) + .ok_or_else(|| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant))? + .iter() + .map(|relay| { + relay + .pointer("/id") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant)) + .and_then(|id| { + MycRateRelayId::new(id).map_err(|_| { + MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant) + }) + }) + }) + .collect::<Result<Vec<_>, _>>()? + .into_boxed_slice(); + MycGovernancePolicies::new( + policy( + MycRateLimitClass::ConnectionAdmission, + "connection_admission", + )?, + policy(MycRateLimitClass::ChallengeCreation, "challenge_creation")?, + policy( + MycRateLimitClass::ChallengeAuthorization, + "challenge_authorization", + )?, + integer("/policy/retention/audit_ms")?, + relays, + ) + .map_err(|_| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant)) +} + fn expected_identities( normalized: &serde_json::Value, ) -> Result<MycExpectedIdentities, MycStateMetadataError> { diff --git a/tests/services_hardening_connection_state.rs b/tests/services_hardening_connection_state.rs @@ -4,14 +4,20 @@ use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path}; use myc::{ - MYC_AUTHORIZATION_CHALLENGE_URL_MAX_BYTES, MYC_CONNECTION_PERMISSION_MAX_COUNT, - MYC_STATE_SCHEMA_VERSION, MycAuthorizationChallengeAdmission, MycAuthorizationChallengeNonce, - MycAuthorizationChallengeRequest, MycAuthorizationChallengeState, MycAuthorizationChallengeUrl, - MycConfigProfile, MycConnectionAdmission, MycConnectionAdmissionPolicy, - MycConnectionAdmissionRequest, MycConnectionNonce, MycConnectionOperatorDecision, - MycConnectionPermission, MycConnectionPermissionSet, MycConnectionPolicyGeneration, - MycConnectionStateErrorKind, MycConnectionStatus, MycConnectionTimeUnixMs, - MycNip46ClientPublicKey, MycNip46EventId, MycNip46RequestId, MycRequestReceivedAtUnixMs, + MYC_AUDIT_PAGE_MAX_ITEMS, MYC_AUDIT_RETENTION_MAX_MS, + MYC_AUTHORIZATION_CHALLENGE_URL_MAX_BYTES, MYC_COMPACTION_MAX_ROWS, + MYC_CONNECTION_PERMISSION_MAX_COUNT, MYC_RATE_MAX_ATTEMPTS, MYC_RATE_MAX_TRACKED_SUBJECTS, + MYC_RATE_RELAY_ID_MAX_BYTES, MYC_RATE_RETENTION_MAX_MS, MYC_RATE_WINDOW_MAX_MS, + MYC_STATE_SCHEMA_VERSION, MycAuditCorrelationId, MycAuditKind, MycAuditOutcome, + MycAuditPageLimit, MycAuditReasonCode, MycAuthorizationChallengeAdmission, + MycAuthorizationChallengeNonce, MycAuthorizationChallengeRequest, + MycAuthorizationChallengeState, MycAuthorizationChallengeUrl, MycConfigProfile, + MycConnectionAdmission, MycConnectionAdmissionPolicy, MycConnectionAdmissionRequest, + MycConnectionNonce, MycConnectionOperatorDecision, MycConnectionPermission, + MycConnectionPermissionSet, MycConnectionPolicyGeneration, MycConnectionStateErrorKind, + MycConnectionStatus, MycConnectionTimeUnixMs, MycGovernanceCompactionPolicy, + MycGovernanceStateErrorKind, MycNip46ClientPublicKey, MycNip46EventId, MycNip46RequestId, + MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, MycRequestReceivedAtUnixMs, MycSignerOperationId, MycSignerOperationNonce, MycSignerRequest, MycSignerRequestDigest, MycSignerRequestMethod, MycStateMetadata, MycStateRepository, MycStateRepositoryErrorKind, RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, initialize_myc_state, @@ -25,6 +31,7 @@ use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; const CONFIG_EXAMPLE: &[u8] = include_bytes!("../contracts/services_hardening/config.v1.example.toml"); const CONNECTION_SOURCE: &str = include_str!("../src/state_connection.rs"); +const GOVERNANCE_SOURCE: &str = include_str!("../src/state_governance.rs"); const CLIENT_PUBLIC_KEY: &str = "2222222222222222222222222222222222222222222222222222222222222222"; fn runtime(root: &Path) -> myc::MycRuntimeContext { @@ -54,8 +61,12 @@ fn prepare_state_directory(runtime: &myc::MycRuntimeContext) { } fn metadata(runtime: &myc::MycRuntimeContext) -> MycStateMetadata { + metadata_from_bytes(runtime, CONFIG_EXAMPLE) +} + +fn metadata_from_bytes(runtime: &myc::MycRuntimeContext, bytes: &[u8]) -> MycStateMetadata { let configuration = - parse_myc_config_v1(CONFIG_EXAMPLE, MycConfigProfile::RepoLocal).expect("configuration"); + parse_myc_config_v1(bytes, MycConfigProfile::RepoLocal).expect("configuration"); MycStateMetadata::new( runtime, &configuration, @@ -100,6 +111,10 @@ fn policy_generation(value: u64) -> MycConnectionPolicyGeneration { MycConnectionPolicyGeneration::new(value).expect("policy generation") } +fn audit_correlation(byte: u8) -> MycAuditCorrelationId { + MycAuditCorrelationId::new([byte; 32]) +} + async fn admit_request( repository: &MycStateRepository<'_>, request_id: &str, @@ -148,6 +163,7 @@ fn connection_request( time(observed_at), authorized_until.map(time), policy, + MycRateRelayId::new("primary").expect("relay ID"), ) .expect("connection request") } @@ -250,6 +266,617 @@ fn connection_inputs_are_closed_bounded_canonical_and_redacted() { } } +#[test] +fn governance_inputs_are_closed_bounded_and_redacted() { + assert!( + MycRateLimitPolicy::new( + MycRateLimitClass::ConnectionAdmission, + MYC_RATE_WINDOW_MAX_MS, + MYC_RATE_MAX_ATTEMPTS, + MYC_RATE_RETENTION_MAX_MS, + MYC_RATE_MAX_TRACKED_SUBJECTS, + ) + .is_ok() + ); + for invalid in [ + MycRateLimitPolicy::new(MycRateLimitClass::ConnectionAdmission, 0, 1, 1, 1), + MycRateLimitPolicy::new( + MycRateLimitClass::ConnectionAdmission, + MYC_RATE_WINDOW_MAX_MS + 1, + 1, + MYC_RATE_WINDOW_MAX_MS + 1, + 1, + ), + MycRateLimitPolicy::new(MycRateLimitClass::ConnectionAdmission, 1, 0, 1, 1), + MycRateLimitPolicy::new( + MycRateLimitClass::ConnectionAdmission, + 1, + MYC_RATE_MAX_ATTEMPTS + 1, + 1, + 1, + ), + MycRateLimitPolicy::new(MycRateLimitClass::ConnectionAdmission, 2, 1, 1, 1), + MycRateLimitPolicy::new( + MycRateLimitClass::ConnectionAdmission, + 1, + 1, + MYC_RATE_RETENTION_MAX_MS + 1, + 1, + ), + MycRateLimitPolicy::new(MycRateLimitClass::ConnectionAdmission, 1, 1, 1, 0), + MycRateLimitPolicy::new( + MycRateLimitClass::ConnectionAdmission, + 1, + 1, + 1, + MYC_RATE_MAX_TRACKED_SUBJECTS + 1, + ), + ] { + assert_eq!( + invalid.expect_err("invalid rate policy").kind(), + MycGovernanceStateErrorKind::InvalidRatePolicy + ); + } + + let maximum_relay = format!("a{}", "1".repeat(MYC_RATE_RELAY_ID_MAX_BYTES - 1)); + assert!(MycRateRelayId::new(&maximum_relay).is_ok()); + for invalid in [ + "", + "A", + "relay-name", + "relay__name", + "relay_", + &format!("a{maximum_relay}"), + ] { + assert_eq!( + MycRateRelayId::new(invalid) + .expect_err("invalid relay ID") + .kind(), + MycGovernanceStateErrorKind::InvalidRelayId + ); + } + assert!(MycAuditPageLimit::new(MYC_AUDIT_PAGE_MAX_ITEMS).is_ok()); + for value in [0, MYC_AUDIT_PAGE_MAX_ITEMS + 1] { + assert_eq!( + MycAuditPageLimit::new(value) + .expect_err("invalid page limit") + .kind(), + MycGovernanceStateErrorKind::InvalidPageLimit + ); + } + assert!( + MycGovernanceCompactionPolicy::new(MYC_AUDIT_RETENTION_MAX_MS, MYC_COMPACTION_MAX_ROWS,) + .is_ok() + ); + for invalid in [ + MycGovernanceCompactionPolicy::new(0, 1), + MycGovernanceCompactionPolicy::new(MYC_AUDIT_RETENTION_MAX_MS + 1, 1), + MycGovernanceCompactionPolicy::new(1, 0), + MycGovernanceCompactionPolicy::new(1, MYC_COMPACTION_MAX_ROWS + 1), + ] { + let error = invalid.expect_err("invalid compaction policy"); + assert_eq!( + error.kind(), + MycGovernanceStateErrorKind::InvalidCompactionPolicy + ); + assert!(Error::source(&error).is_none()); + } + + let relay = MycRateRelayId::new("relay_secret_123").expect("relay ID"); + let policy = + MycRateLimitPolicy::new(MycRateLimitClass::ChallengeCreation, 7, 3, 9, 11).expect("policy"); + let rendered = format!("{relay:?} {policy:?} {:?}", audit_correlation(0x91)); + for secret in ["relay_secret_123", "[145, 145", "window_ms", "retention_ms"] { + assert!(!rendered.contains(secret)); + } +} + +#[tokio::test] +async fn rate_windows_audit_pagination_retention_and_compaction_are_durable_and_bounded() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + prepare_state_directory(&runtime); + let configuration = std::str::from_utf8(CONFIG_EXAMPLE) + .expect("UTF-8 configuration") + .replacen( + "[policy.retention]\nterminal_connections_ms = 604800000\nterminal_challenges_ms = 86400000\nrequest_dedup_ms = 604800000\naudit_ms = 2592000000\ncompleted_outbox_ms = 604800000", + "[policy.retention]\nterminal_connections_ms = 604800000\nterminal_challenges_ms = 86400000\nrequest_dedup_ms = 604800000\naudit_ms = 500\ncompleted_outbox_ms = 604800000", + 1, + ) + .replacen( + "[rate_limits.connection_admission]\nscope = \"global_and_relay\"\nwindow_ms = 60000\nmax_attempts = 10\nretention_ms = 3600000\nmaximum_tracked_subjects = 4096", + "[rate_limits.connection_admission]\nscope = \"global_and_relay\"\nwindow_ms = 100\nmax_attempts = 1\nretention_ms = 200\nmaximum_tracked_subjects = 8", + 1, + ); + let metadata = metadata_from_bytes(&runtime, configuration.as_bytes()); + let (applied_at, build) = migration_evidence(); + initialize_myc_state(&runtime, &metadata, applied_at, &build) + .await + .expect("state initialization"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("writable host"); + let repository = host.repository(); + + let first_operation = admit_request( + &repository, + "rate-first", + 0x80, + MycSignerRequestMethod::Connect, + 0x81, + 100, + ) + .await; + let second_operation = admit_request( + &repository, + "rate-second", + 0x82, + MycSignerRequestMethod::Connect, + 0x83, + 100, + ) + .await; + let first_request = connection_request( + first_operation, + permission_set(&[MycConnectionPermission::Ping]), + 11, + 0x84, + 100, + Some(1_000), + MycConnectionAdmissionPolicy::Trusted, + ); + let unconfigured_relay_request = MycConnectionAdmissionRequest::new( + first_operation, + client(), + permission_set(&[MycConnectionPermission::Ping]), + policy_generation(11), + MycConnectionNonce::from_injected_entropy([0x84; 32]), + time(100), + Some(time(1_000)), + MycConnectionAdmissionPolicy::Trusted, + MycRateRelayId::new("unconfigured").expect("relay ID"), + ) + .expect("unconfigured relay request"); + assert_eq!( + repository + .admit_connection(&unconfigured_relay_request) + .await + .expect_err("unconfigured relay") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + let second_request = connection_request( + second_operation, + permission_set(&[MycConnectionPermission::Ping]), + 11, + 0x85, + 100, + Some(1_000), + MycConnectionAdmissionPolicy::Trusted, + ); + let (first, second) = tokio::join!( + repository.admit_connection(&first_request), + repository.admit_connection(&second_request) + ); + let first = first.expect("first concurrent admission"); + let second = second.expect("second concurrent admission"); + assert_eq!( + usize::from(matches!(first, MycConnectionAdmission::Admitted(_))) + + usize::from(matches!(second, MycConnectionAdmission::Admitted(_))), + 1 + ); + assert_eq!( + usize::from(matches!(first, MycConnectionAdmission::RateLimited)) + + usize::from(matches!(second, MycConnectionAdmission::RateLimited)), + 1 + ); + let (accepted_request, rejected_request) = + if matches!(first, MycConnectionAdmission::Admitted(_)) { + (&first_request, &second_request) + } else { + (&second_request, &first_request) + }; + assert!(matches!( + repository + .admit_connection(rejected_request) + .await + .expect("stable rate-limited replay"), + MycConnectionAdmission::RateLimited + )); + + let boundary_operation = admit_request( + &repository, + "rate-boundary", + 0x86, + MycSignerRequestMethod::Connect, + 0x87, + 200, + ) + .await; + let boundary_request = connection_request( + boundary_operation, + permission_set(&[MycConnectionPermission::Ping]), + 11, + 0x88, + 200, + Some(1_000), + MycConnectionAdmissionPolicy::Trusted, + ); + assert!(matches!( + repository + .admit_connection(&boundary_request) + .await + .expect("exact-window boundary"), + MycConnectionAdmission::RateLimited + )); + + let reset_operation = admit_request( + &repository, + "rate-reset", + 0x89, + MycSignerRequestMethod::Connect, + 0x8a, + 201, + ) + .await; + let reset_request = connection_request( + reset_operation, + permission_set(&[MycConnectionPermission::Ping]), + 11, + 0x8b, + 201, + Some(1_000), + MycConnectionAdmissionPolicy::Trusted, + ); + assert!(matches!( + repository + .admit_connection(&reset_request) + .await + .expect("new window"), + MycConnectionAdmission::Admitted(_) + )); + + let page_one = repository + .read_audit_page(MycAuditPageLimit::new(2).expect("page limit"), None, None) + .await + .expect("first audit page"); + assert_eq!(page_one.snapshot_sequence(), 4); + assert_eq!(page_one.items().len(), 2); + assert!( + page_one + .items() + .iter() + .all(|record| record.operation_id().is_some()) + ); + assert_eq!(page_one.items()[0].outcome(), MycAuditOutcome::Succeeded); + assert_eq!( + page_one.items()[1].reason(), + MycAuditReasonCode::RateLimited + ); + let page_two = repository + .read_audit_page( + MycAuditPageLimit::new(2).expect("page limit"), + Some(page_one.snapshot_sequence()), + page_one.next_before_sequence(), + ) + .await + .expect("second audit page"); + assert_eq!(page_two.items().len(), 2); + assert!(page_two.next_before_sequence().is_none()); + assert_eq!( + page_two + .items() + .iter() + .filter(|record| record.outcome() == MycAuditOutcome::Succeeded) + .count(), + 1 + ); + assert_eq!( + repository + .read_audit_page( + MycAuditPageLimit::new(1).expect("page limit"), + Some(page_one.snapshot_sequence() + 1), + None, + ) + .await + .expect_err("future snapshot") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + assert_eq!( + repository + .compact_governance_evidence( + time(1_000), + MycGovernanceCompactionPolicy::new(499, 1).expect("structural policy"), + audit_correlation(0x8d), + ) + .await + .expect_err("retention policy must match bound configuration") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + + host.close() + .await + .expect("host close before retention pass"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("reopened host"); + let repository = host.repository(); + assert_eq!( + repository + .read_audit_page( + MycAuditPageLimit::new(4).expect("page limit"), + Some(page_one.snapshot_sequence()), + None, + ) + .await + .expect("reopened audit page") + .items() + .len(), + 4 + ); + let compacted = repository + .compact_governance_evidence( + time(1_000), + MycGovernanceCompactionPolicy::new(500, MYC_COMPACTION_MAX_ROWS) + .expect("compaction policy"), + audit_correlation(0x8c), + ) + .await + .expect("bounded compaction"); + assert_eq!(compacted.removed_audit_records(), 4); + assert_eq!(compacted.removed_rate_subjects(), 2); + let replayed_compaction = repository + .compact_governance_evidence( + time(1_000), + MycGovernanceCompactionPolicy::new(500, MYC_COMPACTION_MAX_ROWS) + .expect("compaction policy"), + audit_correlation(0x8c), + ) + .await + .expect("idempotent compaction replay"); + assert_eq!(replayed_compaction.removed_audit_records(), 0); + assert_eq!(replayed_compaction.removed_rate_subjects(), 0); + let retained = repository + .read_audit_page(MycAuditPageLimit::new(2).expect("page limit"), None, None) + .await + .expect("retained compaction audit"); + assert_eq!(retained.items().len(), 1); + assert_eq!( + retained.items()[0].kind(), + MycAuditKind::GovernanceCompaction + ); + assert_eq!( + retained.items()[0].correlation_id(), + audit_correlation(0x8c) + ); + assert!(retained.items()[0].operation_id().is_none()); + assert_eq!(retained.items()[0].reason(), MycAuditReasonCode::Compacted); + assert!(matches!( + repository + .admit_connection(accepted_request) + .await + .expect("authoritative decision replay"), + MycConnectionAdmission::ExactReplay(_) + )); + assert!(matches!( + repository + .admit_connection(&reset_request) + .await + .expect("second authoritative decision replay"), + MycConnectionAdmission::ExactReplay(_) + )); + host.close().await.expect("final host close"); + + let options = SqliteConnectOptions::new() + .filename(runtime.artifacts().state_database()) + .create_if_missing(false) + .read_only(true) + .disable_statement_logging(); + let mut database = sqlx::SqliteConnection::connect_with(&options) + .await + .expect("inspection connection"); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM connections") + .fetch_one(&mut database) + .await + .expect("connection count"), + 2 + ); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_request_decisions") + .fetch_one(&mut database) + .await + .expect("decision count"), + 2 + ); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM operation_audit") + .fetch_one(&mut database) + .await + .expect("audit count"), + 1 + ); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM connection_rate_windows") + .fetch_one(&mut database) + .await + .expect("rate count"), + 0 + ); + database.close().await.expect("inspection close"); +} + +#[tokio::test] +async fn challenge_creation_and_authorization_use_distinct_durable_rate_budgets() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + prepare_state_directory(&runtime); + let configuration = std::str::from_utf8(CONFIG_EXAMPLE) + .expect("UTF-8 configuration") + .replacen( + "[rate_limits.challenge_creation]\nscope = \"connection\"\nwindow_ms = 120000\nmax_attempts = 5\nretention_ms = 86400000\nmaximum_tracked_subjects = 16384", + "[rate_limits.challenge_creation]\nscope = \"connection\"\nwindow_ms = 100\nmax_attempts = 2\nretention_ms = 200\nmaximum_tracked_subjects = 8", + 1, + ) + .replacen( + "[rate_limits.challenge_authorization]\nscope = \"connection\"\nwindow_ms = 120000\nmax_attempts = 5\nretention_ms = 86400000\nmaximum_tracked_subjects = 16384", + "[rate_limits.challenge_authorization]\nscope = \"connection\"\nwindow_ms = 100\nmax_attempts = 1\nretention_ms = 200\nmaximum_tracked_subjects = 8", + 1, + ); + let metadata = metadata_from_bytes(&runtime, configuration.as_bytes()); + let (applied_at, build) = migration_evidence(); + initialize_myc_state(&runtime, &metadata, applied_at, &build) + .await + .expect("state initialization"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("writable host"); + let repository = host.repository(); + let connect = admit_request( + &repository, + "distinct-rates-connect", + 0xa0, + MycSignerRequestMethod::Connect, + 0xa1, + 10, + ) + .await; + let connection = repository + .admit_connection(&connection_request( + connect, + permission_set(&[MycConnectionPermission::Ping]), + 17, + 0xa2, + 11, + Some(1_000), + MycConnectionAdmissionPolicy::Trusted, + )) + .await + .expect("connection admission") + .record() + .expect("admission record") + .connection() + .expect("connection") + .clone(); + let first_operation = admit_request( + &repository, + "distinct-rates-first", + 0xa3, + MycSignerRequestMethod::Ping, + 0xa4, + 20, + ) + .await; + let second_operation = admit_request( + &repository, + "distinct-rates-second", + 0xa5, + MycSignerRequestMethod::Ping, + 0xa6, + 21, + ) + .await; + let first_request = MycAuthorizationChallengeRequest::new( + first_operation, + connection.id(), + policy_generation(17), + MycAuthorizationChallengeUrl::new("https://operator.example/first").expect("URL"), + MycAuthorizationChallengeNonce::from_injected_entropy([0xa7; 32]), + time(22), + time(200), + ) + .expect("first challenge request"); + let second_request = MycAuthorizationChallengeRequest::new( + second_operation, + connection.id(), + policy_generation(17), + MycAuthorizationChallengeUrl::new("https://operator.example/second").expect("URL"), + MycAuthorizationChallengeNonce::from_injected_entropy([0xa8; 32]), + time(23), + time(200), + ) + .expect("second challenge request"); + let first = repository + .issue_authorization_challenge(&first_request) + .await + .expect("first challenge"); + let second = repository + .issue_authorization_challenge(&second_request) + .await + .expect("second challenge"); + assert!(matches!( + first, + MycAuthorizationChallengeAdmission::Created(_) + )); + assert!(matches!( + second, + MycAuthorizationChallengeAdmission::Created(_) + )); + let first_id = first.record().expect("first challenge record").id(); + let second_id = second.record().expect("second challenge record").id(); + assert!( + repository + .authorize_challenge( + first_id, + connection.id(), + first_operation, + policy_generation(17), + time(24), + ) + .await + .expect("first authorization") + .record() + .is_some() + ); + let limited = repository + .authorize_challenge( + second_id, + connection.id(), + second_operation, + policy_generation(17), + time(25), + ) + .await + .expect("bounded authorization"); + assert!(limited.record().is_none()); + assert!(format!("{limited:?}").contains("RateLimited")); + assert!( + repository + .authorize_challenge( + second_id, + connection.id(), + second_operation, + policy_generation(17), + time(125), + ) + .await + .expect("stable authorization rejection") + .record() + .is_none() + ); + host.close().await.expect("host close"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("reopened host"); + assert!( + host.repository() + .authorize_challenge( + second_id, + connection.id(), + second_operation, + policy_generation(17), + time(225), + ) + .await + .expect("reopened stable rejection") + .record() + .is_none() + ); + host.close().await.expect("final close"); +} + #[tokio::test] async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_denial_direct() { let directory = tempfile::tempdir().expect("temporary root"); @@ -308,10 +935,14 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ .expect("trusted admission"); assert!(matches!(trusted, MycConnectionAdmission::Admitted(_))); assert_eq!( - trusted.record().decision(), + trusted.record().expect("admission record").decision(), myc::MycConnectionDecision::Allowed ); - let trusted_connection = trusted.record().connection().expect("connection"); + let trusted_connection = trusted + .record() + .expect("admission record") + .connection() + .expect("connection"); assert_eq!(trusted_connection.status(), MycConnectionStatus::Active); assert_eq!(trusted_connection.granted_permissions(), &requested); let trusted_id = trusted_connection.id(); @@ -338,6 +969,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ assert_eq!( trusted_replay .record() + .expect("admission record") .connection() .expect("connection") .id(), @@ -366,11 +998,12 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ .await .expect("pending admission"); assert_eq!( - pending.record().decision(), + pending.record().expect("admission record").decision(), myc::MycConnectionDecision::PendingApproval ); let pending_id = pending .record() + .expect("admission record") .connection() .expect("pending connection") .id(); @@ -382,6 +1015,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ pending_id, policy_generation(2), time(120), + audit_correlation(0x70), MycConnectionOperatorDecision::Approve { granted_permissions: granted.clone(), authorized_until: Some(time(900)), @@ -398,6 +1032,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ pending_id, policy_generation(2), time(130), + audit_correlation(0x71), MycConnectionOperatorDecision::Approve { granted_permissions: granted.clone(), authorized_until: Some(time(900)), @@ -413,6 +1048,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ pending_id, policy_generation(2), time(131), + audit_correlation(0x71), MycConnectionOperatorDecision::Approve { granted_permissions: granted.clone(), authorized_until: Some(time(900)), @@ -438,7 +1074,10 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ MycConnectionAdmission::ExactReplay(_) )); assert_eq!( - admission_after_approval.record().decision(), + admission_after_approval + .record() + .expect("admission record") + .decision(), myc::MycConnectionDecision::Allowed ); @@ -465,6 +1104,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ .expect("operator-denied pending admission"); let operator_denied_id = operator_pending .record() + .expect("admission record") .connection() .expect("operator-denied connection") .id(); @@ -474,6 +1114,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ operator_denied_id, policy_generation(2), time(137), + audit_correlation(0x72), MycConnectionOperatorDecision::Deny, ) .await @@ -486,6 +1127,7 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ operator_denied_id, policy_generation(2), time(138), + audit_correlation(0x72), MycConnectionOperatorDecision::Deny, ) .await @@ -515,10 +1157,16 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ .await .expect("direct denial"); assert_eq!( - denied.record().decision(), + denied.record().expect("admission record").decision(), myc::MycConnectionDecision::Denied ); - assert!(denied.record().connection().is_none()); + assert!( + denied + .record() + .expect("admission record") + .connection() + .is_none() + ); let denied_replay = repository .admit_connection(&connection_request( denied_operation, @@ -535,7 +1183,13 @@ async fn connection_admission_and_operator_decisions_are_atomic_replay_safe_and_ denied_replay, MycConnectionAdmission::ExactReplay(_) )); - assert!(denied_replay.record().connection().is_none()); + assert!( + denied_replay + .record() + .expect("admission record") + .connection() + .is_none() + ); let mismatch = repository .admit_connection(&connection_request( @@ -616,6 +1270,7 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound .await .expect("connection") .record() + .expect("admission record") .connection() .expect("active connection") .clone(); @@ -680,10 +1335,10 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound MycAuthorizationChallengeAdmission::Created(_) )); assert_eq!( - challenge.record().state(), + challenge.record().expect("challenge record").state(), MycAuthorizationChallengeState::Pending ); - let challenge_id = challenge.record().id(); + let challenge_id = challenge.record().expect("challenge record").id(); assert_eq!( hex(challenge_id.as_bytes()), "b0f43d00c70db2bf058d461a2c1ac940ef52cfc76a4d271b7bbc89464fb25abb" @@ -707,7 +1362,10 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound replay, MycAuthorizationChallengeAdmission::ExactReplay(_) )); - assert_eq!(replay.record().id(), challenge_id); + assert_eq!( + replay.record().expect("challenge record").id(), + challenge_id + ); assert_eq!( repository @@ -735,10 +1393,16 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound .await .expect("authorization at exact deadline"); assert_eq!( - authorized.state(), + authorized.record().expect("authorization record").state(), MycAuthorizationChallengeState::Authorized ); - assert_eq!(authorized.resolved_at(), Some(time(300))); + assert_eq!( + authorized + .record() + .expect("authorization record") + .resolved_at(), + Some(time(300)) + ); let terminal_replay = repository .authorize_challenge( challenge_id, @@ -749,7 +1413,10 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound ) .await .expect("terminal replay"); - assert_eq!(terminal_replay, authorized); + assert_eq!( + terminal_replay.record().expect("authorization record"), + authorized.record().expect("authorization record") + ); let deadline_operation = admit_request( &repository, @@ -757,7 +1424,7 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound 0x3a, MycSignerRequestMethod::Ping, 0x3b, - 215, + 410, ) .await; let deadline_request = MycAuthorizationChallengeRequest::new( @@ -766,8 +1433,8 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound policy_generation(7), MycAuthorizationChallengeUrl::new("https://operator.example/deadline").expect("URL"), MycAuthorizationChallengeNonce::from_injected_entropy([0x3c; 32]), - time(216), - time(240), + time(416), + time(440), ) .expect("deadline request"); let deadline_challenge = repository @@ -776,16 +1443,19 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound .expect("deadline challenge"); let deadline_expired = repository .authorize_challenge( - deadline_challenge.record().id(), + deadline_challenge.record().expect("challenge record").id(), connection.id(), deadline_operation, policy_generation(7), - time(241), + time(441), ) .await .expect("challenge expiry"); assert_eq!( - deadline_expired.state(), + deadline_expired + .record() + .expect("authorization record") + .state(), MycAuthorizationChallengeState::Expired ); @@ -795,7 +1465,7 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound 0x37, MycSignerRequestMethod::Ping, 0x38, - 220, + 450, ) .await; let expiring_request = MycAuthorizationChallengeRequest::new( @@ -804,7 +1474,7 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound policy_generation(7), MycAuthorizationChallengeUrl::new("https://operator.example/expiry").expect("URL"), MycAuthorizationChallengeNonce::from_injected_entropy([0x39; 32]), - time(221), + time(451), time(600), ) .expect("expiring request"); @@ -814,7 +1484,7 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound .expect("expiring challenge"); let expired = repository .authorize_challenge( - expiring.record().id(), + expiring.record().expect("challenge record").id(), connection.id(), expiring_operation, policy_generation(7), @@ -822,16 +1492,29 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound ) .await .expect("connection-expired challenge"); - assert_eq!(expired.state(), MycAuthorizationChallengeState::Expired); + assert_eq!( + expired.record().expect("authorization record").state(), + MycAuthorizationChallengeState::Expired + ); let expired_connection = repository - .expire_connection(connection.id(), policy_generation(7), time(501)) + .expire_connection( + connection.id(), + policy_generation(7), + time(501), + audit_correlation(0x73), + ) .await .expect("connection expiry"); assert_eq!(expired_connection.status(), MycConnectionStatus::Expired); assert_eq!( repository - .expire_connection(connection.id(), policy_generation(7), time(502)) + .expire_connection( + connection.id(), + policy_generation(7), + time(502), + audit_correlation(0x73), + ) .await .expect("expiry replay"), expired_connection @@ -845,8 +1528,10 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound MycAuthorizationChallengeAdmission::ExactReplay(_) )); assert_eq!( - challenge_replay_after_connection_expiry.record(), - &authorized + challenge_replay_after_connection_expiry + .record() + .expect("challenge record"), + authorized.record().expect("authorization record") ); let wrong_binding = repository @@ -862,7 +1547,7 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound assert_eq!(wrong_binding.kind(), MycStateRepositoryErrorKind::Binding); let rendered = format!( "{challenge:?} {:?} {wrong_binding} {wrong_binding:?}", - challenge.record() + challenge.record().expect("challenge record") ); for secret in [ "operator.example", @@ -889,7 +1574,10 @@ async fn challenge_authorization_expiry_and_terminal_replay_remain_exactly_bound ) .await .expect("restart replay"); - assert_eq!(replay, authorized); + assert_eq!( + replay.record().expect("authorization record"), + authorized.record().expect("authorization record") + ); host.close().await.expect("final close"); } @@ -913,6 +1601,10 @@ fn connection_state_source_has_no_ambient_or_external_authority() { !CONNECTION_SOURCE.contains(forbidden), "found forbidden connection-state authority `{forbidden}`" ); + assert!( + !GOVERNANCE_SOURCE.contains(forbidden), + "found forbidden governance-state authority `{forbidden}`" + ); } for required in [ "ServiceSqliteTransaction", @@ -926,4 +1618,7 @@ fn connection_state_source_has_no_ambient_or_external_authority() { "missing governed connection-state boundary `{required}`" ); } + assert!(GOVERNANCE_SOURCE.contains("ServiceSqliteTransaction")); + assert!(GOVERNANCE_SOURCE.contains("LIMIT ?")); + assert!(!GOVERNANCE_SOURCE.contains("SELECT *")); } diff --git a/tests/services_hardening_signer_request_state.rs b/tests/services_hardening_signer_request_state.rs @@ -458,7 +458,7 @@ async fn concurrent_identical_admission_creates_one_request_and_bounded_replay_e } #[tokio::test] -async fn exact_schema_v3_state_advances_to_v4_before_request_admission() { +async fn exact_schema_v3_state_advances_to_v5_before_request_admission() { let directory = tempfile::tempdir().expect("temporary root"); let runtime = runtime(directory.path()); prepare_state_directory(&runtime); @@ -500,6 +500,11 @@ async fn exact_schema_v3_state_advances_to_v4_before_request_admission() { "DROP TRIGGER radroots_service_metadata_guard_update", "DROP TRIGGER myc_state_metadata_no_update", "DROP TRIGGER schema_migrations_no_delete", + "DROP TRIGGER connection_rate_windows_guard_update", + "DROP TRIGGER nip46_request_audit_no_update", + "DROP TRIGGER operation_audit_no_update", + "DROP TRIGGER myc_audit_state_no_delete", + "DROP TRIGGER myc_audit_state_guard_update", "DROP TRIGGER connection_auth_challenges_no_delete", "DROP TRIGGER connection_auth_challenges_guard_update", "DROP TRIGGER nip46_request_decisions_no_delete", @@ -508,13 +513,17 @@ async fn exact_schema_v3_state_advances_to_v4_before_request_admission() { "DROP TRIGGER connection_permissions_no_update", "DROP TRIGGER connections_no_delete", "DROP TRIGGER connections_guard_update", + "DROP TABLE nip46_request_audit", + "DROP TABLE operation_audit", + "DROP TABLE connection_rate_windows", + "DROP TABLE myc_audit_state", "DROP TABLE nip46_request_decisions", "DROP TABLE connection_auth_challenges", "DROP TABLE connection_permissions", "DROP TABLE connections", "UPDATE radroots_service_metadata SET state_schema_version = 3 WHERE singleton = 1", "UPDATE myc_state_metadata SET state_contract_version = 3 WHERE singleton = 1", - "DELETE FROM schema_migrations WHERE version = 4", + "DELETE FROM schema_migrations WHERE version IN (4, 5)", ] { sqlx::query(sql) .execute(&mut connection) diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs @@ -10,8 +10,9 @@ use myc::{ MYC_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_3_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_3_SHA256, MYC_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_4_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_4_SHA256, - MycStateCatalogErrorKind, myc_migration_catalog, myc_schema_catalog, - validate_myc_state_catalogs, + MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT, + MYC_STATE_SCHEMA_VERSION_5_SHA256, MycStateCatalogErrorKind, myc_migration_catalog, + myc_schema_catalog, validate_myc_state_catalogs, }; use radroots_service_sqlite::{ MigrationCatalog, MigrationChecksum, MigrationDescriptor, SchemaCatalog, SchemaDigest, @@ -23,13 +24,13 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); #[test] -fn schema_v1_through_v4_and_all_migrations_have_exact_literal_identities() { +fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { let migrations = myc_migration_catalog().expect("Myc migration catalog"); let schema = myc_schema_catalog().expect("Myc schema catalog"); assert_eq!(MYC_STATE_BASE_SCHEMA_VERSION, 1); - assert_eq!(MYC_STATE_SCHEMA_VERSION, 4); - assert_eq!(migrations.descriptors().len(), 3); + assert_eq!(MYC_STATE_SCHEMA_VERSION, 5); + assert_eq!(migrations.descriptors().len(), 4); let metadata = &migrations.descriptors()[0]; assert_eq!(metadata.target_version(), 2); assert_eq!(metadata.name().as_str(), "create_myc_state_metadata"); @@ -54,13 +55,23 @@ fn schema_v1_through_v4_and_all_migrations_have_exact_literal_identities() { connection.checksum().as_bytes(), &MYC_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256 ); - assert_eq!(migrations.current_version(), 4); + let governance = &migrations.descriptors()[3]; + assert_eq!(governance.target_version(), 5); + assert_eq!( + governance.name().as_str(), + "create_bounded_governance_state" + ); + assert_eq!( + governance.checksum().as_bytes(), + &MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256 + ); + assert_eq!(migrations.current_version(), 5); assert_eq!( migrations.digest().as_bytes(), &MYC_MIGRATION_CATALOG_SHA256 ); - assert_eq!(schema.versions().len(), 4); + assert_eq!(schema.versions().len(), 5); assert_eq!(schema.versions()[0].version(), 1); assert_eq!( schema.versions()[0].object_count(), @@ -100,6 +111,16 @@ fn schema_v1_through_v4_and_all_migrations_have_exact_literal_identities() { schema.versions()[3].digest().as_bytes(), &MYC_STATE_SCHEMA_VERSION_4_SHA256 ); + assert_eq!(schema.versions()[4].version(), 5); + assert_eq!( + schema.versions()[4].object_count(), + MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT + ); + assert_eq!(schema.versions()[4].object_count(), 34); + assert_eq!( + schema.versions()[4].digest().as_bytes(), + &MYC_STATE_SCHEMA_VERSION_5_SHA256 + ); assert_eq!(schema.digest().as_bytes(), &MYC_STATE_SCHEMA_CATALOG_SHA256); assert_eq!(schema.migration_catalog_digest(), migrations.digest()); validate_myc_state_catalogs(&migrations, &schema).expect("exact catalogs"); @@ -110,7 +131,7 @@ fn schema_v1_through_v4_and_all_migrations_have_exact_literal_identities() { ); assert_eq!( hex::encode(MYC_MIGRATION_CATALOG_SHA256), - "453d99f4c19c094c592f1a3fe7e28e82c1dea674514a9e7d063bc466c20bd4bd" + "1d069996435560217dbd3692d5430612ad3a5667ee0b12b1919f0b5edf2cbe7d" ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_VERSION_1_SHA256), @@ -122,7 +143,7 @@ fn schema_v1_through_v4_and_all_migrations_have_exact_literal_identities() { ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_CATALOG_SHA256), - "a47d9f0af804202bfa07268e08fb484ded590b785cc42a01094d1a8b9f37eeca" + "2119efefcfd4ac5477609655341aa4c5b9b9a1befb88cc168a2b9373e99ebbd5" ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256), @@ -140,6 +161,14 @@ fn schema_v1_through_v4_and_all_migrations_have_exact_literal_identities() { hex::encode(MYC_STATE_SCHEMA_VERSION_4_SHA256), "479863d37d91e6c269fa3573db6c2e767cdd3b24a93ab482d774bcec0218c174" ); + assert_eq!( + hex::encode(MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256), + "0e004fcb5d0bc7c951b16f4334533d10efcb18a826f624de0922d86dc8504b33" + ); + assert_eq!( + hex::encode(MYC_STATE_SCHEMA_VERSION_5_SHA256), + "fe89d4af7de4eded78a3f062dc9d300f06cf0f4cfca5fa5cffc93af1146f21c5" + ); } #[test] @@ -184,8 +213,19 @@ fn independent_validator_rejects_migration_or_schema_drift() { let v4_digest = SchemaVersionCatalog::computed_digest(4, [object.clone()]).expect("schema-v4 digest"); let v4 = SchemaVersionCatalog::new(4, [object], v4_digest).expect("schema v4"); - let schema = - SchemaCatalog::new(&expected_migrations, [v1, v2, v3, v4]).expect("drift schema catalog"); + let object = SchemaObject::new( + SchemaObjectKind::Table, + "unexpected", + "unexpected", + SQL, + object_digest, + ) + .expect("schema object"); + let v5_digest = + SchemaVersionCatalog::computed_digest(5, [object.clone()]).expect("schema-v5 digest"); + let v5 = SchemaVersionCatalog::new(5, [object], v5_digest).expect("schema v5"); + let schema = SchemaCatalog::new(&expected_migrations, [v1, v2, v3, v4, v5]) + .expect("drift schema catalog"); assert_eq!( validate_myc_state_catalogs(&expected_migrations, &schema) .expect_err("schema drift") diff --git a/tests/services_hardening_state_repository.rs b/tests/services_hardening_state_repository.rs @@ -119,7 +119,7 @@ async fn initialization_migrates_and_binds_exact_metadata_before_inspection() { .fetch_all(&mut connection) .await .expect("migration rows"); - assert_eq!(migrations.len(), 3); + assert_eq!(migrations.len(), 4); assert_eq!(migrations[0].get::<i64, _>(0), 2); assert_eq!( migrations[0].get::<String, _>(1), @@ -135,6 +135,11 @@ async fn initialization_migrates_and_binds_exact_metadata_before_inspection() { migrations[2].get::<String, _>(1), "create_connection_authorization_state" ); + assert_eq!(migrations[3].get::<i64, _>(0), 5); + assert_eq!( + migrations[3].get::<String, _>(1), + "create_bounded_governance_state" + ); let binding = sqlx::query( "SELECT normalized_config_sha256, transport_public_key, user_public_key, \ discovery_public_key, config_contract_version, state_contract_version, \