myc

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

commit 6207ddfc098349429fa48916ac212e08292597d6
parent fdde56bc7f851660efc46d0e995044eb43427af7
Author: triesap <tyson@radroots.org>
Date:   Fri, 21 Aug 2026 18:17:54 +0000

state: persist delivery workflow evidence

Diffstat:
MREADME | 21+++++++++++++++++----
Msrc/lib.rs | 16++++++++++++++--
Msrc/state_catalog.rs | 418+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Asrc/state_delivery.rs | 2010+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/state_host.rs | 11++++++-----
Msrc/state_metadata.rs | 64++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Atests/services_hardening_delivery_state.rs | 539+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_state_catalog.rs | 62++++++++++++++++++++++++++++++++++++++++++--------------------
Mtests/services_hardening_state_repository.rs | 7++++++-
9 files changed, 3107 insertions(+), 41 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-v5 migrations, binds the normalized configuration, expected +through schema-v6 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 any exact v1 through v4 prefix; read-only inspection requires the current +resume any exact v1 through v5 prefix; read-only inspection requires the current catalog and exact immutable Myc binding. The public Myc repository exposes no raw pool, connection, transaction-control @@ -60,8 +60,10 @@ 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. 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. +audit, request-audit binding, and bounded rate-window objects. Schema v6 adds +checksum-pinned publication-job, target, attempt, transition-guard, and +no-delete objects; discovery and exact signed-event storage remain owned by +their later ordered repository steps. Signer-request admission validates bounded client, request, event, method, canonical request, injected operation entropy, and injected time evidence @@ -98,6 +100,17 @@ 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. +Delivery admission copies the exact normalized write-relay inventory, +required-target flags, acknowledgement policy, retry ceiling, backoff, and +attempt deadline into one immutable publication job. Target leases and attempt +records are advanced only through the shared transaction runner; relay I/O is +never awaited while a transaction is open. Exact retries are idempotent, stale +leases recover through injected time evidence, and an acknowledgement whose +outcome is not known remains `unknown` rather than being relabeled as failure +or delivery. This checkpoint binds an artifact digest and precommit identity; +the exact committed signed-event bytes and discovery projection are owned by +their subsequent state checkpoints. + 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_delivery; mod state_governance; mod state_host; mod state_maintenance; @@ -121,8 +122,10 @@ pub use state_catalog::{ 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, 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, + MYC_STATE_SCHEMA_VERSION_5_SHA256, MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256, + MYC_STATE_SCHEMA_VERSION_6_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_6_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, @@ -135,6 +138,15 @@ pub use state_connection::{ MycConnectionPolicyGeneration, MycConnectionRecord, MycConnectionStateError, MycConnectionStateErrorKind, MycConnectionStatus, MycConnectionTimeUnixMs, }; +pub use state_delivery::{ + MYC_DELIVERY_ATTEMPT_MAX_COUNT, MYC_DELIVERY_RELAY_ID_MAX_BYTES, MYC_DELIVERY_TARGET_MAX_COUNT, + MycDeliveryArtifactDigest, MycDeliveryAttemptId, MycDeliveryAttemptNonce, + MycDeliveryAttemptOutcome, MycDeliveryAttemptRecord, MycDeliveryAttemptStatus, + MycDeliveryClaim, MycDeliveryJobAdmission, MycDeliveryJobId, MycDeliveryJobRecord, + MycDeliveryJobRequest, MycDeliveryJobStatus, MycDeliveryPolicyMode, MycDeliveryRelayId, + MycDeliveryStateError, MycDeliveryStateErrorKind, MycDeliveryTargetRecord, + MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, +}; 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, 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 = 5; +pub const MYC_STATE_SCHEMA_VERSION: u32 = 6; /// The shared metadata and migration-ledger objects present at schema v1. pub const MYC_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6; @@ -29,6 +29,9 @@ 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; +/// The shared objects plus Myc metadata, request, connection, governance, and delivery objects. +pub const MYC_STATE_SCHEMA_VERSION_6_OBJECT_COUNT: u32 = 43; + /// 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, @@ -43,14 +46,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] = [ - 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, + 0xa6, 0x9b, 0xbc, 0x7f, 0x3c, 0x22, 0x75, 0x0d, 0xdd, 0xa5, 0x5f, 0x36, 0xdb, 0x12, 0x56, 0x22, + 0xc4, 0x9f, 0x26, 0xb4, 0x4c, 0xdc, 0x47, 0x21, 0x0f, 0xad, 0xb7, 0xb5, 0xfc, 0x4f, 0x31, 0xe5, ]; /// SHA-256 identity of the schema catalog bound to the migration catalog. pub const MYC_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [ - 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, + 0x94, 0xf7, 0xad, 0xfb, 0xd6, 0x2b, 0xe0, 0x3e, 0x6f, 0x3a, 0xf9, 0xc9, 0x75, 0x4b, 0xf2, 0xf4, + 0x7c, 0xc7, 0x04, 0xa8, 0x7b, 0x27, 0x08, 0xf6, 0x82, 0x6c, 0x44, 0x32, 0x55, 0x6b, 0x6f, 0x06, ]; /// SHA-256 identity of the schema-v2 migration content. @@ -95,6 +98,18 @@ pub const MYC_STATE_SCHEMA_VERSION_5_SHA256: [u8; 32] = [ 0x06, 0xcf, 0x0f, 0x4c, 0xfc, 0xa5, 0xfa, 0x5c, 0xff, 0xc9, 0x3a, 0xf1, 0x14, 0x6f, 0x21, 0xc5, ]; +/// SHA-256 identity of the schema-v6 delivery-state migration. +pub const MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256: [u8; 32] = [ + 0x46, 0x49, 0xb9, 0xaf, 0xd0, 0x3f, 0xc0, 0x7f, 0x89, 0xfe, 0x18, 0x4f, 0x02, 0x79, 0x25, 0x67, + 0x5a, 0x02, 0x65, 0x95, 0x1e, 0xa7, 0xcf, 0x75, 0xe0, 0x6a, 0x43, 0x12, 0xc5, 0x5a, 0x82, 0xbd, +]; + +/// SHA-256 identity of the schema-v6 object snapshot. +pub const MYC_STATE_SCHEMA_VERSION_6_SHA256: [u8; 32] = [ + 0x55, 0x72, 0x73, 0x30, 0x6e, 0x68, 0xe9, 0x30, 0x6d, 0xc1, 0xd7, 0xc7, 0x00, 0x9e, 0x1c, 0xb7, + 0xc9, 0xd1, 0x5b, 0x8d, 0x52, 0xe0, 0x7d, 0xb7, 0xd0, 0xbd, 0x51, 0xe8, 0x3f, 0x05, 0x44, 0x92, +]; + /// 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, @@ -794,6 +809,261 @@ const CREATE_GOVERNANCE_STATE_MIGRATION_SQL: &str = concat!( connection_rate_windows_guard_update_sql!(), ); +macro_rules! publication_outbox_table_sql { + () => { + r#"CREATE TABLE publication_outbox ( + job_id BLOB NOT NULL PRIMARY KEY CHECK (length(job_id) = 32), + signer_operation_id BLOB NOT NULL UNIQUE CHECK (length(signer_operation_id) = 32) + REFERENCES nip46_requests(operation_id), + artifact_sha256 BLOB NOT NULL CHECK (length(artifact_sha256) = 32), + policy_mode TEXT NOT NULL CHECK (policy_mode IN ( + 'at_least_one_required', 'all_required', 'required_quorum' + )), + required_acknowledgements INTEGER NOT NULL + CHECK (required_acknowledgements BETWEEN 1 AND 32), + max_attempts INTEGER NOT NULL CHECK (max_attempts BETWEEN 1 AND 32), + initial_backoff_ms INTEGER NOT NULL CHECK (initial_backoff_ms BETWEEN 1 AND 30000), + maximum_backoff_ms INTEGER NOT NULL CHECK (maximum_backoff_ms BETWEEN initial_backoff_ms AND 300000), + attempt_deadline_ms INTEGER NOT NULL CHECK (attempt_deadline_ms BETWEEN 1 AND 30000), + status TEXT NOT NULL CHECK (status IN ('pending', 'active', 'delivered', 'failed', 'unknown')), + created_at_unix_ms INTEGER NOT NULL + CHECK (created_at_unix_ms BETWEEN 1 AND 9223372036854775807), + updated_at_unix_ms INTEGER NOT NULL + CHECK (updated_at_unix_ms BETWEEN created_at_unix_ms AND 9223372036854775807), + finalized_at_unix_ms INTEGER + CHECK (finalized_at_unix_ms IS NULL OR + finalized_at_unix_ms BETWEEN created_at_unix_ms AND 9223372036854775807), + CHECK ((status IN ('pending', 'active') AND finalized_at_unix_ms IS NULL) + OR (status IN ('delivered', 'failed', 'unknown') AND finalized_at_unix_ms IS NOT NULL)) +) STRICT"# + }; +} + +macro_rules! publication_targets_table_sql { + () => { + r#"CREATE TABLE publication_targets ( + job_id BLOB NOT NULL CHECK (length(job_id) = 32) + REFERENCES publication_outbox(job_id), + target_index INTEGER NOT NULL CHECK (target_index BETWEEN 0 AND 31), + relay_id TEXT NOT NULL CHECK (length(CAST(relay_id AS BLOB)) BETWEEN 1 AND 64), + required INTEGER NOT NULL CHECK (required IN (0, 1)), + attempt_count INTEGER NOT NULL CHECK (attempt_count BETWEEN 0 AND 32), + status TEXT NOT NULL CHECK (status IN ( + 'pending', 'leased', 'submitted', 'delivered', 'retryable', 'unknown', 'exhausted' + )), + active_attempt_id BLOB CHECK (active_attempt_id IS NULL OR length(active_attempt_id) = 32), + next_attempt_at_unix_ms INTEGER + CHECK (next_attempt_at_unix_ms IS NULL OR + next_attempt_at_unix_ms BETWEEN 1 AND 9223372036854775807), + updated_at_unix_ms INTEGER NOT NULL + CHECK (updated_at_unix_ms BETWEEN 1 AND 9223372036854775807), + CHECK ((status IN ('leased', 'submitted') AND active_attempt_id IS NOT NULL + AND next_attempt_at_unix_ms IS NULL) + OR (status = 'retryable' AND active_attempt_id IS NULL + AND next_attempt_at_unix_ms IS NOT NULL) + OR (status = 'unknown' AND active_attempt_id IS NULL) + OR (status IN ('pending', 'delivered', 'exhausted') AND active_attempt_id IS NULL + AND next_attempt_at_unix_ms IS NULL)), + PRIMARY KEY (job_id, target_index), + UNIQUE (job_id, relay_id) +) STRICT"# + }; +} + +macro_rules! publication_attempts_table_sql { + () => { + r#"CREATE TABLE publication_attempts ( + attempt_id BLOB NOT NULL PRIMARY KEY CHECK (length(attempt_id) = 32), + job_id BLOB NOT NULL CHECK (length(job_id) = 32), + target_index INTEGER NOT NULL CHECK (target_index BETWEEN 0 AND 31), + attempt_number INTEGER NOT NULL CHECK (attempt_number BETWEEN 1 AND 32), + attempt_nonce BLOB NOT NULL CHECK (length(attempt_nonce) = 32), + status TEXT NOT NULL CHECK (status IN ('leased', 'submitted', 'delivered', 'failed', 'unknown')), + leased_at_unix_ms INTEGER NOT NULL + CHECK (leased_at_unix_ms BETWEEN 1 AND 9223372036854775807), + lease_expires_at_unix_ms INTEGER NOT NULL + CHECK (lease_expires_at_unix_ms BETWEEN leased_at_unix_ms + 1 AND 9223372036854775807), + submitted_at_unix_ms INTEGER + CHECK (submitted_at_unix_ms IS NULL OR + submitted_at_unix_ms BETWEEN leased_at_unix_ms AND lease_expires_at_unix_ms), + resolved_at_unix_ms INTEGER + CHECK (resolved_at_unix_ms IS NULL OR + resolved_at_unix_ms BETWEEN leased_at_unix_ms AND 9223372036854775807), + reason_code TEXT CHECK (reason_code IS NULL OR reason_code IN ( + 'accepted', 'relay_rejected', 'transport_failed', + 'lease_expired_before_submit', 'acknowledgement_lost' + )), + CHECK ((status = 'leased' AND submitted_at_unix_ms IS NULL + AND resolved_at_unix_ms IS NULL AND reason_code IS NULL) + OR (status = 'submitted' AND submitted_at_unix_ms IS NOT NULL + AND resolved_at_unix_ms IS NULL AND reason_code IS NULL) + OR (status = 'delivered' AND submitted_at_unix_ms IS NOT NULL + AND resolved_at_unix_ms IS NOT NULL AND reason_code = 'accepted') + OR (status = 'failed' AND resolved_at_unix_ms IS NOT NULL + AND reason_code IN ('relay_rejected', 'transport_failed', 'lease_expired_before_submit')) + OR (status = 'unknown' AND submitted_at_unix_ms IS NOT NULL + AND resolved_at_unix_ms IS NOT NULL AND reason_code = 'acknowledgement_lost')), + FOREIGN KEY (job_id, target_index) + REFERENCES publication_targets(job_id, target_index), + UNIQUE (job_id, target_index, attempt_number), + UNIQUE (job_id, target_index, attempt_nonce) +) STRICT"# + }; +} + +macro_rules! publication_outbox_guard_update_sql { + () => { + r#"CREATE TRIGGER publication_outbox_guard_update +BEFORE UPDATE ON publication_outbox +WHEN NEW.job_id != OLD.job_id + OR NEW.signer_operation_id != OLD.signer_operation_id + OR NEW.artifact_sha256 != OLD.artifact_sha256 + OR NEW.policy_mode != OLD.policy_mode + OR NEW.required_acknowledgements != OLD.required_acknowledgements + OR NEW.max_attempts != OLD.max_attempts + OR NEW.initial_backoff_ms != OLD.initial_backoff_ms + OR NEW.maximum_backoff_ms != OLD.maximum_backoff_ms + OR NEW.attempt_deadline_ms != OLD.attempt_deadline_ms + OR NEW.created_at_unix_ms != OLD.created_at_unix_ms + OR NEW.updated_at_unix_ms < OLD.updated_at_unix_ms + OR NOT ((OLD.status = 'pending' AND NEW.status = 'active' + AND NEW.finalized_at_unix_ms IS NULL) + OR (OLD.status IN ('pending', 'active') AND NEW.status IN ('delivered', 'failed', 'unknown') + AND NEW.finalized_at_unix_ms IS NOT NULL)) +BEGIN + SELECT RAISE(ABORT, 'publication job transition is invalid'); +END"# + }; +} + +macro_rules! publication_targets_guard_update_sql { + () => { + r#"CREATE TRIGGER publication_targets_guard_update +BEFORE UPDATE ON publication_targets +WHEN NEW.job_id != OLD.job_id + OR NEW.target_index != OLD.target_index + OR NEW.relay_id != OLD.relay_id + OR NEW.required != OLD.required + OR NEW.updated_at_unix_ms < OLD.updated_at_unix_ms + OR NOT ((OLD.status IN ('pending', 'retryable', 'unknown') AND NEW.status = 'leased' + AND NEW.attempt_count = OLD.attempt_count + 1 + AND NEW.active_attempt_id IS NOT NULL AND NEW.next_attempt_at_unix_ms IS NULL) + OR (OLD.status = 'leased' AND NEW.status = 'submitted' + AND NEW.attempt_count = OLD.attempt_count + AND NEW.active_attempt_id = OLD.active_attempt_id + AND NEW.next_attempt_at_unix_ms IS NULL) + OR (OLD.status IN ('leased', 'submitted') + AND NEW.status IN ('delivered', 'retryable', 'unknown', 'exhausted') + AND NEW.attempt_count = OLD.attempt_count + AND NEW.active_attempt_id IS NULL + AND ((NEW.status = 'retryable' + AND NEW.next_attempt_at_unix_ms IS NOT NULL) + OR NEW.status = 'unknown' + OR (NEW.status IN ('delivered', 'exhausted') + AND NEW.next_attempt_at_unix_ms IS NULL)))) +BEGIN + SELECT RAISE(ABORT, 'publication target transition is invalid'); +END"# + }; +} + +macro_rules! publication_attempts_guard_update_sql { + () => { + r#"CREATE TRIGGER publication_attempts_guard_update +BEFORE UPDATE ON publication_attempts +WHEN NEW.attempt_id != OLD.attempt_id + OR NEW.job_id != OLD.job_id + OR NEW.target_index != OLD.target_index + OR NEW.attempt_number != OLD.attempt_number + OR NEW.attempt_nonce != OLD.attempt_nonce + OR NEW.leased_at_unix_ms != OLD.leased_at_unix_ms + OR NEW.lease_expires_at_unix_ms != OLD.lease_expires_at_unix_ms + OR NOT ((OLD.status = 'leased' AND NEW.status = 'submitted' + AND NEW.submitted_at_unix_ms IS NOT NULL + AND NEW.resolved_at_unix_ms IS NULL AND NEW.reason_code IS NULL) + OR (OLD.status = 'leased' AND NEW.status = 'failed' + AND NEW.submitted_at_unix_ms IS NULL + AND NEW.resolved_at_unix_ms IS NOT NULL + AND NEW.reason_code IN ('transport_failed', 'lease_expired_before_submit')) + OR (OLD.status = 'submitted' AND NEW.status IN ('delivered', 'failed', 'unknown') + AND NEW.submitted_at_unix_ms = OLD.submitted_at_unix_ms + AND NEW.resolved_at_unix_ms IS NOT NULL + AND ((NEW.status = 'delivered' AND NEW.reason_code = 'accepted') + OR (NEW.status = 'failed' + AND NEW.reason_code IN ('relay_rejected', 'transport_failed')) + OR (NEW.status = 'unknown' + AND NEW.reason_code = 'acknowledgement_lost')))) +BEGIN + SELECT RAISE(ABORT, 'publication attempt transition is invalid'); +END"# + }; +} + +macro_rules! publication_outbox_no_delete_sql { + () => { + r#"CREATE TRIGGER publication_outbox_no_delete +BEFORE DELETE ON publication_outbox +BEGIN + SELECT RAISE(ABORT, 'publication jobs are immutable'); +END"# + }; +} + +macro_rules! publication_targets_no_delete_sql { + () => { + r#"CREATE TRIGGER publication_targets_no_delete +BEFORE DELETE ON publication_targets +BEGIN + SELECT RAISE(ABORT, 'publication targets are immutable'); +END"# + }; +} + +macro_rules! publication_attempts_no_delete_sql { + () => { + r#"CREATE TRIGGER publication_attempts_no_delete +BEFORE DELETE ON publication_attempts +BEGIN + SELECT RAISE(ABORT, 'publication attempts are immutable'); +END"# + }; +} + +const CREATE_PUBLICATION_OUTBOX_TABLE_SQL: &str = publication_outbox_table_sql!(); +const CREATE_PUBLICATION_TARGETS_TABLE_SQL: &str = publication_targets_table_sql!(); +const CREATE_PUBLICATION_ATTEMPTS_TABLE_SQL: &str = publication_attempts_table_sql!(); +const CREATE_PUBLICATION_OUTBOX_GUARD_UPDATE_SQL: &str = publication_outbox_guard_update_sql!(); +const CREATE_PUBLICATION_TARGETS_GUARD_UPDATE_SQL: &str = publication_targets_guard_update_sql!(); +const CREATE_PUBLICATION_ATTEMPTS_GUARD_UPDATE_SQL: &str = publication_attempts_guard_update_sql!(); +const CREATE_PUBLICATION_OUTBOX_NO_DELETE_SQL: &str = publication_outbox_no_delete_sql!(); +const CREATE_PUBLICATION_TARGETS_NO_DELETE_SQL: &str = publication_targets_no_delete_sql!(); +const CREATE_PUBLICATION_ATTEMPTS_NO_DELETE_SQL: &str = publication_attempts_no_delete_sql!(); + +const CREATE_DELIVERY_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 = 5 THEN 6 ELSE 0 END WHERE singleton = 1;\n", + myc_state_metadata_no_update_sql!(), + ";\n", + publication_outbox_table_sql!(), + ";\n", + publication_targets_table_sql!(), + ";\n", + publication_attempts_table_sql!(), + ";\n", + publication_outbox_guard_update_sql!(), + ";\n", + publication_targets_guard_update_sql!(), + ";\n", + publication_attempts_guard_update_sql!(), + ";\n", + publication_outbox_no_delete_sql!(), + ";\n", + publication_targets_no_delete_sql!(), + ";\n", + publication_attempts_no_delete_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, @@ -878,6 +1148,42 @@ 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, ]; +const PUBLICATION_OUTBOX_TABLE_SHA256: [u8; 32] = [ + 0x26, 0xbd, 0xa8, 0x77, 0x39, 0x07, 0xdb, 0x05, 0xa8, 0x7c, 0x38, 0x21, 0x9b, 0x58, 0xd7, 0xf8, + 0xc8, 0x01, 0x80, 0xb3, 0x50, 0xae, 0x69, 0x87, 0xd2, 0x13, 0xd1, 0x84, 0xc6, 0x1d, 0x15, 0x0c, +]; +const PUBLICATION_TARGETS_TABLE_SHA256: [u8; 32] = [ + 0x16, 0xd5, 0xa3, 0x85, 0x71, 0x0d, 0xb1, 0x02, 0xaa, 0x25, 0x02, 0x80, 0xbb, 0xc4, 0x23, 0x44, + 0x20, 0xc1, 0xa2, 0x8a, 0x33, 0xaf, 0xd0, 0xe9, 0x6b, 0x65, 0x7a, 0x79, 0xd7, 0xe3, 0x1b, 0x43, +]; +const PUBLICATION_ATTEMPTS_TABLE_SHA256: [u8; 32] = [ + 0x25, 0xed, 0xa7, 0xb2, 0xd5, 0x07, 0xc2, 0xe5, 0x29, 0xef, 0x07, 0x31, 0xe6, 0x84, 0x6a, 0x64, + 0xe9, 0xef, 0x11, 0x47, 0x41, 0x0d, 0xc9, 0x79, 0x43, 0x06, 0x8e, 0x78, 0x47, 0xfb, 0x3b, 0x3d, +]; +const PUBLICATION_OUTBOX_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x6f, 0xb5, 0x23, 0x36, 0x36, 0x1d, 0xea, 0x76, 0x83, 0x4c, 0x55, 0xe3, 0xdb, 0x90, 0xf1, 0xd1, + 0x84, 0xcb, 0x11, 0x23, 0x17, 0x1f, 0x09, 0xd7, 0xee, 0x9c, 0xdc, 0x70, 0x28, 0xbe, 0x24, 0xc4, +]; +const PUBLICATION_TARGETS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x3e, 0x74, 0xf3, 0x95, 0x67, 0x81, 0x3e, 0x55, 0xfd, 0xc6, 0x73, 0x50, 0xa4, 0xa4, 0x4a, 0xb1, + 0x77, 0xf0, 0x12, 0x8a, 0xe7, 0xd0, 0x07, 0x3d, 0x23, 0x92, 0x89, 0x1c, 0x97, 0x9d, 0x2f, 0x34, +]; +const PUBLICATION_ATTEMPTS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x5f, 0xda, 0xc3, 0xa9, 0x01, 0x48, 0x93, 0xac, 0xb4, 0x33, 0x3f, 0xac, 0x16, 0xe9, 0x63, 0x22, + 0xcb, 0x74, 0xff, 0xe6, 0x2c, 0xbf, 0x89, 0xfc, 0x59, 0x4f, 0x75, 0xc6, 0x39, 0xa0, 0x40, 0xc9, +]; +const PUBLICATION_OUTBOX_NO_DELETE_SHA256: [u8; 32] = [ + 0xe9, 0x0e, 0x86, 0x86, 0xe8, 0xf0, 0x51, 0xb7, 0x89, 0x97, 0x92, 0x92, 0x8c, 0x6a, 0xd6, 0x64, + 0xb9, 0x79, 0xaa, 0xbf, 0xbe, 0x08, 0x5d, 0xad, 0xc0, 0x69, 0x02, 0xa7, 0x5e, 0x58, 0xa5, 0x28, +]; +const PUBLICATION_TARGETS_NO_DELETE_SHA256: [u8; 32] = [ + 0x19, 0x0b, 0x6b, 0x34, 0xd0, 0x69, 0x45, 0x24, 0xd8, 0x57, 0x17, 0xf7, 0x27, 0x9e, 0xd9, 0x0c, + 0x5a, 0xc0, 0xc2, 0xa6, 0x63, 0xd6, 0x95, 0x2e, 0x77, 0x55, 0x8f, 0x6f, 0x7c, 0x81, 0xb0, 0x48, +]; +const PUBLICATION_ATTEMPTS_NO_DELETE_SHA256: [u8; 32] = [ + 0x72, 0x0a, 0x1f, 0x3a, 0x40, 0xfb, 0xf8, 0x63, 0xeb, 0x1f, 0xd4, 0x54, 0x7f, 0xa0, 0xbc, 0x43, + 0x50, 0x65, 0x2b, 0xb2, 0xf0, 0x98, 0x9f, 0xda, 0x1f, 0xc6, 0x49, 0xcf, 0xd8, 0xb5, 0xd4, 0xe5, +]; /// Stable classes for invalid embedded Myc catalog definitions. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -978,10 +1284,17 @@ pub fn myc_migration_catalog() -> Result<MigrationCatalog, MycStateCatalogError> MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256), ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; - let catalog = MigrationCatalog::new([metadata, requests, connections, governance]) + let delivery = MigrationDescriptor::sql( + 6, + "create_delivery_evidence_state", + CREATE_DELIVERY_STATE_MIGRATION_SQL, + MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256), + ) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; + let catalog = MigrationCatalog::new([metadata, requests, connections, governance, delivery]) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; if catalog.current_version() != MYC_STATE_SCHEMA_VERSION - || catalog.descriptors().len() != 4 + || catalog.descriptors().len() != 5 || catalog.digest().as_bytes() != &MYC_MIGRATION_CATALOG_SHA256 { return Err(MycStateCatalogError::new( @@ -1024,6 +1337,12 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> { SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_5_SHA256), ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; + let version_six = SchemaVersionCatalog::new( + 6, + myc_state_delivery_objects()?, + SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_6_SHA256), + ) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; let catalog = SchemaCatalog::new( &migrations, [ @@ -1032,6 +1351,7 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> { version_three, version_four, version_five, + version_six, ], ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; @@ -1291,6 +1611,80 @@ fn myc_state_governance_objects() -> Result<Vec<SchemaObject>, MycStateCatalogEr Ok(objects) } +fn myc_state_delivery_objects() -> Result<Vec<SchemaObject>, MycStateCatalogError> { + let mut objects = myc_state_governance_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, + "publication_outbox", + "publication_outbox", + CREATE_PUBLICATION_OUTBOX_TABLE_SQL, + PUBLICATION_OUTBOX_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "publication_targets", + "publication_targets", + CREATE_PUBLICATION_TARGETS_TABLE_SQL, + PUBLICATION_TARGETS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "publication_attempts", + "publication_attempts", + CREATE_PUBLICATION_ATTEMPTS_TABLE_SQL, + PUBLICATION_ATTEMPTS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "publication_outbox_guard_update", + "publication_outbox", + CREATE_PUBLICATION_OUTBOX_GUARD_UPDATE_SQL, + PUBLICATION_OUTBOX_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "publication_targets_guard_update", + "publication_targets", + CREATE_PUBLICATION_TARGETS_GUARD_UPDATE_SQL, + PUBLICATION_TARGETS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "publication_attempts_guard_update", + "publication_attempts", + CREATE_PUBLICATION_ATTEMPTS_GUARD_UPDATE_SQL, + PUBLICATION_ATTEMPTS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "publication_outbox_no_delete", + "publication_outbox", + CREATE_PUBLICATION_OUTBOX_NO_DELETE_SQL, + PUBLICATION_OUTBOX_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "publication_targets_no_delete", + "publication_targets", + CREATE_PUBLICATION_TARGETS_NO_DELETE_SQL, + PUBLICATION_TARGETS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "publication_attempts_no_delete", + "publication_attempts", + CREATE_PUBLICATION_ATTEMPTS_NO_DELETE_SQL, + PUBLICATION_ATTEMPTS_NO_DELETE_SHA256, + )?, + ]); + Ok(objects) +} + /// Independently validates exact catalog versions, counts, and digests. pub fn validate_myc_state_catalogs( migrations: &MigrationCatalog, @@ -1299,7 +1693,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() == 4 + && descriptors.len() == 5 && 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 @@ -1312,9 +1706,12 @@ pub fn validate_myc_state_catalogs( && 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 + && descriptors[4].target_version() == 6 + && descriptors[4].name().as_str() == "create_delivery_evidence_state" + && descriptors[4].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256 && migrations.digest().as_bytes() == &MYC_MIGRATION_CATALOG_SHA256 && schema.migration_catalog_digest() == migrations.digest() - && versions.len() == 5 + && versions.len() == 6 && 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 @@ -1330,6 +1727,9 @@ pub fn validate_myc_state_catalogs( && 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 + && versions[5].version() == 6 + && versions[5].object_count() == MYC_STATE_SCHEMA_VERSION_6_OBJECT_COUNT + && versions[5].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_6_SHA256 && schema.digest().as_bytes() == &MYC_STATE_SCHEMA_CATALOG_SHA256; if valid { Ok(()) diff --git a/src/state_delivery.rs b/src/state_delivery.rs @@ -0,0 +1,2010 @@ +//! Durable delivery-job identity, bounded target attempts, and lease evidence. + +use core::fmt; +use std::error::Error; + +use radroots_service_sqlite::{ + ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, +}; +use sha2::{Digest, Sha256}; +use sqlx::Row; + +use crate::MycSignerOperationId; +use crate::state_repository::{ + MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, + RepositoryOperationError, require_expected_metadata, +}; + +/// Maximum immutable relay targets on one delivery job. +pub const MYC_DELIVERY_TARGET_MAX_COUNT: usize = 32; +/// Maximum durable attempts for one target. +pub const MYC_DELIVERY_ATTEMPT_MAX_COUNT: u32 = 32; +/// Maximum UTF-8 byte length of a canonical delivery relay identifier. +pub const MYC_DELIVERY_RELAY_ID_MAX_BYTES: usize = 64; + +const JOB_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_job.v1\0"; +const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_attempt.v1\0"; + +const READ_SIGNER_OPERATION_SQL: &str = r#"SELECT COUNT(*) AS row_count +FROM nip46_requests +WHERE operation_id = ?"#; + +const READ_JOB_SQL: &str = r#"SELECT + CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 + THEN job_id ELSE NULL END AS job_id, + CASE WHEN typeof(signer_operation_id) = 'blob' AND length(signer_operation_id) = 32 + THEN signer_operation_id ELSE NULL END AS signer_operation_id, + CASE WHEN typeof(artifact_sha256) = 'blob' AND length(artifact_sha256) = 32 + THEN artifact_sha256 ELSE NULL END AS artifact_sha256, + CASE WHEN typeof(policy_mode) = 'text' AND length(CAST(policy_mode AS BLOB)) <= 32 + THEN policy_mode ELSE NULL END AS policy_mode, + required_acknowledgements, max_attempts, initial_backoff_ms, + maximum_backoff_ms, attempt_deadline_ms, + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + created_at_unix_ms, updated_at_unix_ms, finalized_at_unix_ms, + typeof(finalized_at_unix_ms) AS finalized_at_type +FROM publication_outbox +WHERE job_id = ? +LIMIT 2"#; + +const READ_JOB_BY_OPERATION_SQL: &str = r#"SELECT + CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 + THEN job_id ELSE NULL END AS job_id, + CASE WHEN typeof(signer_operation_id) = 'blob' AND length(signer_operation_id) = 32 + THEN signer_operation_id ELSE NULL END AS signer_operation_id, + CASE WHEN typeof(artifact_sha256) = 'blob' AND length(artifact_sha256) = 32 + THEN artifact_sha256 ELSE NULL END AS artifact_sha256, + CASE WHEN typeof(policy_mode) = 'text' AND length(CAST(policy_mode AS BLOB)) <= 32 + THEN policy_mode ELSE NULL END AS policy_mode, + required_acknowledgements, max_attempts, initial_backoff_ms, + maximum_backoff_ms, attempt_deadline_ms, + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + created_at_unix_ms, updated_at_unix_ms, finalized_at_unix_ms, + typeof(finalized_at_unix_ms) AS finalized_at_type +FROM publication_outbox +WHERE signer_operation_id = ? +LIMIT 2"#; + +const READ_TARGETS_SQL: &str = r#"SELECT + target_index, + CASE WHEN typeof(relay_id) = 'text' + AND length(CAST(relay_id AS BLOB)) BETWEEN 1 AND 64 + THEN relay_id ELSE NULL END AS relay_id, + required, attempt_count, + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + CASE WHEN typeof(active_attempt_id) = 'blob' AND length(active_attempt_id) = 32 + THEN active_attempt_id ELSE NULL END AS active_attempt_id, + typeof(active_attempt_id) AS active_attempt_id_type, + next_attempt_at_unix_ms, + typeof(next_attempt_at_unix_ms) AS next_attempt_at_type, + updated_at_unix_ms +FROM publication_targets +WHERE job_id = ? +ORDER BY target_index +LIMIT 33"#; + +const READ_ATTEMPTS_SQL: &str = r#"SELECT + CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 + THEN attempt_id ELSE NULL END AS attempt_id, + attempt_number, + CASE WHEN typeof(attempt_nonce) = 'blob' AND length(attempt_nonce) = 32 + THEN attempt_nonce ELSE NULL END AS attempt_nonce, + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + leased_at_unix_ms, lease_expires_at_unix_ms, + submitted_at_unix_ms, typeof(submitted_at_unix_ms) AS submitted_at_type, + resolved_at_unix_ms, typeof(resolved_at_unix_ms) AS resolved_at_type, + CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 + THEN reason_code ELSE NULL END AS reason_code, + typeof(reason_code) AS reason_code_type +FROM publication_attempts +WHERE job_id = ? AND target_index = ? +ORDER BY attempt_number +LIMIT 33"#; + +const READ_ATTEMPT_SQL: &str = r#"SELECT + CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 + THEN attempt_id ELSE NULL END AS attempt_id, + attempt_number, + CASE WHEN typeof(attempt_nonce) = 'blob' AND length(attempt_nonce) = 32 + THEN attempt_nonce ELSE NULL END AS attempt_nonce, + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + leased_at_unix_ms, lease_expires_at_unix_ms, + submitted_at_unix_ms, typeof(submitted_at_unix_ms) AS submitted_at_type, + resolved_at_unix_ms, typeof(resolved_at_unix_ms) AS resolved_at_type, + CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 + THEN reason_code ELSE NULL END AS reason_code, + typeof(reason_code) AS reason_code_type +FROM publication_attempts +WHERE job_id = ? AND target_index = ? AND attempt_id = ? +LIMIT 2"#; + +const READ_ATTEMPT_BY_NONCE_SQL: &str = r#"SELECT + CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 + THEN attempt_id ELSE NULL END AS attempt_id, + attempt_number, + CASE WHEN typeof(attempt_nonce) = 'blob' AND length(attempt_nonce) = 32 + THEN attempt_nonce ELSE NULL END AS attempt_nonce, + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + leased_at_unix_ms, lease_expires_at_unix_ms, + submitted_at_unix_ms, typeof(submitted_at_unix_ms) AS submitted_at_type, + resolved_at_unix_ms, typeof(resolved_at_unix_ms) AS resolved_at_type, + CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 + THEN reason_code ELSE NULL END AS reason_code, + typeof(reason_code) AS reason_code_type +FROM publication_attempts +WHERE job_id = ? AND target_index = ? AND attempt_nonce = ? +LIMIT 2"#; + +const INSERT_JOB_SQL: &str = r#"INSERT INTO publication_outbox ( + job_id, signer_operation_id, artifact_sha256, policy_mode, + required_acknowledgements, max_attempts, initial_backoff_ms, + maximum_backoff_ms, attempt_deadline_ms, status, + created_at_unix_ms, updated_at_unix_ms, finalized_at_unix_ms +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?, NULL)"#; + +const INSERT_TARGET_SQL: &str = r#"INSERT INTO publication_targets ( + job_id, target_index, relay_id, required, attempt_count, status, + active_attempt_id, next_attempt_at_unix_ms, updated_at_unix_ms +) VALUES (?, ?, ?, ?, 0, 'pending', NULL, NULL, ?)"#; + +const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO publication_attempts ( + attempt_id, job_id, target_index, attempt_number, attempt_nonce, status, + leased_at_unix_ms, lease_expires_at_unix_ms, submitted_at_unix_ms, + resolved_at_unix_ms, reason_code +) VALUES (?, ?, ?, ?, ?, 'leased', ?, ?, NULL, NULL, NULL)"#; + +const CLAIM_TARGET_SQL: &str = r#"UPDATE publication_targets +SET attempt_count = ?, status = 'leased', active_attempt_id = ?, + next_attempt_at_unix_ms = NULL, updated_at_unix_ms = ? +WHERE job_id = ? AND target_index = ? AND attempt_count = ? + AND status IN ('pending', 'retryable', 'unknown') + AND active_attempt_id IS NULL"#; + +const MARK_JOB_ACTIVE_SQL: &str = r#"UPDATE publication_outbox +SET status = 'active', updated_at_unix_ms = ? +WHERE job_id = ? AND status = 'pending'"#; + +const MARK_ATTEMPT_SUBMITTED_SQL: &str = r#"UPDATE publication_attempts +SET status = 'submitted', submitted_at_unix_ms = ? +WHERE attempt_id = ? AND job_id = ? AND target_index = ? + AND status = 'leased' AND lease_expires_at_unix_ms >= ?"#; + +const MARK_TARGET_SUBMITTED_SQL: &str = r#"UPDATE publication_targets +SET status = 'submitted', updated_at_unix_ms = ? +WHERE job_id = ? AND target_index = ? AND active_attempt_id = ? AND status = 'leased'"#; + +const RESOLVE_ATTEMPT_SQL: &str = r#"UPDATE publication_attempts +SET status = ?, resolved_at_unix_ms = ?, reason_code = ? +WHERE attempt_id = ? AND job_id = ? AND target_index = ? AND status = ?"#; + +const RESOLVE_TARGET_SQL: &str = r#"UPDATE publication_targets +SET status = ?, active_attempt_id = NULL, next_attempt_at_unix_ms = ?, + updated_at_unix_ms = ? +WHERE job_id = ? AND target_index = ? AND active_attempt_id = ? AND status = ?"#; + +const FINALIZE_JOB_SQL: &str = r#"UPDATE publication_outbox +SET status = ?, updated_at_unix_ms = ?, finalized_at_unix_ms = ? +WHERE job_id = ? AND status IN ('pending', 'active')"#; + +/// Stable source-free construction failure classes. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliveryStateErrorKind { + InvalidTime, + InvalidRelayId, + InvalidPolicy, +} + +impl MycDeliveryStateErrorKind { + /// Returns the stable machine-readable classification. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::InvalidTime => "delivery_time_invalid", + Self::InvalidRelayId => "delivery_relay_id_invalid", + Self::InvalidPolicy => "delivery_policy_invalid", + } + } +} + +/// Source-free delivery-state construction failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct MycDeliveryStateError { + kind: MycDeliveryStateErrorKind, +} + +impl MycDeliveryStateError { + const fn new(kind: MycDeliveryStateErrorKind) -> Self { + Self { kind } + } + + /// Returns the stable failure class. + #[must_use] + pub const fn kind(self) -> MycDeliveryStateErrorKind { + self.kind + } + + /// Returns the stable machine-readable classification. + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Display for MycDeliveryStateError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + MycDeliveryStateErrorKind::InvalidTime => "delivery time is invalid", + MycDeliveryStateErrorKind::InvalidRelayId => "delivery relay identity is invalid", + MycDeliveryStateErrorKind::InvalidPolicy => "delivery policy is invalid", + }) + } +} + +impl fmt::Debug for MycDeliveryStateError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliveryStateError") + .field("kind", &self.kind) + .finish() + } +} + +impl Error for MycDeliveryStateError {} + +macro_rules! redacted_id { + ($name:ident, $debug:literal) => { + #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] + pub struct $name([u8; 32]); + + impl $name { + /// Returns the exact identity bytes. + #[must_use] + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } + } + + impl fmt::Debug for $name { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str($debug) + } + } + }; +} + +redacted_id!(MycDeliveryJobId, "MycDeliveryJobId([redacted])"); +redacted_id!(MycDeliveryAttemptId, "MycDeliveryAttemptId([redacted])"); +redacted_id!( + MycDeliveryArtifactDigest, + "MycDeliveryArtifactDigest([redacted])" +); + +impl MycDeliveryArtifactDigest { + /// Wraps an independently verified exact-artifact SHA-256 identity. + #[must_use] + pub const fn from_bytes(bytes: [u8; 32]) -> Self { + Self(bytes) + } +} + +/// One-use injected entropy for a new delivery attempt lease. +pub struct MycDeliveryAttemptNonce([u8; 32]); + +impl MycDeliveryAttemptNonce { + /// Wraps exact entropy supplied by the caller's injected boundary. + #[must_use] + pub const fn from_injected_entropy(bytes: [u8; 32]) -> Self { + Self(bytes) + } +} + +impl fmt::Debug for MycDeliveryAttemptNonce { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycDeliveryAttemptNonce([redacted])") + } +} + +/// Positive UTC millisecond evidence representable by SQLite. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct MycDeliveryTimeUnixMs(u64); + +impl MycDeliveryTimeUnixMs { + /// Validates one positive UTC millisecond instant. + pub fn new(value: u64) -> Result<Self, MycDeliveryStateError> { + if value == 0 || i64::try_from(value).is_err() { + return Err(MycDeliveryStateError::new( + MycDeliveryStateErrorKind::InvalidTime, + )); + } + Ok(Self(value)) + } + + /// Returns the validated instant. + #[must_use] + pub const fn get(self) -> u64 { + self.0 + } + + fn sqlite_value(self) -> i64 { + i64::try_from(self.0).expect("validated delivery time fits SQLite") + } +} + +/// Canonical configured relay identifier used as a durable target identity. +#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct MycDeliveryRelayId(Box<str>); + +impl MycDeliveryRelayId { + /// Validates the frozen lower-snake relay grammar before allocation. + pub fn new(value: &str) -> Result<Self, MycDeliveryStateError> { + let bytes = value.as_bytes(); + let valid = !bytes.is_empty() + && bytes.len() <= MYC_DELIVERY_RELAY_ID_MAX_BYTES + && bytes[0].is_ascii_lowercase() + && bytes[bytes.len() - 1].is_ascii_alphanumeric() + && bytes + .iter() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_') + && !bytes.windows(2).any(|window| window == b"__"); + if !valid { + return Err(MycDeliveryStateError::new( + MycDeliveryStateErrorKind::InvalidRelayId, + )); + } + Ok(Self(value.into())) + } + + /// Returns the canonical identifier. + #[must_use] + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl fmt::Debug for MycDeliveryRelayId { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycDeliveryRelayId([redacted])") + } +} + +/// Closed delivery-policy mode copied from normalized configuration. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliveryPolicyMode { + AtLeastOneRequired, + AllRequired, + RequiredQuorum, +} + +impl MycDeliveryPolicyMode { + /// Returns the exact durable spelling. + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::AtLeastOneRequired => "at_least_one_required", + Self::AllRequired => "all_required", + Self::RequiredQuorum => "required_quorum", + } + } + + pub(crate) fn parse(value: &str) -> Option<Self> { + match value { + "at_least_one_required" => Some(Self::AtLeastOneRequired), + "all_required" => Some(Self::AllRequired), + "required_quorum" => Some(Self::RequiredQuorum), + _ => None, + } + } +} + +/// Immutable configured target identity. +#[derive(Clone, PartialEq, Eq)] +pub(crate) struct MycDeliveryTargetPolicy { + relay_id: MycDeliveryRelayId, + required: bool, +} + +/// Immutable configured delivery and retry authority. +#[derive(Clone, PartialEq, Eq)] +pub(crate) struct MycDeliveryPolicies { + mode: MycDeliveryPolicyMode, + required_acknowledgements: u32, + max_attempts: u32, + initial_backoff_ms: u64, + maximum_backoff_ms: u64, + attempt_deadline_ms: u64, + targets: Box<[MycDeliveryTargetPolicy]>, +} + +impl MycDeliveryPolicies { + #[allow(clippy::too_many_arguments)] + pub(crate) fn new( + mode: MycDeliveryPolicyMode, + configured_quorum: Option<u32>, + max_attempts: u32, + initial_backoff_ms: u64, + maximum_backoff_ms: u64, + attempt_deadline_ms: u64, + mut targets: Vec<(MycDeliveryRelayId, bool)>, + ) -> Result<Self, MycDeliveryStateError> { + targets.sort_by(|left, right| left.0.cmp(&right.0)); + let required_count = targets.iter().filter(|(_, required)| *required).count(); + let required_acknowledgements = match mode { + MycDeliveryPolicyMode::AtLeastOneRequired => 1, + MycDeliveryPolicyMode::AllRequired => u32::try_from(required_count).unwrap_or(u32::MAX), + MycDeliveryPolicyMode::RequiredQuorum => configured_quorum.unwrap_or(0), + }; + let valid = !targets.is_empty() + && targets.len() <= MYC_DELIVERY_TARGET_MAX_COUNT + && !targets.windows(2).any(|window| window[0].0 == window[1].0) + && required_acknowledgements != 0 + && usize::try_from(required_acknowledgements) + .is_ok_and(|required| required <= targets.len()) + && usize::try_from(required_acknowledgements) + .is_ok_and(|required| required <= required_count) + && (1..=MYC_DELIVERY_ATTEMPT_MAX_COUNT).contains(&max_attempts) + && initial_backoff_ms != 0 + && initial_backoff_ms <= maximum_backoff_ms + && maximum_backoff_ms <= 300_000 + && attempt_deadline_ms != 0 + && attempt_deadline_ms <= 30_000; + if !valid { + return Err(MycDeliveryStateError::new( + MycDeliveryStateErrorKind::InvalidPolicy, + )); + } + Ok(Self { + mode, + required_acknowledgements, + max_attempts, + initial_backoff_ms, + maximum_backoff_ms, + attempt_deadline_ms, + targets: targets + .into_iter() + .map(|(relay_id, required)| MycDeliveryTargetPolicy { relay_id, required }) + .collect(), + }) + } +} + +impl fmt::Debug for MycDeliveryPolicies { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliveryPolicies") + .field("mode", &self.mode) + .field("target_count", &self.targets.len()) + .finish_non_exhaustive() + } +} + +/// Immutable signer-response delivery job input. +pub struct MycDeliveryJobRequest { + operation_id: MycSignerOperationId, + artifact_digest: MycDeliveryArtifactDigest, + created_at: MycDeliveryTimeUnixMs, +} + +impl MycDeliveryJobRequest { + /// Binds one admitted signer operation to one already-verified artifact digest. + #[must_use] + pub const fn signer_response( + operation_id: MycSignerOperationId, + artifact_digest: MycDeliveryArtifactDigest, + created_at: MycDeliveryTimeUnixMs, + ) -> Self { + Self { + operation_id, + artifact_digest, + created_at, + } + } + + fn owned(&self) -> Self { + Self { + operation_id: self.operation_id, + artifact_digest: self.artifact_digest, + created_at: self.created_at, + } + } +} + +impl fmt::Debug for MycDeliveryJobRequest { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycDeliveryJobRequest([redacted])") + } +} + +/// Durable job state. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliveryJobStatus { + Pending, + Active, + Delivered, + Failed, + Unknown, +} + +impl MycDeliveryJobStatus { + /// Returns the exact durable spelling. + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::Active => "active", + Self::Delivered => "delivered", + Self::Failed => "failed", + Self::Unknown => "unknown", + } + } + + fn parse(value: &str) -> Option<Self> { + match value { + "pending" => Some(Self::Pending), + "active" => Some(Self::Active), + "delivered" => Some(Self::Delivered), + "failed" => Some(Self::Failed), + "unknown" => Some(Self::Unknown), + _ => None, + } + } +} + +/// Durable per-target state. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliveryTargetStatus { + Pending, + Leased, + Submitted, + Delivered, + Retryable, + Unknown, + Exhausted, +} + +impl MycDeliveryTargetStatus { + /// Returns the exact durable spelling. + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::Leased => "leased", + Self::Submitted => "submitted", + Self::Delivered => "delivered", + Self::Retryable => "retryable", + Self::Unknown => "unknown", + Self::Exhausted => "exhausted", + } + } + + fn parse(value: &str) -> Option<Self> { + match value { + "pending" => Some(Self::Pending), + "leased" => Some(Self::Leased), + "submitted" => Some(Self::Submitted), + "delivered" => Some(Self::Delivered), + "retryable" => Some(Self::Retryable), + "unknown" => Some(Self::Unknown), + "exhausted" => Some(Self::Exhausted), + _ => None, + } + } +} + +/// Durable attempt state; `Unknown` is distinct from proof of failure. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliveryAttemptStatus { + Leased, + Submitted, + Delivered, + Failed, + Unknown, +} + +impl MycDeliveryAttemptStatus { + /// Returns the exact durable spelling. + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::Leased => "leased", + Self::Submitted => "submitted", + Self::Delivered => "delivered", + Self::Failed => "failed", + Self::Unknown => "unknown", + } + } + + fn parse(value: &str) -> Option<Self> { + match value { + "leased" => Some(Self::Leased), + "submitted" => Some(Self::Submitted), + "delivered" => Some(Self::Delivered), + "failed" => Some(Self::Failed), + "unknown" => Some(Self::Unknown), + _ => None, + } + } +} + +/// Closed evidence for a terminal attempt observation. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliveryAttemptOutcome { + Delivered, + RelayRejected, + TransportFailed, + UnknownAcknowledgement, +} + +impl MycDeliveryAttemptOutcome { + const fn status(self) -> MycDeliveryAttemptStatus { + match self { + Self::Delivered => MycDeliveryAttemptStatus::Delivered, + Self::RelayRejected | Self::TransportFailed => MycDeliveryAttemptStatus::Failed, + Self::UnknownAcknowledgement => MycDeliveryAttemptStatus::Unknown, + } + } + + const fn reason(self) -> &'static str { + match self { + Self::Delivered => "accepted", + Self::RelayRejected => "relay_rejected", + Self::TransportFailed => "transport_failed", + Self::UnknownAcknowledgement => "acknowledgement_lost", + } + } +} + +/// Immutable summary of one retained delivery job. +#[derive(Clone, PartialEq, Eq)] +pub struct MycDeliveryJobRecord { + id: MycDeliveryJobId, + operation_id: MycSignerOperationId, + artifact_digest: MycDeliveryArtifactDigest, + policy_mode: MycDeliveryPolicyMode, + required_acknowledgements: u32, + max_attempts: u32, + initial_backoff_ms: u64, + maximum_backoff_ms: u64, + attempt_deadline_ms: u64, + status: MycDeliveryJobStatus, + created_at: MycDeliveryTimeUnixMs, + updated_at: MycDeliveryTimeUnixMs, + finalized_at: Option<MycDeliveryTimeUnixMs>, + targets: Box<[MycDeliveryTargetRecord]>, +} + +impl MycDeliveryJobRecord { + #[must_use] + pub const fn id(&self) -> MycDeliveryJobId { + self.id + } + #[must_use] + pub const fn operation_id(&self) -> MycSignerOperationId { + self.operation_id + } + #[must_use] + pub const fn artifact_digest(&self) -> MycDeliveryArtifactDigest { + self.artifact_digest + } + #[must_use] + pub const fn policy_mode(&self) -> MycDeliveryPolicyMode { + self.policy_mode + } + #[must_use] + pub const fn required_acknowledgements(&self) -> u32 { + self.required_acknowledgements + } + #[must_use] + pub const fn max_attempts(&self) -> u32 { + self.max_attempts + } + #[must_use] + pub const fn initial_backoff_ms(&self) -> u64 { + self.initial_backoff_ms + } + #[must_use] + pub const fn maximum_backoff_ms(&self) -> u64 { + self.maximum_backoff_ms + } + #[must_use] + pub const fn attempt_deadline_ms(&self) -> u64 { + self.attempt_deadline_ms + } + #[must_use] + pub const fn status(&self) -> MycDeliveryJobStatus { + self.status + } + #[must_use] + pub const fn targets(&self) -> &[MycDeliveryTargetRecord] { + &self.targets + } + #[must_use] + pub const fn created_at(&self) -> MycDeliveryTimeUnixMs { + self.created_at + } + #[must_use] + pub const fn updated_at(&self) -> MycDeliveryTimeUnixMs { + self.updated_at + } + #[must_use] + pub const fn finalized_at(&self) -> Option<MycDeliveryTimeUnixMs> { + self.finalized_at + } +} + +impl fmt::Debug for MycDeliveryJobRecord { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliveryJobRecord") + .field("status", &self.status) + .field("target_count", &self.targets.len()) + .finish() + } +} + +/// Immutable summary of one configured target and its current state. +#[derive(Clone, PartialEq, Eq)] +pub struct MycDeliveryTargetRecord { + index: u32, + relay_id: MycDeliveryRelayId, + required: bool, + attempt_count: u32, + status: MycDeliveryTargetStatus, + active_attempt_id: Option<MycDeliveryAttemptId>, + next_attempt_at: Option<MycDeliveryTimeUnixMs>, + updated_at: MycDeliveryTimeUnixMs, +} + +impl MycDeliveryTargetRecord { + #[must_use] + pub const fn index(&self) -> u32 { + self.index + } + #[must_use] + pub const fn relay_id(&self) -> &MycDeliveryRelayId { + &self.relay_id + } + #[must_use] + pub const fn required(&self) -> bool { + self.required + } + #[must_use] + pub const fn attempt_count(&self) -> u32 { + self.attempt_count + } + #[must_use] + pub const fn status(&self) -> MycDeliveryTargetStatus { + self.status + } + #[must_use] + pub const fn active_attempt_id(&self) -> Option<MycDeliveryAttemptId> { + self.active_attempt_id + } + #[must_use] + pub const fn next_attempt_at(&self) -> Option<MycDeliveryTimeUnixMs> { + self.next_attempt_at + } + #[must_use] + pub const fn updated_at(&self) -> MycDeliveryTimeUnixMs { + self.updated_at + } +} + +impl fmt::Debug for MycDeliveryTargetRecord { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliveryTargetRecord") + .field("required", &self.required) + .field("attempt_count", &self.attempt_count) + .field("status", &self.status) + .finish() + } +} + +/// Immutable identity plus append-only state for one bounded delivery attempt. +#[derive(Clone, PartialEq, Eq)] +pub struct MycDeliveryAttemptRecord { + id: MycDeliveryAttemptId, + number: u32, + status: MycDeliveryAttemptStatus, + leased_at: MycDeliveryTimeUnixMs, + lease_expires_at: MycDeliveryTimeUnixMs, + submitted_at: Option<MycDeliveryTimeUnixMs>, + resolved_at: Option<MycDeliveryTimeUnixMs>, + reason: Option<&'static str>, +} + +impl MycDeliveryAttemptRecord { + #[must_use] + pub const fn id(&self) -> MycDeliveryAttemptId { + self.id + } + #[must_use] + pub const fn number(&self) -> u32 { + self.number + } + #[must_use] + pub const fn status(&self) -> MycDeliveryAttemptStatus { + self.status + } + #[must_use] + pub const fn leased_at(&self) -> MycDeliveryTimeUnixMs { + self.leased_at + } + #[must_use] + pub const fn lease_expires_at(&self) -> MycDeliveryTimeUnixMs { + self.lease_expires_at + } + #[must_use] + pub const fn submitted_at(&self) -> Option<MycDeliveryTimeUnixMs> { + self.submitted_at + } + #[must_use] + pub const fn resolved_at(&self) -> Option<MycDeliveryTimeUnixMs> { + self.resolved_at + } + #[must_use] + pub const fn reason(&self) -> Option<&'static str> { + self.reason + } +} + +impl fmt::Debug for MycDeliveryAttemptRecord { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliveryAttemptRecord") + .field("number", &self.number) + .field("status", &self.status) + .finish() + } +} + +/// New or idempotently replayed delivery-job creation. +#[derive(Clone, PartialEq, Eq)] +pub enum MycDeliveryJobAdmission { + Created(MycDeliveryJobRecord), + ExactReplay(MycDeliveryJobRecord), +} + +impl MycDeliveryJobAdmission { + #[must_use] + pub const fn record(&self) -> &MycDeliveryJobRecord { + match self { + Self::Created(record) | Self::ExactReplay(record) => record, + } + } +} + +impl fmt::Debug for MycDeliveryJobAdmission { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::Created(_) => "MycDeliveryJobAdmission::Created([redacted])", + Self::ExactReplay(_) => "MycDeliveryJobAdmission::ExactReplay([redacted])", + }) + } +} + +/// Result of attempting to claim one target. +#[derive(Clone, PartialEq, Eq)] +pub enum MycDeliveryClaim { + Claimed(MycDeliveryAttemptRecord), + ExactReplay(MycDeliveryAttemptRecord), + NotReady, + Terminal, +} + +impl fmt::Debug for MycDeliveryClaim { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::Claimed(_) => "MycDeliveryClaim::Claimed([redacted])", + Self::ExactReplay(_) => "MycDeliveryClaim::ExactReplay([redacted])", + Self::NotReady => "MycDeliveryClaim::NotReady", + Self::Terminal => "MycDeliveryClaim::Terminal", + }) + } +} + +impl MycStateRepository<'_> { + /// Atomically creates or exactly replays one config-bound immutable job and target set. + pub async fn create_delivery_job( + &self, + request: &MycDeliveryJobRequest, + ) -> Result<MycDeliveryJobAdmission, MycStateRepositoryError> { + let request = request.owned(); + let policy = self.expected().delivery_policies().clone(); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + create_job(transaction, &request, &policy).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Claims one eligible target under a bounded expiring attempt lease. + pub async fn claim_delivery_target( + &self, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + nonce: MycDeliveryAttemptNonce, + claimed_at: MycDeliveryTimeUnixMs, + ) -> Result<MycDeliveryClaim, MycStateRepositoryError> { + let relay_id = relay_id.clone(); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + claim_target(transaction, job_id, &relay_id, &nonce, claimed_at).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Records that the exact leased attempt reached the external submission boundary. + pub async fn mark_delivery_attempt_submitted( + &self, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: MycDeliveryAttemptId, + submitted_at: MycDeliveryTimeUnixMs, + ) -> Result<MycDeliveryAttemptRecord, MycStateRepositoryError> { + let relay_id = relay_id.clone(); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + mark_submitted(transaction, job_id, &relay_id, attempt_id, submitted_at).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Persists delivered, proven-failed, or unknown acknowledgement evidence. + pub async fn record_delivery_attempt_outcome( + &self, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: MycDeliveryAttemptId, + outcome: MycDeliveryAttemptOutcome, + observed_at: MycDeliveryTimeUnixMs, + ) -> Result<MycDeliveryJobRecord, MycStateRepositoryError> { + let relay_id = relay_id.clone(); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + record_outcome( + transaction, + job_id, + &relay_id, + attempt_id, + outcome, + observed_at, + ) + .await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Converts an expired pre-submit lease to failure or a submitted lease to unknown. + pub async fn recover_expired_delivery_lease( + &self, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: MycDeliveryAttemptId, + observed_at: MycDeliveryTimeUnixMs, + ) -> Result<MycDeliveryJobRecord, MycStateRepositoryError> { + let relay_id = relay_id.clone(); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + recover_expired(transaction, job_id, &relay_id, attempt_id, observed_at).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Reads one bounded job snapshot and its immutable target set. + pub async fn read_delivery_job( + &self, + job_id: MycDeliveryJobId, + ) -> Result<Option<MycDeliveryJobRecord>, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + read_job(transaction, job_id).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Reads at most the configured 32 attempts for one retained target. + pub async fn read_delivery_attempts( + &self, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + ) -> Result<Box<[MycDeliveryAttemptRecord]>, MycStateRepositoryError> { + let relay_id = relay_id.clone(); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let job = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + let target = target_by_relay(&job, &relay_id)?; + read_attempts(transaction, job_id, target.index).await + }) + }) + .await + .map_err(map_transaction_error) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum DeliveryOperationError { + Binding, + Storage, +} + +async fn verify_metadata( + transaction: &mut ServiceSqliteTransaction<'_>, + expected: &PersistedMetadata, +) -> Result<(), DeliveryOperationError> { + require_expected_metadata(transaction, expected) + .await + .map_err(|error| match error { + RepositoryOperationError::Binding => DeliveryOperationError::Binding, + RepositoryOperationError::Storage => DeliveryOperationError::Storage, + }) +} + +async fn create_job( + transaction: &mut ServiceSqliteTransaction<'_>, + request: &MycDeliveryJobRequest, + policy: &MycDeliveryPolicies, +) -> Result<MycDeliveryJobAdmission, DeliveryOperationError> { + require_signer_operation(transaction, request.operation_id).await?; + if let Some(existing) = read_job_by_operation(transaction, request.operation_id).await? { + return exact_job(&existing, request, policy) + .then_some(MycDeliveryJobAdmission::ExactReplay(existing)) + .ok_or(DeliveryOperationError::Binding); + } + let job_id = derive_job_id(request.operation_id, request.artifact_digest); + let result = sqlx::query(INSERT_JOB_SQL) + .bind(job_id.as_bytes().as_slice()) + .bind(request.operation_id.as_bytes().as_slice()) + .bind(request.artifact_digest.as_bytes().as_slice()) + .bind(policy.mode.as_str()) + .bind(i64::from(policy.required_acknowledgements)) + .bind(i64::from(policy.max_attempts)) + .bind( + i64::try_from(policy.initial_backoff_ms) + .map_err(|_| DeliveryOperationError::Binding)?, + ) + .bind( + i64::try_from(policy.maximum_backoff_ms) + .map_err(|_| DeliveryOperationError::Binding)?, + ) + .bind( + i64::try_from(policy.attempt_deadline_ms) + .map_err(|_| DeliveryOperationError::Binding)?, + ) + .bind(request.created_at.sqlite_value()) + .bind(request.created_at.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + for (index, target) in policy.targets.iter().enumerate() { + let index = u32::try_from(index).map_err(|_| DeliveryOperationError::Binding)?; + let result = sqlx::query(INSERT_TARGET_SQL) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(index)) + .bind(target.relay_id.as_str()) + .bind(target.required) + .bind(request.created_at.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + } + let record = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + Ok(MycDeliveryJobAdmission::Created(record)) +} + +async fn claim_target( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + nonce: &MycDeliveryAttemptNonce, + claimed_at: MycDeliveryTimeUnixMs, +) -> Result<MycDeliveryClaim, DeliveryOperationError> { + let job = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + if matches!( + job.status, + MycDeliveryJobStatus::Delivered + | MycDeliveryJobStatus::Failed + | MycDeliveryJobStatus::Unknown + ) { + return Ok(MycDeliveryClaim::Terminal); + } + let target = target_by_relay(&job, relay_id)?.clone(); + if let Some(existing) = read_attempt_by_nonce(transaction, job_id, target.index, nonce).await? { + return Ok(MycDeliveryClaim::ExactReplay(existing)); + } + if target.active_attempt_id.is_some() + || target.next_attempt_at.is_some_and(|next| next > claimed_at) + { + return Ok(MycDeliveryClaim::NotReady); + } + if matches!( + target.status, + MycDeliveryTargetStatus::Delivered | MycDeliveryTargetStatus::Exhausted + ) || target.attempt_count >= job.max_attempts + { + return Ok(MycDeliveryClaim::Terminal); + } + let attempt_number = target.attempt_count + 1; + let lease_expires_value = claimed_at + .get() + .checked_add(job.attempt_deadline_ms) + .ok_or(DeliveryOperationError::Binding)?; + let lease_expires = MycDeliveryTimeUnixMs::new(lease_expires_value) + .map_err(|_| DeliveryOperationError::Binding)?; + let attempt_id = derive_attempt_id(job_id, target.index, attempt_number, nonce); + let result = sqlx::query(INSERT_ATTEMPT_SQL) + .bind(attempt_id.as_bytes().as_slice()) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target.index)) + .bind(i64::from(attempt_number)) + .bind(nonce.0.as_slice()) + .bind(claimed_at.sqlite_value()) + .bind(lease_expires.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + let result = sqlx::query(CLAIM_TARGET_SQL) + .bind(i64::from(attempt_number)) + .bind(attempt_id.as_bytes().as_slice()) + .bind(claimed_at.sqlite_value()) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target.index)) + .bind(i64::from(target.attempt_count)) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + let result = sqlx::query(MARK_JOB_ACTIVE_SQL) + .bind(claimed_at.sqlite_value()) + .bind(job_id.as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + if result.rows_affected() > 1 { + return Err(DeliveryOperationError::Storage); + } + let attempt = read_attempt(transaction, job_id, target.index, attempt_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + Ok(MycDeliveryClaim::Claimed(attempt)) +} + +async fn mark_submitted( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: MycDeliveryAttemptId, + submitted_at: MycDeliveryTimeUnixMs, +) -> Result<MycDeliveryAttemptRecord, DeliveryOperationError> { + let job = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + let target = target_by_relay(&job, relay_id)?; + let attempt = read_attempt(transaction, job_id, target.index, attempt_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + if attempt.status == MycDeliveryAttemptStatus::Submitted { + return (attempt.submitted_at == Some(submitted_at)) + .then_some(attempt) + .ok_or(DeliveryOperationError::Binding); + } + if attempt.status != MycDeliveryAttemptStatus::Leased + || submitted_at < attempt.leased_at + || submitted_at > attempt.lease_expires_at + || target.active_attempt_id != Some(attempt_id) + || target.status != MycDeliveryTargetStatus::Leased + { + return Err(DeliveryOperationError::Binding); + } + let result = sqlx::query(MARK_ATTEMPT_SUBMITTED_SQL) + .bind(submitted_at.sqlite_value()) + .bind(attempt_id.as_bytes().as_slice()) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target.index)) + .bind(submitted_at.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + let result = sqlx::query(MARK_TARGET_SUBMITTED_SQL) + .bind(submitted_at.sqlite_value()) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target.index)) + .bind(attempt_id.as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + read_attempt(transaction, job_id, target.index, attempt_id) + .await? + .ok_or(DeliveryOperationError::Binding) +} + +async fn record_outcome( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: MycDeliveryAttemptId, + outcome: MycDeliveryAttemptOutcome, + observed_at: MycDeliveryTimeUnixMs, +) -> Result<MycDeliveryJobRecord, DeliveryOperationError> { + let job = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + let target = target_by_relay(&job, relay_id)?.clone(); + let attempt = read_attempt(transaction, job_id, target.index, attempt_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + if matches!( + attempt.status, + MycDeliveryAttemptStatus::Delivered + | MycDeliveryAttemptStatus::Failed + | MycDeliveryAttemptStatus::Unknown + ) { + return (attempt.status == outcome.status() + && attempt.reason == Some(outcome.reason()) + && attempt.resolved_at == Some(observed_at)) + .then_some(job) + .ok_or(DeliveryOperationError::Binding); + } + let required_prior = match outcome { + MycDeliveryAttemptOutcome::Delivered + | MycDeliveryAttemptOutcome::UnknownAcknowledgement => MycDeliveryAttemptStatus::Submitted, + MycDeliveryAttemptOutcome::RelayRejected => MycDeliveryAttemptStatus::Submitted, + MycDeliveryAttemptOutcome::TransportFailed => attempt.status, + }; + if attempt.status != required_prior + || !matches!( + required_prior, + MycDeliveryAttemptStatus::Leased | MycDeliveryAttemptStatus::Submitted + ) + || target.active_attempt_id != Some(attempt_id) + || observed_at < attempt.leased_at + || attempt + .submitted_at + .is_some_and(|submitted_at| observed_at < submitted_at) + || observed_at > attempt.lease_expires_at + { + return Err(DeliveryOperationError::Binding); + } + resolve_attempt_and_target( + transaction, + &job, + &target, + &attempt, + required_prior, + outcome.status(), + outcome.reason(), + observed_at, + ) + .await?; + finalize_job_if_terminal(transaction, job_id, observed_at).await +} + +async fn recover_expired( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: MycDeliveryAttemptId, + observed_at: MycDeliveryTimeUnixMs, +) -> Result<MycDeliveryJobRecord, DeliveryOperationError> { + let job = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + let target = target_by_relay(&job, relay_id)?.clone(); + let attempt = read_attempt(transaction, job_id, target.index, attempt_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + if matches!( + attempt.status, + MycDeliveryAttemptStatus::Delivered + | MycDeliveryAttemptStatus::Failed + | MycDeliveryAttemptStatus::Unknown + ) { + return Ok(job); + } + if observed_at <= attempt.lease_expires_at || target.active_attempt_id != Some(attempt_id) { + return Err(DeliveryOperationError::Binding); + } + let (terminal, reason) = match attempt.status { + MycDeliveryAttemptStatus::Leased => ( + MycDeliveryAttemptStatus::Failed, + "lease_expired_before_submit", + ), + MycDeliveryAttemptStatus::Submitted => { + (MycDeliveryAttemptStatus::Unknown, "acknowledgement_lost") + } + MycDeliveryAttemptStatus::Delivered + | MycDeliveryAttemptStatus::Failed + | MycDeliveryAttemptStatus::Unknown => unreachable!(), + }; + resolve_attempt_and_target( + transaction, + &job, + &target, + &attempt, + attempt.status, + terminal, + reason, + observed_at, + ) + .await?; + finalize_job_if_terminal(transaction, job_id, observed_at).await +} + +#[allow(clippy::too_many_arguments)] +async fn resolve_attempt_and_target( + transaction: &mut ServiceSqliteTransaction<'_>, + job: &MycDeliveryJobRecord, + target: &MycDeliveryTargetRecord, + attempt: &MycDeliveryAttemptRecord, + prior: MycDeliveryAttemptStatus, + terminal: MycDeliveryAttemptStatus, + reason: &'static str, + observed_at: MycDeliveryTimeUnixMs, +) -> Result<(), DeliveryOperationError> { + let result = sqlx::query(RESOLVE_ATTEMPT_SQL) + .bind(terminal.as_str()) + .bind(observed_at.sqlite_value()) + .bind(reason) + .bind(attempt.id.as_bytes().as_slice()) + .bind(job.id.as_bytes().as_slice()) + .bind(i64::from(target.index)) + .bind(prior.as_str()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected())?; + let attempts_remaining = target.attempt_count < job.max_attempts; + let (target_status, next_attempt) = match terminal { + MycDeliveryAttemptStatus::Delivered => (MycDeliveryTargetStatus::Delivered, None), + MycDeliveryAttemptStatus::Failed if attempts_remaining => ( + MycDeliveryTargetStatus::Retryable, + Some(next_attempt_time(job, attempt.number, observed_at)?), + ), + MycDeliveryAttemptStatus::Unknown => ( + MycDeliveryTargetStatus::Unknown, + attempts_remaining + .then(|| next_attempt_time(job, attempt.number, observed_at)) + .transpose()?, + ), + MycDeliveryAttemptStatus::Failed => (MycDeliveryTargetStatus::Exhausted, None), + MycDeliveryAttemptStatus::Leased | MycDeliveryAttemptStatus::Submitted => { + return Err(DeliveryOperationError::Binding); + } + }; + let result = sqlx::query(RESOLVE_TARGET_SQL) + .bind(target_status.as_str()) + .bind(next_attempt.map(MycDeliveryTimeUnixMs::sqlite_value)) + .bind(observed_at.sqlite_value()) + .bind(job.id.as_bytes().as_slice()) + .bind(i64::from(target.index)) + .bind(attempt.id.as_bytes().as_slice()) + .bind(match prior { + MycDeliveryAttemptStatus::Leased => "leased", + MycDeliveryAttemptStatus::Submitted => "submitted", + _ => return Err(DeliveryOperationError::Binding), + }) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + require_one(result.rows_affected()) +} + +async fn finalize_job_if_terminal( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + observed_at: MycDeliveryTimeUnixMs, +) -> Result<MycDeliveryJobRecord, DeliveryOperationError> { + let job = read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding)?; + let delivered_required = job + .targets + .iter() + .filter(|target| target.required && target.status == MycDeliveryTargetStatus::Delivered) + .count(); + let possible_required = job + .targets + .iter() + .filter(|target| { + target.required + && target.status != MycDeliveryTargetStatus::Exhausted + && !(target.status == MycDeliveryTargetStatus::Unknown + && target.attempt_count >= job.max_attempts) + }) + .count(); + let required = usize::try_from(job.required_acknowledgements) + .map_err(|_| DeliveryOperationError::Binding)?; + let delivered = match job.policy_mode { + MycDeliveryPolicyMode::AtLeastOneRequired => delivered_required >= required, + MycDeliveryPolicyMode::AllRequired | MycDeliveryPolicyMode::RequiredQuorum => { + delivered_required >= required + } + }; + let possible = match job.policy_mode { + MycDeliveryPolicyMode::AtLeastOneRequired => possible_required >= required, + MycDeliveryPolicyMode::AllRequired | MycDeliveryPolicyMode::RequiredQuorum => { + possible_required >= required + } + }; + if delivered || !possible { + let terminal = if delivered { + MycDeliveryJobStatus::Delivered + } else if job.targets.iter().any(|target| { + target.required + && target.status == MycDeliveryTargetStatus::Unknown + && target.attempt_count >= job.max_attempts + }) { + MycDeliveryJobStatus::Unknown + } else { + MycDeliveryJobStatus::Failed + }; + let result = sqlx::query(FINALIZE_JOB_SQL) + .bind(terminal.as_str()) + .bind(observed_at.sqlite_value()) + .bind(observed_at.sqlite_value()) + .bind(job_id.as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + if result.rows_affected() > 1 { + return Err(DeliveryOperationError::Storage); + } + } + read_job(transaction, job_id) + .await? + .ok_or(DeliveryOperationError::Binding) +} + +fn next_attempt_time( + job: &MycDeliveryJobRecord, + attempt_number: u32, + observed_at: MycDeliveryTimeUnixMs, +) -> Result<MycDeliveryTimeUnixMs, DeliveryOperationError> { + let exponent = attempt_number.saturating_sub(1).min(31); + let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX); + let delay = job + .initial_backoff_ms + .saturating_mul(factor) + .min(job.maximum_backoff_ms); + let value = observed_at + .get() + .checked_add(delay) + .ok_or(DeliveryOperationError::Binding)?; + MycDeliveryTimeUnixMs::new(value).map_err(|_| DeliveryOperationError::Binding) +} + +async fn require_signer_operation( + transaction: &mut ServiceSqliteTransaction<'_>, + operation_id: MycSignerOperationId, +) -> Result<(), DeliveryOperationError> { + let row = sqlx::query(READ_SIGNER_OPERATION_SQL) + .bind(operation_id.as_bytes().as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + let count = row + .try_get::<i64, _>("row_count") + .map_err(|_| DeliveryOperationError::Binding)?; + (count == 1) + .then_some(()) + .ok_or(DeliveryOperationError::Binding) +} + +async fn read_job_by_operation( + transaction: &mut ServiceSqliteTransaction<'_>, + operation_id: MycSignerOperationId, +) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { + let rows = sqlx::query(READ_JOB_BY_OPERATION_SQL) + .bind(operation_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + read_job_rows(transaction, rows).await +} + +async fn read_job( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, +) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { + let rows = sqlx::query(READ_JOB_SQL) + .bind(job_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + read_job_rows(transaction, rows).await +} + +async fn read_job_rows( + transaction: &mut ServiceSqliteTransaction<'_>, + rows: Vec<sqlx::sqlite::SqliteRow>, +) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { + if rows.len() > 1 { + return Err(DeliveryOperationError::Binding); + } + let Some(row) = rows.first() else { + return Ok(None); + }; + let id = MycDeliveryJobId(blob32(row, "job_id")?); + let operation_id = MycSignerOperationId::from_persisted(blob32(row, "signer_operation_id")?); + let artifact_digest = MycDeliveryArtifactDigest(blob32(row, "artifact_sha256")?); + let policy_mode = MycDeliveryPolicyMode::parse(text(row, "policy_mode")?) + .ok_or(DeliveryOperationError::Binding)?; + let required_acknowledgements = positive_u32(row, "required_acknowledgements")?; + let max_attempts = positive_u32(row, "max_attempts")?; + if max_attempts > MYC_DELIVERY_ATTEMPT_MAX_COUNT { + return Err(DeliveryOperationError::Binding); + } + let initial_backoff_ms = positive_u64(row, "initial_backoff_ms")?; + let maximum_backoff_ms = positive_u64(row, "maximum_backoff_ms")?; + let attempt_deadline_ms = positive_u64(row, "attempt_deadline_ms")?; + if initial_backoff_ms > maximum_backoff_ms + || maximum_backoff_ms > 300_000 + || attempt_deadline_ms > 30_000 + { + return Err(DeliveryOperationError::Binding); + } + let status = + MycDeliveryJobStatus::parse(text(row, "status")?).ok_or(DeliveryOperationError::Binding)?; + let created_at = time(row, "created_at_unix_ms")?; + let updated_at = time(row, "updated_at_unix_ms")?; + let finalized_at = optional_time(row, "finalized_at_unix_ms", "finalized_at_type")?; + let targets = read_targets(transaction, id).await?; + let valid = !targets.is_empty() + && targets.len() <= MYC_DELIVERY_TARGET_MAX_COUNT + && targets + .iter() + .enumerate() + .all(|(index, target)| usize::try_from(target.index) == Ok(index)) + && matches!( + status, + MycDeliveryJobStatus::Delivered + | MycDeliveryJobStatus::Failed + | MycDeliveryJobStatus::Unknown + ) == finalized_at.is_some() + && targets.iter().all(|target| { + target.status != MycDeliveryTargetStatus::Unknown + || target.next_attempt_at.is_some() + || target.attempt_count == max_attempts + }); + if !valid { + return Err(DeliveryOperationError::Binding); + } + Ok(Some(MycDeliveryJobRecord { + id, + operation_id, + artifact_digest, + policy_mode, + required_acknowledgements, + max_attempts, + initial_backoff_ms, + maximum_backoff_ms, + attempt_deadline_ms, + status, + created_at, + updated_at, + finalized_at, + targets, + })) +} + +async fn read_targets( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, +) -> Result<Box<[MycDeliveryTargetRecord]>, DeliveryOperationError> { + let rows = sqlx::query(READ_TARGETS_SQL) + .bind(job_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + if rows.len() > MYC_DELIVERY_TARGET_MAX_COUNT { + return Err(DeliveryOperationError::Binding); + } + rows.iter() + .map(parse_target) + .collect::<Result<Vec<_>, _>>() + .map(Vec::into_boxed_slice) +} + +fn parse_target( + row: &sqlx::sqlite::SqliteRow, +) -> Result<MycDeliveryTargetRecord, DeliveryOperationError> { + let index = nonnegative_u32(row, "target_index")?; + let relay_id = MycDeliveryRelayId::new(text(row, "relay_id")?) + .map_err(|_| DeliveryOperationError::Binding)?; + let required = bool_value(row, "required")?; + let attempt_count = nonnegative_u32(row, "attempt_count")?; + if attempt_count > MYC_DELIVERY_ATTEMPT_MAX_COUNT { + return Err(DeliveryOperationError::Binding); + } + let status = MycDeliveryTargetStatus::parse(text(row, "status")?) + .ok_or(DeliveryOperationError::Binding)?; + let active_attempt_id = optional_blob32(row, "active_attempt_id", "active_attempt_id_type")? + .map(MycDeliveryAttemptId); + let next_attempt_at = optional_time(row, "next_attempt_at_unix_ms", "next_attempt_at_type")?; + let updated_at = time(row, "updated_at_unix_ms")?; + let active = matches!( + status, + MycDeliveryTargetStatus::Leased | MycDeliveryTargetStatus::Submitted + ); + let retryable = status == MycDeliveryTargetStatus::Retryable; + if active != active_attempt_id.is_some() + || (retryable && next_attempt_at.is_none()) + || (!retryable && status != MycDeliveryTargetStatus::Unknown && next_attempt_at.is_some()) + { + return Err(DeliveryOperationError::Binding); + } + Ok(MycDeliveryTargetRecord { + index, + relay_id, + required, + attempt_count, + status, + active_attempt_id, + next_attempt_at, + updated_at, + }) +} + +async fn read_attempts( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + target_index: u32, +) -> Result<Box<[MycDeliveryAttemptRecord]>, DeliveryOperationError> { + let rows = sqlx::query(READ_ATTEMPTS_SQL) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target_index)) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + if rows.len() > usize::try_from(MYC_DELIVERY_ATTEMPT_MAX_COUNT).unwrap_or(usize::MAX) { + return Err(DeliveryOperationError::Binding); + } + rows.iter() + .map(parse_attempt) + .collect::<Result<Vec<_>, _>>() + .map(Vec::into_boxed_slice) +} + +async fn read_attempt( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + target_index: u32, + attempt_id: MycDeliveryAttemptId, +) -> Result<Option<MycDeliveryAttemptRecord>, DeliveryOperationError> { + let rows = sqlx::query(READ_ATTEMPT_SQL) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target_index)) + .bind(attempt_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + one_attempt(rows) +} + +async fn read_attempt_by_nonce( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + target_index: u32, + nonce: &MycDeliveryAttemptNonce, +) -> Result<Option<MycDeliveryAttemptRecord>, DeliveryOperationError> { + let rows = sqlx::query(READ_ATTEMPT_BY_NONCE_SQL) + .bind(job_id.as_bytes().as_slice()) + .bind(i64::from(target_index)) + .bind(nonce.0.as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + one_attempt(rows) +} + +fn one_attempt( + rows: Vec<sqlx::sqlite::SqliteRow>, +) -> Result<Option<MycDeliveryAttemptRecord>, DeliveryOperationError> { + if rows.len() > 1 { + return Err(DeliveryOperationError::Binding); + } + rows.first().map(parse_attempt).transpose() +} + +fn parse_attempt( + row: &sqlx::sqlite::SqliteRow, +) -> Result<MycDeliveryAttemptRecord, DeliveryOperationError> { + let id = MycDeliveryAttemptId(blob32(row, "attempt_id")?); + let number = positive_u32(row, "attempt_number")?; + if number > MYC_DELIVERY_ATTEMPT_MAX_COUNT { + return Err(DeliveryOperationError::Binding); + } + let _nonce = blob32(row, "attempt_nonce")?; + let status = MycDeliveryAttemptStatus::parse(text(row, "status")?) + .ok_or(DeliveryOperationError::Binding)?; + let leased_at = time(row, "leased_at_unix_ms")?; + let lease_expires_at = time(row, "lease_expires_at_unix_ms")?; + let submitted_at = optional_time(row, "submitted_at_unix_ms", "submitted_at_type")?; + let resolved_at = optional_time(row, "resolved_at_unix_ms", "resolved_at_type")?; + let reason = optional_reason(row)?; + let valid = lease_expires_at > leased_at + && match status { + MycDeliveryAttemptStatus::Leased => { + submitted_at.is_none() && resolved_at.is_none() && reason.is_none() + } + MycDeliveryAttemptStatus::Submitted => { + submitted_at.is_some() && resolved_at.is_none() && reason.is_none() + } + MycDeliveryAttemptStatus::Delivered => { + submitted_at.is_some() && resolved_at.is_some() && reason == Some("accepted") + } + MycDeliveryAttemptStatus::Failed => resolved_at.is_some() && reason.is_some(), + MycDeliveryAttemptStatus::Unknown => { + submitted_at.is_some() + && resolved_at.is_some() + && reason == Some("acknowledgement_lost") + } + }; + valid + .then_some(MycDeliveryAttemptRecord { + id, + number, + status, + leased_at, + lease_expires_at, + submitted_at, + resolved_at, + reason, + }) + .ok_or(DeliveryOperationError::Binding) +} + +fn optional_reason( + row: &sqlx::sqlite::SqliteRow, +) -> Result<Option<&'static str>, DeliveryOperationError> { + let kind = row + .try_get::<&str, _>("reason_code_type") + .map_err(|_| DeliveryOperationError::Binding)?; + if kind == "null" { + return Ok(None); + } + if kind != "text" { + return Err(DeliveryOperationError::Binding); + } + match row + .try_get::<Option<&str>, _>("reason_code") + .map_err(|_| DeliveryOperationError::Binding)? + .ok_or(DeliveryOperationError::Binding)? + { + "accepted" => Ok(Some("accepted")), + "relay_rejected" => Ok(Some("relay_rejected")), + "transport_failed" => Ok(Some("transport_failed")), + "lease_expired_before_submit" => Ok(Some("lease_expired_before_submit")), + "acknowledgement_lost" => Ok(Some("acknowledgement_lost")), + _ => Err(DeliveryOperationError::Binding), + } +} + +fn target_by_relay<'a>( + job: &'a MycDeliveryJobRecord, + relay_id: &MycDeliveryRelayId, +) -> Result<&'a MycDeliveryTargetRecord, DeliveryOperationError> { + job.targets + .iter() + .find(|target| target.relay_id == *relay_id) + .ok_or(DeliveryOperationError::Binding) +} + +fn exact_job( + job: &MycDeliveryJobRecord, + request: &MycDeliveryJobRequest, + policy: &MycDeliveryPolicies, +) -> bool { + job.operation_id == request.operation_id + && job.artifact_digest == request.artifact_digest + && job.policy_mode == policy.mode + && job.required_acknowledgements == policy.required_acknowledgements + && job.max_attempts == policy.max_attempts + && job.initial_backoff_ms == policy.initial_backoff_ms + && job.maximum_backoff_ms == policy.maximum_backoff_ms + && job.attempt_deadline_ms == policy.attempt_deadline_ms + && job.created_at == request.created_at + && job.targets.len() == policy.targets.len() + && job + .targets + .iter() + .zip(policy.targets.iter()) + .all(|(actual, expected)| { + actual.relay_id == expected.relay_id && actual.required == expected.required + }) +} + +fn derive_job_id( + operation_id: MycSignerOperationId, + artifact: MycDeliveryArtifactDigest, +) -> MycDeliveryJobId { + let mut hasher = Sha256::new(); + hasher.update(JOB_ID_DOMAIN); + hasher.update(operation_id.as_bytes()); + hasher.update(artifact.as_bytes()); + MycDeliveryJobId(hasher.finalize().into()) +} + +fn derive_attempt_id( + job_id: MycDeliveryJobId, + target_index: u32, + attempt_number: u32, + nonce: &MycDeliveryAttemptNonce, +) -> MycDeliveryAttemptId { + let mut hasher = Sha256::new(); + hasher.update(ATTEMPT_ID_DOMAIN); + hasher.update(job_id.as_bytes()); + hasher.update(target_index.to_be_bytes()); + hasher.update(attempt_number.to_be_bytes()); + hasher.update(nonce.0); + MycDeliveryAttemptId(hasher.finalize().into()) +} + +fn text<'a>( + row: &'a sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<&'a str, DeliveryOperationError> { + row.try_get::<Option<&str>, _>(column) + .map_err(|_| DeliveryOperationError::Binding)? + .ok_or(DeliveryOperationError::Binding) +} + +fn blob32(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<[u8; 32], DeliveryOperationError> { + row.try_get::<Option<Vec<u8>>, _>(column) + .map_err(|_| DeliveryOperationError::Binding)? + .ok_or(DeliveryOperationError::Binding)? + .try_into() + .map_err(|_| DeliveryOperationError::Binding) +} + +fn optional_blob32( + row: &sqlx::sqlite::SqliteRow, + column: &str, + type_column: &str, +) -> Result<Option<[u8; 32]>, DeliveryOperationError> { + match row + .try_get::<&str, _>(type_column) + .map_err(|_| DeliveryOperationError::Binding)? + { + "null" => Ok(None), + "blob" => blob32(row, column).map(Some), + _ => Err(DeliveryOperationError::Binding), + } +} + +fn positive_u32( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u32, DeliveryOperationError> { + nonnegative_u32(row, column).and_then(|value| { + (value != 0) + .then_some(value) + .ok_or(DeliveryOperationError::Binding) + }) +} + +fn nonnegative_u32( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u32, DeliveryOperationError> { + let value = row + .try_get::<i64, _>(column) + .map_err(|_| DeliveryOperationError::Binding)?; + u32::try_from(value).map_err(|_| DeliveryOperationError::Binding) +} + +fn positive_u64( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u64, DeliveryOperationError> { + let value = row + .try_get::<i64, _>(column) + .map_err(|_| DeliveryOperationError::Binding)?; + u64::try_from(value) + .ok() + .filter(|value| *value != 0) + .ok_or(DeliveryOperationError::Binding) +} + +fn bool_value(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<bool, DeliveryOperationError> { + match row + .try_get::<i64, _>(column) + .map_err(|_| DeliveryOperationError::Binding)? + { + 0 => Ok(false), + 1 => Ok(true), + _ => Err(DeliveryOperationError::Binding), + } +} + +fn time( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<MycDeliveryTimeUnixMs, DeliveryOperationError> { + let value = positive_u64(row, column)?; + MycDeliveryTimeUnixMs::new(value).map_err(|_| DeliveryOperationError::Binding) +} + +fn optional_time( + row: &sqlx::sqlite::SqliteRow, + column: &str, + type_column: &str, +) -> Result<Option<MycDeliveryTimeUnixMs>, DeliveryOperationError> { + match row + .try_get::<&str, _>(type_column) + .map_err(|_| DeliveryOperationError::Binding)? + { + "null" => Ok(None), + "integer" => time(row, column).map(Some), + _ => Err(DeliveryOperationError::Binding), + } +} + +fn require_one(rows: u64) -> Result<(), DeliveryOperationError> { + (rows == 1) + .then_some(()) + .ok_or(DeliveryOperationError::Storage) +} + +fn map_transaction_error( + error: ServiceSqliteTransactionError<DeliveryOperationError>, +) -> MycStateRepositoryError { + if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { + return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); + } + let kind = match error.operation_error() { + Some(DeliveryOperationError::Binding) => MycStateRepositoryErrorKind::Binding, + Some(DeliveryOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, + }; + MycStateRepositoryError::new(kind) +} diff --git a/src/state_host.rs b/src/state_host.rs @@ -377,17 +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() == 4 + && outcome.applied_count() == 5 } 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, 4) - | (2, 3) - | (3, 2) - | (4, 1) + (MYC_STATE_BASE_SCHEMA_VERSION, 5) + | (2, 4) + | (3, 3) + | (4, 2) + | (5, 1) | (MYC_STATE_SCHEMA_VERSION, 0) ) } diff --git a/src/state_metadata.rs b/src/state_metadata.rs @@ -11,6 +11,7 @@ use radroots_service_sqlite::{ use radroots_storage::event::SourceGeneration; use sha2::{Digest, Sha256}; +use crate::state_delivery::{MycDeliveryPolicies, MycDeliveryPolicyMode, MycDeliveryRelayId}; use crate::state_governance::{ MycGovernancePolicies, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, }; @@ -171,6 +172,7 @@ pub struct MycStateMetadata { configuration: MycNormalizedConfigDigest, identities: MycExpectedIdentities, governance: MycGovernancePolicies, + delivery: MycDeliveryPolicies, policy_versions: MycStatePolicyVersions, } @@ -209,6 +211,7 @@ impl MycStateMetadata { ); let normalized = configuration.normalized(); let governance = governance_policies(normalized)?; + let delivery = delivery_policies(normalized)?; let configuration = normalized_config_digest(configuration.profile(), normalized)?; let identities = expected_identities(normalized)?; let policy_versions = MycStatePolicyVersions::governed(); @@ -231,6 +234,7 @@ impl MycStateMetadata { configuration, identities, governance, + delivery, policy_versions, }) } @@ -280,6 +284,10 @@ impl MycStateMetadata { self.governance.audit_retention_ms() } + pub(crate) const fn delivery_policies(&self) -> &MycDeliveryPolicies { + &self.delivery + } + pub(crate) fn matches_runtime(&self, runtime: &MycRuntimeContext) -> bool { ServiceSqlitePaths::from_runtime_context(runtime.context()) .is_ok_and(|paths| paths == self.paths) @@ -295,6 +303,7 @@ impl fmt::Debug for MycStateMetadata { .field("configuration", &self.configuration) .field("identities", &self.identities) .field("governance", &"[redacted]") + .field("delivery", &"[redacted]") .field("policy_versions", &self.policy_versions) .field("paths", &"[redacted]") .finish() @@ -445,6 +454,61 @@ fn governance_policies( .map_err(|_| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant)) } +fn delivery_policies( + normalized: &serde_json::Value, +) -> Result<MycDeliveryPolicies, MycStateMetadataError> { + let invalid = || MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant); + let integer = |pointer: &str| { + normalized + .pointer(pointer) + .and_then(serde_json::Value::as_u64) + .ok_or_else(invalid) + }; + let mode = normalized + .pointer("/transport/delivery_policy/mode") + .and_then(serde_json::Value::as_str) + .and_then(MycDeliveryPolicyMode::parse) + .ok_or_else(invalid)?; + let configured_quorum = normalized + .pointer("/transport/delivery_policy/required_acknowledgements") + .map(|value| { + value + .as_u64() + .and_then(|number| u32::try_from(number).ok()) + .ok_or_else(invalid) + }) + .transpose()?; + let targets = normalized + .pointer("/relays") + .and_then(serde_json::Value::as_array) + .ok_or_else(invalid)? + .iter() + .filter(|relay| relay.pointer("/write").and_then(serde_json::Value::as_bool) == Some(true)) + .map(|relay| { + let id = relay + .pointer("/id") + .and_then(serde_json::Value::as_str) + .ok_or_else(invalid)?; + let required = relay + .pointer("/required") + .and_then(serde_json::Value::as_bool) + .ok_or_else(invalid)?; + let id = MycDeliveryRelayId::new(id).map_err(|_| invalid())?; + Ok((id, required)) + }) + .collect::<Result<Vec<_>, MycStateMetadataError>>()?; + MycDeliveryPolicies::new( + mode, + configured_quorum, + u32::try_from(integer("/transport/publish_retry/max_attempts")?).map_err(|_| invalid())?, + integer("/transport/publish_retry/initial_backoff_ms")?, + integer("/transport/publish_retry/maximum_backoff_ms")?, + integer("/transport/publish_retry/attempt_deadline_ms")?, + targets, + ) + .map_err(|_| invalid()) +} + fn expected_identities( normalized: &serde_json::Value, ) -> Result<MycExpectedIdentities, MycStateMetadataError> { diff --git a/tests/services_hardening_delivery_state.rs b/tests/services_hardening_delivery_state.rs @@ -0,0 +1,539 @@ +#![forbid(unsafe_code)] +#![cfg(any(target_os = "linux", target_os = "macos"))] + +use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path}; + +use myc::{ + MYC_DELIVERY_RELAY_ID_MAX_BYTES, MYC_STATE_SCHEMA_VERSION, MycConfigProfile, + MycDeliveryArtifactDigest, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, + MycDeliveryAttemptStatus, MycDeliveryClaim, MycDeliveryJobAdmission, MycDeliveryJobRequest, + MycDeliveryJobStatus, MycDeliveryPolicyMode, MycDeliveryRelayId, MycDeliveryStateErrorKind, + MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, MycNip46ClientPublicKey, MycNip46EventId, + MycNip46RequestId, MycRequestReceivedAtUnixMs, MycSignerOperationNonce, MycSignerRequest, + MycSignerRequestDigest, MycSignerRequestMethod, MycStateMetadata, MycStateRepositoryErrorKind, + RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, initialize_myc_state, + open_myc_state_read_write, parse_myc_cli_v1_from, parse_myc_config_v1, + resolve_myc_runtime_context, +}; +use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; +use radroots_storage::event::SourceGeneration; +use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; + +const CONFIG_EXAMPLE: &[u8] = + include_bytes!("../contracts/services_hardening/config.v1.example.toml"); +const DELIVERY_SOURCE: &str = include_str!("../src/state_delivery.rs"); +const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); +const LIB_SOURCE: &str = include_str!("../src/lib.rs"); +const CLIENT_PUBLIC_KEY: &str = "2222222222222222222222222222222222222222222222222222222222222222"; + +fn runtime(root: &Path) -> myc::MycRuntimeContext { + let root = root.to_str().expect("UTF-8 temporary root"); + let invocation = parse_myc_cli_v1_from([ + "myc", + "--profile", + "repo-local", + "--instance", + "primary", + "--repo-local-root", + root, + "run", + ]) + .expect("valid invocation"); + resolve_myc_runtime_context( + &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), + &invocation, + ) + .expect("runtime context") +} + +fn prepare_state_directory(runtime: &myc::MycRuntimeContext) { + let directory = runtime.context().paths().state(); + fs::create_dir_all(directory).expect("state directory"); + fs::set_permissions(directory, fs::Permissions::from_mode(0o700)).expect("state mode"); +} + +fn metadata(runtime: &myc::MycRuntimeContext, source: &[u8]) -> MycStateMetadata { + let configuration = + parse_myc_config_v1(source, MycConfigProfile::RepoLocal).expect("configuration"); + MycStateMetadata::new( + runtime, + &configuration, + SourceGeneration::new([0x5a; 32]).expect("generation"), + 1_725_000_000_000, + ) + .expect("metadata") +} + +fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) { + let applied_at = MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time"); + let build = MigrationBuildIdentity::new( + env!("CARGO_PKG_VERSION"), + "1111111111111111111111111111111111111111", + "b44119fbac5985be8127ad1bf56d2950e6399427", + "rustc-test", + "test-target", + "service-host", + 1, + MYC_STATE_SCHEMA_VERSION, + 1, + 1, + 1, + ) + .expect("build identity"); + (applied_at, build) +} + +fn signer_request(request_id: &str, nonce: u8, received_at: u64) -> MycSignerRequest { + let canonical = format!(r#"{{"id":"{request_id}","method":"ping","params":[]}}"#); + MycSignerRequest::new( + MycNip46ClientPublicKey::new(CLIENT_PUBLIC_KEY).expect("client identity"), + MycNip46RequestId::new(request_id).expect("request ID"), + MycNip46EventId::from_bytes([nonce; 32]), + MycSignerRequestMethod::Ping, + MycSignerRequestDigest::for_canonical_request(canonical.as_bytes()).expect("digest"), + MycSignerOperationNonce::from_injected_entropy([nonce; 32]), + MycRequestReceivedAtUnixMs::new(received_at).expect("time"), + ) +} + +fn time(value: u64) -> MycDeliveryTimeUnixMs { + MycDeliveryTimeUnixMs::new(value).expect("delivery time") +} + +#[test] +fn delivery_inputs_and_diagnostics_are_closed_bounded_and_redacted() { + let maximum = format!("a{}", "1".repeat(MYC_DELIVERY_RELAY_ID_MAX_BYTES - 1)); + assert!(MycDeliveryRelayId::new(&maximum).is_ok()); + for invalid in [ + "", + "Primary", + "relay-name", + "relay__name", + "relay_", + &format!("a{maximum}"), + &"x".repeat(1024 * 1024), + ] { + assert_eq!( + MycDeliveryRelayId::new(invalid) + .expect_err("invalid relay") + .kind(), + MycDeliveryStateErrorKind::InvalidRelayId + ); + } + for invalid in [0, i64::MAX.unsigned_abs() + 1] { + assert_eq!( + MycDeliveryTimeUnixMs::new(invalid) + .expect_err("invalid time") + .kind(), + MycDeliveryStateErrorKind::InvalidTime + ); + } + assert!(MycDeliveryTimeUnixMs::new(i64::MAX.unsigned_abs()).is_ok()); + assert_eq!( + [ + MycDeliveryPolicyMode::AtLeastOneRequired, + MycDeliveryPolicyMode::AllRequired, + MycDeliveryPolicyMode::RequiredQuorum, + ] + .map(MycDeliveryPolicyMode::as_str), + ["at_least_one_required", "all_required", "required_quorum"] + ); + + let relay = MycDeliveryRelayId::new("relay_secret_123").expect("relay"); + let nonce = MycDeliveryAttemptNonce::from_injected_entropy([0x91; 32]); + let digest = MycDeliveryArtifactDigest::from_bytes([0x92; 32]); + let error = MycDeliveryRelayId::new("secret-invalid").expect_err("invalid"); + assert!(Error::source(&error).is_none()); + let rendered = format!("{relay:?} {nonce:?} {digest:?} {error} {error:?}"); + for secret in ["relay_secret_123", "secret-invalid", "145, 145", "146, 146"] { + assert!(!rendered.contains(secret)); + } +} + +#[tokio::test] +async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_aware() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + prepare_state_directory(&runtime); + let metadata = metadata(&runtime, CONFIG_EXAMPLE); + let (applied_at, build) = migration_evidence(); + initialize_myc_state(&runtime, &metadata, applied_at, &build) + .await + .expect("initialization"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("writer"); + let admitted = host + .repository() + .admit_signer_request(&signer_request("delivery-01", 0x31, 100)) + .await + .expect("signer request"); + let request = MycDeliveryJobRequest::signer_response( + admitted.record().operation_id(), + MycDeliveryArtifactDigest::from_bytes([0x44; 32]), + time(110), + ); + let job = host + .repository() + .create_delivery_job(&request) + .await + .expect("job"); + assert!(matches!(job, MycDeliveryJobAdmission::Created(_))); + let job_id = job.record().id(); + assert_eq!( + job.record().policy_mode(), + MycDeliveryPolicyMode::AllRequired + ); + assert_eq!(job.record().required_acknowledgements(), 2); + assert_eq!(job.record().max_attempts(), 5); + assert_eq!(job.record().initial_backoff_ms(), 250); + assert_eq!(job.record().maximum_backoff_ms(), 30_000); + assert_eq!(job.record().attempt_deadline_ms(), 15_000); + assert_eq!(job.record().targets().len(), 2); + assert_eq!(job.record().targets()[0].relay_id().as_str(), "primary"); + assert_eq!(job.record().targets()[1].relay_id().as_str(), "secondary"); + assert!( + job.record() + .targets() + .iter() + .all(|target| target.required()) + ); + + let replay = host + .repository() + .create_delivery_job(&request) + .await + .expect("exact replay"); + assert!(matches!(replay, MycDeliveryJobAdmission::ExactReplay(_))); + let conflicting = MycDeliveryJobRequest::signer_response( + admitted.record().operation_id(), + MycDeliveryArtifactDigest::from_bytes([0x45; 32]), + time(110), + ); + assert_eq!( + host.repository() + .create_delivery_job(&conflicting) + .await + .expect_err("conflicting job") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + + let primary = MycDeliveryRelayId::new("primary").expect("primary"); + let left_repository = host.repository(); + let right_repository = host.repository(); + let (left, right) = tokio::join!( + left_repository.claim_delivery_target( + job_id, + &primary, + MycDeliveryAttemptNonce::from_injected_entropy([0x51; 32]), + time(120), + ), + right_repository.claim_delivery_target( + job_id, + &primary, + MycDeliveryAttemptNonce::from_injected_entropy([0x51; 32]), + time(120), + ) + ); + let (claimed, replayed) = match (left.expect("left claim"), right.expect("right claim")) { + (MycDeliveryClaim::Claimed(claimed), MycDeliveryClaim::ExactReplay(replayed)) + | (MycDeliveryClaim::ExactReplay(replayed), MycDeliveryClaim::Claimed(claimed)) => { + (claimed, replayed) + } + outcome => panic!("unexpected concurrent outcome: {outcome:?}"), + }; + assert_eq!(claimed.id(), replayed.id()); + assert_eq!(claimed.number(), 1); + assert_eq!( + host.repository() + .record_delivery_attempt_outcome( + job_id, + &primary, + claimed.id(), + MycDeliveryAttemptOutcome::Delivered, + time(121), + ) + .await + .expect_err("delivery before submission") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + let submitted = host + .repository() + .mark_delivery_attempt_submitted(job_id, &primary, claimed.id(), time(121)) + .await + .expect("submitted"); + assert_eq!(submitted.status(), MycDeliveryAttemptStatus::Submitted); + let active = host + .repository() + .record_delivery_attempt_outcome( + job_id, + &primary, + claimed.id(), + MycDeliveryAttemptOutcome::Delivered, + time(122), + ) + .await + .expect("primary delivered"); + assert_eq!(active.status(), MycDeliveryJobStatus::Active); + + let secondary = MycDeliveryRelayId::new("secondary").expect("secondary"); + let first = match host + .repository() + .claim_delivery_target( + job_id, + &secondary, + MycDeliveryAttemptNonce::from_injected_entropy([0x61; 32]), + time(123), + ) + .await + .expect("secondary claim") + { + MycDeliveryClaim::Claimed(attempt) => attempt, + other => panic!("unexpected claim: {other:?}"), + }; + host.repository() + .mark_delivery_attempt_submitted(job_id, &secondary, first.id(), time(124)) + .await + .expect("secondary submitted"); + let unknown = host + .repository() + .record_delivery_attempt_outcome( + job_id, + &secondary, + first.id(), + MycDeliveryAttemptOutcome::UnknownAcknowledgement, + time(125), + ) + .await + .expect("unknown acknowledgement"); + assert_eq!(unknown.status(), MycDeliveryJobStatus::Active); + assert_eq!( + unknown.targets()[1].status(), + MycDeliveryTargetStatus::Unknown + ); + assert_eq!(unknown.targets()[1].next_attempt_at(), Some(time(375))); + assert!(matches!( + host.repository() + .claim_delivery_target( + job_id, + &secondary, + MycDeliveryAttemptNonce::from_injected_entropy([0x62; 32]), + time(374), + ) + .await + .expect("not ready"), + MycDeliveryClaim::NotReady + )); + let second = match host + .repository() + .claim_delivery_target( + job_id, + &secondary, + MycDeliveryAttemptNonce::from_injected_entropy([0x62; 32]), + time(375), + ) + .await + .expect("retry claim") + { + MycDeliveryClaim::Claimed(attempt) => attempt, + other => panic!("unexpected retry: {other:?}"), + }; + let retry = host + .repository() + .recover_expired_delivery_lease(job_id, &secondary, second.id(), time(15_376)) + .await + .expect("expired pre-submit lease"); + assert_eq!( + retry.targets()[1].status(), + MycDeliveryTargetStatus::Retryable + ); + assert_eq!(retry.targets()[1].next_attempt_at(), Some(time(15_876))); + let third = match host + .repository() + .claim_delivery_target( + job_id, + &secondary, + MycDeliveryAttemptNonce::from_injected_entropy([0x63; 32]), + time(15_876), + ) + .await + .expect("third claim") + { + MycDeliveryClaim::Claimed(attempt) => attempt, + other => panic!("unexpected third claim: {other:?}"), + }; + host.repository() + .mark_delivery_attempt_submitted(job_id, &secondary, third.id(), time(15_877)) + .await + .expect("third submitted"); + let delivered = host + .repository() + .record_delivery_attempt_outcome( + job_id, + &secondary, + third.id(), + MycDeliveryAttemptOutcome::Delivered, + time(15_878), + ) + .await + .expect("job delivered"); + assert_eq!(delivered.status(), MycDeliveryJobStatus::Delivered); + assert_eq!(delivered.finalized_at(), Some(time(15_878))); + let attempts = host + .repository() + .read_delivery_attempts(job_id, &secondary) + .await + .expect("attempt history"); + assert_eq!(attempts.len(), 3); + assert_eq!(attempts[0].status(), MycDeliveryAttemptStatus::Unknown); + assert_eq!(attempts[0].reason(), Some("acknowledgement_lost")); + assert_eq!(attempts[1].status(), MycDeliveryAttemptStatus::Failed); + assert_eq!(attempts[1].reason(), Some("lease_expired_before_submit")); + assert_eq!(attempts[2].status(), MycDeliveryAttemptStatus::Delivered); + host.close().await.expect("first close"); + + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("reopen"); + assert_eq!( + host.repository() + .read_delivery_job(job_id) + .await + .expect("read after reopen") + .expect("job") + .status(), + MycDeliveryJobStatus::Delivered + ); + host.close().await.expect("final close"); +} + +#[tokio::test] +async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_evidence() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + prepare_state_directory(&runtime); + let source = String::from_utf8(CONFIG_EXAMPLE.to_vec()) + .expect("UTF-8 config") + .replace("max_attempts = 5", "max_attempts = 1"); + let metadata = metadata(&runtime, source.as_bytes()); + let (applied_at, build) = migration_evidence(); + initialize_myc_state(&runtime, &metadata, applied_at, &build) + .await + .expect("initialization"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("writer"); + let admitted = host + .repository() + .admit_signer_request(&signer_request("delivery-unknown", 0x71, 100)) + .await + .expect("request"); + let job = host + .repository() + .create_delivery_job(&MycDeliveryJobRequest::signer_response( + admitted.record().operation_id(), + MycDeliveryArtifactDigest::from_bytes([0x72; 32]), + time(110), + )) + .await + .expect("job"); + let job_id = job.record().id(); + let primary = MycDeliveryRelayId::new("primary").expect("primary"); + let attempt = match host + .repository() + .claim_delivery_target( + job_id, + &primary, + MycDeliveryAttemptNonce::from_injected_entropy([0x73; 32]), + time(120), + ) + .await + .expect("claim") + { + MycDeliveryClaim::Claimed(attempt) => attempt, + other => panic!("unexpected claim: {other:?}"), + }; + host.repository() + .mark_delivery_attempt_submitted(job_id, &primary, attempt.id(), time(121)) + .await + .expect("submit"); + let job = host + .repository() + .record_delivery_attempt_outcome( + job_id, + &primary, + attempt.id(), + MycDeliveryAttemptOutcome::UnknownAcknowledgement, + time(122), + ) + .await + .expect("unknown"); + assert_eq!(job.status(), MycDeliveryJobStatus::Unknown); + assert_eq!(job.targets()[0].status(), MycDeliveryTargetStatus::Unknown); + assert_eq!(job.targets()[0].next_attempt_at(), None); + assert_eq!(job.finalized_at(), Some(time(122))); + host.close().await.expect("close"); + + let options = SqliteConnectOptions::new() + .filename(runtime.artifacts().state_database()) + .create_if_missing(false) + .disable_statement_logging(); + let mut connection = sqlx::SqliteConnection::connect_with(&options) + .await + .expect("inspection connection"); + for statement in [ + "DELETE FROM publication_attempts", + "DELETE FROM publication_targets", + "DELETE FROM publication_outbox", + "UPDATE publication_targets SET relay_id = 'changed'", + "UPDATE publication_outbox SET artifact_sha256 = zeroblob(32)", + ] { + assert!( + sqlx::query(statement) + .execute(&mut connection) + .await + .is_err(), + "guard accepted {statement}" + ); + } + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM publication_attempts") + .fetch_one(&mut connection) + .await + .expect("attempt count"), + 1 + ); + connection.close().await.expect("connection close"); +} + +#[test] +fn delivery_boundary_is_typed_sqlx_only_and_has_no_external_wait_or_legacy_authority() { + assert!(LIB_SOURCE.contains("mod state_delivery;")); + assert!(!LIB_SOURCE.contains("pub mod state_delivery;")); + assert!(DELIVERY_SOURCE.contains("ServiceSqliteTransaction<'_>")); + assert!(DELIVERY_SOURCE.contains("publication_attempts")); + assert!(DELIVERY_SOURCE.contains("UnknownAcknowledgement")); + assert!(CATALOG_SOURCE.contains("publication_outbox_no_delete")); + assert!(CATALOG_SOURCE.contains("publication_targets_no_delete")); + assert!(CATALOG_SOURCE.contains("publication_attempts_no_delete")); + for forbidden in [ + "SqlitePool", + "SqliteConnection", + "BEGIN ", + "COMMIT", + "ROLLBACK", + "relay_url", + "reqwest", + "tokio::spawn", + "std::env", + "std::time", + "outbox_sqlite", + ] { + assert!( + !DELIVERY_SOURCE.contains(forbidden), + "found forbidden delivery authority `{forbidden}`" + ); + } +} diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs @@ -11,8 +11,10 @@ use myc::{ 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, 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, + MYC_STATE_SCHEMA_VERSION_5_SHA256, MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256, + MYC_STATE_SCHEMA_VERSION_6_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_6_SHA256, + MycStateCatalogErrorKind, myc_migration_catalog, myc_schema_catalog, + validate_myc_state_catalogs, }; use radroots_service_sqlite::{ MigrationCatalog, MigrationChecksum, MigrationDescriptor, SchemaCatalog, SchemaDigest, @@ -24,13 +26,13 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); #[test] -fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { +fn schema_v1_through_v6_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, 5); - assert_eq!(migrations.descriptors().len(), 4); + assert_eq!(MYC_STATE_SCHEMA_VERSION, 6); + assert_eq!(migrations.descriptors().len(), 5); let metadata = &migrations.descriptors()[0]; assert_eq!(metadata.target_version(), 2); assert_eq!(metadata.name().as_str(), "create_myc_state_metadata"); @@ -65,13 +67,20 @@ fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { governance.checksum().as_bytes(), &MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256 ); - assert_eq!(migrations.current_version(), 5); + let delivery = &migrations.descriptors()[4]; + assert_eq!(delivery.target_version(), 6); + assert_eq!(delivery.name().as_str(), "create_delivery_evidence_state"); + assert_eq!( + delivery.checksum().as_bytes(), + &MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256 + ); + assert_eq!(migrations.current_version(), 6); assert_eq!( migrations.digest().as_bytes(), &MYC_MIGRATION_CATALOG_SHA256 ); - assert_eq!(schema.versions().len(), 5); + assert_eq!(schema.versions().len(), 6); assert_eq!(schema.versions()[0].version(), 1); assert_eq!( schema.versions()[0].object_count(), @@ -121,6 +130,16 @@ fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { schema.versions()[4].digest().as_bytes(), &MYC_STATE_SCHEMA_VERSION_5_SHA256 ); + assert_eq!(schema.versions()[5].version(), 6); + assert_eq!( + schema.versions()[5].object_count(), + MYC_STATE_SCHEMA_VERSION_6_OBJECT_COUNT + ); + assert_eq!(schema.versions()[5].object_count(), 43); + assert_eq!( + schema.versions()[5].digest().as_bytes(), + &MYC_STATE_SCHEMA_VERSION_6_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"); @@ -131,7 +150,7 @@ fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { ); assert_eq!( hex::encode(MYC_MIGRATION_CATALOG_SHA256), - "1d069996435560217dbd3692d5430612ad3a5667ee0b12b1919f0b5edf2cbe7d" + "a69bbc7f3c22750ddda55f36db125622c49f26b44cdc47210fadb7b5fc4f31e5" ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_VERSION_1_SHA256), @@ -143,7 +162,7 @@ fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_CATALOG_SHA256), - "2119efefcfd4ac5477609655341aa4c5b9b9a1befb88cc168a2b9373e99ebbd5" + "94f7adfbd62be03e6f3af9c9754bf2f47cc704a87b2708f6826c4432556b6f06" ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256), @@ -169,6 +188,14 @@ fn schema_v1_through_v5_and_all_migrations_have_exact_literal_identities() { hex::encode(MYC_STATE_SCHEMA_VERSION_5_SHA256), "fe89d4af7de4eded78a3f062dc9d300f06cf0f4cfca5fa5cffc93af1146f21c5" ); + assert_eq!( + hex::encode(MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256), + "4649b9afd03fc07f89fe184f027925675a0265951ea7cf75e06a4312c55a82bd" + ); + assert_eq!( + hex::encode(MYC_STATE_SCHEMA_VERSION_6_SHA256), + "557273306e68e9306dc1d7c7009e1cb7c9d15b8d52e07db7d0bd51e83f054492" + ); } #[test] @@ -212,19 +239,14 @@ fn independent_validator_rejects_migration_or_schema_drift() { let v3 = SchemaVersionCatalog::new(3, [object.clone()], v3_digest).expect("schema v3"); 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 object = SchemaObject::new( - SchemaObjectKind::Table, - "unexpected", - "unexpected", - SQL, - object_digest, - ) - .expect("schema object"); + let v4 = SchemaVersionCatalog::new(4, [object.clone()], v4_digest).expect("schema v4"); 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]) + let v5 = SchemaVersionCatalog::new(5, [object.clone()], v5_digest).expect("schema v5"); + let v6_digest = + SchemaVersionCatalog::computed_digest(6, [object.clone()]).expect("schema-v6 digest"); + let v6 = SchemaVersionCatalog::new(6, [object], v6_digest).expect("schema v6"); + let schema = SchemaCatalog::new(&expected_migrations, [v1, v2, v3, v4, v5, v6]) .expect("drift schema catalog"); assert_eq!( validate_myc_state_catalogs(&expected_migrations, &schema) 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(), 4); + assert_eq!(migrations.len(), 5); assert_eq!(migrations[0].get::<i64, _>(0), 2); assert_eq!( migrations[0].get::<String, _>(1), @@ -140,6 +140,11 @@ async fn initialization_migrates_and_binds_exact_metadata_before_inspection() { migrations[3].get::<String, _>(1), "create_bounded_governance_state" ); + assert_eq!(migrations[4].get::<i64, _>(0), 6); + assert_eq!( + migrations[4].get::<String, _>(1), + "create_delivery_evidence_state" + ); let binding = sqlx::query( "SELECT normalized_config_sha256, transport_public_key, user_public_key, \ discovery_public_key, config_contract_version, state_contract_version, \