myc

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

commit 961787013babf953e8dbdf2dc059a94081b0fc8d
parent 6207ddfc098349429fa48916ac212e08292597d6
Author: triesap <tyson@radroots.org>
Date:   Fri, 21 Aug 2026 19:20:43 +0000

state: govern discovery publication evidence

- persist canonical signed discovery events and NIP-05 projections atomically
- generalize immutable delivery evidence with source-bound schema guards
- pin schema v7 and cover replay, restart, promotion, and corruption boundaries

Diffstat:
MREADME | 9+++++++++
Msrc/lib.rs | 16++++++++++++----
Msrc/state_catalog.rs | 753+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Msrc/state_delivery.rs | 223+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Asrc/state_discovery.rs | 1235+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/state_host.rs | 13+++++++------
Msrc/state_metadata.rs | 12+++++++++++-
Mtests/services_hardening_delivery_state.rs | 20++++++++++----------
Atests/services_hardening_discovery_state.rs | 522+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_state_catalog.rs | 61++++++++++++++++++++++++++++++++++++++++++++++++++-----------
10 files changed, 2767 insertions(+), 97 deletions(-)

diff --git a/README b/README @@ -111,6 +111,15 @@ 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. +Discovery desired state is committed atomically with signature-verified, +canonical NIP-89 event bytes, deterministic NIP-05 projection inputs, the +immutable relay targets, and the initial delivery job. The desired generation +may advance independently, while current generation advances only after the +corresponding desired job has proven the configured delivery policy. Restart +reads return the same committed event bytes and projection inputs; publication +and hosted NIP-05 responses remain later runtime concerns and never run inside +the SQLite transaction. + 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 @@ -28,6 +28,7 @@ pub mod sql; mod state_catalog; mod state_connection; mod state_delivery; +mod state_discovery; mod state_governance; mod state_host; mod state_maintenance; @@ -124,8 +125,9 @@ pub use state_catalog::{ MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT, 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, + MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT, + MYC_STATE_SCHEMA_VERSION_7_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, @@ -144,8 +146,14 @@ pub use state_delivery::{ MycDeliveryAttemptOutcome, MycDeliveryAttemptRecord, MycDeliveryAttemptStatus, MycDeliveryClaim, MycDeliveryJobAdmission, MycDeliveryJobId, MycDeliveryJobRecord, MycDeliveryJobRequest, MycDeliveryJobStatus, MycDeliveryPolicyMode, MycDeliveryRelayId, - MycDeliveryStateError, MycDeliveryStateErrorKind, MycDeliveryTargetRecord, - MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, + MycDeliverySourceKind, MycDeliveryStateError, MycDeliveryStateErrorKind, + MycDeliveryTargetRecord, MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, +}; +pub use state_discovery::{ + MYC_DISCOVERY_DOCUMENT_MAX_BYTES, MYC_NIP05_PROJECTION_MAX_BYTES, MycDiscoveryCommitAdmission, + MycDiscoveryCommitRecord, MycDiscoveryCommitRequest, MycDiscoveryDocumentDigest, + MycDiscoveryDocumentRecord, MycDiscoveryGenerationId, MycDiscoveryPublicationState, + MycDiscoveryStateError, MycDiscoveryStateErrorKind, MycNip05ProjectionDigest, }; pub use state_governance::{ MYC_AUDIT_PAGE_MAX_ITEMS, MYC_AUDIT_RETENTION_MAX_MS, MYC_COMPACTION_MAX_ROWS, 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 = 6; +pub const MYC_STATE_SCHEMA_VERSION: u32 = 7; /// The shared metadata and migration-ledger objects present at schema v1. pub const MYC_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6; @@ -32,6 +32,9 @@ 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; +/// The shared objects plus Myc discovery desired/current state and exact documents. +pub const MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT: u32 = 53; + /// 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, @@ -46,14 +49,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] = [ - 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, + 0x6f, 0x47, 0xd1, 0xa6, 0x61, 0x42, 0x93, 0xb8, 0xd8, 0xd3, 0x10, 0x2a, 0xb5, 0x87, 0xc4, 0x2e, + 0x2b, 0xde, 0x40, 0xb2, 0x49, 0xa1, 0x65, 0xe2, 0x50, 0x4c, 0xe4, 0xa6, 0xf5, 0x08, 0x16, 0x77, ]; /// SHA-256 identity of the schema catalog bound to the migration catalog. pub const MYC_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [ - 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, + 0x4f, 0x8d, 0x6e, 0xe9, 0x87, 0x59, 0xcc, 0xad, 0x98, 0x42, 0xb1, 0xbe, 0xb6, 0xa1, 0xdc, 0x96, + 0xc2, 0x16, 0x49, 0x00, 0x8e, 0xd6, 0x47, 0xff, 0x73, 0x11, 0xcf, 0x22, 0x06, 0xc1, 0xf6, 0x20, ]; /// SHA-256 identity of the schema-v2 migration content. @@ -110,6 +113,18 @@ pub const MYC_STATE_SCHEMA_VERSION_6_SHA256: [u8; 32] = [ 0xc9, 0xd1, 0x5b, 0x8d, 0x52, 0xe0, 0x7d, 0xb7, 0xd0, 0xbd, 0x51, 0xe8, 0x3f, 0x05, 0x44, 0x92, ]; +/// SHA-256 identity of the schema-v7 discovery-state migration. +pub const MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256: [u8; 32] = [ + 0x60, 0x97, 0xc4, 0x07, 0x76, 0xa5, 0x7d, 0xd4, 0xbd, 0xdc, 0x04, 0xe6, 0x52, 0x98, 0x72, 0x17, + 0xd6, 0xf5, 0x99, 0x1f, 0x2d, 0x5e, 0x94, 0xd3, 0x22, 0xd1, 0x20, 0x86, 0x19, 0x25, 0x4e, 0x6c, +]; + +/// SHA-256 identity of the schema-v7 object snapshot. +pub const MYC_STATE_SCHEMA_VERSION_7_SHA256: [u8; 32] = [ + 0x3b, 0x35, 0x28, 0x91, 0x1a, 0x29, 0x34, 0x99, 0xd9, 0x72, 0x1d, 0xa7, 0x1b, 0x6a, 0xd0, 0x7a, + 0x5e, 0x72, 0x3e, 0x97, 0xb6, 0xeb, 0xf5, 0xb4, 0x0a, 0x19, 0xb9, 0x05, 0x89, 0x58, 0xbf, 0x79, +]; + /// 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, @@ -1064,6 +1079,463 @@ const CREATE_DELIVERY_STATE_MIGRATION_SQL: &str = concat!( publication_attempts_no_delete_sql!(), ); +macro_rules! delivery_jobs_table_sql { + () => { + r#"CREATE TABLE delivery_jobs ( + job_id BLOB NOT NULL PRIMARY KEY CHECK (length(job_id) = 32), + source_kind TEXT NOT NULL CHECK (source_kind IN ('signer_response', 'discovery_handler')), + source_id BLOB NOT NULL CHECK (length(source_id) = 32), + 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), + UNIQUE (source_kind, source_id), + 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! delivery_targets_table_sql { + () => { + r#"CREATE TABLE delivery_targets ( + job_id BLOB NOT NULL CHECK (length(job_id) = 32) + REFERENCES delivery_jobs(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! delivery_attempts_table_sql { + () => { + r#"CREATE TABLE delivery_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 delivery_targets(job_id, target_index), + UNIQUE (job_id, target_index, attempt_number), + UNIQUE (job_id, target_index, attempt_nonce) +) STRICT"# + }; +} + +macro_rules! delivery_jobs_guard_update_sql { + () => { + r#"CREATE TRIGGER delivery_jobs_guard_update +BEFORE UPDATE ON delivery_jobs +WHEN NEW.job_id != OLD.job_id + OR NEW.source_kind != OLD.source_kind + OR NEW.source_id != OLD.source_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! delivery_jobs_guard_insert_sql { + () => { + r#"CREATE TRIGGER delivery_jobs_guard_insert +BEFORE INSERT ON delivery_jobs +WHEN (NEW.source_kind = 'signer_response' + AND NOT EXISTS ( + SELECT 1 FROM nip46_requests WHERE operation_id = NEW.source_id + )) + OR (NEW.source_kind = 'discovery_handler' + AND NOT EXISTS ( + SELECT 1 FROM discovery_desired_state WHERE generation_id = NEW.source_id + )) +BEGIN + SELECT RAISE(ABORT, 'delivery job source is unknown'); +END"# + }; +} + +macro_rules! delivery_targets_guard_update_sql { + () => { + r#"CREATE TRIGGER delivery_targets_guard_update +BEFORE UPDATE ON delivery_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! delivery_attempts_guard_update_sql { + () => { + r#"CREATE TRIGGER delivery_attempts_guard_update +BEFORE UPDATE ON delivery_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! discovery_desired_state_table_sql { + () => { + r#"CREATE TABLE discovery_desired_state ( + generation_id BLOB NOT NULL PRIMARY KEY CHECK (length(generation_id) = 32), + normalized_config_sha256 BLOB NOT NULL CHECK (length(normalized_config_sha256) = 32), + desired_sha256 BLOB NOT NULL UNIQUE CHECK (length(desired_sha256) = 32), + created_at_unix_ms INTEGER NOT NULL + CHECK (created_at_unix_ms BETWEEN 1 AND 9223372036854775807) +) STRICT"# + }; +} + +macro_rules! discovery_documents_table_sql { + () => { + r#"CREATE TABLE discovery_documents ( + generation_id BLOB NOT NULL PRIMARY KEY CHECK (length(generation_id) = 32) + REFERENCES discovery_desired_state(generation_id), + event_id BLOB NOT NULL UNIQUE CHECK (length(event_id) = 32), + event_sha256 BLOB NOT NULL UNIQUE CHECK (length(event_sha256) = 32), + event_bytes BLOB NOT NULL CHECK (length(event_bytes) BETWEEN 1 AND 524288), + nip05_projection_sha256 BLOB NOT NULL CHECK (length(nip05_projection_sha256) = 32), + nip05_projection_bytes BLOB NOT NULL + CHECK (length(nip05_projection_bytes) BETWEEN 1 AND 524288) +) STRICT"# + }; +} + +macro_rules! discovery_publication_state_table_sql { + () => { + r#"CREATE TABLE discovery_publication_state ( + singleton INTEGER NOT NULL PRIMARY KEY CHECK (singleton = 1), + desired_generation_id BLOB NOT NULL UNIQUE CHECK (length(desired_generation_id) = 32) + REFERENCES discovery_desired_state(generation_id), + desired_job_id BLOB NOT NULL UNIQUE CHECK (length(desired_job_id) = 32) + REFERENCES delivery_jobs(job_id), + current_generation_id BLOB UNIQUE + CHECK (current_generation_id IS NULL OR length(current_generation_id) = 32) + REFERENCES discovery_desired_state(generation_id), + current_job_id BLOB UNIQUE + CHECK (current_job_id IS NULL OR length(current_job_id) = 32) + REFERENCES delivery_jobs(job_id), + updated_at_unix_ms INTEGER NOT NULL + CHECK (updated_at_unix_ms BETWEEN 1 AND 9223372036854775807), + CHECK ((current_generation_id IS NULL AND current_job_id IS NULL) + OR (current_generation_id IS NOT NULL AND current_job_id IS NOT NULL)) +) STRICT"# + }; +} + +macro_rules! immutable_no_update_sql { + ($trigger:literal, $table:literal, $message:literal) => { + concat!( + "CREATE TRIGGER ", + $trigger, + " BEFORE UPDATE ON ", + $table, + " BEGIN SELECT RAISE(ABORT, '", + $message, + "'); END" + ) + }; +} + +macro_rules! immutable_no_delete_sql { + ($trigger:literal, $table:literal, $message:literal) => { + concat!( + "CREATE TRIGGER ", + $trigger, + " BEFORE DELETE ON ", + $table, + " BEGIN SELECT RAISE(ABORT, '", + $message, + "'); END" + ) + }; +} + +macro_rules! discovery_publication_state_guard_update_sql { + () => { + r#"CREATE TRIGGER discovery_publication_state_guard_update +BEFORE UPDATE ON discovery_publication_state +WHEN NEW.singleton != OLD.singleton + OR NEW.updated_at_unix_ms < OLD.updated_at_unix_ms + OR NOT ( + (NEW.desired_generation_id != OLD.desired_generation_id + AND NEW.desired_job_id != OLD.desired_job_id + AND NEW.current_generation_id IS OLD.current_generation_id + AND NEW.current_job_id IS OLD.current_job_id) + OR (NEW.desired_generation_id = OLD.desired_generation_id + AND NEW.desired_job_id = OLD.desired_job_id + AND NEW.current_generation_id = OLD.desired_generation_id + AND NEW.current_job_id = OLD.desired_job_id) + ) +BEGIN + SELECT RAISE(ABORT, 'discovery publication transition is invalid'); +END"# + }; +} + +const CREATE_DELIVERY_JOBS_TABLE_SQL: &str = delivery_jobs_table_sql!(); +const CREATE_DELIVERY_TARGETS_TABLE_SQL: &str = delivery_targets_table_sql!(); +const CREATE_DELIVERY_ATTEMPTS_TABLE_SQL: &str = delivery_attempts_table_sql!(); +const CREATE_DELIVERY_JOBS_GUARD_INSERT_SQL: &str = delivery_jobs_guard_insert_sql!(); +const CREATE_DELIVERY_JOBS_GUARD_UPDATE_SQL: &str = delivery_jobs_guard_update_sql!(); +const CREATE_DELIVERY_TARGETS_GUARD_UPDATE_SQL: &str = delivery_targets_guard_update_sql!(); +const CREATE_DELIVERY_ATTEMPTS_GUARD_UPDATE_SQL: &str = delivery_attempts_guard_update_sql!(); +const CREATE_DELIVERY_JOBS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "delivery_jobs_no_delete", + "delivery_jobs", + "publication jobs are immutable" +); +const CREATE_DELIVERY_TARGETS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "delivery_targets_no_delete", + "delivery_targets", + "publication targets are immutable" +); +const CREATE_DELIVERY_ATTEMPTS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "delivery_attempts_no_delete", + "delivery_attempts", + "publication attempts are immutable" +); +const CREATE_DISCOVERY_DESIRED_STATE_TABLE_SQL: &str = discovery_desired_state_table_sql!(); +const CREATE_DISCOVERY_DOCUMENTS_TABLE_SQL: &str = discovery_documents_table_sql!(); +const CREATE_DISCOVERY_PUBLICATION_STATE_TABLE_SQL: &str = discovery_publication_state_table_sql!(); +const CREATE_DISCOVERY_DESIRED_STATE_NO_UPDATE_SQL: &str = immutable_no_update_sql!( + "discovery_desired_state_no_update", + "discovery_desired_state", + "discovery desired state is immutable" +); +const CREATE_DISCOVERY_DESIRED_STATE_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "discovery_desired_state_no_delete", + "discovery_desired_state", + "discovery desired state is immutable" +); +const CREATE_DISCOVERY_DOCUMENTS_NO_UPDATE_SQL: &str = immutable_no_update_sql!( + "discovery_documents_no_update", + "discovery_documents", + "discovery documents are immutable" +); +const CREATE_DISCOVERY_DOCUMENTS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "discovery_documents_no_delete", + "discovery_documents", + "discovery documents are immutable" +); +const CREATE_DISCOVERY_PUBLICATION_STATE_GUARD_UPDATE_SQL: &str = + discovery_publication_state_guard_update_sql!(); +const CREATE_DISCOVERY_PUBLICATION_STATE_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "discovery_publication_state_no_delete", + "discovery_publication_state", + "discovery publication state is retained" +); + +const CREATE_DISCOVERY_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 = 6 THEN 7 ELSE 0 END WHERE singleton = 1;\n", + myc_state_metadata_no_update_sql!(), + ";\n", + "DROP TRIGGER publication_attempts_no_delete;\n", + "DROP TRIGGER publication_targets_no_delete;\n", + "DROP TRIGGER publication_outbox_no_delete;\n", + "DROP TRIGGER publication_attempts_guard_update;\n", + "DROP TRIGGER publication_targets_guard_update;\n", + "DROP TRIGGER publication_outbox_guard_update;\n", + delivery_jobs_table_sql!(), + ";\n", + delivery_targets_table_sql!(), + ";\n", + delivery_attempts_table_sql!(), + ";\n", + "INSERT INTO delivery_jobs (job_id, source_kind, source_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) SELECT job_id, 'signer_response', ", + "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 FROM publication_outbox;\n", + "INSERT INTO delivery_targets SELECT * FROM publication_targets;\n", + "INSERT INTO delivery_attempts SELECT * FROM publication_attempts;\n", + "DROP TABLE publication_attempts;\n", + "DROP TABLE publication_targets;\n", + "DROP TABLE publication_outbox;\n", + delivery_jobs_guard_update_sql!(), + ";\n", + delivery_targets_guard_update_sql!(), + ";\n", + delivery_attempts_guard_update_sql!(), + ";\n", + immutable_no_delete_sql!( + "delivery_jobs_no_delete", + "delivery_jobs", + "publication jobs are immutable" + ), + ";\n", + immutable_no_delete_sql!( + "delivery_targets_no_delete", + "delivery_targets", + "publication targets are immutable" + ), + ";\n", + immutable_no_delete_sql!( + "delivery_attempts_no_delete", + "delivery_attempts", + "publication attempts are immutable" + ), + ";\n", + discovery_desired_state_table_sql!(), + ";\n", + discovery_documents_table_sql!(), + ";\n", + discovery_publication_state_table_sql!(), + ";\n", + delivery_jobs_guard_insert_sql!(), + ";\n", + immutable_no_update_sql!( + "discovery_desired_state_no_update", + "discovery_desired_state", + "discovery desired state is immutable" + ), + ";\n", + immutable_no_delete_sql!( + "discovery_desired_state_no_delete", + "discovery_desired_state", + "discovery desired state is immutable" + ), + ";\n", + immutable_no_update_sql!( + "discovery_documents_no_update", + "discovery_documents", + "discovery documents are immutable" + ), + ";\n", + immutable_no_delete_sql!( + "discovery_documents_no_delete", + "discovery_documents", + "discovery documents are immutable" + ), + ";\n", + discovery_publication_state_guard_update_sql!(), + ";\n", + immutable_no_delete_sql!( + "discovery_publication_state_no_delete", + "discovery_publication_state", + "discovery publication state is retained" + ), +); + 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, @@ -1184,6 +1656,82 @@ 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, ]; +const DELIVERY_JOBS_TABLE_SHA256: [u8; 32] = [ + 0xbb, 0xb1, 0xcf, 0xa2, 0x5e, 0x1e, 0xdd, 0x70, 0x30, 0x61, 0xd8, 0x43, 0xbe, 0x8d, 0x8c, 0x2b, + 0x16, 0x8c, 0x35, 0x6c, 0x9a, 0x9c, 0x0b, 0xe7, 0x88, 0x64, 0x18, 0xde, 0xca, 0xdd, 0x05, 0x01, +]; +const DELIVERY_TARGETS_TABLE_SHA256: [u8; 32] = [ + 0xce, 0xa5, 0x3a, 0xa4, 0x64, 0x18, 0xb1, 0x8e, 0xaf, 0x61, 0xe0, 0x90, 0xf9, 0x51, 0xe7, 0x76, + 0x79, 0xba, 0x33, 0x97, 0xc5, 0xa9, 0x5b, 0x3c, 0xa7, 0xeb, 0x83, 0x81, 0x4a, 0x8c, 0xaf, 0x69, +]; +const DELIVERY_ATTEMPTS_TABLE_SHA256: [u8; 32] = [ + 0xae, 0x8c, 0x2a, 0xa4, 0x61, 0xeb, 0xb2, 0x04, 0x69, 0x7b, 0xdc, 0x57, 0xc9, 0x34, 0x9c, 0x1b, + 0x86, 0x4b, 0x0f, 0xcd, 0x7b, 0xba, 0x9a, 0x45, 0xa7, 0x72, 0x2b, 0x49, 0xa0, 0x95, 0x94, 0xb8, +]; +const DELIVERY_JOBS_GUARD_INSERT_SHA256: [u8; 32] = [ + 0x29, 0x60, 0xe9, 0x9e, 0x32, 0xcd, 0x50, 0x0a, 0xfb, 0xa2, 0x5a, 0xc6, 0xbc, 0x2b, 0xfc, 0xa0, + 0x1a, 0x71, 0x95, 0x68, 0xa6, 0x2c, 0xa7, 0xac, 0x63, 0xd4, 0x79, 0x30, 0xc8, 0x5a, 0x8f, 0x7f, +]; +const DELIVERY_JOBS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0xc9, 0x5f, 0x5c, 0xac, 0x67, 0xa5, 0xe8, 0x1e, 0x88, 0xca, 0x5d, 0xf2, 0xc4, 0x8e, 0x71, 0x39, + 0xb6, 0x11, 0xc4, 0x4f, 0xa9, 0x59, 0x09, 0xea, 0xf7, 0x5e, 0x34, 0x61, 0xab, 0x5a, 0x5c, 0xe8, +]; +const DELIVERY_TARGETS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x1b, 0xf7, 0x81, 0xe9, 0x9e, 0xe2, 0xdd, 0x1a, 0xb9, 0xc7, 0xd5, 0xd4, 0x25, 0x96, 0x65, 0x11, + 0x79, 0x50, 0x1f, 0x74, 0xbf, 0x1a, 0x09, 0xdb, 0x99, 0x9c, 0x19, 0x9f, 0x25, 0xff, 0x98, 0x97, +]; +const DELIVERY_ATTEMPTS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x24, 0x5f, 0x0a, 0x53, 0xa0, 0xd5, 0x96, 0xa9, 0x62, 0x82, 0x70, 0xca, 0x9b, 0x89, 0x9e, 0xc2, + 0xd8, 0xe7, 0xf1, 0x45, 0x81, 0x61, 0x6f, 0x61, 0x62, 0x66, 0x98, 0xf5, 0x48, 0x29, 0x76, 0x9c, +]; +const DELIVERY_JOBS_NO_DELETE_SHA256: [u8; 32] = [ + 0x19, 0x9b, 0x4a, 0xc3, 0x1e, 0x53, 0x87, 0xed, 0x6a, 0xbc, 0x15, 0x69, 0x51, 0x38, 0x0d, 0xf2, + 0x61, 0xf7, 0x3c, 0x97, 0x0c, 0xc3, 0x7b, 0x16, 0xfc, 0x9e, 0xca, 0x47, 0x1a, 0x64, 0x3d, 0xb2, +]; +const DELIVERY_TARGETS_NO_DELETE_SHA256: [u8; 32] = [ + 0xa6, 0xa1, 0xd0, 0xb7, 0x42, 0xfd, 0x4f, 0xd7, 0x1d, 0x14, 0x0c, 0x64, 0xfd, 0xda, 0x1b, 0x1a, + 0x60, 0xfb, 0xea, 0xc4, 0xea, 0x6e, 0x22, 0xd8, 0x12, 0x79, 0xc3, 0xee, 0xbb, 0xf8, 0xb1, 0x29, +]; +const DELIVERY_ATTEMPTS_NO_DELETE_SHA256: [u8; 32] = [ + 0x8b, 0x6d, 0x86, 0x04, 0xe3, 0x6b, 0xf6, 0xb9, 0xa2, 0x21, 0x13, 0x39, 0xda, 0xd1, 0xf0, 0x9a, + 0x6a, 0xf3, 0x31, 0xef, 0x46, 0xbc, 0xdd, 0xd1, 0xe0, 0xe3, 0xef, 0x35, 0x07, 0x40, 0x08, 0xf1, +]; +const DISCOVERY_DESIRED_STATE_TABLE_SHA256: [u8; 32] = [ + 0xad, 0x85, 0x4a, 0x61, 0x50, 0xa4, 0x0f, 0xf9, 0x99, 0xe7, 0x2e, 0xf8, 0xa7, 0x63, 0xa8, 0xc8, + 0xcc, 0x57, 0x9d, 0x37, 0x23, 0xa2, 0x20, 0x67, 0x3c, 0xb1, 0xee, 0xbe, 0x9a, 0xda, 0xec, 0xce, +]; +const DISCOVERY_DOCUMENTS_TABLE_SHA256: [u8; 32] = [ + 0x0b, 0x44, 0x57, 0x52, 0xbd, 0xdc, 0x28, 0x9c, 0x90, 0x96, 0x42, 0x80, 0x00, 0x8e, 0x00, 0x2a, + 0x96, 0x81, 0x1a, 0x2a, 0x67, 0x96, 0x9c, 0x6c, 0x55, 0x44, 0x1b, 0xa2, 0xff, 0x0e, 0x01, 0x28, +]; +const DISCOVERY_PUBLICATION_STATE_TABLE_SHA256: [u8; 32] = [ + 0xc1, 0x61, 0x4f, 0xc6, 0x3e, 0x56, 0xae, 0xcd, 0xcf, 0x81, 0xe5, 0xea, 0x45, 0x93, 0x16, 0xc1, + 0x44, 0x71, 0x97, 0x11, 0xda, 0x97, 0x9c, 0xd5, 0xe4, 0x2a, 0x62, 0x54, 0x01, 0x24, 0xa1, 0x62, +]; +const DISCOVERY_DESIRED_STATE_NO_UPDATE_SHA256: [u8; 32] = [ + 0x5a, 0x4e, 0x12, 0x67, 0x2a, 0xb5, 0x75, 0xb7, 0x91, 0x35, 0xb0, 0xd2, 0x87, 0xcb, 0x33, 0x83, + 0x89, 0x8a, 0x51, 0x32, 0xf2, 0x34, 0xc1, 0x20, 0xd4, 0x5a, 0x56, 0x10, 0x6d, 0xe6, 0x92, 0x4c, +]; +const DISCOVERY_DESIRED_STATE_NO_DELETE_SHA256: [u8; 32] = [ + 0x5e, 0x1e, 0xcc, 0x9f, 0x37, 0x97, 0x38, 0xdf, 0x66, 0x0f, 0x89, 0xb2, 0xb6, 0x5d, 0x85, 0xd4, + 0xd2, 0x2c, 0xa1, 0x47, 0x7f, 0x20, 0x2e, 0x9b, 0xcb, 0x52, 0xb8, 0x95, 0x55, 0x8f, 0xeb, 0xf0, +]; +const DISCOVERY_DOCUMENTS_NO_UPDATE_SHA256: [u8; 32] = [ + 0x80, 0xf0, 0x11, 0x68, 0xe6, 0x3e, 0x94, 0xd1, 0xa6, 0xcd, 0x63, 0xf6, 0x16, 0xed, 0xdf, 0x05, + 0x7d, 0x58, 0xb4, 0x77, 0x8b, 0x83, 0xac, 0xc6, 0x07, 0xd3, 0x98, 0x11, 0xbf, 0x0d, 0x83, 0xf7, +]; +const DISCOVERY_DOCUMENTS_NO_DELETE_SHA256: [u8; 32] = [ + 0x23, 0xd4, 0x66, 0xe2, 0x21, 0xd3, 0x7a, 0xa1, 0xe0, 0xb7, 0xa6, 0x4e, 0x01, 0x5b, 0x1c, 0xf0, + 0xd5, 0x8d, 0xad, 0x1c, 0x31, 0x5b, 0x1a, 0x9e, 0x88, 0x17, 0x5c, 0x6b, 0x59, 0xd4, 0x47, 0xec, +]; +const DISCOVERY_PUBLICATION_STATE_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0xfc, 0xe8, 0x4a, 0xa7, 0x55, 0x00, 0x1c, 0x00, 0x12, 0x99, 0xc7, 0x58, 0xa1, 0xc6, 0x6a, 0x70, + 0x99, 0xba, 0x07, 0x77, 0xc7, 0x8e, 0xce, 0x5a, 0xeb, 0xa7, 0xf8, 0x5a, 0x27, 0x52, 0xb7, 0x5f, +]; +const DISCOVERY_PUBLICATION_STATE_NO_DELETE_SHA256: [u8; 32] = [ + 0xcd, 0x1b, 0xb2, 0x56, 0x6e, 0x3e, 0x46, 0x14, 0x4a, 0x4d, 0x38, 0xf1, 0xcb, 0xf9, 0xa0, 0xc4, + 0xb8, 0x94, 0xb8, 0x96, 0x71, 0x75, 0x49, 0x9c, 0x12, 0xb0, 0x15, 0x1d, 0xa0, 0x06, 0x1c, 0xa0, +]; /// Stable classes for invalid embedded Myc catalog definitions. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -1291,10 +1839,24 @@ pub fn myc_migration_catalog() -> Result<MigrationCatalog, MycStateCatalogError> 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))?; + let discovery = MigrationDescriptor::sql( + 7, + "create_discovery_desired_state", + CREATE_DISCOVERY_STATE_MIGRATION_SQL, + MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256), + ) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; + let catalog = MigrationCatalog::new([ + metadata, + requests, + connections, + governance, + delivery, + discovery, + ]) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?; if catalog.current_version() != MYC_STATE_SCHEMA_VERSION - || catalog.descriptors().len() != 5 + || catalog.descriptors().len() != 6 || catalog.digest().as_bytes() != &MYC_MIGRATION_CATALOG_SHA256 { return Err(MycStateCatalogError::new( @@ -1343,6 +1905,12 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> { SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_6_SHA256), ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; + let version_seven = SchemaVersionCatalog::new( + 7, + myc_state_discovery_objects()?, + SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_7_SHA256), + ) + .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; let catalog = SchemaCatalog::new( &migrations, [ @@ -1352,6 +1920,7 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> { version_four, version_five, version_six, + version_seven, ], ) .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?; @@ -1685,6 +2254,164 @@ fn myc_state_delivery_objects() -> Result<Vec<SchemaObject>, MycStateCatalogErro Ok(objects) } +fn myc_state_discovery_objects() -> Result<Vec<SchemaObject>, MycStateCatalogError> { + let mut objects = myc_state_delivery_objects()?; + objects.retain(|object| { + !matches!( + object.name(), + "publication_outbox" + | "publication_targets" + | "publication_attempts" + | "publication_outbox_guard_update" + | "publication_targets_guard_update" + | "publication_attempts_guard_update" + | "publication_outbox_no_delete" + | "publication_targets_no_delete" + | "publication_attempts_no_delete" + ) + }); + 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, + "delivery_jobs", + "delivery_jobs", + CREATE_DELIVERY_JOBS_TABLE_SQL, + DELIVERY_JOBS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "delivery_targets", + "delivery_targets", + CREATE_DELIVERY_TARGETS_TABLE_SQL, + DELIVERY_TARGETS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "delivery_attempts", + "delivery_attempts", + CREATE_DELIVERY_ATTEMPTS_TABLE_SQL, + DELIVERY_ATTEMPTS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_jobs_guard_update", + "delivery_jobs", + CREATE_DELIVERY_JOBS_GUARD_UPDATE_SQL, + DELIVERY_JOBS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_targets_guard_update", + "delivery_targets", + CREATE_DELIVERY_TARGETS_GUARD_UPDATE_SQL, + DELIVERY_TARGETS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_attempts_guard_update", + "delivery_attempts", + CREATE_DELIVERY_ATTEMPTS_GUARD_UPDATE_SQL, + DELIVERY_ATTEMPTS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_jobs_no_delete", + "delivery_jobs", + CREATE_DELIVERY_JOBS_NO_DELETE_SQL, + DELIVERY_JOBS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_targets_no_delete", + "delivery_targets", + CREATE_DELIVERY_TARGETS_NO_DELETE_SQL, + DELIVERY_TARGETS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_attempts_no_delete", + "delivery_attempts", + CREATE_DELIVERY_ATTEMPTS_NO_DELETE_SQL, + DELIVERY_ATTEMPTS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "discovery_desired_state", + "discovery_desired_state", + CREATE_DISCOVERY_DESIRED_STATE_TABLE_SQL, + DISCOVERY_DESIRED_STATE_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "discovery_documents", + "discovery_documents", + CREATE_DISCOVERY_DOCUMENTS_TABLE_SQL, + DISCOVERY_DOCUMENTS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "discovery_publication_state", + "discovery_publication_state", + CREATE_DISCOVERY_PUBLICATION_STATE_TABLE_SQL, + DISCOVERY_PUBLICATION_STATE_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "delivery_jobs_guard_insert", + "delivery_jobs", + CREATE_DELIVERY_JOBS_GUARD_INSERT_SQL, + DELIVERY_JOBS_GUARD_INSERT_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "discovery_desired_state_no_update", + "discovery_desired_state", + CREATE_DISCOVERY_DESIRED_STATE_NO_UPDATE_SQL, + DISCOVERY_DESIRED_STATE_NO_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "discovery_desired_state_no_delete", + "discovery_desired_state", + CREATE_DISCOVERY_DESIRED_STATE_NO_DELETE_SQL, + DISCOVERY_DESIRED_STATE_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "discovery_documents_no_update", + "discovery_documents", + CREATE_DISCOVERY_DOCUMENTS_NO_UPDATE_SQL, + DISCOVERY_DOCUMENTS_NO_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "discovery_documents_no_delete", + "discovery_documents", + CREATE_DISCOVERY_DOCUMENTS_NO_DELETE_SQL, + DISCOVERY_DOCUMENTS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "discovery_publication_state_guard_update", + "discovery_publication_state", + CREATE_DISCOVERY_PUBLICATION_STATE_GUARD_UPDATE_SQL, + DISCOVERY_PUBLICATION_STATE_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "discovery_publication_state_no_delete", + "discovery_publication_state", + CREATE_DISCOVERY_PUBLICATION_STATE_NO_DELETE_SQL, + DISCOVERY_PUBLICATION_STATE_NO_DELETE_SHA256, + )?, + ]); + Ok(objects) +} + /// Independently validates exact catalog versions, counts, and digests. pub fn validate_myc_state_catalogs( migrations: &MigrationCatalog, @@ -1693,7 +2420,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() == 5 + && descriptors.len() == 6 && 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 @@ -1709,9 +2436,12 @@ pub fn validate_myc_state_catalogs( && 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 + && descriptors[5].target_version() == 7 + && descriptors[5].name().as_str() == "create_discovery_desired_state" + && descriptors[5].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256 && migrations.digest().as_bytes() == &MYC_MIGRATION_CATALOG_SHA256 && schema.migration_catalog_digest() == migrations.digest() - && versions.len() == 6 + && versions.len() == 7 && 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 @@ -1730,6 +2460,9 @@ pub fn validate_myc_state_catalogs( && 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 + && versions[6].version() == 7 + && versions[6].object_count() == MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT + && versions[6].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_7_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 @@ -32,8 +32,10 @@ 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(source_kind) = 'text' AND length(CAST(source_kind AS BLOB)) <= 32 + THEN source_kind ELSE NULL END AS source_kind, + CASE WHEN typeof(source_id) = 'blob' AND length(source_id) = 32 + THEN source_id ELSE NULL END AS source_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 @@ -44,15 +46,17 @@ const READ_JOB_SQL: &str = r#"SELECT 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 +FROM delivery_jobs WHERE job_id = ? LIMIT 2"#; -const READ_JOB_BY_OPERATION_SQL: &str = r#"SELECT +const READ_JOB_BY_SOURCE_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(source_kind) = 'text' AND length(CAST(source_kind AS BLOB)) <= 32 + THEN source_kind ELSE NULL END AS source_kind, + CASE WHEN typeof(source_id) = 'blob' AND length(source_id) = 32 + THEN source_id ELSE NULL END AS source_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 @@ -63,8 +67,8 @@ const READ_JOB_BY_OPERATION_SQL: &str = r#"SELECT 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 = ? +FROM delivery_jobs +WHERE source_kind = ? AND source_id = ? LIMIT 2"#; const READ_TARGETS_SQL: &str = r#"SELECT @@ -81,7 +85,7 @@ const READ_TARGETS_SQL: &str = r#"SELECT next_attempt_at_unix_ms, typeof(next_attempt_at_unix_ms) AS next_attempt_at_type, updated_at_unix_ms -FROM publication_targets +FROM delivery_targets WHERE job_id = ? ORDER BY target_index LIMIT 33"#; @@ -100,7 +104,7 @@ const READ_ATTEMPTS_SQL: &str = r#"SELECT 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 +FROM delivery_attempts WHERE job_id = ? AND target_index = ? ORDER BY attempt_number LIMIT 33"#; @@ -119,7 +123,7 @@ const READ_ATTEMPT_SQL: &str = r#"SELECT 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 +FROM delivery_attempts WHERE job_id = ? AND target_index = ? AND attempt_id = ? LIMIT 2"#; @@ -137,58 +141,58 @@ const READ_ATTEMPT_BY_NONCE_SQL: &str = r#"SELECT 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 +FROM delivery_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, +const INSERT_JOB_SQL: &str = r#"INSERT INTO delivery_jobs ( + job_id, source_kind, source_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)"#; +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?, NULL)"#; -const INSERT_TARGET_SQL: &str = r#"INSERT INTO publication_targets ( +const INSERT_TARGET_SQL: &str = r#"INSERT INTO delivery_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 ( +const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO delivery_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 +const CLAIM_TARGET_SQL: &str = r#"UPDATE delivery_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 +const MARK_JOB_ACTIVE_SQL: &str = r#"UPDATE delivery_jobs SET status = 'active', updated_at_unix_ms = ? WHERE job_id = ? AND status = 'pending'"#; -const MARK_ATTEMPT_SUBMITTED_SQL: &str = r#"UPDATE publication_attempts +const MARK_ATTEMPT_SUBMITTED_SQL: &str = r#"UPDATE delivery_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 +const MARK_TARGET_SUBMITTED_SQL: &str = r#"UPDATE delivery_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 +const RESOLVE_ATTEMPT_SQL: &str = r#"UPDATE delivery_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 +const RESOLVE_TARGET_SQL: &str = r#"UPDATE delivery_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 +const FINALIZE_JOB_SQL: &str = r#"UPDATE delivery_jobs SET status = ?, updated_at_unix_ms = ?, finalized_at_unix_ms = ? WHERE job_id = ? AND status IN ('pending', 'active')"#; @@ -285,6 +289,12 @@ redacted_id!( "MycDeliveryArtifactDigest([redacted])" ); +impl MycDeliveryJobId { + pub(crate) const fn from_persisted(bytes: [u8; 32]) -> Self { + Self(bytes) + } +} + impl MycDeliveryArtifactDigest { /// Wraps an independently verified exact-artifact SHA-256 identity. #[must_use] @@ -331,7 +341,7 @@ impl MycDeliveryTimeUnixMs { self.0 } - fn sqlite_value(self) -> i64 { + pub(crate) fn sqlite_value(self) -> i64 { i64::try_from(self.0).expect("validated delivery time fits SQLite") } } @@ -483,6 +493,72 @@ impl fmt::Debug for MycDeliveryPolicies { } } +/// Closed authority that produced the exact bytes retained by a delivery job. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDeliverySourceKind { + SignerResponse, + DiscoveryHandler, +} + +impl MycDeliverySourceKind { + /// Returns the exact durable spelling. + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::SignerResponse => "signer_response", + Self::DiscoveryHandler => "discovery_handler", + } + } + + pub(crate) fn parse(value: &str) -> Option<Self> { + match value { + "signer_response" => Some(Self::SignerResponse), + "discovery_handler" => Some(Self::DiscoveryHandler), + _ => None, + } + } +} + +#[derive(Clone, Copy, PartialEq, Eq)] +pub(crate) struct MycDeliverySource { + kind: MycDeliverySourceKind, + id: [u8; 32], +} + +impl MycDeliverySource { + pub(crate) const fn signer_response(operation_id: MycSignerOperationId) -> Self { + Self { + kind: MycDeliverySourceKind::SignerResponse, + id: *operation_id.as_bytes(), + } + } + + pub(crate) const fn discovery_handler(generation_id: [u8; 32]) -> Self { + Self { + kind: MycDeliverySourceKind::DiscoveryHandler, + id: generation_id, + } + } + + pub(crate) const fn kind(self) -> MycDeliverySourceKind { + self.kind + } + + pub(crate) const fn id(self) -> [u8; 32] { + self.id + } +} + +impl fmt::Debug for MycDeliverySource { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliverySource") + .field("kind", &self.kind) + .field("id", &"[redacted]") + .finish() + } +} + /// Immutable signer-response delivery job input. pub struct MycDeliveryJobRequest { operation_id: MycSignerOperationId, @@ -543,7 +619,7 @@ impl MycDeliveryJobStatus { } } - fn parse(value: &str) -> Option<Self> { + pub(crate) fn parse(value: &str) -> Option<Self> { match value { "pending" => Some(Self::Pending), "active" => Some(Self::Active), @@ -663,7 +739,7 @@ impl MycDeliveryAttemptOutcome { #[derive(Clone, PartialEq, Eq)] pub struct MycDeliveryJobRecord { id: MycDeliveryJobId, - operation_id: MycSignerOperationId, + source: MycDeliverySource, artifact_digest: MycDeliveryArtifactDigest, policy_mode: MycDeliveryPolicyMode, required_acknowledgements: u32, @@ -684,8 +760,17 @@ impl MycDeliveryJobRecord { self.id } #[must_use] - pub const fn operation_id(&self) -> MycSignerOperationId { - self.operation_id + pub const fn operation_id(&self) -> Option<MycSignerOperationId> { + match self.source.kind { + MycDeliverySourceKind::SignerResponse => { + Some(MycSignerOperationId::from_persisted(self.source.id)) + } + MycDeliverySourceKind::DiscoveryHandler => None, + } + } + #[must_use] + pub const fn source_kind(&self) -> MycDeliverySourceKind { + self.source.kind } #[must_use] pub const fn artifact_digest(&self) -> MycDeliveryArtifactDigest { @@ -741,6 +826,7 @@ impl fmt::Debug for MycDeliveryJobRecord { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter .debug_struct("MycDeliveryJobRecord") + .field("source_kind", &self.source.kind) .field("status", &self.status) .field("target_count", &self.targets.len()) .finish() @@ -916,13 +1002,22 @@ impl MycStateRepository<'_> { request: &MycDeliveryJobRequest, ) -> Result<MycDeliveryJobAdmission, MycStateRepositoryError> { let request = request.owned(); + let source = MycDeliverySource::signer_response(request.operation_id); 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 + require_signer_operation(transaction, request.operation_id).await?; + create_job( + transaction, + source, + request.artifact_digest, + request.created_at, + &policy, + ) + .await }) }) .await @@ -1064,7 +1159,7 @@ impl MycStateRepository<'_> { } #[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum DeliveryOperationError { +pub(crate) enum DeliveryOperationError { Binding, Storage, } @@ -1081,22 +1176,24 @@ async fn verify_metadata( }) } -async fn create_job( +pub(crate) async fn create_job( transaction: &mut ServiceSqliteTransaction<'_>, - request: &MycDeliveryJobRequest, + source: MycDeliverySource, + artifact_digest: MycDeliveryArtifactDigest, + created_at: MycDeliveryTimeUnixMs, 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) + if let Some(existing) = read_job_by_source(transaction, source).await? { + return exact_job(&existing, source, artifact_digest, created_at, policy) .then_some(MycDeliveryJobAdmission::ExactReplay(existing)) .ok_or(DeliveryOperationError::Binding); } - let job_id = derive_job_id(request.operation_id, request.artifact_digest); + let job_id = derive_job_id(source, 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(source.kind().as_str()) + .bind(source.id().as_slice()) + .bind(artifact_digest.as_bytes().as_slice()) .bind(policy.mode.as_str()) .bind(i64::from(policy.required_acknowledgements)) .bind(i64::from(policy.max_attempts)) @@ -1112,8 +1209,8 @@ async fn create_job( i64::try_from(policy.attempt_deadline_ms) .map_err(|_| DeliveryOperationError::Binding)?, ) - .bind(request.created_at.sqlite_value()) - .bind(request.created_at.sqlite_value()) + .bind(created_at.sqlite_value()) + .bind(created_at.sqlite_value()) .execute(&mut *transaction) .await .map_err(|_| DeliveryOperationError::Storage)?; @@ -1125,7 +1222,7 @@ async fn create_job( .bind(i64::from(index)) .bind(target.relay_id.as_str()) .bind(target.required) - .bind(request.created_at.sqlite_value()) + .bind(created_at.sqlite_value()) .execute(&mut *transaction) .await .map_err(|_| DeliveryOperationError::Storage)?; @@ -1540,16 +1637,22 @@ async fn require_signer_operation( .ok_or(DeliveryOperationError::Binding) } -async fn read_job_by_operation( +async fn read_job_by_source( transaction: &mut ServiceSqliteTransaction<'_>, - operation_id: MycSignerOperationId, + source: MycDeliverySource, ) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { - let rows = sqlx::query(READ_JOB_BY_OPERATION_SQL) - .bind(operation_id.as_bytes().as_slice()) + let rows = sqlx::query(READ_JOB_BY_SOURCE_SQL) + .bind(source.kind().as_str()) + .bind(source.id().as_slice()) .fetch_all(&mut *transaction) .await .map_err(|_| DeliveryOperationError::Storage)?; - read_job_rows(transaction, rows).await + let job = read_job_rows(transaction, rows).await?; + match job { + Some(job) if job.source == source => Ok(Some(job)), + Some(_) => Err(DeliveryOperationError::Binding), + None => Ok(None), + } } async fn read_job( @@ -1575,7 +1678,12 @@ async fn read_job_rows( return Ok(None); }; let id = MycDeliveryJobId(blob32(row, "job_id")?); - let operation_id = MycSignerOperationId::from_persisted(blob32(row, "signer_operation_id")?); + let source_kind = MycDeliverySourceKind::parse(text(row, "source_kind")?) + .ok_or(DeliveryOperationError::Binding)?; + let source = MycDeliverySource { + kind: source_kind, + id: blob32(row, "source_id")?, + }; let artifact_digest = MycDeliveryArtifactDigest(blob32(row, "artifact_sha256")?); let policy_mode = MycDeliveryPolicyMode::parse(text(row, "policy_mode")?) .ok_or(DeliveryOperationError::Binding)?; @@ -1621,7 +1729,7 @@ async fn read_job_rows( } Ok(Some(MycDeliveryJobRecord { id, - operation_id, + source, artifact_digest, policy_mode, required_acknowledgements, @@ -1842,18 +1950,20 @@ fn target_by_relay<'a>( fn exact_job( job: &MycDeliveryJobRecord, - request: &MycDeliveryJobRequest, + source: MycDeliverySource, + artifact_digest: MycDeliveryArtifactDigest, + created_at: MycDeliveryTimeUnixMs, policy: &MycDeliveryPolicies, ) -> bool { - job.operation_id == request.operation_id - && job.artifact_digest == request.artifact_digest + job.source == source + && job.artifact_digest == 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.created_at == created_at && job.targets.len() == policy.targets.len() && job .targets @@ -1865,12 +1975,15 @@ fn exact_job( } fn derive_job_id( - operation_id: MycSignerOperationId, + source: MycDeliverySource, artifact: MycDeliveryArtifactDigest, ) -> MycDeliveryJobId { let mut hasher = Sha256::new(); hasher.update(JOB_ID_DOMAIN); - hasher.update(operation_id.as_bytes()); + if source.kind() == MycDeliverySourceKind::DiscoveryHandler { + hasher.update(b"discovery_handler\0"); + } + hasher.update(source.id()); hasher.update(artifact.as_bytes()); MycDeliveryJobId(hasher.finalize().into()) } diff --git a/src/state_discovery.rs b/src/state_discovery.rs @@ -0,0 +1,1235 @@ +//! Durable discovery desired/current state and exact committed publication bytes. + +use core::fmt; +use std::{collections::BTreeMap, error::Error}; + +use radroots_service_sqlite::{ + ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, +}; +use serde::Serialize; +use sha2::{Digest, Sha256}; +use sqlx::Row; + +use crate::nostr_contract::{ + RadrootsNostrEvent, RadrootsNostrKind, RadrootsNostrMetadata, RadrootsNostrRelayUrl, +}; +use crate::state_delivery::{ + DeliveryOperationError, MycDeliveryArtifactDigest, MycDeliveryJobId, MycDeliveryJobRecord, + MycDeliveryJobStatus, MycDeliverySource, MycDeliveryTimeUnixMs, create_job, +}; +use crate::state_repository::{ + MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, + RepositoryOperationError, require_expected_metadata, +}; +use crate::{MycExpectedIdentities, MycStateMetadata}; + +/// Maximum exact signed-event bytes admitted from the configured event bound. +pub const MYC_DISCOVERY_DOCUMENT_MAX_BYTES: usize = 524_288; +/// Maximum deterministic NIP-05 projection input bytes. +pub const MYC_NIP05_PROJECTION_MAX_BYTES: usize = 524_288; + +const NIP46_RPC_KIND: u32 = 24_133; +const NIP89_HANDLER_KIND: u16 = 31_990; +const DESIRED_DIGEST_DOMAIN: &[u8] = b"radroots.myc.discovery_desired.v1\0"; +const GENERATION_ID_DOMAIN: &[u8] = b"radroots.myc.discovery_generation.v1\0"; + +const READ_STATE_SQL: &str = r#"SELECT + CASE WHEN typeof(desired_generation_id) = 'blob' AND length(desired_generation_id) = 32 + THEN desired_generation_id ELSE NULL END AS desired_generation_id, + CASE WHEN typeof(desired_job_id) = 'blob' AND length(desired_job_id) = 32 + THEN desired_job_id ELSE NULL END AS desired_job_id, + typeof(current_generation_id) AS current_generation_id_type, + CASE WHEN typeof(current_generation_id) = 'blob' AND length(current_generation_id) = 32 + THEN current_generation_id ELSE NULL END AS current_generation_id, + typeof(current_job_id) AS current_job_id_type, + CASE WHEN typeof(current_job_id) = 'blob' AND length(current_job_id) = 32 + THEN current_job_id ELSE NULL END AS current_job_id, + updated_at_unix_ms +FROM discovery_publication_state +WHERE singleton = 1 +LIMIT 2"#; + +const READ_DOCUMENT_BY_GENERATION_SQL: &str = r#"SELECT + CASE WHEN typeof(generation_id) = 'blob' AND length(generation_id) = 32 + THEN generation_id ELSE NULL END AS generation_id, + CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32 + THEN event_id ELSE NULL END AS event_id, + CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 + THEN event_sha256 ELSE NULL END AS event_sha256, + CASE WHEN typeof(event_bytes) = 'blob' AND length(event_bytes) BETWEEN 1 AND 524288 + THEN event_bytes ELSE NULL END AS event_bytes, + CASE WHEN typeof(nip05_projection_sha256) = 'blob' + AND length(nip05_projection_sha256) = 32 + THEN nip05_projection_sha256 ELSE NULL END AS nip05_projection_sha256, + CASE WHEN typeof(nip05_projection_bytes) = 'blob' + AND length(nip05_projection_bytes) BETWEEN 1 AND 524288 + THEN nip05_projection_bytes ELSE NULL END AS nip05_projection_bytes +FROM discovery_documents +WHERE generation_id = ? +LIMIT 2"#; + +const READ_DESIRED_SQL: &str = r#"SELECT + CASE WHEN typeof(generation_id) = 'blob' AND length(generation_id) = 32 + THEN generation_id ELSE NULL END AS generation_id, + CASE WHEN typeof(normalized_config_sha256) = 'blob' + AND length(normalized_config_sha256) = 32 + THEN normalized_config_sha256 ELSE NULL END AS normalized_config_sha256, + CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 + THEN desired_sha256 ELSE NULL END AS desired_sha256, + created_at_unix_ms +FROM discovery_desired_state +WHERE generation_id = ? +LIMIT 2"#; + +const READ_JOB_SQL: &str = r#"SELECT status +FROM delivery_jobs +WHERE job_id = ? AND source_kind = 'discovery_handler' AND source_id = ? +LIMIT 2"#; + +const INSERT_DESIRED_SQL: &str = r#"INSERT INTO discovery_desired_state ( + generation_id, normalized_config_sha256, desired_sha256, created_at_unix_ms +) VALUES (?, ?, ?, ?)"#; + +const INSERT_DOCUMENT_SQL: &str = r#"INSERT INTO discovery_documents ( + generation_id, event_id, event_sha256, event_bytes, + nip05_projection_sha256, nip05_projection_bytes +) VALUES (?, ?, ?, ?, ?, ?)"#; + +const INSERT_STATE_SQL: &str = r#"INSERT INTO discovery_publication_state ( + singleton, desired_generation_id, desired_job_id, current_generation_id, + current_job_id, updated_at_unix_ms +) VALUES (1, ?, ?, NULL, NULL, ?)"#; + +const REPLACE_DESIRED_SQL: &str = r#"UPDATE discovery_publication_state +SET desired_generation_id = ?, desired_job_id = ?, updated_at_unix_ms = ? +WHERE singleton = 1 AND updated_at_unix_ms <= ? + AND desired_generation_id = ? AND desired_job_id = ?"#; + +const PROMOTE_CURRENT_SQL: &str = r#"UPDATE discovery_publication_state +SET current_generation_id = desired_generation_id, + current_job_id = desired_job_id, + updated_at_unix_ms = ? +WHERE singleton = 1 AND desired_generation_id = ? AND desired_job_id = ? + AND updated_at_unix_ms <= ? + AND (current_generation_id IS NULL OR current_generation_id != desired_generation_id + OR current_job_id IS NULL OR current_job_id != desired_job_id)"#; + +/// Stable source-free construction and admission failure classes. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycDiscoveryStateErrorKind { + Disabled, + TooLarge, + InvalidProjection, + InvalidEvent, + IdentityMismatch, +} + +impl MycDiscoveryStateErrorKind { + /// Returns the stable machine-readable code. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::Disabled => "discovery_state_disabled", + Self::TooLarge => "discovery_state_too_large", + Self::InvalidProjection => "discovery_projection_invalid", + Self::InvalidEvent => "discovery_event_invalid", + Self::IdentityMismatch => "discovery_identity_mismatch", + } + } +} + +/// Source-free failure at the discovery-state boundary. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct MycDiscoveryStateError { + kind: MycDiscoveryStateErrorKind, +} + +impl MycDiscoveryStateError { + const fn new(kind: MycDiscoveryStateErrorKind) -> Self { + Self { kind } + } + + /// Returns the stable failure class. + #[must_use] + pub const fn kind(self) -> MycDiscoveryStateErrorKind { + self.kind + } + + /// Returns the stable machine-readable code. + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Display for MycDiscoveryStateError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + MycDiscoveryStateErrorKind::Disabled => "discovery state is not enabled", + MycDiscoveryStateErrorKind::TooLarge => "discovery state exceeds its configured bound", + MycDiscoveryStateErrorKind::InvalidProjection => { + "discovery projection inputs are invalid" + } + MycDiscoveryStateErrorKind::InvalidEvent => "discovery event is invalid", + MycDiscoveryStateErrorKind::IdentityMismatch => { + "discovery event identity does not match configuration" + } + }) + } +} + +impl fmt::Debug for MycDiscoveryStateError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDiscoveryStateError") + .field("kind", &self.kind) + .finish() + } +} + +impl Error for MycDiscoveryStateError {} + +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!( + MycDiscoveryGenerationId, + "MycDiscoveryGenerationId([redacted])" +); +redacted_id!( + MycDiscoveryDocumentDigest, + "MycDiscoveryDocumentDigest([redacted])" +); +redacted_id!( + MycNip05ProjectionDigest, + "MycNip05ProjectionDigest([redacted])" +); + +#[derive(Clone, PartialEq, Eq)] +pub(crate) struct MycDiscoveryPolicies { + domain: Box<str>, + handler_identifier: Box<str>, + author_public_key: Box<str>, + public_relays: Box<[Box<str>]>, + nostrconnect_url: Option<Box<str>>, + metadata_json: Box<str>, + event_max_bytes: usize, +} + +impl MycDiscoveryPolicies { + pub(crate) fn from_normalized( + normalized: &serde_json::Value, + identities: &MycExpectedIdentities, + ) -> Result<Option<Self>, MycDiscoveryStateError> { + let discovery = normalized + .pointer("/discovery") + .and_then(serde_json::Value::as_object) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })?; + let enabled = discovery + .get("enabled") + .and_then(serde_json::Value::as_bool) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })?; + if !enabled { + return Ok(None); + } + let author = identities.discovery().ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::IdentityMismatch) + })?; + let string = |name: &str| { + discovery + .get(name) + .and_then(serde_json::Value::as_str) + .filter(|value| !value.is_empty()) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + }) + }; + let public_ids = discovery + .get("public_relay_ids") + .and_then(serde_json::Value::as_array) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })?; + let relay_map = normalized + .pointer("/relays") + .and_then(serde_json::Value::as_array) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })? + .iter() + .map(|relay| { + let id = relay.pointer("/id").and_then(serde_json::Value::as_str)?; + let url = relay.pointer("/url").and_then(serde_json::Value::as_str)?; + Some((id, url)) + }) + .collect::<Option<BTreeMap<_, _>>>() + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })?; + let public_relays = public_ids + .iter() + .map(|id| { + let id = id.as_str()?; + let url = *relay_map.get(id)?; + RadrootsNostrRelayUrl::parse(url).ok()?; + Some(Box::<str>::from(url)) + }) + .collect::<Option<Vec<_>>>() + .filter(|relays| !relays.is_empty()) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })? + .into_boxed_slice(); + let metadata = discovery + .get("metadata") + .and_then(serde_json::Value::as_object) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })?; + let optional = |name: &str| { + metadata + .get(name) + .and_then(serde_json::Value::as_str) + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(ToOwned::to_owned) + }; + let metadata = RadrootsNostrMetadata { + name: optional("name"), + display_name: optional("display_name"), + about: optional("about"), + website: optional("website"), + picture: optional("picture"), + ..RadrootsNostrMetadata::default() + }; + let metadata_json = if metadata.name.is_none() + && metadata.display_name.is_none() + && metadata.about.is_none() + && metadata.website.is_none() + && metadata.picture.is_none() + { + Box::<str>::from("") + } else { + serde_json::to_string(&metadata) + .map_err(|_| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })? + .into_boxed_str() + }; + let template = discovery + .get("nostrconnect_url_template") + .and_then(serde_json::Value::as_str) + .filter(|value| !value.is_empty()); + let signer = identities.user().as_hex(); + let nostrconnect_url = template + .map(|template| render_nostrconnect_url(template, signer, &public_relays)) + .transpose()? + .map(String::into_boxed_str); + let event_max_bytes = normalized + .pointer("/resource_limits/events/wire_bytes") + .and_then(serde_json::Value::as_u64) + .and_then(|value| usize::try_from(value).ok()) + .filter(|value| (1..=MYC_DISCOVERY_DOCUMENT_MAX_BYTES).contains(value)) + .ok_or_else(|| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })?; + Ok(Some(Self { + domain: string("domain")?.into(), + handler_identifier: string("handler_identifier")?.into(), + author_public_key: author.as_hex().into(), + public_relays, + nostrconnect_url, + metadata_json, + event_max_bytes, + })) + } +} + +impl fmt::Debug for MycDiscoveryPolicies { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDiscoveryPolicies") + .field("relay_count", &self.public_relays.len()) + .field("values", &"[redacted]") + .finish() + } +} + +/// Exact signature-verified discovery generation prepared outside SQLite. +pub struct MycDiscoveryCommitRequest { + generation_id: MycDiscoveryGenerationId, + desired_digest: MycDiscoveryDocumentDigest, + configuration_digest: [u8; 32], + event_id: [u8; 32], + event_digest: MycDeliveryArtifactDigest, + event_bytes: Box<[u8]>, + projection_digest: MycNip05ProjectionDigest, + projection_bytes: Box<[u8]>, + created_at: MycDeliveryTimeUnixMs, +} + +impl MycDiscoveryCommitRequest { + /// Validates exact canonical event bytes, signature, identity, NIP-89 semantics, and bounds. + pub fn new( + metadata: &MycStateMetadata, + event_bytes: &[u8], + created_at: MycDeliveryTimeUnixMs, + ) -> Result<Self, MycDiscoveryStateError> { + let policy = metadata + .discovery_policies() + .ok_or_else(|| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::Disabled))?; + if event_bytes.is_empty() || event_bytes.len() > policy.event_max_bytes { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::TooLarge, + )); + } + let event: RadrootsNostrEvent = serde_json::from_slice(event_bytes) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidEvent))?; + let canonical = serde_json::to_vec(&event) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidEvent))?; + if canonical != event_bytes || event.verify().is_err() { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::InvalidEvent, + )); + } + validate_event(&event, policy)?; + let projection_bytes = projection_bytes(policy)?; + if projection_bytes.is_empty() || projection_bytes.len() > MYC_NIP05_PROJECTION_MAX_BYTES { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::TooLarge, + )); + } + let event_digest_bytes: [u8; 32] = Sha256::digest(event_bytes).into(); + let projection_digest_bytes: [u8; 32] = Sha256::digest(&projection_bytes).into(); + let configuration_digest = *metadata.configuration_digest().as_bytes(); + let mut desired_hasher = Sha256::new(); + desired_hasher.update(DESIRED_DIGEST_DOMAIN); + desired_hasher.update(configuration_digest); + desired_hasher.update(event_digest_bytes); + desired_hasher.update(projection_digest_bytes); + let desired_digest_bytes: [u8; 32] = desired_hasher.finalize().into(); + let mut generation_hasher = Sha256::new(); + generation_hasher.update(GENERATION_ID_DOMAIN); + generation_hasher.update(desired_digest_bytes); + let generation_id = MycDiscoveryGenerationId(generation_hasher.finalize().into()); + Ok(Self { + generation_id, + desired_digest: MycDiscoveryDocumentDigest(desired_digest_bytes), + configuration_digest, + event_id: *event.id.as_bytes(), + event_digest: MycDeliveryArtifactDigest::from_bytes(event_digest_bytes), + event_bytes: event_bytes.into(), + projection_digest: MycNip05ProjectionDigest(projection_digest_bytes), + projection_bytes: projection_bytes.into_boxed_slice(), + created_at, + }) + } + + /// Returns the deterministic desired-state generation identity. + #[must_use] + pub const fn generation_id(&self) -> MycDiscoveryGenerationId { + self.generation_id + } + + /// Returns the exact digest binding normalized configuration, event, and projection inputs. + #[must_use] + pub const fn desired_digest(&self) -> MycDiscoveryDocumentDigest { + self.desired_digest + } + + fn owned(&self) -> Self { + Self { + generation_id: self.generation_id, + desired_digest: self.desired_digest, + configuration_digest: self.configuration_digest, + event_id: self.event_id, + event_digest: self.event_digest, + event_bytes: self.event_bytes.clone(), + projection_digest: self.projection_digest, + projection_bytes: self.projection_bytes.clone(), + created_at: self.created_at, + } + } +} + +impl fmt::Debug for MycDiscoveryCommitRequest { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycDiscoveryCommitRequest([redacted])") + } +} + +/// Immutable exact document bytes retained for a discovery generation. +#[derive(Clone, PartialEq, Eq)] +pub struct MycDiscoveryDocumentRecord { + generation_id: MycDiscoveryGenerationId, + event_id: [u8; 32], + event_digest: MycDeliveryArtifactDigest, + event_bytes: Box<[u8]>, + projection_digest: MycNip05ProjectionDigest, + projection_bytes: Box<[u8]>, +} + +impl MycDiscoveryDocumentRecord { + #[must_use] + pub const fn generation_id(&self) -> MycDiscoveryGenerationId { + self.generation_id + } + + /// Returns the sole exact publishable event bytes. + #[must_use] + pub fn event_bytes(&self) -> &[u8] { + &self.event_bytes + } + + /// Returns deterministic NIP-05 projection inputs, not a hosted response. + #[must_use] + pub fn nip05_projection_bytes(&self) -> &[u8] { + &self.projection_bytes + } + + #[must_use] + pub const fn event_digest(&self) -> MycDeliveryArtifactDigest { + self.event_digest + } + + /// Returns the verified NIP-01 event identity. + #[must_use] + pub const fn event_id(&self) -> &[u8; 32] { + &self.event_id + } + + /// Returns the exact digest of the deterministic NIP-05 projection inputs. + #[must_use] + pub const fn nip05_projection_digest(&self) -> MycNip05ProjectionDigest { + self.projection_digest + } +} + +impl fmt::Debug for MycDiscoveryDocumentRecord { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDiscoveryDocumentRecord") + .field("event_bytes", &"[redacted]") + .field("nip05_projection_bytes", &"[redacted]") + .finish() + } +} + +/// Current desired and proven-delivered discovery generations. +#[derive(Clone, PartialEq, Eq)] +pub struct MycDiscoveryPublicationState { + desired_generation_id: MycDiscoveryGenerationId, + desired_job_id: MycDeliveryJobId, + current_generation_id: Option<MycDiscoveryGenerationId>, + current_job_id: Option<MycDeliveryJobId>, + updated_at: MycDeliveryTimeUnixMs, +} + +impl MycDiscoveryPublicationState { + #[must_use] + pub const fn desired_generation_id(&self) -> MycDiscoveryGenerationId { + self.desired_generation_id + } + #[must_use] + pub const fn desired_job_id(&self) -> MycDeliveryJobId { + self.desired_job_id + } + #[must_use] + pub const fn current_generation_id(&self) -> Option<MycDiscoveryGenerationId> { + self.current_generation_id + } + #[must_use] + pub const fn current_job_id(&self) -> Option<MycDeliveryJobId> { + self.current_job_id + } +} + +impl fmt::Debug for MycDiscoveryPublicationState { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDiscoveryPublicationState") + .field("has_current", &self.current_generation_id.is_some()) + .field("updated_at", &self.updated_at) + .finish() + } +} + +/// Atomic result of one discovery desired-state commit. +#[derive(Clone, PartialEq, Eq)] +pub struct MycDiscoveryCommitRecord { + state: MycDiscoveryPublicationState, + document: MycDiscoveryDocumentRecord, + job: MycDeliveryJobRecord, +} + +impl MycDiscoveryCommitRecord { + #[must_use] + pub const fn state(&self) -> &MycDiscoveryPublicationState { + &self.state + } + #[must_use] + pub const fn document(&self) -> &MycDiscoveryDocumentRecord { + &self.document + } + #[must_use] + pub const fn job(&self) -> &MycDeliveryJobRecord { + &self.job + } +} + +impl fmt::Debug for MycDiscoveryCommitRecord { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycDiscoveryCommitRecord([redacted])") + } +} + +/// Created or exact replay outcome of one atomic discovery commit. +#[derive(Clone, PartialEq, Eq)] +pub enum MycDiscoveryCommitAdmission { + Created(MycDiscoveryCommitRecord), + ExactReplay(MycDiscoveryCommitRecord), +} + +impl MycDiscoveryCommitAdmission { + #[must_use] + pub const fn record(&self) -> &MycDiscoveryCommitRecord { + match self { + Self::Created(record) | Self::ExactReplay(record) => record, + } + } +} + +impl fmt::Debug for MycDiscoveryCommitAdmission { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::Created(_) => "MycDiscoveryCommitAdmission::Created([redacted])", + Self::ExactReplay(_) => "MycDiscoveryCommitAdmission::ExactReplay([redacted])", + }) + } +} + +impl MycStateRepository<'_> { + /// Atomically commits desired state, exact documents, targets, and initial delivery evidence. + pub async fn commit_discovery_desired_state( + &self, + request: &MycDiscoveryCommitRequest, + ) -> Result<MycDiscoveryCommitAdmission, MycStateRepositoryError> { + let request = request.owned(); + let expected = PersistedMetadata::from(self.expected()); + let delivery_policy = self.expected().delivery_policies().clone(); + let discovery_policy = self.expected().discovery_policies().cloned(); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let discovery_policy = discovery_policy + .as_ref() + .ok_or(DiscoveryOperationError::Binding)?; + commit_desired(transaction, &request, &delivery_policy, discovery_policy).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Reads the bounded desired/current discovery state, when configured and committed. + pub async fn read_discovery_publication_state( + &self, + ) -> Result<Option<MycDiscoveryPublicationState>, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + read_state(transaction).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Reads exact committed bytes for one discovery delivery job. + pub async fn read_discovery_document_for_job( + &self, + job_id: MycDeliveryJobId, + ) -> Result<Option<MycDiscoveryDocumentRecord>, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + let policy = self.expected().discovery_policies().cloned(); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let generation = discovery_generation_for_job(transaction, job_id).await?; + match (generation, policy.as_ref()) { + (Some(generation), Some(policy)) => { + read_document(transaction, generation, policy).await + } + (Some(_), None) => Err(DiscoveryOperationError::Binding), + (None, _) => Ok(None), + } + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Advances current state only from the current desired job's proven-delivered evidence. + pub async fn promote_delivered_discovery_state( + &self, + job_id: MycDeliveryJobId, + observed_at: MycDeliveryTimeUnixMs, + ) -> Result<MycDiscoveryPublicationState, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + promote_current(transaction, job_id, observed_at).await + }) + }) + .await + .map_err(map_transaction_error) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum DiscoveryOperationError { + Binding, + Storage, +} + +impl From<DeliveryOperationError> for DiscoveryOperationError { + fn from(error: DeliveryOperationError) -> Self { + match error { + DeliveryOperationError::Binding => Self::Binding, + DeliveryOperationError::Storage => Self::Storage, + } + } +} + +async fn verify_metadata( + transaction: &mut ServiceSqliteTransaction<'_>, + expected: &PersistedMetadata, +) -> Result<(), DiscoveryOperationError> { + require_expected_metadata(transaction, expected) + .await + .map_err(|error| match error { + RepositoryOperationError::Binding => DiscoveryOperationError::Binding, + RepositoryOperationError::Storage => DiscoveryOperationError::Storage, + }) +} + +async fn commit_desired( + transaction: &mut ServiceSqliteTransaction<'_>, + request: &MycDiscoveryCommitRequest, + delivery_policy: &crate::state_delivery::MycDeliveryPolicies, + discovery_policy: &MycDiscoveryPolicies, +) -> Result<MycDiscoveryCommitAdmission, DiscoveryOperationError> { + let existing_desired = read_desired(transaction, request.generation_id).await?; + let existing_document = + read_document(transaction, request.generation_id, discovery_policy).await?; + let exact_existing = match (&existing_desired, &existing_document) { + (Some(desired), Some(document)) => { + desired.configuration_digest == request.configuration_digest + && desired.desired_digest == request.desired_digest + && desired.created_at == request.created_at + && exact_document(document, request) + } + (None, None) => false, + _ => return Err(DiscoveryOperationError::Binding), + }; + if existing_desired.is_some() && !exact_existing { + return Err(DiscoveryOperationError::Binding); + } + if existing_desired.is_none() { + insert_desired(transaction, request).await?; + insert_document(transaction, request).await?; + } + let source = MycDeliverySource::discovery_handler(*request.generation_id.as_bytes()); + let job = create_job( + transaction, + source, + request.event_digest, + request.created_at, + delivery_policy, + ) + .await?; + let job_record = job.record().clone(); + let prior = read_state(transaction).await?; + let exact_pointer = prior.as_ref().is_some_and(|state| { + state.desired_generation_id == request.generation_id + && state.desired_job_id == job_record.id() + }); + let created = existing_desired.is_none(); + match prior { + None => { + let result = sqlx::query(INSERT_STATE_SQL) + .bind(request.generation_id.as_bytes().as_slice()) + .bind(job_record.id().as_bytes().as_slice()) + .bind(request.created_at.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + require_one(result.rows_affected())?; + } + Some(state) if exact_pointer => { + if !exact_existing { + return Err(DiscoveryOperationError::Binding); + } + let document = existing_document.ok_or(DiscoveryOperationError::Binding)?; + return Ok(MycDiscoveryCommitAdmission::ExactReplay( + MycDiscoveryCommitRecord { + state, + document, + job: job_record, + }, + )); + } + Some(state) => { + if request.created_at <= state.updated_at { + return Err(DiscoveryOperationError::Binding); + } + let result = sqlx::query(REPLACE_DESIRED_SQL) + .bind(request.generation_id.as_bytes().as_slice()) + .bind(job_record.id().as_bytes().as_slice()) + .bind(request.created_at.sqlite_value()) + .bind(request.created_at.sqlite_value()) + .bind(state.desired_generation_id.as_bytes().as_slice()) + .bind(state.desired_job_id.as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + require_one(result.rows_affected())?; + } + } + let state = read_state(transaction) + .await? + .ok_or(DiscoveryOperationError::Binding)?; + let document = read_document(transaction, request.generation_id, discovery_policy) + .await? + .ok_or(DiscoveryOperationError::Binding)?; + let record = MycDiscoveryCommitRecord { + state, + document, + job: job_record, + }; + if created { + Ok(MycDiscoveryCommitAdmission::Created(record)) + } else { + Err(DiscoveryOperationError::Binding) + } +} + +#[derive(Clone, PartialEq, Eq)] +struct DesiredRecord { + configuration_digest: [u8; 32], + desired_digest: MycDiscoveryDocumentDigest, + created_at: MycDeliveryTimeUnixMs, +} + +async fn insert_desired( + transaction: &mut ServiceSqliteTransaction<'_>, + request: &MycDiscoveryCommitRequest, +) -> Result<(), DiscoveryOperationError> { + let result = sqlx::query(INSERT_DESIRED_SQL) + .bind(request.generation_id.as_bytes().as_slice()) + .bind(request.configuration_digest.as_slice()) + .bind(request.desired_digest.as_bytes().as_slice()) + .bind(request.created_at.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + require_one(result.rows_affected()) +} + +async fn insert_document( + transaction: &mut ServiceSqliteTransaction<'_>, + request: &MycDiscoveryCommitRequest, +) -> Result<(), DiscoveryOperationError> { + let result = sqlx::query(INSERT_DOCUMENT_SQL) + .bind(request.generation_id.as_bytes().as_slice()) + .bind(request.event_id.as_slice()) + .bind(request.event_digest.as_bytes().as_slice()) + .bind(request.event_bytes.as_ref()) + .bind(request.projection_digest.as_bytes().as_slice()) + .bind(request.projection_bytes.as_ref()) + .execute(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + require_one(result.rows_affected()) +} + +async fn read_desired( + transaction: &mut ServiceSqliteTransaction<'_>, + generation: MycDiscoveryGenerationId, +) -> Result<Option<DesiredRecord>, DiscoveryOperationError> { + let rows = sqlx::query(READ_DESIRED_SQL) + .bind(generation.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + if rows.len() > 1 { + return Err(DiscoveryOperationError::Binding); + } + rows.first() + .map(|row| { + if blob32(row, "generation_id")? != *generation.as_bytes() { + return Err(DiscoveryOperationError::Binding); + } + Ok(DesiredRecord { + configuration_digest: blob32(row, "normalized_config_sha256")?, + desired_digest: MycDiscoveryDocumentDigest(blob32(row, "desired_sha256")?), + created_at: time(row, "created_at_unix_ms")?, + }) + }) + .transpose() +} + +async fn read_document( + transaction: &mut ServiceSqliteTransaction<'_>, + generation: MycDiscoveryGenerationId, + policy: &MycDiscoveryPolicies, +) -> Result<Option<MycDiscoveryDocumentRecord>, DiscoveryOperationError> { + let rows = sqlx::query(READ_DOCUMENT_BY_GENERATION_SQL) + .bind(generation.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + if rows.len() > 1 { + return Err(DiscoveryOperationError::Binding); + } + rows.first() + .map(|row| parse_document(row, policy)) + .transpose() +} + +fn parse_document( + row: &sqlx::sqlite::SqliteRow, + policy: &MycDiscoveryPolicies, +) -> Result<MycDiscoveryDocumentRecord, DiscoveryOperationError> { + let generation_id = MycDiscoveryGenerationId(blob32(row, "generation_id")?); + let event_id = blob32(row, "event_id")?; + let event_digest = MycDeliveryArtifactDigest::from_bytes(blob32(row, "event_sha256")?); + let event_bytes = bounded_blob(row, "event_bytes", MYC_DISCOVERY_DOCUMENT_MAX_BYTES)?; + let projection_digest = MycNip05ProjectionDigest(blob32(row, "nip05_projection_sha256")?); + let stored_projection_bytes = bounded_blob( + row, + "nip05_projection_bytes", + MYC_NIP05_PROJECTION_MAX_BYTES, + )?; + let actual_event_digest: [u8; 32] = Sha256::digest(&event_bytes).into(); + let actual_projection_digest: [u8; 32] = Sha256::digest(&stored_projection_bytes).into(); + let event: RadrootsNostrEvent = + serde_json::from_slice(&event_bytes).map_err(|_| DiscoveryOperationError::Binding)?; + let canonical = serde_json::to_vec(&event).map_err(|_| DiscoveryOperationError::Binding)?; + let expected_projection = + projection_bytes(policy).map_err(|_| DiscoveryOperationError::Binding)?; + if actual_event_digest != *event_digest.as_bytes() + || actual_projection_digest != *projection_digest.as_bytes() + || canonical.as_slice() != event_bytes.as_ref() + || event.verify().is_err() + || event.id.as_bytes() != &event_id + || validate_event(&event, policy).is_err() + || expected_projection.as_slice() != stored_projection_bytes.as_ref() + { + return Err(DiscoveryOperationError::Binding); + } + Ok(MycDiscoveryDocumentRecord { + generation_id, + event_id, + event_digest, + event_bytes, + projection_digest, + projection_bytes: stored_projection_bytes, + }) +} + +async fn read_state( + transaction: &mut ServiceSqliteTransaction<'_>, +) -> Result<Option<MycDiscoveryPublicationState>, DiscoveryOperationError> { + let rows = sqlx::query(READ_STATE_SQL) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + if rows.len() > 1 { + return Err(DiscoveryOperationError::Binding); + } + rows.first().map(parse_state).transpose() +} + +fn parse_state( + row: &sqlx::sqlite::SqliteRow, +) -> Result<MycDiscoveryPublicationState, DiscoveryOperationError> { + let current_generation_id = + optional_blob32(row, "current_generation_id", "current_generation_id_type")? + .map(MycDiscoveryGenerationId); + let current_job_id = optional_blob32(row, "current_job_id", "current_job_id_type")? + .map(MycDeliveryJobId::from_persisted); + if current_generation_id.is_some() != current_job_id.is_some() { + return Err(DiscoveryOperationError::Binding); + } + Ok(MycDiscoveryPublicationState { + desired_generation_id: MycDiscoveryGenerationId(blob32(row, "desired_generation_id")?), + desired_job_id: MycDeliveryJobId::from_persisted(blob32(row, "desired_job_id")?), + current_generation_id, + current_job_id, + updated_at: time(row, "updated_at_unix_ms")?, + }) +} + +async fn discovery_generation_for_job( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, +) -> Result<Option<MycDiscoveryGenerationId>, DiscoveryOperationError> { + let rows = sqlx::query( + "SELECT source_id FROM delivery_jobs \ + WHERE job_id = ? AND source_kind = 'discovery_handler' LIMIT 2", + ) + .bind(job_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + if rows.len() > 1 { + return Err(DiscoveryOperationError::Binding); + } + rows.first() + .map(|row| blob32(row, "source_id").map(MycDiscoveryGenerationId)) + .transpose() +} + +async fn promote_current( + transaction: &mut ServiceSqliteTransaction<'_>, + job_id: MycDeliveryJobId, + observed_at: MycDeliveryTimeUnixMs, +) -> Result<MycDiscoveryPublicationState, DiscoveryOperationError> { + let state = read_state(transaction) + .await? + .ok_or(DiscoveryOperationError::Binding)?; + if state.desired_job_id != job_id || observed_at < state.updated_at { + return Err(DiscoveryOperationError::Binding); + } + let rows = sqlx::query(READ_JOB_SQL) + .bind(job_id.as_bytes().as_slice()) + .bind(state.desired_generation_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + if rows.len() != 1 + || rows[0] + .try_get::<&str, _>("status") + .ok() + .and_then(MycDeliveryJobStatus::parse) + != Some(MycDeliveryJobStatus::Delivered) + { + return Err(DiscoveryOperationError::Binding); + } + if state.current_generation_id == Some(state.desired_generation_id) + && state.current_job_id == Some(state.desired_job_id) + { + return Ok(state); + } + let result = sqlx::query(PROMOTE_CURRENT_SQL) + .bind(observed_at.sqlite_value()) + .bind(state.desired_generation_id.as_bytes().as_slice()) + .bind(job_id.as_bytes().as_slice()) + .bind(observed_at.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| DiscoveryOperationError::Storage)?; + require_one(result.rows_affected())?; + read_state(transaction) + .await? + .ok_or(DiscoveryOperationError::Binding) +} + +fn exact_document( + record: &MycDiscoveryDocumentRecord, + request: &MycDiscoveryCommitRequest, +) -> bool { + record.generation_id == request.generation_id + && record.event_id == request.event_id + && record.event_digest == request.event_digest + && record.event_bytes.as_ref() == request.event_bytes.as_ref() + && record.projection_digest == request.projection_digest + && record.projection_bytes.as_ref() == request.projection_bytes.as_ref() +} + +fn validate_event( + event: &RadrootsNostrEvent, + policy: &MycDiscoveryPolicies, +) -> Result<(), MycDiscoveryStateError> { + if event.pubkey.to_hex() != policy.author_public_key.as_ref() { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::IdentityMismatch, + )); + } + if event.kind != RadrootsNostrKind::Custom(NIP89_HANDLER_KIND) { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::InvalidEvent, + )); + } + let mut expected_tags = vec![ + vec!["d".to_owned(), policy.handler_identifier.to_string()], + vec!["k".to_owned(), NIP46_RPC_KIND.to_string()], + ]; + expected_tags.extend( + policy + .public_relays + .iter() + .map(|relay| vec!["relay".to_owned(), relay.to_string()]), + ); + if let Some(url) = policy.nostrconnect_url.as_ref() { + expected_tags.push(vec!["nostrconnect_url".to_owned(), url.to_string()]); + } + let actual_tags = event + .tags + .iter() + .map(|tag| tag.as_slice().to_vec()) + .collect::<Vec<_>>(); + if actual_tags != expected_tags || event.content != policy.metadata_json.as_ref() { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::InvalidEvent, + )); + } + Ok(()) +} + +#[derive(Serialize)] +struct Nip05Projection<'a> { + schema: &'static str, + schema_version: u32, + domain: &'a str, + name: &'static str, + public_key: &'a str, + relays: Vec<&'a str>, + #[serde(skip_serializing_if = "Option::is_none")] + nostrconnect_url: Option<&'a str>, +} + +fn projection_bytes(policy: &MycDiscoveryPolicies) -> Result<Vec<u8>, MycDiscoveryStateError> { + serde_json::to_vec(&Nip05Projection { + schema: "radroots.myc.nip05-projection-input.v1", + schema_version: 1, + domain: &policy.domain, + name: "_", + public_key: &policy.author_public_key, + relays: policy.public_relays.iter().map(AsRef::as_ref).collect(), + nostrconnect_url: policy.nostrconnect_url.as_deref(), + }) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection)) +} + +fn render_nostrconnect_url( + template: &str, + signer_public_key: &str, + public_relays: &[Box<str>], +) -> Result<String, MycDiscoveryStateError> { + let mut serializer = url::form_urlencoded::Serializer::new(String::new()); + for relay in public_relays { + serializer.append_pair("relay", relay); + } + let bunker_uri = format!("bunker://{signer_public_key}?{}", serializer.finish()); + let bunker_uri = radroots_nostr_connect::uri::Uri::parse(&bunker_uri) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))? + .to_string(); + let encoded: String = url::form_urlencoded::byte_serialize(bunker_uri.as_bytes()).collect(); + let rendered = template.replace("<nostrconnect>", &encoded); + nostr::Url::parse(&rendered) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))?; + Ok(rendered) +} + +fn bounded_blob( + row: &sqlx::sqlite::SqliteRow, + column: &str, + maximum: usize, +) -> Result<Box<[u8]>, DiscoveryOperationError> { + let value = row + .try_get::<Option<Vec<u8>>, _>(column) + .map_err(|_| DiscoveryOperationError::Binding)? + .ok_or(DiscoveryOperationError::Binding)?; + if value.is_empty() || value.len() > maximum { + return Err(DiscoveryOperationError::Binding); + } + Ok(value.into_boxed_slice()) +} + +fn blob32( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<[u8; 32], DiscoveryOperationError> { + row.try_get::<Option<Vec<u8>>, _>(column) + .map_err(|_| DiscoveryOperationError::Binding)? + .ok_or(DiscoveryOperationError::Binding)? + .try_into() + .map_err(|_| DiscoveryOperationError::Binding) +} + +fn optional_blob32( + row: &sqlx::sqlite::SqliteRow, + column: &str, + type_column: &str, +) -> Result<Option<[u8; 32]>, DiscoveryOperationError> { + match row + .try_get::<&str, _>(type_column) + .map_err(|_| DiscoveryOperationError::Binding)? + { + "null" => Ok(None), + "blob" => blob32(row, column).map(Some), + _ => Err(DiscoveryOperationError::Binding), + } +} + +fn time( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<MycDeliveryTimeUnixMs, DiscoveryOperationError> { + let value = row + .try_get::<i64, _>(column) + .map_err(|_| DiscoveryOperationError::Binding)?; + u64::try_from(value) + .ok() + .and_then(|value| MycDeliveryTimeUnixMs::new(value).ok()) + .ok_or(DiscoveryOperationError::Binding) +} + +fn require_one(rows: u64) -> Result<(), DiscoveryOperationError> { + (rows == 1) + .then_some(()) + .ok_or(DiscoveryOperationError::Storage) +} + +fn map_transaction_error( + error: ServiceSqliteTransactionError<DiscoveryOperationError>, +) -> MycStateRepositoryError { + if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { + return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); + } + let kind = match error.operation_error() { + Some(DiscoveryOperationError::Binding) => MycStateRepositoryErrorKind::Binding, + Some(DiscoveryOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, + }; + MycStateRepositoryError::new(kind) +} diff --git a/src/state_host.rs b/src/state_host.rs @@ -377,18 +377,19 @@ 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() == 5 + && outcome.applied_count() == 6 } 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, 5) - | (2, 4) - | (3, 3) - | (4, 2) - | (5, 1) + (MYC_STATE_BASE_SCHEMA_VERSION, 6) + | (2, 5) + | (3, 4) + | (4, 3) + | (5, 2) + | (6, 1) | (MYC_STATE_SCHEMA_VERSION, 0) ) } diff --git a/src/state_metadata.rs b/src/state_metadata.rs @@ -12,6 +12,7 @@ use radroots_storage::event::SourceGeneration; use sha2::{Digest, Sha256}; use crate::state_delivery::{MycDeliveryPolicies, MycDeliveryPolicyMode, MycDeliveryRelayId}; +use crate::state_discovery::MycDiscoveryPolicies; use crate::state_governance::{ MycGovernancePolicies, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, }; @@ -173,6 +174,7 @@ pub struct MycStateMetadata { identities: MycExpectedIdentities, governance: MycGovernancePolicies, delivery: MycDeliveryPolicies, + discovery: Option<MycDiscoveryPolicies>, policy_versions: MycStatePolicyVersions, } @@ -211,9 +213,11 @@ impl MycStateMetadata { ); let normalized = configuration.normalized(); let governance = governance_policies(normalized)?; + let identities = expected_identities(normalized)?; let delivery = delivery_policies(normalized)?; + let discovery = MycDiscoveryPolicies::from_normalized(normalized, &identities) + .map_err(|_| MycStateMetadataError::new(MycStateMetadataErrorKind::Invariant))?; let configuration = normalized_config_digest(configuration.profile(), normalized)?; - let identities = expected_identities(normalized)?; let policy_versions = MycStatePolicyVersions::governed(); if [ policy_versions.configuration, @@ -235,6 +239,7 @@ impl MycStateMetadata { identities, governance, delivery, + discovery, policy_versions, }) } @@ -288,6 +293,10 @@ impl MycStateMetadata { &self.delivery } + pub(crate) const fn discovery_policies(&self) -> Option<&MycDiscoveryPolicies> { + self.discovery.as_ref() + } + pub(crate) fn matches_runtime(&self, runtime: &MycRuntimeContext) -> bool { ServiceSqlitePaths::from_runtime_context(runtime.context()) .is_ok_and(|paths| paths == self.paths) @@ -304,6 +313,7 @@ impl fmt::Debug for MycStateMetadata { .field("identities", &self.identities) .field("governance", &"[redacted]") .field("delivery", &"[redacted]") + .field("discovery", &self.discovery.as_ref().map(|_| "[redacted]")) .field("policy_versions", &self.policy_versions) .field("paths", &"[redacted]") .finish() diff --git a/tests/services_hardening_delivery_state.rs b/tests/services_hardening_delivery_state.rs @@ -484,11 +484,11 @@ async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_e .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)", + "DELETE FROM delivery_attempts", + "DELETE FROM delivery_targets", + "DELETE FROM delivery_jobs", + "UPDATE delivery_targets SET relay_id = 'changed'", + "UPDATE delivery_jobs SET artifact_sha256 = zeroblob(32)", ] { assert!( sqlx::query(statement) @@ -499,7 +499,7 @@ async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_e ); } assert_eq!( - sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM publication_attempts") + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_attempts") .fetch_one(&mut connection) .await .expect("attempt count"), @@ -513,11 +513,11 @@ fn delivery_boundary_is_typed_sqlx_only_and_has_no_external_wait_or_legacy_autho 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("delivery_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")); + assert!(CATALOG_SOURCE.contains("delivery_jobs_no_delete")); + assert!(CATALOG_SOURCE.contains("delivery_targets_no_delete")); + assert!(CATALOG_SOURCE.contains("delivery_attempts_no_delete")); for forbidden in [ "SqlitePool", "SqliteConnection", diff --git a/tests/services_hardening_discovery_state.rs b/tests/services_hardening_discovery_state.rs @@ -0,0 +1,522 @@ +#![forbid(unsafe_code)] +#![cfg(any(target_os = "linux", target_os = "macos"))] + +use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path}; + +use myc::host_identity::RadrootsIdentity; +use myc::nostr_contract::{ + RadrootsNostrApplicationHandlerSpec, RadrootsNostrMetadata, RadrootsNostrTimestamp, + radroots_nostr_build_application_handler_event, +}; +use myc::{ + MYC_DISCOVERY_DOCUMENT_MAX_BYTES, MYC_STATE_SCHEMA_VERSION, MycConfigProfile, + MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, MycDeliveryClaim, MycDeliveryJobStatus, + MycDeliveryRelayId, MycDeliverySourceKind, MycDeliveryTimeUnixMs, MycDiscoveryCommitAdmission, + MycDiscoveryCommitRequest, MycDiscoveryStateErrorKind, 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 DISCOVERY_SOURCE: &str = include_str!("../src/state_discovery.rs"); +const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); +const LIB_SOURCE: &str = include_str!("../src/lib.rs"); +const DISCOVERY_SECRET: &str = "3333333333333333333333333333333333333333333333333333333333333333"; + +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 discovery_identity() -> RadrootsIdentity { + RadrootsIdentity::from_secret_key_str(DISCOVERY_SECRET).expect("discovery identity") +} + +fn config_source(enabled: bool) -> Vec<u8> { + let identity = discovery_identity(); + let source = String::from_utf8(CONFIG_EXAMPLE.to_vec()) + .expect("UTF-8 configuration") + .replace( + "3333333333333333333333333333333333333333333333333333333333333333", + &identity.public_key_hex(), + ); + if enabled { + return source.into_bytes(); + } + let before_binding = source + .split_once("[identity.discovery.binding]") + .expect("discovery binding") + .0; + let after_binding = source.split_once("[[relays]]").expect("relay inventory").1; + let without_binding = format!( + "{}[[relays]]{}", + before_binding.replace( + "[identity.discovery]\nenabled = true", + "[identity.discovery]\nenabled = false" + ), + after_binding, + ); + let before_discovery = without_binding + .split_once("[discovery]") + .expect("discovery section") + .0; + format!("{before_discovery}[discovery]\nenabled = false\n").into_bytes() +} + +fn state_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 time(value: u64) -> MycDeliveryTimeUnixMs { + MycDeliveryTimeUnixMs::new(value).expect("delivery time") +} + +fn nostrconnect_url() -> String { + let mut query = url::form_urlencoded::Serializer::new(String::new()); + query.append_pair("relay", "wss://relay-primary.example.test/"); + query.append_pair("relay", "wss://relay-secondary.example.test/"); + let bunker = format!( + "bunker://{}?{}", + "2222222222222222222222222222222222222222222222222222222222222222", + query.finish() + ); + let encoded: String = url::form_urlencoded::byte_serialize(bunker.as_bytes()).collect(); + format!("https://myc.example.test/connect?uri={encoded}") +} + +fn signed_handler_event(created_at: u64) -> Vec<u8> { + let metadata = RadrootsNostrMetadata { + name: Some("myc".to_owned()), + display_name: Some("Radroots Myc".to_owned()), + about: Some("NIP-46 signer".to_owned()), + website: Some("https://myc.example.test/".to_owned()), + picture: Some("https://myc.example.test/myc.png".to_owned()), + ..RadrootsNostrMetadata::default() + }; + let spec = RadrootsNostrApplicationHandlerSpec::new(vec![24_133]) + .with_identifier("myc") + .with_relays(vec![ + "wss://relay-primary.example.test/".to_owned(), + "wss://relay-secondary.example.test/".to_owned(), + ]) + .with_nostr_connect_url(nostrconnect_url()) + .with_metadata(metadata); + let event = radroots_nostr_build_application_handler_event(&spec) + .expect("typed handler event") + .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at)) + .sign_with_keys(discovery_identity().keys()) + .expect("signed event"); + serde_json::to_vec(&event).expect("canonical event bytes") +} + +async fn deliver_all_required( + repository: &myc::MycStateRepository<'_>, + job_id: myc::MycDeliveryJobId, + start: u64, +) { + for (offset, relay_name) in ["primary", "secondary"].into_iter().enumerate() { + let relay = MycDeliveryRelayId::new(relay_name).expect("relay"); + let claimed = match repository + .claim_delivery_target( + job_id, + &relay, + MycDeliveryAttemptNonce::from_injected_entropy([0x70 + offset as u8; 32]), + time(start + offset as u64 * 10), + ) + .await + .expect("claim") + { + MycDeliveryClaim::Claimed(attempt) => attempt, + other => panic!("unexpected claim: {other:?}"), + }; + repository + .mark_delivery_attempt_submitted( + job_id, + &relay, + claimed.id(), + time(start + offset as u64 * 10 + 1), + ) + .await + .expect("submitted"); + repository + .record_delivery_attempt_outcome( + job_id, + &relay, + claimed.id(), + MycDeliveryAttemptOutcome::Delivered, + time(start + offset as u64 * 10 + 2), + ) + .await + .expect("delivered"); + } +} + +#[test] +fn discovery_input_is_signature_verified_bounded_and_redacted() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + let metadata = state_metadata(&runtime, &config_source(true)); + let event = signed_handler_event(1_725_000_000); + let request = MycDiscoveryCommitRequest::new(&metadata, &event, time(100)).expect("request"); + assert!( + !request + .generation_id() + .as_bytes() + .iter() + .all(|byte| *byte == 0) + ); + + let mut altered = event.clone(); + let position = altered + .iter() + .position(|byte| *byte == b'm') + .expect("mutable event byte"); + altered[position] = b'n'; + assert_eq!( + MycDiscoveryCommitRequest::new(&metadata, &altered, time(100)) + .expect_err("signature/canonical drift") + .kind(), + MycDiscoveryStateErrorKind::InvalidEvent + ); + let mut whitespace = event.clone(); + whitespace.push(b'\n'); + assert_eq!( + MycDiscoveryCommitRequest::new(&metadata, &whitespace, time(100)) + .expect_err("noncanonical event") + .kind(), + MycDiscoveryStateErrorKind::InvalidEvent + ); + assert_eq!( + MycDiscoveryCommitRequest::new( + &metadata, + &vec![b'x'; MYC_DISCOVERY_DOCUMENT_MAX_BYTES + 1], + time(100), + ) + .expect_err("oversize") + .kind(), + MycDiscoveryStateErrorKind::TooLarge + ); + let disabled = state_metadata(&runtime, &config_source(false)); + let error = MycDiscoveryCommitRequest::new(&disabled, &event, time(100)) + .expect_err("disabled discovery"); + assert_eq!(error.kind(), MycDiscoveryStateErrorKind::Disabled); + assert!(Error::source(&error).is_none()); + let rendered = format!("{request:?} {error} {error:?}"); + assert!(!rendered.contains(DISCOVERY_SECRET)); + assert!(!rendered.contains("myc.example.test")); + assert!(!rendered.contains(&String::from_utf8(event).expect("event text"))); +} + +#[tokio::test] +async fn discovery_desired_current_and_exact_documents_are_atomic_restart_safe_and_delivery_bound() +{ + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + prepare_state_directory(&runtime); + let metadata = state_metadata(&runtime, &config_source(true)); + 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 event = signed_handler_event(1_725_000_000); + let request = MycDiscoveryCommitRequest::new(&metadata, &event, time(100)).expect("request"); + let committed = host + .repository() + .commit_discovery_desired_state(&request) + .await + .expect("commit"); + assert!(matches!(committed, MycDiscoveryCommitAdmission::Created(_))); + let record = committed.record(); + assert_eq!( + record.job().source_kind(), + MycDeliverySourceKind::DiscoveryHandler + ); + assert_eq!(record.job().operation_id(), None); + assert_eq!(record.job().status(), MycDeliveryJobStatus::Pending); + assert!( + !request + .desired_digest() + .as_bytes() + .iter() + .all(|byte| *byte == 0) + ); + assert_eq!(record.document().event_id().len(), 32); + assert!( + !record + .document() + .nip05_projection_digest() + .as_bytes() + .iter() + .all(|byte| *byte == 0) + ); + assert_eq!(record.document().event_bytes(), event); + let projection: serde_json::Value = + serde_json::from_slice(record.document().nip05_projection_bytes()).expect("projection"); + assert_eq!( + projection["schema"], + "radroots.myc.nip05-projection-input.v1" + ); + assert_eq!(projection["domain"], "myc.example.test"); + assert_eq!(projection["name"], "_"); + assert_eq!(projection["relays"].as_array().expect("relays").len(), 2); + assert_eq!(record.state().current_generation_id(), None); + assert_eq!(record.state().current_job_id(), None); + let job_id = record.job().id(); + + let replay = host + .repository() + .commit_discovery_desired_state(&request) + .await + .expect("exact replay"); + assert!(matches!( + replay, + MycDiscoveryCommitAdmission::ExactReplay(_) + )); + assert_eq!(replay.record().job().id(), job_id); + let same_time_event = signed_handler_event(1_725_000_002); + let same_time_request = MycDiscoveryCommitRequest::new(&metadata, &same_time_event, time(100)) + .expect("same-time request"); + assert_eq!( + host.repository() + .commit_discovery_desired_state(&same_time_request) + .await + .expect_err("different desired state at the same time") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + assert_eq!( + host.repository() + .promote_delivered_discovery_state(job_id, time(101)) + .await + .expect_err("pending job cannot become current") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + deliver_all_required(&host.repository(), job_id, 110).await; + let promoted = host + .repository() + .promote_delivered_discovery_state(job_id, time(140)) + .await + .expect("promoted current"); + assert_eq!( + promoted.current_generation_id(), + Some(request.generation_id()) + ); + assert_eq!(promoted.current_job_id(), Some(job_id)); + assert_eq!( + host.repository() + .read_discovery_document_for_job(job_id) + .await + .expect("document read") + .expect("document") + .event_bytes(), + event + ); + let next_event = signed_handler_event(1_725_000_001); + let next_request = + MycDiscoveryCommitRequest::new(&metadata, &next_event, time(150)).expect("next request"); + let next = host + .repository() + .commit_discovery_desired_state(&next_request) + .await + .expect("next desired state"); + assert!(matches!(next, MycDiscoveryCommitAdmission::Created(_))); + assert_eq!( + next.record().state().desired_generation_id(), + next_request.generation_id() + ); + assert_eq!( + next.record().state().current_generation_id(), + Some(request.generation_id()) + ); + assert_eq!(next.record().state().current_job_id(), Some(job_id)); + let next_job_id = next.record().job().id(); + assert_eq!( + host.repository() + .promote_delivered_discovery_state(next_job_id, time(151)) + .await + .expect_err("new desired state is not yet current") + .kind(), + MycStateRepositoryErrorKind::Binding + ); + deliver_all_required(&host.repository(), next_job_id, 160).await; + let next_promoted = host + .repository() + .promote_delivered_discovery_state(next_job_id, time(190)) + .await + .expect("next current state"); + assert_eq!( + next_promoted.current_generation_id(), + Some(next_request.generation_id()) + ); + assert_eq!(next_promoted.current_job_id(), Some(next_job_id)); + host.close().await.expect("close"); + + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("reopen"); + let state = host + .repository() + .read_discovery_publication_state() + .await + .expect("state") + .expect("committed state"); + assert_eq!(state.current_job_id(), Some(next_job_id)); + host.close().await.expect("final close"); +} + +#[tokio::test] +async fn discovery_schema_guards_reject_mutation_and_corrupt_source_kinds() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = runtime(directory.path()); + prepare_state_directory(&runtime); + let metadata = state_metadata(&runtime, &config_source(true)); + 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 request = + MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_000), time(100)) + .expect("request"); + host.repository() + .commit_discovery_desired_state(&request) + .await + .expect("commit"); + 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 discovery_desired_state", + "UPDATE discovery_documents SET event_bytes = X'01'", + "DELETE FROM discovery_publication_state", + "UPDATE discovery_publication_state SET desired_job_id = zeroblob(32)", + "UPDATE delivery_jobs SET source_kind = 'invalid'", + ] { + assert!( + sqlx::query(statement) + .execute(&mut connection) + .await + .is_err(), + "guard accepted {statement}" + ); + } + for (index, source_kind) in ["signer_response", "discovery_handler"] + .into_iter() + .enumerate() + { + let job_id = [0xa0 + u8::try_from(index).expect("bounded index"); 32]; + let source_id = [0xf0 + u8::try_from(index).expect("bounded index"); 32]; + let artifact = [0xb0 + u8::try_from(index).expect("bounded index"); 32]; + let result = sqlx::query( + "INSERT INTO delivery_jobs (job_id, source_kind, source_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 (?, ?, ?, ?, 'all_required', 1, 1, 1, 1, 1, 'pending', 1, 1, NULL)", + ) + .bind(job_id.as_slice()) + .bind(source_kind) + .bind(source_id.as_slice()) + .bind(artifact.as_slice()) + .execute(&mut connection) + .await; + assert!(result.is_err(), "accepted orphaned {source_kind} source"); + } + connection.close().await.expect("connection close"); +} + +#[test] +fn discovery_boundary_is_typed_sqlx_only_and_commits_before_any_publication() { + assert!(LIB_SOURCE.contains("mod state_discovery;")); + assert!(!LIB_SOURCE.contains("pub mod state_discovery;")); + assert!(DISCOVERY_SOURCE.contains("event.verify()")); + assert!(DISCOVERY_SOURCE.contains("ServiceSqliteTransaction<'_>")); + assert!(DISCOVERY_SOURCE.contains("source_kind = 'discovery_handler'")); + assert!(CATALOG_SOURCE.contains("discovery_desired_state_no_update")); + assert!(CATALOG_SOURCE.contains("discovery_documents_no_delete")); + assert!(CATALOG_SOURCE.contains("delivery_jobs_guard_insert")); + for forbidden in [ + "SqlitePool", + "SqliteConnection", + "BEGIN ", + "COMMIT", + "ROLLBACK", + "reqwest", + "send_event", + "publish_nip89_event", + "tokio::spawn", + "std::env", + "std::time", + "rusqlite", + ] { + assert!( + !DISCOVERY_SOURCE.contains(forbidden), + "found forbidden discovery authority `{forbidden}`" + ); + } +} diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs @@ -13,8 +13,9 @@ use myc::{ MYC_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_5_OBJECT_COUNT, 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, + MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT, + MYC_STATE_SCHEMA_VERSION_7_SHA256, MycStateCatalogErrorKind, myc_migration_catalog, + myc_schema_catalog, validate_myc_state_catalogs, }; use radroots_service_sqlite::{ MigrationCatalog, MigrationChecksum, MigrationDescriptor, SchemaCatalog, SchemaDigest, @@ -26,13 +27,13 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); #[test] -fn schema_v1_through_v6_and_all_migrations_have_exact_literal_identities() { +fn schema_v1_through_v7_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, 6); - assert_eq!(migrations.descriptors().len(), 5); + assert_eq!(MYC_STATE_SCHEMA_VERSION, 7); + assert_eq!(migrations.descriptors().len(), 6); let metadata = &migrations.descriptors()[0]; assert_eq!(metadata.target_version(), 2); assert_eq!(metadata.name().as_str(), "create_myc_state_metadata"); @@ -74,13 +75,20 @@ fn schema_v1_through_v6_and_all_migrations_have_exact_literal_identities() { delivery.checksum().as_bytes(), &MYC_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256 ); - assert_eq!(migrations.current_version(), 6); + let discovery = &migrations.descriptors()[5]; + assert_eq!(discovery.target_version(), 7); + assert_eq!(discovery.name().as_str(), "create_discovery_desired_state"); + assert_eq!( + discovery.checksum().as_bytes(), + &MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256 + ); + assert_eq!(migrations.current_version(), 7); assert_eq!( migrations.digest().as_bytes(), &MYC_MIGRATION_CATALOG_SHA256 ); - assert_eq!(schema.versions().len(), 6); + assert_eq!(schema.versions().len(), 7); assert_eq!(schema.versions()[0].version(), 1); assert_eq!( schema.versions()[0].object_count(), @@ -140,6 +148,16 @@ fn schema_v1_through_v6_and_all_migrations_have_exact_literal_identities() { schema.versions()[5].digest().as_bytes(), &MYC_STATE_SCHEMA_VERSION_6_SHA256 ); + assert_eq!(schema.versions()[6].version(), 7); + assert_eq!( + schema.versions()[6].object_count(), + MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT + ); + assert_eq!(schema.versions()[6].object_count(), 53); + assert_eq!( + schema.versions()[6].digest().as_bytes(), + &MYC_STATE_SCHEMA_VERSION_7_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"); @@ -150,7 +168,7 @@ fn schema_v1_through_v6_and_all_migrations_have_exact_literal_identities() { ); assert_eq!( hex::encode(MYC_MIGRATION_CATALOG_SHA256), - "a69bbc7f3c22750ddda55f36db125622c49f26b44cdc47210fadb7b5fc4f31e5" + "6f47d1a6614293b8d8d3102ab587c42e2bde40b249a165e2504ce4a6f5081677" ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_VERSION_1_SHA256), @@ -162,7 +180,7 @@ fn schema_v1_through_v6_and_all_migrations_have_exact_literal_identities() { ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_CATALOG_SHA256), - "94f7adfbd62be03e6f3af9c9754bf2f47cc704a87b2708f6826c4432556b6f06" + "4f8d6ee98759ccad9842b1beb6a1dc96c21649008ed647ff7311cf2206c1f620" ); assert_eq!( hex::encode(MYC_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256), @@ -196,6 +214,14 @@ fn schema_v1_through_v6_and_all_migrations_have_exact_literal_identities() { hex::encode(MYC_STATE_SCHEMA_VERSION_6_SHA256), "557273306e68e9306dc1d7c7009e1cb7c9d15b8d52e07db7d0bd51e83f054492" ); + assert_eq!( + hex::encode(MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256), + "6097c40776a57dd4bddc04e652987217d6f5991f2d5e94d322d1208619254e6c" + ); + assert_eq!( + hex::encode(MYC_STATE_SCHEMA_VERSION_7_SHA256), + "3b3528911a293499d9721da71b6ad07a5e723e97b6ebf5b40a19b9058958bf79" + ); } #[test] @@ -245,8 +271,11 @@ fn independent_validator_rejects_migration_or_schema_drift() { 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]) + let v6 = SchemaVersionCatalog::new(6, [object.clone()], v6_digest).expect("schema v6"); + let v7_digest = + SchemaVersionCatalog::computed_digest(7, [object.clone()]).expect("schema-v7 digest"); + let v7 = SchemaVersionCatalog::new(7, [object], v7_digest).expect("schema v7"); + let schema = SchemaCatalog::new(&expected_migrations, [v1, v2, v3, v4, v5, v6, v7]) .expect("drift schema catalog"); assert_eq!( validate_myc_state_catalogs(&expected_migrations, &schema) @@ -284,7 +313,17 @@ fn catalog_source_is_pure_pinned_and_uses_only_the_shared_authority() { assert!(CATALOG_SOURCE.contains("MigrationDescriptor::sql(")); assert!(CATALOG_SOURCE.contains("MigrationChecksum::from_bytes(")); assert!(CATALOG_SOURCE.contains("SchemaDigest::from_bytes(")); + assert!(CATALOG_SOURCE.contains("FROM publication_outbox;")); + assert!( + CATALOG_SOURCE.contains("INSERT INTO delivery_targets SELECT * FROM publication_targets;") + ); + assert!( + CATALOG_SOURCE + .contains("INSERT INTO delivery_attempts SELECT * FROM publication_attempts;") + ); + assert!(CATALOG_SOURCE.contains("delivery_jobs_guard_insert")); assert!(!CATALOG_SOURCE.contains("computed_digest")); + assert!(!CATALOG_SOURCE.contains("ALTER TABLE")); for forbidden in [ "sqlx::", "rusqlite",