lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 086041a2ad73228468954103bf2d307f7dd9c945
parent 8f9cd897c4d3c5395f2e569372950ad82694264d
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 04:58:54 +0000

outbox: persist typed Phase 1 publications

- add authenticated schema-v2 typed publication storage
- fence signing and target transitions with revisioned leases
- bind operation, dispatch, receipt, and repair identities
- cover idempotency, corruption, races, rollback, and reopen

Diffstat:
MCHANGELOG.md | 9+++++++++
MCargo.lock | 3+++
Acontracts/conformance/vectors/outbox/phase1_publication.v1.json | 34++++++++++++++++++++++++++++++++++
Mcontracts/outbox_feature_matrix.toml | 9++++++++-
Mcontracts/releases/1.0.0-alpha.1.toml | 13+++++++++++++
Mcrates/outbox/Cargo.toml | 15++++++++++++++-
Mcrates/outbox/contracts/migration_authority_v1.manifest.json | 92++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
Mcrates/outbox/contracts/migration_authority_v1.manifest.sha256 | 2+-
Mcrates/outbox/contracts/migration_registry.v1.json | 28++++++++++++++++++++++++++++
Acrates/outbox/contracts/phase1_publication_v1.descriptor.json | 115+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/outbox/contracts/phase1_publication_v1.manifest.json | 152+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/outbox/contracts/phase1_publication_v1.manifest.schema.json | 177+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/outbox/contracts/phase1_publication_v1.manifest.sha256 | 1+
Acrates/outbox/migrations/0002_phase1_publication.down.sql | 5+++++
Acrates/outbox/migrations/0002_phase1_publication.up.sql | 95+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/outbox/src/generated/outbox_migration_registry.rs | 32+++++++++++++++++++++++++++++++-
Mcrates/outbox/src/lib.rs | 15+++++++++++++++
Mcrates/outbox/src/migrations.rs | 5+++--
Acrates/outbox/src/phase1_publication.rs | 2814+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/outbox/src/schema.rs | 83++++++++++++++++++++++++++++++++++++++++++++++---------------------------------
Mcrates/outbox/src/store.rs | 2+-
Acrates/outbox/tests/fixtures/phase1_publication.v1.json | 34++++++++++++++++++++++++++++++++++
Acrates/outbox/tests/phase1_publication_v1_result_vector.rs | 251+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtools/xtask/src/contract.rs | 16++++++++++++++--
Mtools/xtask/src/contract/outbox_migration.rs | 24+++++++++++++++++++-----
Atools/xtask/src/contract/outbox_phase1_publication.rs | 638+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtools/xtask/src/main.rs | 19+++++++++++++++++++
27 files changed, 4607 insertions(+), 76 deletions(-)

diff --git a/CHANGELOG.md b/CHANGELOG.md @@ -30,6 +30,15 @@ publish policy both pass for the same source revision. authenticated schema status and owned open-time migration. A machine-readable matrix now governs no-default, SQLite, Tokio, event-store-adapter, and all-feature builds. +<!-- release-change: outbox-phase1-publication-state --> +- Outbox schema version `2` adds an isolated typed Phase 1 publication state + machine. Enqueue accepts only sealed allowlisted artifacts with complete + media-readiness bindings, persists exact canonical bytes and raw fixed-width + operation identities, and uses bounded canonical relay targets. Revision-CAS + transitions and opaque expiring claims fence signing and target workers; + signed-event bytes are immutable, while dispatch intent, uncertain results, + receipts, and accepted-observation repair identities remain durable without + performing network publication. - Event-store schema initialization now uses a transactional, checksummed migration authority with exact legacy-baseline adoption, shared-database catalog scoping, tamper-evident fail-closed managed history, exact catalog diff --git a/Cargo.lock b/Cargo.lock @@ -4888,8 +4888,10 @@ dependencies = [ name = "radroots_outbox" version = "1.0.0-alpha.1" dependencies = [ + "getrandom 0.2.17", "hex", "radroots_event", + "radroots_event_codec", "radroots_event_store", "radroots_nostr", "radroots_transport", @@ -4900,6 +4902,7 @@ dependencies = [ "tempfile", "thiserror 1.0.69", "tokio", + "url", ] [[package]] diff --git a/contracts/conformance/vectors/outbox/phase1_publication.v1.json b/contracts/conformance/vectors/outbox/phase1_publication.v1.json @@ -0,0 +1,34 @@ +{ + "schema_version": 1, + "contract_id": "radroots_outbox.phase1_publication.v1", + "executor": { + "id": "radroots_outbox.phase1_publication.v1.result_vector_executor.v1", + "path": "crates/outbox/tests/phase1_publication_v1_result_vector.rs", + "test": "phase1_publication_v1_result_vector" + }, + "identity_vector": { + "fixture": "update", + "target_uri": "wss://relay.example/", + "required_target_count": 1, + "artifact_digest": "9ad318496bd4a710fbc7f3e0f6d5a01de808d352db91ca56a12e710904784cc5", + "readiness_digest": "bf9c1ff2a2c26d62d45b9f5a1eea727ce8f8d7ee92e749364900c12ff1546be5", + "endpoint_fingerprint": "59bf3d289a2344bdc61c426484877df2f122d3936115f45e68e812d2189f21d2", + "target_policy_digest": "e7577fbe841b140993f9cf317e1f562d926b64eb19625220a7c0bfaa8bce4fbb", + "operation_digest": "01a16b215292e90af4c7a76556a86b54d7910ed01757acc8429bda5f1c5fb72b", + "dispatch_digest": "0e35ef0a2769586300ecb3b4760e6a034a60ee68344126dc875168e3876618a4" + }, + "cases": [ + { "id": "identity_preimages", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "typed_enqueue", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "duplicate_enqueue", "execution": "direct_executor", "expected_outcome": "accepted_idempotent" }, + { "id": "target_count_exact", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "target_count_one_over", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_target_count" }, + { "id": "target_uri_exact", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "target_uri_one_over", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_target_uri_too_large" }, + { "id": "empty_required_policy", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_required_target_count" }, + { "id": "two_worker_claim_race", "execution": "direct_executor", "expected_outcome": "one_winner" }, + { "id": "expired_lease_reclaim", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "stale_claim_rejected", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_claim_invalid" }, + { "id": "migration_rollback_reopen", "execution": "direct_executor", "expected_outcome": "accepted" } + ] +} diff --git a/contracts/outbox_feature_matrix.toml b/contracts/outbox_feature_matrix.toml @@ -3,7 +3,14 @@ package = "radroots_outbox" [feature_edges] default = ["event-store-adapter"] -sqlite = ["dep:sqlx", "sqlx/sqlite-bundled"] +sqlite = [ + "dep:getrandom", + "dep:radroots_event_codec", + "dep:sqlx", + "dep:url", + "radroots_event/signature", + "sqlx/sqlite-bundled", +] runtime-tokio = ["sqlite", "sqlx/runtime-tokio", "dep:tokio"] event-store-adapter = [ "runtime-tokio", diff --git a/contracts/releases/1.0.0-alpha.1.toml b/contracts/releases/1.0.0-alpha.1.toml @@ -560,3 +560,16 @@ semver_impacts = [ "change_exported_algorithm_behavior", ] summary = "Replace raw outbox migration SQL and live migrate-down with an append-only generated registry, immutable checksums, exact unledgered-0001 adoption, a tamper-evident ledger and catalog fingerprint, registry-derived version bounds, governed feature profiles, and an executable conformance contract while freezing the 0001 migration bytes." + +[[changes]] +id = "outbox-phase1-publication-state" +classification = "breaking" +semver_impacts = [ + "add_exported_type", + "add_exported_function", + "add_exported_constant", + "add_enum_variant", + "add_conformance_vector", + "change_exported_algorithm_behavior", +] +summary = "Add authenticated schema version 2 and a sealed Phase 1 publication state machine that persists exact artifact, readiness, signed-event, target-policy, dispatch-intent, receipt, and observation-repair authority under bounded canonical targets, stable raw fixed-width identities, revision compare-and-swap, and opaque expiring claims without performing network publication." diff --git a/crates/outbox/Cargo.toml b/crates/outbox/Cargo.toml @@ -13,7 +13,14 @@ readme = "README" [features] default = ["event-store-adapter"] -sqlite = ["dep:sqlx", "sqlx/sqlite-bundled"] +sqlite = [ + "dep:getrandom", + "dep:radroots_event_codec", + "dep:sqlx", + "dep:url", + "radroots_event/signature", + "sqlx/sqlite-bundled", +] runtime-tokio = ["sqlite", "sqlx/runtime-tokio", "dep:tokio"] event-store-adapter = [ "runtime-tokio", @@ -28,7 +35,12 @@ radroots_event = { workspace = true, default-features = false, features = [ "serde", ] } radroots_event_store = { workspace = true, default-features = false, optional = true } +radroots_event_codec = { workspace = true, default-features = false, features = [ + "serde_json", + "std", +], optional = true } radroots_transport = { workspace = true, default-features = false } +getrandom = { workspace = true, optional = true, features = ["std"] } hex = { workspace = true } serde = { workspace = true, features = ["std"] } serde_json = { workspace = true, features = ["std"] } @@ -36,6 +48,7 @@ sha2 = { workspace = true } sqlx = { workspace = true, optional = true, features = ["derive"] } thiserror = { workspace = true } tokio = { workspace = true, optional = true, features = ["time"] } +url = { workspace = true, optional = true } [dev-dependencies] radroots_nostr = { workspace = true, default-features = false, features = [ diff --git a/crates/outbox/contracts/migration_authority_v1.manifest.json b/crates/outbox/contracts/migration_authority_v1.manifest.json @@ -28,7 +28,11 @@ }, { "enables": [ + "dep:getrandom", + "dep:radroots_event_codec", "dep:sqlx", + "dep:url", + "radroots_event/signature", "sqlx/sqlite-bundled" ], "feature": "sqlite" @@ -79,17 +83,17 @@ } ], "source": { - "byte_length": 902, + "byte_length": 1001, "hash_algorithm": "sha256_bytes_v1", "path": "contracts/outbox_feature_matrix.toml", - "sha256": "846a9c0a41f538e5f55f352602ca898f4508e0b6d657ddb11b77eb0f76d19e55" + "sha256": "bc6578853c07f8901019e007e1cd9d7e87c41e0ef92425ad17781e528bcd15a1" } }, "generated_runtime": { - "byte_length": 1452, + "byte_length": 2664, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/generated/outbox_migration_registry.rs", - "sha256": "616be98c5ba8054e08a6109357a713a39a0086040272a1f6165f54530b469a91" + "sha256": "23c6438b11f7576b64848cc645bdd844727771a26bea73a20243ea7db38d5386" }, "ledger": { "adoption": "exact_unledgered_0001_catalog_only_v1", @@ -144,13 +148,47 @@ "sha256": "a7ee775d32c2b9f845961425362e1b1e558ce0d025f7d22dd58f118ba4dab4fa" }, "version": 1 + }, + { + "down": { + "byte_length": 208, + "hash_algorithm": "sha256_bytes_v1", + "path": "crates/outbox/migrations/0002_phase1_publication.down.sql", + "sha256": "57a5a00ca4257973097acf7f5cc64494dd0bc73fcfa00af2e6c8bb1f61823928" + }, + "name": "phase1_publication", + "owned_objects": [ + "outbox_phase1_delivery_target", + "outbox_phase1_delivery_target_ready_idx", + "outbox_phase1_dispatch_intent", + "outbox_phase1_observation_repair", + "outbox_phase1_publication", + "outbox_phase1_publication_event_idx", + "outbox_phase1_publication_ready_idx", + "outbox_phase1_target_receipt" + ], + "owned_tables": [ + "outbox_phase1_delivery_target", + "outbox_phase1_dispatch_intent", + "outbox_phase1_observation_repair", + "outbox_phase1_publication", + "outbox_phase1_target_receipt" + ], + "schema_sha256": "a56af9ba400fd51c97d48886fbb3f3733adb97458d7109fa8989c1b7e0c8bcaf", + "up": { + "byte_length": 5733, + "hash_algorithm": "sha256_bytes_v1", + "path": "crates/outbox/migrations/0002_phase1_publication.up.sql", + "sha256": "84f0c9897cff8d002961cb6ad9dee53edcf28853d1407483519b00bdbf029308" + }, + "version": 2 } ], "registry_source": { - "byte_length": 1329, + "byte_length": 2497, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/contracts/migration_registry.v1.json", - "sha256": "c46e65d7274ea58d0574efaa6694e6f65702f9cf2cc49000af1278ae7c123e0e" + "sha256": "a93768922b973ab4a98dcaf371cfea1d9c307d0a1eb9a3909899a9e679de31a5" }, "release": { "change_id": "outbox-versioned-migration-authority", @@ -178,19 +216,19 @@ "source_files": [ { "file": { - "byte_length": 1532, + "byte_length": 1873, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/Cargo.toml", - "sha256": "144c3da533519c9e5cf830dfd9ca7e7c8309cd8b581245141cf869f6ae41a47b" + "sha256": "95364a40f818d5d7aa4589a42774e650b590e1e1e7d4c3ed4dc8220cc8cadcbf" }, "role": "outbox_package_manifest" }, { "file": { - "byte_length": 1684, + "byte_length": 2604, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/lib.rs", - "sha256": "8a5ac45f4b88eee78078773d6400ae658f5180771a8f51c3ba1c6c6a95d2fe76" + "sha256": "a97deba8ec374102ba6b3c243e414cf99f14ab4fa374dd6f245b5339793ca037" }, "role": "outbox_public_surface" }, @@ -214,19 +252,19 @@ }, { "file": { - "byte_length": 22839, + "byte_length": 22911, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/migrations.rs", - "sha256": "f207cab712cbc4bdf8b9fe1892e553ced3512adab6b7b6883c8cbe79c7205255" + "sha256": "360d74823f943d541585d0fe41296c0ef1cb5a5a898f8c2957d105081787018d" }, "role": "migration_registry_runtime" }, { "file": { - "byte_length": 63715, + "byte_length": 64349, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/schema.rs", - "sha256": "135341d8d47171bf587bd17bb5ec5085c00648290ba84cdfc9feb348e6e214dd" + "sha256": "4c06effa3c57fd6f56e733a6a61a920ed00c63938cc7122da762a3135d2402c5" }, "role": "schema_runtime" }, @@ -241,10 +279,10 @@ }, { "file": { - "byte_length": 318687, + "byte_length": 318698, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/store.rs", - "sha256": "6aa76ba0e86554a92dd3e21b6b74853bd5c061edef69c5024e740dcbf777d193" + "sha256": "7f324a44658eb670ab82d30eaa80fe290a93f1cc3203e94ec62517ab015182c9" }, "role": "store_integration" }, @@ -259,28 +297,28 @@ }, { "file": { - "byte_length": 46145, + "byte_length": 46807, "hash_algorithm": "sha256_bytes_v1", "path": "tools/xtask/src/contract/outbox_migration.rs", - "sha256": "2ee39f355148e42a80ed73f7d292a7f03c42b12e9bdd03d849c01dda17254784" + "sha256": "feecba3ecb9ac05891877ef1de1b1e17f85fd349dbcc5b0f4cd7a09d2d6e0dda" }, "role": "contract_governance" }, { "file": { - "byte_length": 485783, + "byte_length": 486359, "hash_algorithm": "sha256_bytes_v1", "path": "tools/xtask/src/contract.rs", - "sha256": "cb3cbb45378551db7a3b1309dbc742d81f90af0c8a067065bb4e8698b744c5ef" + "sha256": "96468090685099e9d3722d45cf05170f3da83376150cbf58fee776157e6be649" }, "role": "contract_dispatch" }, { "file": { - "byte_length": 23735, + "byte_length": 24752, "hash_algorithm": "sha256_bytes_v1", "path": "tools/xtask/src/main.rs", - "sha256": "5f08b60c2d35a0a2a2b2743c9906538963282a5ed93a03a76087a9bfd282c07d" + "sha256": "4defaf823f7ab58322b59a568eb07d68e09c66830a9624b043529e29a3032daf" }, "role": "xtask_dispatch" }, @@ -313,25 +351,25 @@ }, { "file": { - "byte_length": 25312, + "byte_length": 25965, "hash_algorithm": "sha256_bytes_v1", "path": "contracts/releases/1.0.0-alpha.1.toml", - "sha256": "1ae79c394b13c64d785436e98a0743789eba7786adb2ec4edff7df1344c4f2c6" + "sha256": "a23dae6f24e7d4202826970b2f8c6a7d942d79e068c3ab1f17e5fc96e0ca1603" }, "role": "release_record" }, { "file": { - "byte_length": 36323, + "byte_length": 36955, "hash_algorithm": "sha256_bytes_v1", "path": "CHANGELOG.md", - "sha256": "3ee8e0d37238113448ddcf44ecd7f60e38d74e87a45ba712a1a3730d157510a3" + "sha256": "68e74a02a28a5c24643befca7038d143c5f5e6f941e314e98b1d6dfc1094e28d" }, "role": "release_notes" } ], "version_bounds": { - "current": 1, + "current": 2, "derivation": "first_and_last_ordered_registry_entries_v1", "minimum": 1 } diff --git a/crates/outbox/contracts/migration_authority_v1.manifest.sha256 b/crates/outbox/contracts/migration_authority_v1.manifest.sha256 @@ -1 +1 @@ -cbc5f03a84e21694878a6e3cd3173cbb0da222f8258a5f25a565fdebf00d3c49 +94a67ed7f4726eb45fc8481cf231a8129a776d4a73918a2bb371811763fa2589 diff --git a/crates/outbox/contracts/migration_registry.v1.json b/crates/outbox/contracts/migration_registry.v1.json @@ -34,6 +34,34 @@ "outbox_event", "outbox_operations" ] + }, + { + "version": 2, + "name": "phase1_publication", + "up_path": "crates/outbox/migrations/0002_phase1_publication.up.sql", + "down_path": "crates/outbox/migrations/0002_phase1_publication.down.sql", + "up_byte_length": 5733, + "down_byte_length": 208, + "up_sha256": "84f0c9897cff8d002961cb6ad9dee53edcf28853d1407483519b00bdbf029308", + "down_sha256": "57a5a00ca4257973097acf7f5cc64494dd0bc73fcfa00af2e6c8bb1f61823928", + "schema_sha256": "a56af9ba400fd51c97d48886fbb3f3733adb97458d7109fa8989c1b7e0c8bcaf", + "owned_objects": [ + "outbox_phase1_delivery_target", + "outbox_phase1_delivery_target_ready_idx", + "outbox_phase1_dispatch_intent", + "outbox_phase1_observation_repair", + "outbox_phase1_publication", + "outbox_phase1_publication_event_idx", + "outbox_phase1_publication_ready_idx", + "outbox_phase1_target_receipt" + ], + "owned_tables": [ + "outbox_phase1_delivery_target", + "outbox_phase1_dispatch_intent", + "outbox_phase1_observation_repair", + "outbox_phase1_publication", + "outbox_phase1_target_receipt" + ] } ] } diff --git a/crates/outbox/contracts/phase1_publication_v1.descriptor.json b/crates/outbox/contracts/phase1_publication_v1.descriptor.json @@ -0,0 +1,115 @@ +{ + "schema_version": 1, + "contract_id": "radroots_outbox.phase1_publication.v1", + "migration": { + "version": 2, + "name": "phase1_publication", + "up_path": "crates/outbox/migrations/0002_phase1_publication.up.sql", + "up_sha256": "84f0c9897cff8d002961cb6ad9dee53edcf28853d1407483519b00bdbf029308", + "down_path": "crates/outbox/migrations/0002_phase1_publication.down.sql", + "down_sha256": "57a5a00ca4257973097acf7f5cc64494dd0bc73fcfa00af2e6c8bb1f61823928", + "schema_sha256": "a56af9ba400fd51c97d48886fbb3f3733adb97458d7109fa8989c1b7e0c8bcaf" + }, + "resource_limits": { + "target_count": 16, + "target_uri_bytes": 2048, + "diagnostic_bytes": 4096, + "claim_lease_millis": 300000 + }, + "operation_identity": { + "algorithm": "sha256_raw_fixed_width_v1", + "domain": "radroots.phase1.publication-operation.v1", + "domain_terminator_hex": "00", + "preimage": [ + "artifact_digest_32", + "media_readiness_binding_digest_32", + "expected_author_32", + "target_policy_digest_32" + ] + }, + "dispatch_identity": { + "algorithm": "sha256_raw_fixed_width_v1", + "domain": "radroots.phase1.relay-dispatch.v1", + "domain_terminator_hex": "00", + "preimage": [ + "event_id_32", + "target_policy_digest_32", + "endpoint_fingerprint_32" + ] + }, + "event_states": [ + "ready", + "claimed-for-signing", + "signed-ready", + "dispatching", + "published", + "failed-retryable", + "failed-terminal", + "quarantined", + "cancelled" + ], + "target_states": [ + "pending", + "in-flight", + "accepted-observation-pending", + "accepted-observed", + "failed-retryable", + "failed-terminal", + "uncertain", + "cancelled" + ], + "transitions": [ + { "id": "claim-ready", "scope": "event", "from": "ready", "to": "claimed-for-signing", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "replace-claim", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, + { "id": "reclaim-signing", "scope": "event", "from": "claimed-for-signing", "to": "claimed-for-signing", "revision_cas": true, "lease_predicate": "expired", "durable_side_effect": "replace-claim", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "claim-sign-retry", "scope": "event", "from": "failed-retryable", "to": "claimed-for-signing", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "replace-claim", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "renew-signing", "scope": "event", "from": "claimed-for-signing", "to": "claimed-for-signing", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "extend-claim", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, + { "id": "release-signing", "scope": "event", "from": "claimed-for-signing", "to": "ready", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "clear-claim", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, + { "id": "complete-signing", "scope": "event", "from": "claimed-for-signing", "to": "signed-ready", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-immutable-signed-bytes", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, + { "id": "retry-signing", "scope": "event", "from": "claimed-for-signing", "to": "failed-retryable", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "fail-signing", "scope": "event", "from": "claimed-for-signing", "to": "failed-terminal", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "quarantine-signing", "scope": "event", "from": "claimed-for-signing", "to": "quarantined", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "cancel-signing", "scope": "event", "from": "claimed-for-signing", "to": "cancelled", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "clear-claim", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "begin-dispatch", "scope": "event", "from": "signed-ready", "to": "dispatching", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-dispatch-intent", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, + { "id": "continue-dispatch", "scope": "event", "from": "dispatching", "to": "dispatching", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-dispatch-intent", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "dispatch-published", "scope": "event", "from": "dispatching", "to": "published", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-receipt", "retry_class": "none", "repair_edge": false, "terminal_destination": true }, + { "id": "dispatch-waiting", "scope": "event", "from": "dispatching", "to": "signed-ready", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-result", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "dispatch-exhausted", "scope": "event", "from": "dispatching", "to": "failed-terminal", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-result", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "claim-target", "scope": "target", "from": "pending", "to": "in-flight", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "persist-dispatch-intent", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, + { "id": "retry-target", "scope": "target", "from": "failed-retryable", "to": "in-flight", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "reuse-dispatch-intent", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "repair-uncertain-target", "scope": "target", "from": "uncertain", "to": "in-flight", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "reuse-dispatch-intent", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, + { "id": "target-accepted-pending", "scope": "target", "from": "in-flight", "to": "accepted-observation-pending", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-receipt-and-repair", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, + { "id": "target-accepted", "scope": "target", "from": "in-flight", "to": "accepted-observed", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-receipt", "retry_class": "none", "repair_edge": false, "terminal_destination": true }, + { "id": "target-retryable", "scope": "target", "from": "in-flight", "to": "failed-retryable", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, + { "id": "target-terminal", "scope": "target", "from": "in-flight", "to": "failed-terminal", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "target-uncertain", "scope": "target", "from": "in-flight", "to": "uncertain", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, + { "id": "target-cancelled", "scope": "target", "from": "in-flight", "to": "cancelled", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "clear-claim", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "observation-repaired", "scope": "target", "from": "accepted-observation-pending", "to": "accepted-observed", "revision_cas": true, "lease_predicate": "repair-revision-cas", "durable_side_effect": "complete-repair-and-receipt", "retry_class": "repair", "repair_edge": true, "terminal_destination": true } + ], + "stable_errors": [ + "phase1_publication_artifact_invalid", + "phase1_publication_claim_invalid", + "phase1_publication_diagnostic_too_large", + "phase1_publication_entropy_unavailable", + "phase1_publication_idempotency_conflict", + "phase1_publication_integer_range", + "phase1_publication_lease_invalid", + "phase1_publication_not_found", + "phase1_publication_readiness_invalid", + "phase1_publication_required_target_count", + "phase1_publication_revision_conflict", + "phase1_publication_signed_event_invalid", + "phase1_publication_signed_event_mismatch", + "phase1_publication_sqlite", + "phase1_publication_state_conflict", + "phase1_publication_stored_authority_invalid", + "phase1_publication_stored_digest_invalid", + "phase1_publication_stored_state_invalid", + "phase1_publication_stored_value_too_large", + "phase1_publication_target_count", + "phase1_publication_target_duplicate", + "phase1_publication_target_not_found", + "phase1_publication_target_uri_invalid", + "phase1_publication_target_uri_too_large", + "phase1_publication_time_invalid" + ] +} diff --git a/crates/outbox/contracts/phase1_publication_v1.manifest.json b/crates/outbox/contracts/phase1_publication_v1.manifest.json @@ -0,0 +1,152 @@ +{ + "contract_id": "radroots_outbox.phase1_publication.v1", + "descriptor": { + "byte_length": 10119, + "path": "crates/outbox/contracts/phase1_publication_v1.descriptor.json", + "sha256": "2a85ff1ff4ef45dcfbac2e255b818000b61a218277b97931dc2af7a489716067" + }, + "manifest_schema": { + "byte_length": 3806, + "path": "crates/outbox/contracts/phase1_publication_v1.manifest.schema.json", + "sha256": "48587292d857b56a97fa053f68712251b16de6e814d3697f13e7e5d5a4497d77" + }, + "migration": { + "down": { + "byte_length": 208, + "path": "crates/outbox/migrations/0002_phase1_publication.down.sql", + "sha256": "57a5a00ca4257973097acf7f5cc64494dd0bc73fcfa00af2e6c8bb1f61823928" + }, + "schema_sha256": "a56af9ba400fd51c97d48886fbb3f3733adb97458d7109fa8989c1b7e0c8bcaf", + "up": { + "byte_length": 5733, + "path": "crates/outbox/migrations/0002_phase1_publication.up.sql", + "sha256": "84f0c9897cff8d002961cb6ad9dee53edcf28853d1407483519b00bdbf029308" + }, + "version": 2 + }, + "release": { + "change_id": "outbox-phase1-publication-state", + "changelog": "CHANGELOG.md", + "record": "contracts/releases/1.0.0-alpha.1.toml" + }, + "result_vector": { + "canonical": { + "byte_length": 2464, + "path": "contracts/conformance/vectors/outbox/phase1_publication.v1.json", + "sha256": "0a21d59ce5050a0e27569bbde4dd958483a7d91099eac3a482bc989c80a385eb" + }, + "executor": { + "byte_length": 8528, + "path": "crates/outbox/tests/phase1_publication_v1_result_vector.rs", + "sha256": "8e656a6cf7cdf9366ec815846634998c7ae1e10e33ade31044748cfd6a33ba90" + }, + "executor_id": "radroots_outbox.phase1_publication.v1.result_vector_executor.v1", + "executor_test": "phase1_publication_v1_result_vector", + "mirror_path": "crates/outbox/tests/fixtures/phase1_publication.v1.json" + }, + "schema_version": 1, + "source_files": [ + { + "file": { + "byte_length": 1873, + "path": "crates/outbox/Cargo.toml", + "sha256": "95364a40f818d5d7aa4589a42774e650b590e1e1e7d4c3ed4dc8220cc8cadcbf" + }, + "role": "outbox_package_manifest" + }, + { + "file": { + "byte_length": 2604, + "path": "crates/outbox/src/lib.rs", + "sha256": "a97deba8ec374102ba6b3c243e414cf99f14ab4fa374dd6f245b5339793ca037" + }, + "role": "outbox_public_surface" + }, + { + "file": { + "byte_length": 101112, + "path": "crates/outbox/src/phase1_publication.rs", + "sha256": "9019a23d593b27d9a0e1b67871d8da53c91b5a0b0a2a37698c64cbbfe9111272" + }, + "role": "phase1_publication_runtime" + }, + { + "file": { + "byte_length": 64349, + "path": "crates/outbox/src/schema.rs", + "sha256": "4c06effa3c57fd6f56e733a6a61a920ed00c63938cc7122da762a3135d2402c5" + }, + "role": "outbox_schema_runtime" + }, + { + "file": { + "byte_length": 2664, + "path": "crates/outbox/src/generated/outbox_migration_registry.rs", + "sha256": "23c6438b11f7576b64848cc645bdd844727771a26bea73a20243ea7db38d5386" + }, + "role": "migration_registry_runtime" + }, + { + "file": { + "byte_length": 2497, + "path": "crates/outbox/contracts/migration_registry.v1.json", + "sha256": "a93768922b973ab4a98dcaf371cfea1d9c307d0a1eb9a3909899a9e679de31a5" + }, + "role": "migration_registry_source" + }, + { + "file": { + "byte_length": 8528, + "path": "crates/outbox/tests/phase1_publication_v1_result_vector.rs", + "sha256": "8e656a6cf7cdf9366ec815846634998c7ae1e10e33ade31044748cfd6a33ba90" + }, + "role": "vector_executor" + }, + { + "file": { + "byte_length": 24248, + "path": "tools/xtask/src/contract/outbox_phase1_publication.rs", + "sha256": "4c81dea379e674066f0822d52f7d90e601a2d5311d5f88b809f760682c42bfa0" + }, + "role": "contract_governance" + }, + { + "file": { + "byte_length": 486359, + "path": "tools/xtask/src/contract.rs", + "sha256": "96468090685099e9d3722d45cf05170f3da83376150cbf58fee776157e6be649" + }, + "role": "contract_dispatch" + }, + { + "file": { + "byte_length": 24752, + "path": "tools/xtask/src/main.rs", + "sha256": "4defaf823f7ab58322b59a568eb07d68e09c66830a9624b043529e29a3032daf" + }, + "role": "xtask_dispatch" + }, + { + "file": { + "byte_length": 25965, + "path": "contracts/releases/1.0.0-alpha.1.toml", + "sha256": "a23dae6f24e7d4202826970b2f8c6a7d942d79e068c3ab1f17e5fc96e0ca1603" + }, + "role": "release_record" + }, + { + "file": { + "byte_length": 36955, + "path": "CHANGELOG.md", + "sha256": "68e74a02a28a5c24643befca7038d143c5f5e6f941e314e98b1d6dfc1094e28d" + }, + "role": "release_notes" + } + ], + "state_machine": { + "event_state_count": 9, + "stable_error_count": 25, + "target_state_count": 8, + "transition_count": 25 + } +} diff --git a/crates/outbox/contracts/phase1_publication_v1.manifest.schema.json b/crates/outbox/contracts/phase1_publication_v1.manifest.schema.json @@ -0,0 +1,177 @@ +{ + "$defs": { + "file": { + "additionalProperties": false, + "properties": { + "byte_length": { + "minimum": 1, + "type": "integer" + }, + "path": { + "minLength": 1, + "type": "string" + }, + "sha256": { + "$ref": "#/$defs/sha256" + } + }, + "required": [ + "path", + "byte_length", + "sha256" + ], + "type": "object" + }, + "sha256": { + "pattern": "^[0-9a-f]{64}$", + "type": "string" + } + }, + "$id": "https://radroots.org/contracts/outbox/phase1_publication_v1.manifest.schema.json", + "$schema": "https://json-schema.org/draft/2020-12/schema", + "additionalProperties": false, + "properties": { + "contract_id": { + "const": "radroots_outbox.phase1_publication.v1" + }, + "descriptor": { + "$ref": "#/$defs/file" + }, + "manifest_schema": { + "$ref": "#/$defs/file" + }, + "migration": { + "additionalProperties": false, + "properties": { + "down": { + "$ref": "#/$defs/file" + }, + "schema_sha256": { + "$ref": "#/$defs/sha256" + }, + "up": { + "$ref": "#/$defs/file" + }, + "version": { + "const": 2 + } + }, + "required": [ + "version", + "schema_sha256", + "up", + "down" + ], + "type": "object" + }, + "release": { + "additionalProperties": false, + "properties": { + "change_id": { + "const": "outbox-phase1-publication-state" + }, + "changelog": { + "const": "CHANGELOG.md" + }, + "record": { + "const": "contracts/releases/1.0.0-alpha.1.toml" + } + }, + "required": [ + "change_id", + "record", + "changelog" + ], + "type": "object" + }, + "result_vector": { + "additionalProperties": false, + "properties": { + "canonical": { + "$ref": "#/$defs/file" + }, + "executor": { + "$ref": "#/$defs/file" + }, + "executor_id": { + "const": "radroots_outbox.phase1_publication.v1.result_vector_executor.v1" + }, + "executor_test": { + "const": "phase1_publication_v1_result_vector" + }, + "mirror_path": { + "const": "crates/outbox/tests/fixtures/phase1_publication.v1.json" + } + }, + "required": [ + "canonical", + "mirror_path", + "executor", + "executor_id", + "executor_test" + ], + "type": "object" + }, + "schema_version": { + "const": 1 + }, + "source_files": { + "items": { + "additionalProperties": false, + "properties": { + "file": { + "$ref": "#/$defs/file" + }, + "role": { + "minLength": 1, + "type": "string" + } + }, + "required": [ + "role", + "file" + ], + "type": "object" + }, + "maxItems": 12, + "minItems": 12, + "type": "array" + }, + "state_machine": { + "additionalProperties": false, + "properties": { + "event_state_count": { + "const": 9 + }, + "stable_error_count": { + "const": 25 + }, + "target_state_count": { + "const": 8 + }, + "transition_count": { + "const": 25 + } + }, + "required": [ + "event_state_count", + "target_state_count", + "transition_count", + "stable_error_count" + ], + "type": "object" + } + }, + "required": [ + "schema_version", + "contract_id", + "manifest_schema", + "descriptor", + "migration", + "state_machine", + "result_vector", + "source_files", + "release" + ], + "type": "object" +} diff --git a/crates/outbox/contracts/phase1_publication_v1.manifest.sha256 b/crates/outbox/contracts/phase1_publication_v1.manifest.sha256 @@ -0,0 +1 @@ +99a1a2466aa97d3adb02def1c2cd20d7cc95fd364221c4e694ac6aa6e3c272d3 diff --git a/crates/outbox/migrations/0002_phase1_publication.down.sql b/crates/outbox/migrations/0002_phase1_publication.down.sql @@ -0,0 +1,5 @@ +DROP TABLE outbox_phase1_observation_repair; +DROP TABLE outbox_phase1_target_receipt; +DROP TABLE outbox_phase1_dispatch_intent; +DROP TABLE outbox_phase1_delivery_target; +DROP TABLE outbox_phase1_publication; diff --git a/crates/outbox/migrations/0002_phase1_publication.up.sql b/crates/outbox/migrations/0002_phase1_publication.up.sql @@ -0,0 +1,95 @@ +CREATE TABLE outbox_phase1_publication ( + publication_id INTEGER PRIMARY KEY AUTOINCREMENT, + operation_digest BLOB NOT NULL UNIQUE CHECK (length(operation_digest) = 32), + artifact_schema_version INTEGER NOT NULL CHECK (artifact_schema_version = 1), + artifact_json BLOB NOT NULL CHECK (length(artifact_json) BETWEEN 1 AND 2097152), + artifact_digest BLOB NOT NULL CHECK (length(artifact_digest) = 32), + readiness_schema_version INTEGER NOT NULL CHECK (readiness_schema_version = 1), + readiness_json BLOB NOT NULL CHECK (length(readiness_json) BETWEEN 1 AND 4194304), + readiness_digest BLOB NOT NULL CHECK (length(readiness_digest) = 32), + semantic_role TEXT NOT NULL CHECK (semantic_role IN ('profile', 'update', 'photo_update', 'ask', 'event_date', 'event_time', 'food_availability')), + expected_author BLOB NOT NULL CHECK (length(expected_author) = 32), + expected_event_id BLOB NOT NULL CHECK (length(expected_event_id) = 32), + target_policy_digest BLOB NOT NULL CHECK (length(target_policy_digest) = 32), + target_count INTEGER NOT NULL CHECK (target_count BETWEEN 1 AND 16), + required_target_count INTEGER NOT NULL CHECK (required_target_count BETWEEN 1 AND target_count), + state TEXT NOT NULL CHECK (state IN ('ready', 'claimed-for-signing', 'signed-ready', 'dispatching', 'published', 'failed-retryable', 'failed-terminal', 'quarantined', 'cancelled')), + state_revision INTEGER NOT NULL CHECK (state_revision >= 0), + claim_token BLOB CHECK (claim_token IS NULL OR length(claim_token) = 32), + claim_expires_at_ms INTEGER, + next_attempt_after_ms INTEGER NOT NULL, + signed_event_json BLOB CHECK (signed_event_json IS NULL OR length(signed_event_json) BETWEEN 1 AND 1048576), + signed_event_digest BLOB CHECK (signed_event_digest IS NULL OR length(signed_event_digest) = 32), + signed_event_id BLOB CHECK (signed_event_id IS NULL OR length(signed_event_id) = 32), + last_error TEXT CHECK (last_error IS NULL OR length(CAST(last_error AS BLOB)) <= 4096), + created_at_ms INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL, + CHECK ( + (claim_token IS NULL AND claim_expires_at_ms IS NULL) + OR (claim_token IS NOT NULL AND claim_expires_at_ms IS NOT NULL AND state = 'claimed-for-signing') + ), + CHECK ( + (signed_event_json IS NULL AND signed_event_digest IS NULL AND signed_event_id IS NULL) + OR (signed_event_json IS NOT NULL AND signed_event_digest IS NOT NULL AND signed_event_id IS NOT NULL) + ), + CHECK (state NOT IN ('signed-ready', 'dispatching', 'published') OR signed_event_json IS NOT NULL) +) STRICT; + +CREATE INDEX outbox_phase1_publication_ready_idx +ON outbox_phase1_publication(state, next_attempt_after_ms, claim_expires_at_ms, created_at_ms, publication_id); + +CREATE INDEX outbox_phase1_publication_event_idx +ON outbox_phase1_publication(expected_event_id, publication_id); + +CREATE TABLE outbox_phase1_delivery_target ( + target_id INTEGER PRIMARY KEY AUTOINCREMENT, + publication_id INTEGER NOT NULL REFERENCES outbox_phase1_publication(publication_id) ON DELETE CASCADE, + target_ordinal INTEGER NOT NULL CHECK (target_ordinal BETWEEN 0 AND 15), + endpoint_uri TEXT NOT NULL CHECK (length(CAST(endpoint_uri AS BLOB)) BETWEEN 1 AND 2048), + endpoint_fingerprint BLOB NOT NULL CHECK (length(endpoint_fingerprint) = 32), + dispatch_digest BLOB NOT NULL CHECK (length(dispatch_digest) = 32), + state TEXT NOT NULL CHECK (state IN ('pending', 'in-flight', 'accepted-observation-pending', 'accepted-observed', 'failed-retryable', 'failed-terminal', 'uncertain', 'cancelled')), + state_revision INTEGER NOT NULL CHECK (state_revision >= 0), + claim_token BLOB CHECK (claim_token IS NULL OR length(claim_token) = 32), + claim_expires_at_ms INTEGER, + next_attempt_after_ms INTEGER NOT NULL, + last_error TEXT CHECK (last_error IS NULL OR length(CAST(last_error AS BLOB)) <= 4096), + updated_at_ms INTEGER NOT NULL, + UNIQUE(publication_id, target_ordinal), + UNIQUE(publication_id, endpoint_uri), + UNIQUE(dispatch_digest), + CHECK ( + (claim_token IS NULL AND claim_expires_at_ms IS NULL) + OR (claim_token IS NOT NULL AND claim_expires_at_ms IS NOT NULL AND state = 'in-flight') + ) +) STRICT; + +CREATE INDEX outbox_phase1_delivery_target_ready_idx +ON outbox_phase1_delivery_target(state, next_attempt_after_ms, claim_expires_at_ms, publication_id, target_id); + +CREATE TABLE outbox_phase1_dispatch_intent ( + intent_digest BLOB PRIMARY KEY CHECK (length(intent_digest) = 32), + target_id INTEGER NOT NULL UNIQUE REFERENCES outbox_phase1_delivery_target(target_id) ON DELETE CASCADE, + signed_event_digest BLOB NOT NULL CHECK (length(signed_event_digest) = 32), + state TEXT NOT NULL CHECK (state IN ('in-flight', 'uncertain', 'completed', 'cancelled')), + state_revision INTEGER NOT NULL CHECK (state_revision >= 0), + created_at_ms INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL +) STRICT, WITHOUT ROWID; + +CREATE TABLE outbox_phase1_target_receipt ( + receipt_digest BLOB PRIMARY KEY CHECK (length(receipt_digest) = 32), + target_id INTEGER NOT NULL REFERENCES outbox_phase1_delivery_target(target_id) ON DELETE CASCADE, + observation_kind TEXT NOT NULL CHECK (observation_kind IN ('accepted-pending', 'accepted-observed')), + observed_at_ms INTEGER NOT NULL, + UNIQUE(target_id, observation_kind) +) STRICT, WITHOUT ROWID; + +CREATE TABLE outbox_phase1_observation_repair ( + repair_digest BLOB PRIMARY KEY CHECK (length(repair_digest) = 32), + target_id INTEGER NOT NULL UNIQUE REFERENCES outbox_phase1_delivery_target(target_id) ON DELETE CASCADE, + state TEXT NOT NULL CHECK (state IN ('pending', 'complete', 'failed-terminal')), + state_revision INTEGER NOT NULL CHECK (state_revision >= 0), + created_at_ms INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL +) STRICT, WITHOUT ROWID; diff --git a/crates/outbox/src/generated/outbox_migration_registry.rs b/crates/outbox/src/generated/outbox_migration_registry.rs @@ -36,4 +36,34 @@ const OUTBOX_MIGRATION_0001: OutboxMigration = OutboxMigration { ], }; -pub(crate) const OUTBOX_MIGRATIONS: &[OutboxMigration] = &[OUTBOX_MIGRATION_0001]; +const OUTBOX_MIGRATION_0002: OutboxMigration = OutboxMigration { + version: 2, + name: "phase1_publication", + up_sql: include_str!("../../migrations/0002_phase1_publication.up.sql"), + down_sql: include_str!("../../migrations/0002_phase1_publication.down.sql"), + up_len: 5733, + down_len: 208, + up_sha256: "84f0c9897cff8d002961cb6ad9dee53edcf28853d1407483519b00bdbf029308", + down_sha256: "57a5a00ca4257973097acf7f5cc64494dd0bc73fcfa00af2e6c8bb1f61823928", + schema_sha256: "a56af9ba400fd51c97d48886fbb3f3733adb97458d7109fa8989c1b7e0c8bcaf", + owned_object_names: &[ + "outbox_phase1_delivery_target", + "outbox_phase1_delivery_target_ready_idx", + "outbox_phase1_dispatch_intent", + "outbox_phase1_observation_repair", + "outbox_phase1_publication", + "outbox_phase1_publication_event_idx", + "outbox_phase1_publication_ready_idx", + "outbox_phase1_target_receipt", + ], + owned_table_names: &[ + "outbox_phase1_delivery_target", + "outbox_phase1_dispatch_intent", + "outbox_phase1_observation_repair", + "outbox_phase1_publication", + "outbox_phase1_target_receipt", + ], +}; + +pub(crate) const OUTBOX_MIGRATIONS: &[OutboxMigration] = + &[OUTBOX_MIGRATION_0001, OUTBOX_MIGRATION_0002]; diff --git a/crates/outbox/src/lib.rs b/crates/outbox/src/lib.rs @@ -8,6 +8,8 @@ mod generated; mod migrations; mod model; #[cfg(feature = "sqlite")] +mod phase1_publication; +#[cfg(feature = "sqlite")] mod schema; #[cfg(feature = "sqlite")] mod sqlite_lifecycle; @@ -30,6 +32,19 @@ pub use model::{ RadrootsOutboxTradeMutationInput, }; #[cfg(feature = "sqlite")] +pub use phase1_publication::{ + RADROOTS_PHASE1_PUBLICATION_CLAIM_LEASE_MAX_MILLIS, + RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES, RADROOTS_PHASE1_PUBLICATION_ERROR_CODES, + RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES, + RADROOTS_PHASE1_PUBLICATION_TRANSITIONS, RadrootsPhase1PublicationClaim, + RadrootsPhase1PublicationEnqueueReceipt, RadrootsPhase1PublicationEnqueueStatus, + RadrootsPhase1PublicationError, RadrootsPhase1PublicationEventState, + RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationTarget, + RadrootsPhase1PublicationTargetClaim, RadrootsPhase1PublicationTargetPolicy, + RadrootsPhase1PublicationTargetState, RadrootsPhase1PublicationTransition, + RadrootsPhase1PublicationTransitionRetryClass, RadrootsPhase1PublicationTransitionScope, +}; +#[cfg(feature = "sqlite")] pub use schema::{RadrootsOutboxSchemaStatus, inspect_outbox_schema_status}; #[cfg(feature = "sqlite")] pub use sqlite_lifecycle::{ diff --git a/crates/outbox/src/migrations.rs b/crates/outbox/src/migrations.rs @@ -437,7 +437,8 @@ mod tests { let registry = [legacy]; assert!(migration_for_version(OUTBOX_MIGRATIONS, 1).is_some()); - assert!(migration_for_version(OUTBOX_MIGRATIONS, 2).is_none()); + assert!(migration_for_version(OUTBOX_MIGRATIONS, 2).is_some()); + assert!(migration_for_version(OUTBOX_MIGRATIONS, 3).is_none()); assert!(sqlite_identifier_starts_with("OUTBOX_EVENT", "outbox_")); assert!(!sqlite_identifier_starts_with("short", "outbox_")); assert!(is_outbox_owned_table_name(&registry, "outbox_new")); @@ -529,7 +530,7 @@ mod tests { 2, )); - assert_registry_defect(validate_migration_registry(OUTBOX_MIGRATIONS, 1, 2)); + assert_registry_defect(validate_migration_registry(OUTBOX_MIGRATIONS, 1, 3)); validate_migration_registry(&[OUTBOX_MIGRATIONS[0], future_migration()], 1, 2) .expect("contiguous synthetic registry"); diff --git a/crates/outbox/src/phase1_publication.rs b/crates/outbox/src/phase1_publication.rs @@ -0,0 +1,2814 @@ +#![forbid(unsafe_code)] + +use crate::RadrootsOutbox; +use radroots_event::draft::{RadrootsSignedEvent, RadrootsVerifiedSignedEvent}; +use radroots_event::wire::RadrootsNip01EventWire; +use radroots_event_codec::wire::publication::allowlist::allow_phase1_publication_canonical_json; +use radroots_event_codec::wire::publication::{ + RADROOTS_PHASE1_PUBLICATION_ARTIFACT_MAX_BYTES, + RADROOTS_PHASE1_PUBLICATION_ARTIFACT_SCHEMA_VERSION, + RADROOTS_PHASE1_PUBLICATION_MEDIA_READINESS_BINDING_MAX_BYTES, + RADROOTS_PHASE1_PUBLICATION_MEDIA_READINESS_BINDING_SCHEMA_VERSION, + RADROOTS_PHASE1_PUBLICATION_SIGNED_EVENT_MAX_BYTES, + RadrootsPhase1MediaReadyPublicationArtifact, validate_phase1_publication_media_readiness, +}; +use sha2::{Digest, Sha256}; +use sqlx::{Row, Sqlite, Transaction}; +use std::collections::BTreeSet; +use std::fmt; +use url::Url; + +pub const RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT: usize = 16; +pub const RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES: usize = 2_048; +pub const RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES: usize = 4_096; +pub const RADROOTS_PHASE1_PUBLICATION_CLAIM_LEASE_MAX_MILLIS: i64 = 300_000; +pub const RADROOTS_PHASE1_PUBLICATION_ERROR_CODES: &[&str] = &[ + "phase1_publication_artifact_invalid", + "phase1_publication_claim_invalid", + "phase1_publication_diagnostic_too_large", + "phase1_publication_entropy_unavailable", + "phase1_publication_idempotency_conflict", + "phase1_publication_integer_range", + "phase1_publication_lease_invalid", + "phase1_publication_not_found", + "phase1_publication_readiness_invalid", + "phase1_publication_required_target_count", + "phase1_publication_revision_conflict", + "phase1_publication_signed_event_invalid", + "phase1_publication_signed_event_mismatch", + "phase1_publication_sqlite", + "phase1_publication_state_conflict", + "phase1_publication_stored_authority_invalid", + "phase1_publication_stored_digest_invalid", + "phase1_publication_stored_state_invalid", + "phase1_publication_stored_value_too_large", + "phase1_publication_target_count", + "phase1_publication_target_duplicate", + "phase1_publication_target_not_found", + "phase1_publication_target_uri_invalid", + "phase1_publication_target_uri_too_large", + "phase1_publication_time_invalid", +]; + +const OPERATION_DOMAIN: &[u8] = b"radroots.phase1.publication-operation.v1\0"; +const TARGET_POLICY_DOMAIN: &[u8] = b"radroots.phase1.target-policy.v1\0"; +const ENDPOINT_DOMAIN: &[u8] = b"radroots.phase1.endpoint.v1\0"; +const DISPATCH_DOMAIN: &[u8] = b"radroots.phase1.relay-dispatch.v1\0"; +const RECEIPT_DOMAIN: &[u8] = b"radroots.phase1.target-receipt.v1\0"; +const REPAIR_DOMAIN: &[u8] = b"radroots.phase1.observation-repair.v1\0"; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RadrootsPhase1PublicationEventState { + Ready, + ClaimedForSigning, + SignedReady, + Dispatching, + Published, + FailedRetryable, + FailedTerminal, + Quarantined, + Cancelled, +} + +impl RadrootsPhase1PublicationEventState { + pub const fn as_str(self) -> &'static str { + match self { + Self::Ready => "ready", + Self::ClaimedForSigning => "claimed-for-signing", + Self::SignedReady => "signed-ready", + Self::Dispatching => "dispatching", + Self::Published => "published", + Self::FailedRetryable => "failed-retryable", + Self::FailedTerminal => "failed-terminal", + Self::Quarantined => "quarantined", + Self::Cancelled => "cancelled", + } + } + + fn parse(value: &str) -> Result<Self, RadrootsPhase1PublicationError> { + match value { + "ready" => Ok(Self::Ready), + "claimed-for-signing" => Ok(Self::ClaimedForSigning), + "signed-ready" => Ok(Self::SignedReady), + "dispatching" => Ok(Self::Dispatching), + "published" => Ok(Self::Published), + "failed-retryable" => Ok(Self::FailedRetryable), + "failed-terminal" => Ok(Self::FailedTerminal), + "quarantined" => Ok(Self::Quarantined), + "cancelled" => Ok(Self::Cancelled), + _ => Err(RadrootsPhase1PublicationError::StoredStateInvalid), + } + } + + pub const fn is_terminal(self) -> bool { + matches!( + self, + Self::Published | Self::FailedTerminal | Self::Quarantined | Self::Cancelled + ) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RadrootsPhase1PublicationTargetState { + Pending, + InFlight, + AcceptedObservationPending, + AcceptedObserved, + FailedRetryable, + FailedTerminal, + Uncertain, + Cancelled, +} + +impl RadrootsPhase1PublicationTargetState { + pub const fn as_str(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::InFlight => "in-flight", + Self::AcceptedObservationPending => "accepted-observation-pending", + Self::AcceptedObserved => "accepted-observed", + Self::FailedRetryable => "failed-retryable", + Self::FailedTerminal => "failed-terminal", + Self::Uncertain => "uncertain", + Self::Cancelled => "cancelled", + } + } + + fn parse(value: &str) -> Result<Self, RadrootsPhase1PublicationError> { + match value { + "pending" => Ok(Self::Pending), + "in-flight" => Ok(Self::InFlight), + "accepted-observation-pending" => Ok(Self::AcceptedObservationPending), + "accepted-observed" => Ok(Self::AcceptedObserved), + "failed-retryable" => Ok(Self::FailedRetryable), + "failed-terminal" => Ok(Self::FailedTerminal), + "uncertain" => Ok(Self::Uncertain), + "cancelled" => Ok(Self::Cancelled), + _ => Err(RadrootsPhase1PublicationError::StoredStateInvalid), + } + } + + pub const fn is_terminal(self) -> bool { + matches!( + self, + Self::AcceptedObserved | Self::FailedTerminal | Self::Cancelled + ) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RadrootsPhase1PublicationTransitionScope { + Event, + Target, +} + +impl RadrootsPhase1PublicationTransitionScope { + pub const fn as_str(self) -> &'static str { + match self { + Self::Event => "event", + Self::Target => "target", + } + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RadrootsPhase1PublicationTransitionRetryClass { + None, + Retryable, + Repair, + Terminal, +} + +impl RadrootsPhase1PublicationTransitionRetryClass { + pub const fn as_str(self) -> &'static str { + match self { + Self::None => "none", + Self::Retryable => "retryable", + Self::Repair => "repair", + Self::Terminal => "terminal", + } + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RadrootsPhase1PublicationTransition { + pub id: &'static str, + pub scope: RadrootsPhase1PublicationTransitionScope, + pub from: &'static str, + pub to: &'static str, + pub revision_cas: bool, + pub lease_predicate: &'static str, + pub durable_side_effect: &'static str, + pub retry_class: RadrootsPhase1PublicationTransitionRetryClass, + pub repair_edge: bool, + pub terminal_destination: bool, +} + +macro_rules! transition { + ($id:literal, $scope:ident, $from:literal, $to:literal, $lease:literal, $effect:literal, $retry:ident, $repair:literal, $terminal:literal) => { + RadrootsPhase1PublicationTransition { + id: $id, + scope: RadrootsPhase1PublicationTransitionScope::$scope, + from: $from, + to: $to, + revision_cas: true, + lease_predicate: $lease, + durable_side_effect: $effect, + retry_class: RadrootsPhase1PublicationTransitionRetryClass::$retry, + repair_edge: $repair, + terminal_destination: $terminal, + } + }; +} + +pub const RADROOTS_PHASE1_PUBLICATION_TRANSITIONS: &[RadrootsPhase1PublicationTransition] = &[ + transition!( + "claim-ready", + Event, + "ready", + "claimed-for-signing", + "absent-or-expired", + "replace-claim", + None, + false, + false + ), + transition!( + "reclaim-signing", + Event, + "claimed-for-signing", + "claimed-for-signing", + "expired", + "replace-claim", + Retryable, + false, + false + ), + transition!( + "claim-sign-retry", + Event, + "failed-retryable", + "claimed-for-signing", + "absent-or-expired", + "replace-claim", + Retryable, + false, + false + ), + transition!( + "renew-signing", + Event, + "claimed-for-signing", + "claimed-for-signing", + "matching-live-token", + "extend-claim", + None, + false, + false + ), + transition!( + "release-signing", + Event, + "claimed-for-signing", + "ready", + "matching-live-token", + "clear-claim", + None, + false, + false + ), + transition!( + "complete-signing", + Event, + "claimed-for-signing", + "signed-ready", + "matching-live-token", + "persist-immutable-signed-bytes", + None, + false, + false + ), + transition!( + "retry-signing", + Event, + "claimed-for-signing", + "failed-retryable", + "matching-live-token", + "persist-bounded-error", + Retryable, + false, + false + ), + transition!( + "fail-signing", + Event, + "claimed-for-signing", + "failed-terminal", + "matching-live-token", + "persist-bounded-error", + Terminal, + false, + true + ), + transition!( + "quarantine-signing", + Event, + "claimed-for-signing", + "quarantined", + "matching-live-token", + "persist-bounded-error", + Terminal, + false, + true + ), + transition!( + "cancel-signing", + Event, + "claimed-for-signing", + "cancelled", + "matching-live-token", + "clear-claim", + Terminal, + false, + true + ), + transition!( + "begin-dispatch", + Event, + "signed-ready", + "dispatching", + "target-matching-live-token", + "persist-dispatch-intent", + None, + false, + false + ), + transition!( + "continue-dispatch", + Event, + "dispatching", + "dispatching", + "target-matching-live-token", + "persist-dispatch-intent", + Retryable, + false, + false + ), + transition!( + "dispatch-published", + Event, + "dispatching", + "published", + "target-matching-live-token", + "persist-target-receipt", + None, + false, + true + ), + transition!( + "dispatch-waiting", + Event, + "dispatching", + "signed-ready", + "target-matching-live-token", + "persist-target-result", + Retryable, + false, + false + ), + transition!( + "dispatch-exhausted", + Event, + "dispatching", + "failed-terminal", + "target-matching-live-token", + "persist-target-result", + Terminal, + false, + true + ), + transition!( + "claim-target", + Target, + "pending", + "in-flight", + "absent-or-expired", + "persist-dispatch-intent", + None, + false, + false + ), + transition!( + "retry-target", + Target, + "failed-retryable", + "in-flight", + "absent-or-expired", + "reuse-dispatch-intent", + Retryable, + false, + false + ), + transition!( + "repair-uncertain-target", + Target, + "uncertain", + "in-flight", + "absent-or-expired", + "reuse-dispatch-intent", + Repair, + true, + false + ), + transition!( + "target-accepted-pending", + Target, + "in-flight", + "accepted-observation-pending", + "matching-live-token", + "persist-receipt-and-repair", + Repair, + true, + false + ), + transition!( + "target-accepted", + Target, + "in-flight", + "accepted-observed", + "matching-live-token", + "persist-receipt", + None, + false, + true + ), + transition!( + "target-retryable", + Target, + "in-flight", + "failed-retryable", + "matching-live-token", + "persist-bounded-error", + Retryable, + false, + false + ), + transition!( + "target-terminal", + Target, + "in-flight", + "failed-terminal", + "matching-live-token", + "persist-bounded-error", + Terminal, + false, + true + ), + transition!( + "target-uncertain", + Target, + "in-flight", + "uncertain", + "matching-live-token", + "persist-bounded-error", + Repair, + true, + false + ), + transition!( + "target-cancelled", + Target, + "in-flight", + "cancelled", + "matching-live-token", + "clear-claim", + Terminal, + false, + true + ), + transition!( + "observation-repaired", + Target, + "accepted-observation-pending", + "accepted-observed", + "repair-revision-cas", + "complete-repair-and-receipt", + Repair, + true, + true + ), +]; + +#[derive(Debug)] +#[non_exhaustive] +pub enum RadrootsPhase1PublicationError { + Sqlite(sqlx::Error), + ArtifactInvalid { + source_code: &'static str, + }, + ReadinessInvalid { + source_code: &'static str, + }, + TargetCount { + max: usize, + actual: usize, + }, + RequiredTargetCount { + target_count: usize, + required: usize, + }, + TargetUriTooLarge { + max: usize, + actual: usize, + }, + TargetUriInvalid, + DuplicateTarget, + InvalidTime, + InvalidLease { + max_millis: i64, + actual: i64, + }, + EntropyUnavailable, + PublicationNotFound { + publication_id: i64, + }, + TargetNotFound { + target_id: i64, + }, + IdempotencyConflict, + RevisionConflict, + ClaimInvalid, + StateConflict, + SignedEventMismatch, + SignedEventInvalid, + DiagnosticTooLarge { + max: usize, + actual: usize, + }, + StoredValueTooLarge { + field: &'static str, + max: usize, + actual: usize, + }, + StoredDigestInvalid { + field: &'static str, + }, + StoredStateInvalid, + StoredAuthorityInvalid, + IntegerRange { + field: &'static str, + }, +} + +impl RadrootsPhase1PublicationError { + pub const fn code(&self) -> &'static str { + match self { + Self::Sqlite(_) => "phase1_publication_sqlite", + Self::ArtifactInvalid { .. } => "phase1_publication_artifact_invalid", + Self::ReadinessInvalid { .. } => "phase1_publication_readiness_invalid", + Self::TargetCount { .. } => "phase1_publication_target_count", + Self::RequiredTargetCount { .. } => "phase1_publication_required_target_count", + Self::TargetUriTooLarge { .. } => "phase1_publication_target_uri_too_large", + Self::TargetUriInvalid => "phase1_publication_target_uri_invalid", + Self::DuplicateTarget => "phase1_publication_target_duplicate", + Self::InvalidTime => "phase1_publication_time_invalid", + Self::InvalidLease { .. } => "phase1_publication_lease_invalid", + Self::EntropyUnavailable => "phase1_publication_entropy_unavailable", + Self::PublicationNotFound { .. } => "phase1_publication_not_found", + Self::TargetNotFound { .. } => "phase1_publication_target_not_found", + Self::IdempotencyConflict => "phase1_publication_idempotency_conflict", + Self::RevisionConflict => "phase1_publication_revision_conflict", + Self::ClaimInvalid => "phase1_publication_claim_invalid", + Self::StateConflict => "phase1_publication_state_conflict", + Self::SignedEventMismatch => "phase1_publication_signed_event_mismatch", + Self::SignedEventInvalid => "phase1_publication_signed_event_invalid", + Self::DiagnosticTooLarge { .. } => "phase1_publication_diagnostic_too_large", + Self::StoredValueTooLarge { .. } => "phase1_publication_stored_value_too_large", + Self::StoredDigestInvalid { .. } => "phase1_publication_stored_digest_invalid", + Self::StoredStateInvalid => "phase1_publication_stored_state_invalid", + Self::StoredAuthorityInvalid => "phase1_publication_stored_authority_invalid", + Self::IntegerRange { .. } => "phase1_publication_integer_range", + } + } + + pub fn public_diagnostic(&self) -> String { + let diagnostic = self.to_string(); + if diagnostic.len() <= RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES { + return diagnostic; + } + let suffix = "…"; + let mut end = RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES - suffix.len(); + while !diagnostic.is_char_boundary(end) { + end -= 1; + } + format!("{}{}", &diagnostic[..end], suffix) + } +} + +impl fmt::Display for RadrootsPhase1PublicationError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::TargetCount { max, actual } => { + write!( + formatter, + "target count {actual} exceeds valid range 1..={max}" + ) + } + Self::RequiredTargetCount { + target_count, + required, + } => write!( + formatter, + "required target count {required} is invalid for {target_count} targets" + ), + Self::TargetUriTooLarge { max, actual } => { + write!(formatter, "target URI is {actual} bytes; maximum is {max}") + } + Self::InvalidLease { max_millis, actual } => write!( + formatter, + "claim lease {actual} ms is outside 1..={max_millis}" + ), + Self::PublicationNotFound { publication_id } => { + write!( + formatter, + "Phase 1 publication {publication_id} was not found" + ) + } + Self::TargetNotFound { target_id } => { + write!( + formatter, + "Phase 1 publication target {target_id} was not found" + ) + } + Self::DiagnosticTooLarge { max, actual } => { + write!(formatter, "diagnostic is {actual} bytes; maximum is {max}") + } + Self::StoredValueTooLarge { field, max, actual } => write!( + formatter, + "stored {field} is {actual} bytes; maximum is {max}" + ), + Self::StoredDigestInvalid { field } => write!(formatter, "stored {field} is invalid"), + Self::IntegerRange { field } => write!(formatter, "stored {field} is out of range"), + Self::Sqlite(_) + | Self::ArtifactInvalid { .. } + | Self::ReadinessInvalid { .. } + | Self::TargetUriInvalid + | Self::DuplicateTarget + | Self::InvalidTime + | Self::EntropyUnavailable + | Self::IdempotencyConflict + | Self::RevisionConflict + | Self::ClaimInvalid + | Self::StateConflict + | Self::SignedEventMismatch + | Self::SignedEventInvalid + | Self::StoredStateInvalid + | Self::StoredAuthorityInvalid => formatter.write_str(self.code()), + } + } +} + +impl std::error::Error for RadrootsPhase1PublicationError { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + match self { + Self::Sqlite(error) => Some(error), + _ => None, + } + } +} + +impl From<sqlx::Error> for RadrootsPhase1PublicationError { + fn from(value: sqlx::Error) -> Self { + Self::Sqlite(value) + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsPhase1PublicationTargetPolicy { + targets: Vec<String>, + required_target_count: usize, + digest: [u8; 32], +} + +impl RadrootsPhase1PublicationTargetPolicy { + pub fn new<I, S>( + targets: I, + required_target_count: usize, + ) -> Result<Self, RadrootsPhase1PublicationError> + where + I: IntoIterator<Item = S>, + S: AsRef<str>, + { + let mut canonical = Vec::new(); + let mut targets = targets.into_iter(); + let (minimum, _) = targets.size_hint(); + if minimum > RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT { + return Err(RadrootsPhase1PublicationError::TargetCount { + max: RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, + actual: minimum, + }); + } + canonical.try_reserve_exact(minimum).map_err(|_| { + RadrootsPhase1PublicationError::TargetCount { + max: RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, + actual: minimum, + } + })?; + for target in &mut targets { + if canonical.len() == RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT { + return Err(RadrootsPhase1PublicationError::TargetCount { + max: RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, + actual: canonical.len() + 1, + }); + } + canonical.push(canonical_target_uri(target.as_ref())?); + } + canonical.sort(); + if canonical.windows(2).any(|pair| pair[0] == pair[1]) { + return Err(RadrootsPhase1PublicationError::DuplicateTarget); + } + if canonical.is_empty() + || required_target_count == 0 + || required_target_count > canonical.len() + { + return Err(RadrootsPhase1PublicationError::RequiredTargetCount { + target_count: canonical.len(), + required: required_target_count, + }); + } + let digest = target_policy_digest(&canonical, required_target_count)?; + Ok(Self { + targets: canonical, + required_target_count, + digest, + }) + } + + pub fn targets(&self) -> &[String] { + &self.targets + } + + pub const fn required_target_count(&self) -> usize { + self.required_target_count + } + + pub const fn digest(&self) -> &[u8; 32] { + &self.digest + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RadrootsPhase1PublicationEnqueueStatus { + Inserted, + Existing, +} + +#[derive(Clone, Debug)] +pub struct RadrootsPhase1PublicationEnqueueReceipt { + status: RadrootsPhase1PublicationEnqueueStatus, + record: RadrootsPhase1PublicationRecord, +} + +impl RadrootsPhase1PublicationEnqueueReceipt { + pub const fn status(&self) -> RadrootsPhase1PublicationEnqueueStatus { + self.status + } + + pub const fn record(&self) -> &RadrootsPhase1PublicationRecord { + &self.record + } + + pub fn into_record(self) -> RadrootsPhase1PublicationRecord { + self.record + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsPhase1PublicationTarget { + target_id: i64, + ordinal: usize, + endpoint_uri: String, + endpoint_fingerprint: [u8; 32], + dispatch_digest: [u8; 32], + state: RadrootsPhase1PublicationTargetState, + revision: u64, + next_attempt_after_ms: i64, +} + +impl RadrootsPhase1PublicationTarget { + pub const fn target_id(&self) -> i64 { + self.target_id + } + + pub const fn ordinal(&self) -> usize { + self.ordinal + } + + pub fn endpoint_uri(&self) -> &str { + &self.endpoint_uri + } + + pub const fn endpoint_fingerprint(&self) -> &[u8; 32] { + &self.endpoint_fingerprint + } + + pub const fn dispatch_digest(&self) -> &[u8; 32] { + &self.dispatch_digest + } + + pub const fn state(&self) -> RadrootsPhase1PublicationTargetState { + self.state + } + + pub const fn revision(&self) -> u64 { + self.revision + } + + pub const fn next_attempt_after_ms(&self) -> i64 { + self.next_attempt_after_ms + } +} + +#[derive(Clone, Debug)] +pub struct RadrootsPhase1PublicationRecord { + publication_id: i64, + operation_digest: [u8; 32], + ready_artifact: RadrootsPhase1MediaReadyPublicationArtifact, + target_policy: RadrootsPhase1PublicationTargetPolicy, + targets: Vec<RadrootsPhase1PublicationTarget>, + state: RadrootsPhase1PublicationEventState, + revision: u64, + next_attempt_after_ms: i64, + signed_event: Option<RadrootsVerifiedSignedEvent>, +} + +impl RadrootsPhase1PublicationRecord { + pub const fn publication_id(&self) -> i64 { + self.publication_id + } + + pub const fn operation_digest(&self) -> &[u8; 32] { + &self.operation_digest + } + + pub const fn ready_artifact(&self) -> &RadrootsPhase1MediaReadyPublicationArtifact { + &self.ready_artifact + } + + pub const fn target_policy(&self) -> &RadrootsPhase1PublicationTargetPolicy { + &self.target_policy + } + + pub fn targets(&self) -> &[RadrootsPhase1PublicationTarget] { + &self.targets + } + + pub const fn state(&self) -> RadrootsPhase1PublicationEventState { + self.state + } + + pub const fn revision(&self) -> u64 { + self.revision + } + + pub const fn next_attempt_after_ms(&self) -> i64 { + self.next_attempt_after_ms + } + + pub const fn signed_event(&self) -> Option<&RadrootsVerifiedSignedEvent> { + self.signed_event.as_ref() + } +} + +#[derive(Clone, Debug)] +pub struct RadrootsPhase1PublicationClaim { + publication_id: i64, + revision: u64, + token: [u8; 32], + expires_at_ms: i64, +} + +impl RadrootsPhase1PublicationClaim { + pub const fn publication_id(&self) -> i64 { + self.publication_id + } + + pub const fn revision(&self) -> u64 { + self.revision + } + + pub const fn expires_at_ms(&self) -> i64 { + self.expires_at_ms + } +} + +#[derive(Clone, Debug)] +pub struct RadrootsPhase1PublicationTargetClaim { + publication_id: i64, + publication_revision: u64, + target_id: i64, + target_revision: u64, + token: [u8; 32], + expires_at_ms: i64, +} + +impl RadrootsPhase1PublicationTargetClaim { + pub const fn publication_id(&self) -> i64 { + self.publication_id + } + + pub const fn publication_revision(&self) -> u64 { + self.publication_revision + } + + pub const fn target_id(&self) -> i64 { + self.target_id + } + + pub const fn target_revision(&self) -> u64 { + self.target_revision + } + + pub const fn expires_at_ms(&self) -> i64 { + self.expires_at_ms + } +} + +struct PreparedPublication { + artifact_json: Vec<u8>, + artifact_digest: [u8; 32], + readiness_json: Vec<u8>, + readiness_digest: [u8; 32], + semantic_role: &'static str, + expected_author: [u8; 32], + expected_event_id: [u8; 32], + operation_digest: [u8; 32], +} + +impl RadrootsOutbox { + pub async fn enqueue_phase1_publication( + &self, + ready: &RadrootsPhase1MediaReadyPublicationArtifact, + target_policy: &RadrootsPhase1PublicationTargetPolicy, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationEnqueueReceipt, RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + let prepared = prepare_publication(ready, target_policy)?; + let mut transaction = self.pool.begin().await?; + if let Some(row) = sqlx::query( + "SELECT publication_id FROM outbox_phase1_publication WHERE operation_digest = ?", + ) + .bind(prepared.operation_digest.as_slice()) + .fetch_optional(&mut *transaction) + .await? + { + let publication_id = row.try_get("publication_id")?; + transaction.commit().await?; + let record = self.load_phase1_publication(publication_id).await?; + return Ok(RadrootsPhase1PublicationEnqueueReceipt { + status: RadrootsPhase1PublicationEnqueueStatus::Existing, + record, + }); + } + if sqlx::query( + "SELECT publication_id FROM outbox_phase1_publication WHERE artifact_digest = ? AND expected_author = ? LIMIT 1", + ) + .bind(prepared.artifact_digest.as_slice()) + .bind(prepared.expected_author.as_slice()) + .fetch_optional(&mut *transaction) + .await? + .is_some() + { + return Err(RadrootsPhase1PublicationError::IdempotencyConflict); + } + let result = sqlx::query( + "INSERT INTO outbox_phase1_publication(operation_digest, artifact_schema_version, artifact_json, artifact_digest, readiness_schema_version, readiness_json, readiness_digest, semantic_role, expected_author, expected_event_id, target_policy_digest, target_count, required_target_count, state, state_revision, next_attempt_after_ms, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'ready', 0, ?, ?, ?)", + ) + .bind(prepared.operation_digest.as_slice()) + .bind(i64::from(RADROOTS_PHASE1_PUBLICATION_ARTIFACT_SCHEMA_VERSION)) + .bind(prepared.artifact_json) + .bind(prepared.artifact_digest.as_slice()) + .bind(i64::from( + RADROOTS_PHASE1_PUBLICATION_MEDIA_READINESS_BINDING_SCHEMA_VERSION, + )) + .bind(prepared.readiness_json) + .bind(prepared.readiness_digest.as_slice()) + .bind(prepared.semantic_role) + .bind(prepared.expected_author.as_slice()) + .bind(prepared.expected_event_id.as_slice()) + .bind(target_policy.digest.as_slice()) + .bind(i64::try_from(target_policy.targets.len()).map_err(|_| { + RadrootsPhase1PublicationError::IntegerRange { + field: "target_count", + } + })?) + .bind(i64::try_from(target_policy.required_target_count).map_err(|_| { + RadrootsPhase1PublicationError::IntegerRange { + field: "required_target_count", + } + })?) + .bind(now_ms) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *transaction) + .await?; + let publication_id = result.last_insert_rowid(); + for (ordinal, endpoint_uri) in target_policy.targets.iter().enumerate() { + let endpoint_fingerprint = endpoint_fingerprint(endpoint_uri); + let dispatch_digest = dispatch_digest( + &prepared.expected_event_id, + &target_policy.digest, + &endpoint_fingerprint, + ); + sqlx::query( + "INSERT INTO outbox_phase1_delivery_target(publication_id, target_ordinal, endpoint_uri, endpoint_fingerprint, dispatch_digest, state, state_revision, next_attempt_after_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, 'pending', 0, ?, ?)", + ) + .bind(publication_id) + .bind(i64::try_from(ordinal).map_err(|_| { + RadrootsPhase1PublicationError::IntegerRange { + field: "target_ordinal", + } + })?) + .bind(endpoint_uri) + .bind(endpoint_fingerprint.as_slice()) + .bind(dispatch_digest.as_slice()) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *transaction) + .await?; + } + transaction.commit().await?; + let record = self.load_phase1_publication(publication_id).await?; + Ok(RadrootsPhase1PublicationEnqueueReceipt { + status: RadrootsPhase1PublicationEnqueueStatus::Inserted, + record, + }) + } + + pub async fn load_phase1_publication( + &self, + publication_id: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + let row = sqlx::query( + "SELECT publication_id, operation_digest, + length(artifact_json) AS artifact_bytes, + CASE WHEN length(artifact_json) <= ? THEN artifact_json END AS bounded_artifact_json, + artifact_digest, + length(readiness_json) AS readiness_bytes, + CASE WHEN length(readiness_json) <= ? THEN readiness_json END AS bounded_readiness_json, + readiness_digest, semantic_role, expected_author, expected_event_id, + target_policy_digest, target_count, required_target_count, state, state_revision, + next_attempt_after_ms, + length(signed_event_json) AS signed_event_bytes, + CASE WHEN signed_event_json IS NULL OR length(signed_event_json) <= ? THEN signed_event_json END AS bounded_signed_event_json, + signed_event_digest, signed_event_id + FROM outbox_phase1_publication WHERE publication_id = ?", + ) + .bind(i64::try_from(RADROOTS_PHASE1_PUBLICATION_ARTIFACT_MAX_BYTES).unwrap_or(i64::MAX)) + .bind( + i64::try_from(RADROOTS_PHASE1_PUBLICATION_MEDIA_READINESS_BINDING_MAX_BYTES) + .unwrap_or(i64::MAX), + ) + .bind(i64::try_from(RADROOTS_PHASE1_PUBLICATION_SIGNED_EVENT_MAX_BYTES).unwrap_or(i64::MAX)) + .bind(publication_id) + .fetch_optional(&self.pool) + .await? + .ok_or(RadrootsPhase1PublicationError::PublicationNotFound { publication_id })?; + let artifact_json = bounded_blob( + &row, + "bounded_artifact_json", + "artifact_bytes", + "artifact_json", + RADROOTS_PHASE1_PUBLICATION_ARTIFACT_MAX_BYTES, + )? + .ok_or(RadrootsPhase1PublicationError::StoredAuthorityInvalid)?; + let readiness_json = bounded_blob( + &row, + "bounded_readiness_json", + "readiness_bytes", + "readiness_json", + RADROOTS_PHASE1_PUBLICATION_MEDIA_READINESS_BINDING_MAX_BYTES, + )? + .ok_or(RadrootsPhase1PublicationError::StoredAuthorityInvalid)?; + let allowlisted = + allow_phase1_publication_canonical_json(&artifact_json).map_err(|error| { + RadrootsPhase1PublicationError::ArtifactInvalid { + source_code: error.code(), + } + })?; + let ready_artifact = RadrootsPhase1MediaReadyPublicationArtifact::from_canonical_json( + allowlisted, + &readiness_json, + ) + .map_err(|error| RadrootsPhase1PublicationError::ReadinessInvalid { + source_code: error.code(), + })?; + validate_phase1_publication_media_readiness(&ready_artifact).map_err(|error| { + RadrootsPhase1PublicationError::ReadinessInvalid { + source_code: error.code(), + } + })?; + + let targets = load_targets(&self.pool, publication_id).await?; + let required_target_count = usize_from_i64( + row.try_get("required_target_count")?, + "required_target_count", + )?; + let target_policy = RadrootsPhase1PublicationTargetPolicy::new( + targets.iter().map(|target| target.endpoint_uri.as_str()), + required_target_count, + )?; + if targets.len() != usize_from_i64(row.try_get("target_count")?, "target_count")? { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + let operation_digest = blob32(&row, "operation_digest")?; + let artifact_digest = blob32(&row, "artifact_digest")?; + let readiness_digest = blob32(&row, "readiness_digest")?; + let expected_author = blob32(&row, "expected_author")?; + let expected_event_id = blob32(&row, "expected_event_id")?; + let stored_policy_digest = blob32(&row, "target_policy_digest")?; + let artifact = ready_artifact.artifact(); + if artifact_digest != *artifact.artifact_digest().as_bytes() + || readiness_digest != *ready_artifact.binding_digest().as_bytes() + || expected_author + != decode_hex32(artifact.expected_author().as_str(), "expected_author")? + || expected_event_id + != decode_hex32(artifact.expected_event_id().as_str(), "expected_event_id")? + || stored_policy_digest != target_policy.digest + || row.try_get::<String, _>("semantic_role")? != artifact.semantic_variant().as_str() + { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + let expected_operation = operation_digest_for( + &artifact_digest, + &readiness_digest, + &expected_author, + &stored_policy_digest, + ); + if operation_digest != expected_operation { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + validate_target_identities(&targets, &expected_event_id, &stored_policy_digest)?; + + let signed_json = bounded_blob( + &row, + "bounded_signed_event_json", + "signed_event_bytes", + "signed_event_json", + RADROOTS_PHASE1_PUBLICATION_SIGNED_EVENT_MAX_BYTES, + )?; + let signed_digest: Option<Vec<u8>> = row.try_get("signed_event_digest")?; + let signed_event_id: Option<Vec<u8>> = row.try_get("signed_event_id")?; + let signed_event = + reload_signed_event(signed_json, signed_digest, signed_event_id, &ready_artifact)?; + let state = + RadrootsPhase1PublicationEventState::parse(&row.try_get::<String, _>("state")?)?; + let signed_required = matches!( + state, + RadrootsPhase1PublicationEventState::SignedReady + | RadrootsPhase1PublicationEventState::Dispatching + | RadrootsPhase1PublicationEventState::Published + ); + let signed_forbidden = matches!( + state, + RadrootsPhase1PublicationEventState::Ready + | RadrootsPhase1PublicationEventState::ClaimedForSigning + | RadrootsPhase1PublicationEventState::FailedRetryable + ); + if (signed_required && signed_event.is_none()) + || (signed_forbidden && signed_event.is_some()) + { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + Ok(RadrootsPhase1PublicationRecord { + publication_id, + operation_digest, + ready_artifact, + target_policy, + targets, + state, + revision: u64_from_i64(row.try_get("state_revision")?, "state_revision")?, + next_attempt_after_ms: row.try_get("next_attempt_after_ms")?, + signed_event, + }) + } + + pub async fn claim_phase1_publication_for_signing( + &self, + publication_id: i64, + observed_revision: u64, + now_ms: i64, + lease_millis: i64, + ) -> Result<RadrootsPhase1PublicationClaim, RadrootsPhase1PublicationError> { + let record = self.load_phase1_publication(publication_id).await?; + if record.revision != observed_revision { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let expires_at_ms = validated_expiry(now_ms, lease_millis)?; + let token = new_claim_token()?; + let affected = sqlx::query( + "UPDATE outbox_phase1_publication + SET state = 'claimed-for-signing', state_revision = state_revision + 1, + claim_token = ?, claim_expires_at_ms = ?, updated_at_ms = ? + WHERE publication_id = ? AND state_revision = ? AND signed_event_json IS NULL + AND state IN ('ready', 'claimed-for-signing', 'failed-retryable') + AND (claim_token IS NULL OR claim_expires_at_ms <= ?)", + ) + .bind(token.as_slice()) + .bind(expires_at_ms) + .bind(now_ms) + .bind(publication_id) + .bind(i64_from_u64(observed_revision, "state_revision")?) + .bind(now_ms) + .execute(&self.pool) + .await? + .rows_affected(); + if affected != 1 { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + Ok(RadrootsPhase1PublicationClaim { + publication_id, + revision: observed_revision + 1, + token, + expires_at_ms, + }) + } + + pub async fn renew_phase1_publication_claim( + &self, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + lease_millis: i64, + ) -> Result<RadrootsPhase1PublicationClaim, RadrootsPhase1PublicationError> { + self.load_phase1_publication(claim.publication_id).await?; + let expires_at_ms = validated_expiry(now_ms, lease_millis)?; + let affected = sqlx::query( + "UPDATE outbox_phase1_publication + SET state_revision = state_revision + 1, claim_expires_at_ms = ?, updated_at_ms = ? + WHERE publication_id = ? AND state = 'claimed-for-signing' AND state_revision = ? + AND claim_token = ? AND claim_expires_at_ms > ?", + ) + .bind(expires_at_ms) + .bind(now_ms) + .bind(claim.publication_id) + .bind(i64_from_u64(claim.revision, "state_revision")?) + .bind(claim.token.as_slice()) + .bind(now_ms) + .execute(&self.pool) + .await? + .rows_affected(); + if affected != 1 { + return Err(RadrootsPhase1PublicationError::ClaimInvalid); + } + Ok(RadrootsPhase1PublicationClaim { + publication_id: claim.publication_id, + revision: claim.revision + 1, + token: claim.token, + expires_at_ms, + }) + } + + pub async fn release_phase1_publication_claim( + &self, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + transition_signing_claim(self, claim, now_ms, "ready", None, now_ms).await + } + + pub async fn fail_phase1_publication_signing_retryable( + &self, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + next_attempt_after_ms: i64, + diagnostic: &str, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_time(next_attempt_after_ms)?; + validate_diagnostic(diagnostic)?; + transition_signing_claim( + self, + claim, + now_ms, + "failed-retryable", + Some(diagnostic), + next_attempt_after_ms, + ) + .await + } + + pub async fn fail_phase1_publication_signing_terminal( + &self, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + diagnostic: &str, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_diagnostic(diagnostic)?; + transition_signing_claim( + self, + claim, + now_ms, + "failed-terminal", + Some(diagnostic), + now_ms, + ) + .await + } + + pub async fn quarantine_phase1_publication( + &self, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + diagnostic: &str, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_diagnostic(diagnostic)?; + transition_signing_claim(self, claim, now_ms, "quarantined", Some(diagnostic), now_ms).await + } + + pub async fn cancel_phase1_publication( + &self, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + transition_signing_claim(self, claim, now_ms, "cancelled", None, now_ms).await + } + + pub async fn complete_phase1_publication_signing( + &self, + claim: &RadrootsPhase1PublicationClaim, + verified: &RadrootsVerifiedSignedEvent, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + let record = self.load_phase1_publication(claim.publication_id).await?; + validate_signed_matches_artifact(verified.signed_event(), &record.ready_artifact)?; + let signed_json = verified.signed_event().raw_json().as_bytes(); + if signed_json.is_empty() + || signed_json.len() > RADROOTS_PHASE1_PUBLICATION_SIGNED_EVENT_MAX_BYTES + { + return Err(RadrootsPhase1PublicationError::SignedEventInvalid); + } + let signed_digest: [u8; 32] = Sha256::digest(signed_json).into(); + let signed_event_id = decode_hex32(verified.signed_event().id_str(), "signed_event_id")?; + let affected = sqlx::query( + "UPDATE outbox_phase1_publication + SET state = 'signed-ready', state_revision = state_revision + 1, + claim_token = NULL, claim_expires_at_ms = NULL, + signed_event_json = ?, signed_event_digest = ?, signed_event_id = ?, + last_error = NULL, next_attempt_after_ms = ?, updated_at_ms = ? + WHERE publication_id = ? AND state = 'claimed-for-signing' AND state_revision = ? + AND claim_token = ? AND claim_expires_at_ms > ? AND signed_event_json IS NULL", + ) + .bind(signed_json) + .bind(signed_digest.as_slice()) + .bind(signed_event_id.as_slice()) + .bind(now_ms) + .bind(now_ms) + .bind(claim.publication_id) + .bind(i64_from_u64(claim.revision, "state_revision")?) + .bind(claim.token.as_slice()) + .bind(now_ms) + .execute(&self.pool) + .await? + .rows_affected(); + if affected != 1 { + return Err(RadrootsPhase1PublicationError::ClaimInvalid); + } + self.load_phase1_publication(claim.publication_id).await + } + + pub async fn claim_phase1_publication_target( + &self, + publication_id: i64, + observed_publication_revision: u64, + target_id: i64, + observed_target_revision: u64, + now_ms: i64, + lease_millis: i64, + ) -> Result<RadrootsPhase1PublicationTargetClaim, RadrootsPhase1PublicationError> { + let record = self.load_phase1_publication(publication_id).await?; + if record.revision != observed_publication_revision + || record.signed_event.is_none() + || !matches!( + record.state, + RadrootsPhase1PublicationEventState::SignedReady + | RadrootsPhase1PublicationEventState::Dispatching + ) + { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let target = record + .targets + .iter() + .find(|target| target.target_id == target_id) + .ok_or(RadrootsPhase1PublicationError::TargetNotFound { target_id })?; + if target.revision != observed_target_revision { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let expires_at_ms = validated_expiry(now_ms, lease_millis)?; + let token = new_claim_token()?; + let signed_digest: [u8; 32] = Sha256::digest( + record + .signed_event + .as_ref() + .expect("checked signed event") + .signed_event() + .raw_json() + .as_bytes(), + ) + .into(); + let mut transaction = self.pool.begin().await?; + let target_affected = sqlx::query( + "UPDATE outbox_phase1_delivery_target + SET state = 'in-flight', state_revision = state_revision + 1, + claim_token = ?, claim_expires_at_ms = ?, updated_at_ms = ? + WHERE target_id = ? AND publication_id = ? AND state_revision = ? + AND state IN ('pending', 'failed-retryable', 'uncertain') + AND (claim_token IS NULL OR claim_expires_at_ms <= ?)", + ) + .bind(token.as_slice()) + .bind(expires_at_ms) + .bind(now_ms) + .bind(target_id) + .bind(publication_id) + .bind(i64_from_u64(observed_target_revision, "target_revision")?) + .bind(now_ms) + .execute(&mut *transaction) + .await? + .rows_affected(); + if target_affected != 1 { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let publication_affected = sqlx::query( + "UPDATE outbox_phase1_publication + SET state = 'dispatching', state_revision = state_revision + 1, updated_at_ms = ? + WHERE publication_id = ? AND state_revision = ? AND state IN ('signed-ready', 'dispatching')", + ) + .bind(now_ms) + .bind(publication_id) + .bind(i64_from_u64( + observed_publication_revision, + "publication_revision", + )?) + .execute(&mut *transaction) + .await? + .rows_affected(); + if publication_affected != 1 { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + sqlx::query( + "INSERT INTO outbox_phase1_dispatch_intent(intent_digest, target_id, signed_event_digest, state, state_revision, created_at_ms, updated_at_ms) + VALUES (?, ?, ?, 'in-flight', 0, ?, ?) + ON CONFLICT(intent_digest) DO UPDATE SET + state = 'in-flight', state_revision = outbox_phase1_dispatch_intent.state_revision + 1, + updated_at_ms = excluded.updated_at_ms + WHERE outbox_phase1_dispatch_intent.target_id = excluded.target_id + AND outbox_phase1_dispatch_intent.signed_event_digest = excluded.signed_event_digest", + ) + .bind(target.dispatch_digest.as_slice()) + .bind(target_id) + .bind(signed_digest.as_slice()) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *transaction) + .await?; + transaction.commit().await?; + Ok(RadrootsPhase1PublicationTargetClaim { + publication_id, + publication_revision: observed_publication_revision + 1, + target_id, + target_revision: observed_target_revision + 1, + token, + expires_at_ms, + }) + } + + pub async fn complete_phase1_target_accepted_pending( + &self, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + complete_target( + self, + claim, + now_ms, + RadrootsPhase1PublicationTargetState::AcceptedObservationPending, + None, + now_ms, + ) + .await + } + + pub async fn complete_phase1_target_accepted_observed( + &self, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + complete_target( + self, + claim, + now_ms, + RadrootsPhase1PublicationTargetState::AcceptedObserved, + None, + now_ms, + ) + .await + } + + pub async fn fail_phase1_target_retryable( + &self, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + next_attempt_after_ms: i64, + diagnostic: &str, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_time(next_attempt_after_ms)?; + validate_diagnostic(diagnostic)?; + complete_target( + self, + claim, + now_ms, + RadrootsPhase1PublicationTargetState::FailedRetryable, + Some(diagnostic), + next_attempt_after_ms, + ) + .await + } + + pub async fn fail_phase1_target_terminal( + &self, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + diagnostic: &str, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_diagnostic(diagnostic)?; + complete_target( + self, + claim, + now_ms, + RadrootsPhase1PublicationTargetState::FailedTerminal, + Some(diagnostic), + now_ms, + ) + .await + } + + pub async fn mark_phase1_target_uncertain( + &self, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + diagnostic: &str, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_diagnostic(diagnostic)?; + complete_target( + self, + claim, + now_ms, + RadrootsPhase1PublicationTargetState::Uncertain, + Some(diagnostic), + now_ms, + ) + .await + } + + pub async fn cancel_phase1_target( + &self, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + complete_target( + self, + claim, + now_ms, + RadrootsPhase1PublicationTargetState::Cancelled, + None, + now_ms, + ) + .await + } + + pub async fn complete_phase1_observation_repair( + &self, + publication_id: i64, + target_id: i64, + observed_target_revision: u64, + now_ms: i64, + ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + let record = self.load_phase1_publication(publication_id).await?; + let target = record + .targets + .iter() + .find(|target| target.target_id == target_id) + .ok_or(RadrootsPhase1PublicationError::TargetNotFound { target_id })?; + if target.revision != observed_target_revision + || target.state != RadrootsPhase1PublicationTargetState::AcceptedObservationPending + { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let repair_digest = repair_digest(&target.dispatch_digest); + let receipt_digest = receipt_digest(&target.dispatch_digest, "accepted-observed"); + let mut transaction = self.pool.begin().await?; + let affected = sqlx::query( + "UPDATE outbox_phase1_delivery_target SET state = 'accepted-observed', state_revision = state_revision + 1, updated_at_ms = ? WHERE target_id = ? AND publication_id = ? AND state = 'accepted-observation-pending' AND state_revision = ?", + ) + .bind(now_ms) + .bind(target_id) + .bind(publication_id) + .bind(i64_from_u64(observed_target_revision, "target_revision")?) + .execute(&mut *transaction) + .await? + .rows_affected(); + if affected != 1 { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let repaired = sqlx::query( + "UPDATE outbox_phase1_observation_repair SET state = 'complete', state_revision = state_revision + 1, updated_at_ms = ? WHERE repair_digest = ? AND target_id = ? AND state = 'pending'", + ) + .bind(now_ms) + .bind(repair_digest.as_slice()) + .bind(target_id) + .execute(&mut *transaction) + .await? + .rows_affected(); + if repaired != 1 { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + sqlx::query( + "INSERT INTO outbox_phase1_target_receipt(receipt_digest, target_id, observation_kind, observed_at_ms) VALUES (?, ?, 'accepted-observed', ?)", + ) + .bind(receipt_digest.as_slice()) + .bind(target_id) + .bind(now_ms) + .execute(&mut *transaction) + .await?; + transaction.commit().await?; + self.load_phase1_publication(publication_id).await + } +} + +async fn transition_signing_claim( + outbox: &RadrootsOutbox, + claim: &RadrootsPhase1PublicationClaim, + now_ms: i64, + destination: &'static str, + diagnostic: Option<&str>, + next_attempt_after_ms: i64, +) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + outbox.load_phase1_publication(claim.publication_id).await?; + let affected = sqlx::query( + "UPDATE outbox_phase1_publication + SET state = ?, state_revision = state_revision + 1, claim_token = NULL, + claim_expires_at_ms = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? + WHERE publication_id = ? AND state = 'claimed-for-signing' AND state_revision = ? + AND claim_token = ? AND claim_expires_at_ms > ? AND signed_event_json IS NULL", + ) + .bind(destination) + .bind(diagnostic) + .bind(next_attempt_after_ms) + .bind(now_ms) + .bind(claim.publication_id) + .bind(i64_from_u64(claim.revision, "state_revision")?) + .bind(claim.token.as_slice()) + .bind(now_ms) + .execute(&outbox.pool) + .await? + .rows_affected(); + if affected != 1 { + return Err(RadrootsPhase1PublicationError::ClaimInvalid); + } + outbox.load_phase1_publication(claim.publication_id).await +} + +async fn complete_target( + outbox: &RadrootsOutbox, + claim: &RadrootsPhase1PublicationTargetClaim, + now_ms: i64, + destination: RadrootsPhase1PublicationTargetState, + diagnostic: Option<&str>, + next_attempt_after_ms: i64, +) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + let record = outbox.load_phase1_publication(claim.publication_id).await?; + if record.revision != claim.publication_revision { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let target = record + .targets + .iter() + .find(|target| target.target_id == claim.target_id) + .ok_or(RadrootsPhase1PublicationError::TargetNotFound { + target_id: claim.target_id, + })?; + if target.revision != claim.target_revision { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + let mut transaction = outbox.pool.begin().await?; + let affected = sqlx::query( + "UPDATE outbox_phase1_delivery_target + SET state = ?, state_revision = state_revision + 1, claim_token = NULL, + claim_expires_at_ms = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? + WHERE target_id = ? AND publication_id = ? AND state = 'in-flight' AND state_revision = ? + AND claim_token = ? AND claim_expires_at_ms > ?", + ) + .bind(destination.as_str()) + .bind(diagnostic) + .bind(next_attempt_after_ms) + .bind(now_ms) + .bind(claim.target_id) + .bind(claim.publication_id) + .bind(i64_from_u64(claim.target_revision, "target_revision")?) + .bind(claim.token.as_slice()) + .bind(now_ms) + .execute(&mut *transaction) + .await? + .rows_affected(); + if affected != 1 { + return Err(RadrootsPhase1PublicationError::ClaimInvalid); + } + match destination { + RadrootsPhase1PublicationTargetState::AcceptedObservationPending + | RadrootsPhase1PublicationTargetState::AcceptedObserved => { + let observation_kind = if destination + == RadrootsPhase1PublicationTargetState::AcceptedObservationPending + { + "accepted-pending" + } else { + "accepted-observed" + }; + sqlx::query( + "INSERT INTO outbox_phase1_target_receipt(receipt_digest, target_id, observation_kind, observed_at_ms) VALUES (?, ?, ?, ?)", + ) + .bind(receipt_digest(&target.dispatch_digest, observation_kind).as_slice()) + .bind(claim.target_id) + .bind(observation_kind) + .bind(now_ms) + .execute(&mut *transaction) + .await?; + if destination == RadrootsPhase1PublicationTargetState::AcceptedObservationPending { + sqlx::query( + "INSERT INTO outbox_phase1_observation_repair(repair_digest, target_id, state, state_revision, created_at_ms, updated_at_ms) VALUES (?, ?, 'pending', 0, ?, ?)", + ) + .bind(repair_digest(&target.dispatch_digest).as_slice()) + .bind(claim.target_id) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *transaction) + .await?; + } + sqlx::query( + "UPDATE outbox_phase1_dispatch_intent SET state = 'completed', state_revision = state_revision + 1, updated_at_ms = ? WHERE intent_digest = ? AND target_id = ?", + ) + .bind(now_ms) + .bind(target.dispatch_digest.as_slice()) + .bind(claim.target_id) + .execute(&mut *transaction) + .await?; + } + RadrootsPhase1PublicationTargetState::Uncertain => { + sqlx::query( + "UPDATE outbox_phase1_dispatch_intent SET state = 'uncertain', state_revision = state_revision + 1, updated_at_ms = ? WHERE intent_digest = ? AND target_id = ?", + ) + .bind(now_ms) + .bind(target.dispatch_digest.as_slice()) + .bind(claim.target_id) + .execute(&mut *transaction) + .await?; + } + RadrootsPhase1PublicationTargetState::Cancelled => { + sqlx::query( + "UPDATE outbox_phase1_dispatch_intent SET state = 'cancelled', state_revision = state_revision + 1, updated_at_ms = ? WHERE intent_digest = ? AND target_id = ?", + ) + .bind(now_ms) + .bind(target.dispatch_digest.as_slice()) + .bind(claim.target_id) + .execute(&mut *transaction) + .await?; + } + RadrootsPhase1PublicationTargetState::FailedRetryable + | RadrootsPhase1PublicationTargetState::FailedTerminal => { + sqlx::query( + "UPDATE outbox_phase1_dispatch_intent SET state = 'completed', state_revision = state_revision + 1, updated_at_ms = ? WHERE intent_digest = ? AND target_id = ?", + ) + .bind(now_ms) + .bind(target.dispatch_digest.as_slice()) + .bind(claim.target_id) + .execute(&mut *transaction) + .await?; + } + RadrootsPhase1PublicationTargetState::Pending + | RadrootsPhase1PublicationTargetState::InFlight => { + return Err(RadrootsPhase1PublicationError::StateConflict); + } + } + let destination_event = aggregate_publication_state( + &mut transaction, + claim.publication_id, + record.target_policy.required_target_count, + ) + .await?; + let publication_affected = sqlx::query( + "UPDATE outbox_phase1_publication SET state = ?, state_revision = state_revision + 1, updated_at_ms = ? WHERE publication_id = ? AND state_revision = ? AND state = 'dispatching'", + ) + .bind(destination_event.as_str()) + .bind(now_ms) + .bind(claim.publication_id) + .bind(i64_from_u64( + claim.publication_revision, + "publication_revision", + )?) + .execute(&mut *transaction) + .await? + .rows_affected(); + if publication_affected != 1 { + return Err(RadrootsPhase1PublicationError::RevisionConflict); + } + transaction.commit().await?; + outbox.load_phase1_publication(claim.publication_id).await +} + +async fn aggregate_publication_state( + transaction: &mut Transaction<'_, Sqlite>, + publication_id: i64, + required_target_count: usize, +) -> Result<RadrootsPhase1PublicationEventState, RadrootsPhase1PublicationError> { + let row = sqlx::query( + "SELECT + SUM(CASE WHEN state IN ('accepted-observation-pending', 'accepted-observed') THEN 1 ELSE 0 END) AS accepted, + SUM(CASE WHEN state IN ('pending', 'in-flight', 'failed-retryable', 'uncertain') THEN 1 ELSE 0 END) AS recoverable + FROM outbox_phase1_delivery_target WHERE publication_id = ?", + ) + .bind(publication_id) + .fetch_one(&mut **transaction) + .await?; + let accepted = usize_from_i64(row.try_get("accepted")?, "accepted_target_count")?; + let recoverable = usize_from_i64(row.try_get("recoverable")?, "recoverable_target_count")?; + if accepted >= required_target_count { + Ok(RadrootsPhase1PublicationEventState::Published) + } else if accepted + recoverable < required_target_count { + Ok(RadrootsPhase1PublicationEventState::FailedTerminal) + } else { + Ok(RadrootsPhase1PublicationEventState::SignedReady) + } +} + +async fn load_targets( + pool: &sqlx::SqlitePool, + publication_id: i64, +) -> Result<Vec<RadrootsPhase1PublicationTarget>, RadrootsPhase1PublicationError> { + let rows = sqlx::query( + "SELECT target_id, target_ordinal, + length(CAST(endpoint_uri AS BLOB)) AS endpoint_uri_bytes, + CASE WHEN length(CAST(endpoint_uri AS BLOB)) <= ? THEN endpoint_uri END AS bounded_endpoint_uri, + endpoint_fingerprint, dispatch_digest, state, state_revision, next_attempt_after_ms + FROM outbox_phase1_delivery_target WHERE publication_id = ? + ORDER BY target_ordinal LIMIT 17", + ) + .bind(i64::try_from(RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES).unwrap_or(i64::MAX)) + .bind(publication_id) + .fetch_all(pool) + .await?; + if rows.len() > RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT { + return Err(RadrootsPhase1PublicationError::TargetCount { + max: RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, + actual: rows.len(), + }); + } + rows.into_iter() + .map(|row| { + let actual = usize_from_i64(row.try_get("endpoint_uri_bytes")?, "endpoint_uri_bytes")?; + if actual > RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES { + return Err(RadrootsPhase1PublicationError::StoredValueTooLarge { + field: "endpoint_uri", + max: RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES, + actual, + }); + } + let endpoint_uri: Option<String> = row.try_get("bounded_endpoint_uri")?; + let endpoint_uri = + endpoint_uri.ok_or(RadrootsPhase1PublicationError::StoredAuthorityInvalid)?; + Ok(RadrootsPhase1PublicationTarget { + target_id: row.try_get("target_id")?, + ordinal: usize_from_i64(row.try_get("target_ordinal")?, "target_ordinal")?, + endpoint_uri, + endpoint_fingerprint: blob32(&row, "endpoint_fingerprint")?, + dispatch_digest: blob32(&row, "dispatch_digest")?, + state: RadrootsPhase1PublicationTargetState::parse( + &row.try_get::<String, _>("state")?, + )?, + revision: u64_from_i64(row.try_get("state_revision")?, "target_revision")?, + next_attempt_after_ms: row.try_get("next_attempt_after_ms")?, + }) + }) + .collect() +} + +fn prepare_publication( + ready: &RadrootsPhase1MediaReadyPublicationArtifact, + target_policy: &RadrootsPhase1PublicationTargetPolicy, +) -> Result<PreparedPublication, RadrootsPhase1PublicationError> { + validate_phase1_publication_media_readiness(ready).map_err(|error| { + RadrootsPhase1PublicationError::ReadinessInvalid { + source_code: error.code(), + } + })?; + let artifact_json = ready.artifact().to_canonical_json(); + let allowlisted = allow_phase1_publication_canonical_json(&artifact_json).map_err(|error| { + RadrootsPhase1PublicationError::ArtifactInvalid { + source_code: error.code(), + } + })?; + let readiness_json = ready.to_canonical_json(); + let reloaded = RadrootsPhase1MediaReadyPublicationArtifact::from_canonical_json( + allowlisted, + &readiness_json, + ) + .map_err(|error| RadrootsPhase1PublicationError::ReadinessInvalid { + source_code: error.code(), + })?; + if &reloaded != ready { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + let artifact = ready.artifact(); + let artifact_digest = *artifact.artifact_digest().as_bytes(); + let readiness_digest = *ready.binding_digest().as_bytes(); + let expected_author = decode_hex32(artifact.expected_author().as_str(), "expected_author")?; + let expected_event_id = + decode_hex32(artifact.expected_event_id().as_str(), "expected_event_id")?; + let operation_digest = operation_digest_for( + &artifact_digest, + &readiness_digest, + &expected_author, + &target_policy.digest, + ); + Ok(PreparedPublication { + artifact_json, + artifact_digest, + readiness_json, + readiness_digest, + semantic_role: artifact.semantic_variant().as_str(), + expected_author, + expected_event_id, + operation_digest, + }) +} + +fn validate_signed_matches_artifact( + signed: &RadrootsSignedEvent, + ready: &RadrootsPhase1MediaReadyPublicationArtifact, +) -> Result<(), RadrootsPhase1PublicationError> { + let artifact = ready.artifact(); + let draft = artifact.draft(); + if signed.pubkey_str() != artifact.expected_author().as_str() + || signed.id_str() != artifact.expected_event_id().as_str() + || signed.created_at() != draft.created_at() + || signed.kind() != draft.kind() + || signed.tags_as_vec() != draft.tags() + || signed.content() != draft.content() + { + return Err(RadrootsPhase1PublicationError::SignedEventMismatch); + } + Ok(()) +} + +fn reload_signed_event( + signed_json: Option<Vec<u8>>, + signed_digest: Option<Vec<u8>>, + signed_event_id: Option<Vec<u8>>, + ready: &RadrootsPhase1MediaReadyPublicationArtifact, +) -> Result<Option<RadrootsVerifiedSignedEvent>, RadrootsPhase1PublicationError> { + match (signed_json, signed_digest, signed_event_id) { + (None, None, None) => Ok(None), + (Some(bytes), Some(stored_digest), Some(stored_event_id)) => { + if stored_digest.as_slice() != Sha256::digest(&bytes).as_slice() + || stored_event_id.len() != 32 + { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + let raw = String::from_utf8(bytes) + .map_err(|_| RadrootsPhase1PublicationError::SignedEventInvalid)?; + let wire = RadrootsNip01EventWire::parse_json(&raw) + .map_err(|_| RadrootsPhase1PublicationError::SignedEventInvalid)?; + let signed = RadrootsSignedEvent::from_wire_verified_id(wire, raw) + .map_err(|_| RadrootsPhase1PublicationError::SignedEventInvalid)?; + if decode_hex32(signed.id_str(), "signed_event_id")?.as_slice() + != stored_event_id.as_slice() + { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + validate_signed_matches_artifact(&signed, ready)?; + let verified = signed + .verify_signature() + .map_err(|_| RadrootsPhase1PublicationError::SignedEventInvalid)?; + Ok(Some(verified)) + } + _ => Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid), + } +} + +fn validate_target_identities( + targets: &[RadrootsPhase1PublicationTarget], + expected_event_id: &[u8; 32], + target_policy_digest: &[u8; 32], +) -> Result<(), RadrootsPhase1PublicationError> { + let mut ids = BTreeSet::new(); + for (ordinal, target) in targets.iter().enumerate() { + if target.ordinal != ordinal + || canonical_target_uri(&target.endpoint_uri)? != target.endpoint_uri + || target.endpoint_fingerprint != endpoint_fingerprint(&target.endpoint_uri) + || target.dispatch_digest + != dispatch_digest( + expected_event_id, + target_policy_digest, + &target.endpoint_fingerprint, + ) + || !ids.insert(target.target_id) + { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + } + Ok(()) +} + +fn canonical_target_uri(value: &str) -> Result<String, RadrootsPhase1PublicationError> { + if value.is_empty() || value.len() > RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES { + return Err(RadrootsPhase1PublicationError::TargetUriTooLarge { + max: RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES, + actual: value.len(), + }); + } + let parsed = Url::parse(value).map_err(|_| RadrootsPhase1PublicationError::TargetUriInvalid)?; + if !matches!(parsed.scheme(), "ws" | "wss") + || parsed.host_str().is_none() + || !parsed.username().is_empty() + || parsed.password().is_some() + || parsed.fragment().is_some() + { + return Err(RadrootsPhase1PublicationError::TargetUriInvalid); + } + let canonical = parsed.to_string(); + if canonical.len() > RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES { + return Err(RadrootsPhase1PublicationError::TargetUriTooLarge { + max: RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES, + actual: canonical.len(), + }); + } + Ok(canonical) +} + +fn target_policy_digest( + targets: &[String], + required_target_count: usize, +) -> Result<[u8; 32], RadrootsPhase1PublicationError> { + let mut digest = Sha256::new(); + digest.update(TARGET_POLICY_DOMAIN); + digest.update( + u16::try_from(required_target_count) + .map_err(|_| RadrootsPhase1PublicationError::IntegerRange { + field: "required_target_count", + })? + .to_be_bytes(), + ); + digest.update( + u16::try_from(targets.len()) + .map_err(|_| RadrootsPhase1PublicationError::IntegerRange { + field: "target_count", + })? + .to_be_bytes(), + ); + for target in targets { + digest.update(endpoint_fingerprint(target)); + } + Ok(digest.finalize().into()) +} + +fn endpoint_fingerprint(target: &str) -> [u8; 32] { + let mut digest = Sha256::new(); + digest.update(ENDPOINT_DOMAIN); + digest.update((target.len() as u64).to_be_bytes()); + digest.update(target.as_bytes()); + digest.finalize().into() +} + +fn operation_digest_for( + artifact_digest: &[u8; 32], + readiness_digest: &[u8; 32], + expected_author: &[u8; 32], + target_policy_digest: &[u8; 32], +) -> [u8; 32] { + let mut digest = Sha256::new(); + digest.update(OPERATION_DOMAIN); + digest.update(artifact_digest); + digest.update(readiness_digest); + digest.update(expected_author); + digest.update(target_policy_digest); + digest.finalize().into() +} + +fn dispatch_digest( + event_id: &[u8; 32], + target_policy_digest: &[u8; 32], + endpoint_fingerprint: &[u8; 32], +) -> [u8; 32] { + let mut digest = Sha256::new(); + digest.update(DISPATCH_DOMAIN); + digest.update(event_id); + digest.update(target_policy_digest); + digest.update(endpoint_fingerprint); + digest.finalize().into() +} + +fn receipt_digest(dispatch_digest: &[u8; 32], observation_kind: &str) -> [u8; 32] { + let mut digest = Sha256::new(); + digest.update(RECEIPT_DOMAIN); + digest.update(dispatch_digest); + digest.update((observation_kind.len() as u64).to_be_bytes()); + digest.update(observation_kind.as_bytes()); + digest.finalize().into() +} + +fn repair_digest(dispatch_digest: &[u8; 32]) -> [u8; 32] { + let mut digest = Sha256::new(); + digest.update(REPAIR_DOMAIN); + digest.update(dispatch_digest); + digest.finalize().into() +} + +fn new_claim_token() -> Result<[u8; 32], RadrootsPhase1PublicationError> { + let mut token = [0_u8; 32]; + getrandom::getrandom(&mut token) + .map_err(|_| RadrootsPhase1PublicationError::EntropyUnavailable)?; + Ok(token) +} + +fn validated_expiry(now_ms: i64, lease_millis: i64) -> Result<i64, RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + if !(1..=RADROOTS_PHASE1_PUBLICATION_CLAIM_LEASE_MAX_MILLIS).contains(&lease_millis) { + return Err(RadrootsPhase1PublicationError::InvalidLease { + max_millis: RADROOTS_PHASE1_PUBLICATION_CLAIM_LEASE_MAX_MILLIS, + actual: lease_millis, + }); + } + now_ms + .checked_add(lease_millis) + .ok_or(RadrootsPhase1PublicationError::InvalidTime) +} + +fn validate_time(value: i64) -> Result<(), RadrootsPhase1PublicationError> { + if value < 0 { + Err(RadrootsPhase1PublicationError::InvalidTime) + } else { + Ok(()) + } +} + +fn validate_diagnostic(value: &str) -> Result<(), RadrootsPhase1PublicationError> { + if value.len() > RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES { + Err(RadrootsPhase1PublicationError::DiagnosticTooLarge { + max: RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES, + actual: value.len(), + }) + } else { + Ok(()) + } +} + +fn decode_hex32( + value: &str, + field: &'static str, +) -> Result<[u8; 32], RadrootsPhase1PublicationError> { + let bytes = hex::decode(value) + .map_err(|_| RadrootsPhase1PublicationError::StoredDigestInvalid { field })?; + bytes + .try_into() + .map_err(|_| RadrootsPhase1PublicationError::StoredDigestInvalid { field }) +} + +fn blob32( + row: &sqlx::sqlite::SqliteRow, + field: &'static str, +) -> Result<[u8; 32], RadrootsPhase1PublicationError> { + let bytes: Vec<u8> = row.try_get(field)?; + bytes + .try_into() + .map_err(|_| RadrootsPhase1PublicationError::StoredDigestInvalid { field }) +} + +fn bounded_blob( + row: &sqlx::sqlite::SqliteRow, + value_field: &'static str, + length_field: &'static str, + authority_field: &'static str, + max: usize, +) -> Result<Option<Vec<u8>>, RadrootsPhase1PublicationError> { + let length: Option<i64> = row.try_get(length_field)?; + let Some(length) = length else { + return Ok(None); + }; + let actual = usize_from_i64(length, length_field)?; + if actual > max { + return Err(RadrootsPhase1PublicationError::StoredValueTooLarge { + field: authority_field, + max, + actual, + }); + } + row.try_get(value_field).map_err(Into::into) +} + +fn usize_from_i64( + value: i64, + field: &'static str, +) -> Result<usize, RadrootsPhase1PublicationError> { + usize::try_from(value).map_err(|_| RadrootsPhase1PublicationError::IntegerRange { field }) +} + +fn u64_from_i64(value: i64, field: &'static str) -> Result<u64, RadrootsPhase1PublicationError> { + u64::try_from(value).map_err(|_| RadrootsPhase1PublicationError::IntegerRange { field }) +} + +fn i64_from_u64(value: u64, field: &'static str) -> Result<i64, RadrootsPhase1PublicationError> { + i64::try_from(value).map_err(|_| RadrootsPhase1PublicationError::IntegerRange { field }) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::RadrootsOutboxRollbackConfirmation; + use radroots_event::post::RadrootsAuthoredUpdate; + use radroots_event_codec::wire::publication::allowlist::allow_phase1_publication_artifact; + use radroots_event_codec::wire::publication::{ + RadrootsPhase1PublicationArtifact, bind_phase1_publication_media_readiness, + }; + use radroots_nostr::prelude::{ + RadrootsNostrKeys, RadrootsNostrSecretKey, RadrootsNostrTimestamp, + radroots_nostr_build_update_event, + }; + + const ALICE_SECRET_KEY: &str = + "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5"; + const ALICE_PUBLIC_KEY: &str = + "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; + const DESCRIPTOR: &[u8] = include_bytes!("../contracts/phase1_publication_v1.descriptor.json"); + + fn ready_update() -> RadrootsPhase1MediaReadyPublicationArtifact { + ready_update_with_content("Carrots harvested in Victoria today") + } + + fn ready_update_with_content(content: &str) -> RadrootsPhase1MediaReadyPublicationArtifact { + let update = RadrootsAuthoredUpdate::new(content).unwrap(); + let artifact = RadrootsPhase1PublicationArtifact::from_update( + &update, + 1_780_000_000, + ALICE_PUBLIC_KEY, + ) + .unwrap(); + let allowlisted = allow_phase1_publication_artifact(artifact).unwrap(); + bind_phase1_publication_media_readiness(allowlisted, Vec::new()).unwrap() + } + + fn signed_update( + ready: &RadrootsPhase1MediaReadyPublicationArtifact, + ) -> RadrootsVerifiedSignedEvent { + let artifact = ready.artifact(); + let draft = artifact.draft(); + let update = RadrootsAuthoredUpdate::new("Carrots harvested in Victoria today").unwrap(); + let secret = RadrootsNostrSecretKey::from_hex(ALICE_SECRET_KEY).unwrap(); + let event = radroots_nostr_build_update_event(&update) + .unwrap() + .custom_created_at(RadrootsNostrTimestamp::from_secs(draft.created_at())) + .sign_with_keys(&RadrootsNostrKeys::new(secret)) + .unwrap(); + let raw = serde_json::to_string(&event).unwrap(); + let wire = RadrootsNip01EventWire::parse_json(&raw).unwrap(); + let signed = RadrootsSignedEvent::from_wire_verified_id(wire, raw).unwrap(); + assert_eq!(signed.id_str(), artifact.expected_event_id().as_str()); + signed.verify_signature().unwrap() + } + + #[test] + fn phase1_publication_transition_matrix_is_closed_and_unique() { + let event_states = BTreeSet::from([ + "ready", + "claimed-for-signing", + "signed-ready", + "dispatching", + "published", + "failed-retryable", + "failed-terminal", + "quarantined", + "cancelled", + ]); + let target_states = BTreeSet::from([ + "pending", + "in-flight", + "accepted-observation-pending", + "accepted-observed", + "failed-retryable", + "failed-terminal", + "uncertain", + "cancelled", + ]); + let mut ids = BTreeSet::new(); + for transition in RADROOTS_PHASE1_PUBLICATION_TRANSITIONS { + assert!(ids.insert(transition.id)); + let states = match transition.scope { + RadrootsPhase1PublicationTransitionScope::Event => &event_states, + RadrootsPhase1PublicationTransitionScope::Target => &target_states, + }; + assert!(states.contains(transition.from)); + assert!(states.contains(transition.to)); + assert!(transition.revision_cas); + assert!(!transition.lease_predicate.is_empty()); + assert!(!transition.durable_side_effect.is_empty()); + } + assert_eq!(ids.len(), 25); + } + + #[test] + fn phase1_publication_descriptor_matches_runtime_authority() { + let descriptor: serde_json::Value = serde_json::from_slice(DESCRIPTOR).unwrap(); + assert_eq!( + descriptor["resource_limits"], + serde_json::json!({ + "target_count": RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, + "target_uri_bytes": RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES, + "diagnostic_bytes": RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES, + "claim_lease_millis": RADROOTS_PHASE1_PUBLICATION_CLAIM_LEASE_MAX_MILLIS, + }) + ); + assert_eq!( + descriptor["stable_errors"] + .as_array() + .unwrap() + .iter() + .map(|value| value.as_str().unwrap()) + .collect::<Vec<_>>(), + RADROOTS_PHASE1_PUBLICATION_ERROR_CODES + ); + let transitions = descriptor["transitions"].as_array().unwrap(); + assert_eq!( + transitions.len(), + RADROOTS_PHASE1_PUBLICATION_TRANSITIONS.len() + ); + for (wire, runtime) in transitions + .iter() + .zip(RADROOTS_PHASE1_PUBLICATION_TRANSITIONS) + { + assert_eq!(wire["id"], runtime.id); + assert_eq!(wire["scope"], runtime.scope.as_str()); + assert_eq!(wire["from"], runtime.from); + assert_eq!(wire["to"], runtime.to); + assert_eq!(wire["revision_cas"], runtime.revision_cas); + assert_eq!(wire["lease_predicate"], runtime.lease_predicate); + assert_eq!(wire["durable_side_effect"], runtime.durable_side_effect); + assert_eq!(wire["retry_class"], runtime.retry_class.as_str()); + assert_eq!(wire["repair_edge"], runtime.repair_edge); + assert_eq!(wire["terminal_destination"], runtime.terminal_destination); + } + } + + #[test] + fn phase1_publication_target_policy_is_canonical_bounded_and_stable() { + let policy = RadrootsPhase1PublicationTargetPolicy::new( + ["wss://B.example:443", "wss://a.example/relay"], + 1, + ) + .unwrap(); + assert_eq!( + policy.targets(), + ["wss://a.example/relay", "wss://b.example/"] + ); + assert_eq!( + policy, + RadrootsPhase1PublicationTargetPolicy::new( + ["wss://a.example/relay", "wss://b.example/"], + 1, + ) + .unwrap() + ); + assert!(matches!( + RadrootsPhase1PublicationTargetPolicy::new(Vec::<String>::new(), 0), + Err(RadrootsPhase1PublicationError::RequiredTargetCount { .. }) + )); + assert!(matches!( + RadrootsPhase1PublicationTargetPolicy::new( + ["wss://same.example", "wss://same.example/"], + 1 + ), + Err(RadrootsPhase1PublicationError::DuplicateTarget) + )); + let one_over = format!("wss://example.com/{}", "a".repeat(2_048)); + assert!(matches!( + RadrootsPhase1PublicationTargetPolicy::new([one_over], 1), + Err(RadrootsPhase1PublicationError::TargetUriTooLarge { .. }) + )); + let targets = (0..=RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT) + .map(|index| format!("wss://relay-{index}.example")) + .collect::<Vec<_>>(); + assert!(matches!( + RadrootsPhase1PublicationTargetPolicy::new(targets, 1), + Err(RadrootsPhase1PublicationError::TargetCount { .. }) + )); + + let prefix = "wss://example.com/"; + let exact = format!( + "{prefix}{}", + "a".repeat(RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES - prefix.len()) + ); + assert_eq!( + exact.len(), + RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES + ); + assert!(RadrootsPhase1PublicationTargetPolicy::new([exact.clone()], 1).is_ok()); + assert!(matches!( + RadrootsPhase1PublicationTargetPolicy::new([format!("{exact}a")], 1), + Err(RadrootsPhase1PublicationError::TargetUriTooLarge { .. }) + )); + let exact_targets = (0..RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT) + .map(|index| format!("wss://relay-{index}.example")) + .collect::<Vec<_>>(); + assert!(RadrootsPhase1PublicationTargetPolicy::new(exact_targets, 16).is_ok()); + } + + #[tokio::test] + async fn phase1_publication_enqueue_is_typed_idempotent_and_revalidated() { + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let ready = ready_update(); + let policy = RadrootsPhase1PublicationTargetPolicy::new( + ["wss://relay-b.example", "wss://relay-a.example"], + 1, + ) + .unwrap(); + let inserted = outbox + .enqueue_phase1_publication(&ready, &policy, 10) + .await + .unwrap(); + assert_eq!( + inserted.status(), + RadrootsPhase1PublicationEnqueueStatus::Inserted + ); + assert_eq!( + inserted.record().state(), + RadrootsPhase1PublicationEventState::Ready + ); + assert_eq!(inserted.record().targets().len(), 2); + assert_eq!( + inserted.record().targets()[0].endpoint_uri(), + "wss://relay-a.example/" + ); + let duplicate = outbox + .enqueue_phase1_publication(&ready, &policy, 99) + .await + .unwrap(); + assert_eq!( + duplicate.status(), + RadrootsPhase1PublicationEnqueueStatus::Existing + ); + assert_eq!( + duplicate.record().operation_digest(), + inserted.record().operation_digest() + ); + let conflicting_policy = + RadrootsPhase1PublicationTargetPolicy::new(["wss://other.example"], 1).unwrap(); + assert_eq!( + outbox + .enqueue_phase1_publication(&ready, &conflicting_policy, 100) + .await + .unwrap_err() + .code(), + "phase1_publication_idempotency_conflict" + ); + + sqlx::query( + "UPDATE outbox_phase1_publication SET readiness_digest = zeroblob(32) WHERE publication_id = ?", + ) + .bind(inserted.record().publication_id()) + .execute(outbox.pool()) + .await + .unwrap(); + assert_eq!( + outbox + .load_phase1_publication(inserted.record().publication_id()) + .await + .unwrap_err() + .code(), + "phase1_publication_stored_authority_invalid" + ); + } + + #[tokio::test] + async fn phase1_publication_enqueue_failure_rolls_back_every_typed_row() { + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + sqlx::query( + "CREATE TEMP TRIGGER phase1_publication_fail_target BEFORE INSERT ON outbox_phase1_delivery_target BEGIN SELECT RAISE(ABORT, 'injected target failure'); END", + ) + .execute(outbox.pool()) + .await + .unwrap(); + let result = outbox + .enqueue_phase1_publication( + &ready_update(), + &RadrootsPhase1PublicationTargetPolicy::new(["wss://relay.example"], 1).unwrap(), + 1, + ) + .await; + assert!(matches!( + result, + Err(RadrootsPhase1PublicationError::Sqlite(_)) + )); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM outbox_phase1_publication") + .fetch_one(outbox.pool()) + .await + .unwrap(), + 0 + ); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM outbox_phase1_delivery_target") + .fetch_one(outbox.pool()) + .await + .unwrap(), + 0 + ); + } + + #[tokio::test] + async fn phase1_publication_reload_rejects_malformed_swapped_and_cross_row_authority() { + for mutation in [ + "UPDATE outbox_phase1_publication SET artifact_json = x'7b7d' WHERE publication_id = ?", + "UPDATE outbox_phase1_publication SET artifact_json = substr(artifact_json, 1, length(artifact_json) - 1) WHERE publication_id = ?", + "UPDATE outbox_phase1_publication SET artifact_json = CAST(artifact_json || x'20' AS BLOB) WHERE publication_id = ?", + "UPDATE outbox_phase1_publication SET artifact_json = readiness_json WHERE publication_id = ?", + "UPDATE outbox_phase1_delivery_target SET endpoint_uri = 'wss://retarget.example/' WHERE publication_id = ?", + ] { + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let receipt = outbox + .enqueue_phase1_publication( + &ready_update(), + &RadrootsPhase1PublicationTargetPolicy::new(["wss://relay.example"], 1) + .unwrap(), + 1, + ) + .await + .unwrap(); + sqlx::query(mutation) + .bind(receipt.record().publication_id()) + .execute(outbox.pool()) + .await + .unwrap(); + assert!( + outbox + .load_phase1_publication(receipt.record().publication_id()) + .await + .is_err(), + "mutation must fail closed: {mutation}" + ); + } + + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let policy = + RadrootsPhase1PublicationTargetPolicy::new(["wss://relay.example"], 1).unwrap(); + let first = outbox + .enqueue_phase1_publication(&ready_update(), &policy, 1) + .await + .unwrap(); + let second = outbox + .enqueue_phase1_publication( + &ready_update_with_content("A second Victoria harvest update"), + &policy, + 2, + ) + .await + .unwrap(); + sqlx::query( + "UPDATE outbox_phase1_publication SET artifact_json = (SELECT artifact_json FROM outbox_phase1_publication WHERE publication_id = ?) WHERE publication_id = ?", + ) + .bind(second.record().publication_id()) + .bind(first.record().publication_id()) + .execute(outbox.pool()) + .await + .unwrap(); + assert_eq!( + outbox + .load_phase1_publication(first.record().publication_id()) + .await + .unwrap_err() + .code(), + "phase1_publication_readiness_invalid" + ); + } + + #[tokio::test] + async fn phase1_publication_claims_are_lease_fenced_and_race_safe() { + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let receipt = outbox + .enqueue_phase1_publication( + &ready_update(), + &RadrootsPhase1PublicationTargetPolicy::new(["wss://relay.example"], 1).unwrap(), + 10, + ) + .await + .unwrap(); + let publication_id = receipt.record().publication_id(); + let first = outbox + .claim_phase1_publication_for_signing(publication_id, 0, 20, 10) + .await + .unwrap(); + assert_eq!( + outbox + .claim_phase1_publication_for_signing(publication_id, 0, 20, 10) + .await + .unwrap_err() + .code(), + "phase1_publication_revision_conflict" + ); + let reclaimed = outbox + .claim_phase1_publication_for_signing(publication_id, 1, 30, 10) + .await + .unwrap(); + assert_eq!( + outbox + .renew_phase1_publication_claim(&first, 31, 10) + .await + .unwrap_err() + .code(), + "phase1_publication_claim_invalid" + ); + let retryable = outbox + .fail_phase1_publication_signing_retryable(&reclaimed, 31, 40, "signer unavailable") + .await + .unwrap(); + assert_eq!( + retryable.state(), + RadrootsPhase1PublicationEventState::FailedRetryable + ); + let final_claim = outbox + .claim_phase1_publication_for_signing(publication_id, retryable.revision(), 40, 10) + .await + .unwrap(); + let cancelled = outbox + .cancel_phase1_publication(&final_claim, 41) + .await + .unwrap(); + assert_eq!( + cancelled.state(), + RadrootsPhase1PublicationEventState::Cancelled + ); + } + + #[tokio::test] + async fn phase1_publication_signed_dispatch_and_observation_repair_are_durable() { + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let ready = ready_update(); + let receipt = outbox + .enqueue_phase1_publication( + &ready, + &RadrootsPhase1PublicationTargetPolicy::new( + ["wss://relay-a.example", "wss://relay-b.example"], + 1, + ) + .unwrap(), + 100, + ) + .await + .unwrap(); + let claim = outbox + .claim_phase1_publication_for_signing( + receipt.record().publication_id(), + receipt.record().revision(), + 101, + 100, + ) + .await + .unwrap(); + let signed = signed_update(&ready); + let signed_record = outbox + .complete_phase1_publication_signing(&claim, &signed, 102) + .await + .unwrap(); + assert_eq!( + signed_record.state(), + RadrootsPhase1PublicationEventState::SignedReady + ); + assert_eq!( + signed_record + .signed_event() + .unwrap() + .signed_event() + .raw_json(), + signed.signed_event().raw_json() + ); + let target = &signed_record.targets()[0]; + let target_claim = outbox + .claim_phase1_publication_target( + signed_record.publication_id(), + signed_record.revision(), + target.target_id(), + target.revision(), + 103, + 100, + ) + .await + .unwrap(); + let pending = outbox + .complete_phase1_target_accepted_pending(&target_claim, 104) + .await + .unwrap(); + assert_eq!( + pending.state(), + RadrootsPhase1PublicationEventState::Published + ); + let repaired_target = pending + .targets() + .iter() + .find(|candidate| candidate.target_id() == target_claim.target_id()) + .unwrap(); + assert_eq!( + repaired_target.state(), + RadrootsPhase1PublicationTargetState::AcceptedObservationPending + ); + let repaired = outbox + .complete_phase1_observation_repair( + pending.publication_id(), + repaired_target.target_id(), + repaired_target.revision(), + 105, + ) + .await + .unwrap(); + assert_eq!( + repaired.state(), + RadrootsPhase1PublicationEventState::Published + ); + assert_eq!( + repaired.targets()[0].state(), + RadrootsPhase1PublicationTargetState::AcceptedObserved + ); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM outbox_phase1_dispatch_intent") + .fetch_one(outbox.pool()) + .await + .unwrap(), + 1 + ); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM outbox_phase1_target_receipt") + .fetch_one(outbox.pool()) + .await + .unwrap(), + 2 + ); + } + + #[tokio::test] + async fn phase1_publication_reopen_upgrade_and_exact_rollback_are_transactional() { + let temp = tempfile::tempdir().unwrap(); + let path = temp.path().join("phase1.sqlite"); + let ready = ready_update(); + let policy = + RadrootsPhase1PublicationTargetPolicy::new(["wss://relay.example"], 1).unwrap(); + let outbox = RadrootsOutbox::open_file(&path).await.unwrap(); + let receipt = outbox + .enqueue_phase1_publication(&ready, &policy, 1) + .await + .unwrap(); + let publication_id = receipt.record().publication_id(); + outbox.close().await; + + let reopened = RadrootsOutbox::open_file(&path).await.unwrap(); + assert_eq!( + reopened + .load_phase1_publication(publication_id) + .await + .unwrap() + .operation_digest(), + receipt.record().operation_digest() + ); + reopened.close().await; + RadrootsOutbox::rollback_file_schema_offline( + &path, + 1, + RadrootsOutboxRollbackConfirmation::acknowledge_data_loss(), + ) + .await + .unwrap(); + let remigrated = RadrootsOutbox::open_file(&path).await.unwrap(); + assert!(matches!( + remigrated.load_phase1_publication(publication_id).await, + Err(RadrootsPhase1PublicationError::PublicationNotFound { .. }) + )); + } + + #[test] + fn phase1_publication_diagnostics_enforce_exact_and_one_over() { + assert!( + validate_diagnostic(&"a".repeat(RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES)) + .is_ok() + ); + assert!(matches!( + validate_diagnostic(&"a".repeat(RADROOTS_PHASE1_PUBLICATION_DIAGNOSTIC_MAX_BYTES + 1)), + Err(RadrootsPhase1PublicationError::DiagnosticTooLarge { .. }) + )); + } +} diff --git a/crates/outbox/src/schema.rs b/crates/outbox/src/schema.rs @@ -917,8 +917,14 @@ pub(crate) async fn validate_outbox_owned_integrity( connection: &mut SqliteConnection, ) -> Result<(), RadrootsOutboxError> { validate_embedded_migration_registry()?; + let current = match inspect_outbox_schema_on_connection(connection).await? { + RadrootsOutboxSchemaStatus::Managed { version } => version, + RadrootsOutboxSchemaStatus::UnledgeredBaseline => RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + RadrootsOutboxSchemaStatus::Uninitialized => return Ok(()), + }; let tables = OUTBOX_MIGRATIONS .iter() + .filter(|migration| migration.version <= current) .flat_map(|migration| migration.owned_table_names.iter().copied()) .collect::<BTreeSet<_>>(); for table in tables { @@ -1068,7 +1074,9 @@ mod tests { inspect_outbox_schema_status(&pool) .await .expect("managed schema"), - RadrootsOutboxSchemaStatus::Managed { version: 1 } + RadrootsOutboxSchemaStatus::Managed { + version: RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + } ); } @@ -1089,15 +1097,15 @@ mod tests { .expect("ledger count") } - async fn synthetic_v2_registry(pool: &SqlitePool) -> [OutboxMigration; 2] { + async fn synthetic_v3_registry(pool: &SqlitePool) -> [OutboxMigration; 3] { const UP_SQL: &str = "CREATE TABLE outbox_future (value TEXT NOT NULL) STRICT;\n"; const DOWN_SQL: &str = "DROP TABLE outbox_future;\n"; - migrate_outbox_schema(pool).await.expect("version 1"); + migrate_outbox_schema(pool).await.expect("version 2"); sqlx::raw_sql(UP_SQL) .execute(pool) .await - .expect("synthetic version 2 schema"); + .expect("synthetic version 3 schema"); let catalog = sqlx::query( "SELECT type, name, tbl_name, sql FROM main.sqlite_schema WHERE lower(substr(name, 1, 7)) != 'sqlite_' @@ -1110,7 +1118,7 @@ mod tests { .bind(OUTBOX_RESERVED_PREFIX) .fetch_all(pool) .await - .expect("synthetic version 2 catalog") + .expect("synthetic version 3 catalog") .into_iter() .map(|row| CatalogRow { object_type: row.try_get("type").expect("catalog type"), @@ -1123,7 +1131,7 @@ mod tests { sqlx::raw_sql(DOWN_SQL) .execute(pool) .await - .expect("remove synthetic version 2 schema"); + .expect("remove synthetic version 3 schema"); let up_sha256 = Box::leak(crate::migrations::sha256_hex(UP_SQL.as_bytes()).into_boxed_str()); @@ -1131,8 +1139,9 @@ mod tests { Box::leak(crate::migrations::sha256_hex(DOWN_SQL.as_bytes()).into_boxed_str()); [ OUTBOX_MIGRATIONS[0], + OUTBOX_MIGRATIONS[1], OutboxMigration { - version: 2, + version: 3, name: "future", up_sql: UP_SQL, down_sql: DOWN_SQL, @@ -1168,16 +1177,16 @@ mod tests { .execute(&pool) .await .expect("caller state"); - let registry = synthetic_v2_registry(&pool).await; + let registry = synthetic_v3_registry(&pool).await; - migrate_outbox_schema_with_registry(&pool, &registry, 1, 2) + migrate_outbox_schema_with_registry(&pool, &registry, 1, 3) .await - .expect("advance to version 2"); + .expect("advance to version 3"); assert_eq!( - inspect_with_registry(&pool, &registry, 2) + inspect_with_registry(&pool, &registry, 3) .await - .expect("managed version 2"), - RadrootsOutboxSchemaStatus::Managed { version: 2 } + .expect("managed version 3"), + RadrootsOutboxSchemaStatus::Managed { version: 3 } ); let future_objects: i64 = sqlx::query_scalar( "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'outbox_future'", @@ -1191,15 +1200,15 @@ mod tests { .begin_with("BEGIN EXCLUSIVE") .await .expect("rollback transaction"); - let result = rollback_schema_on_connection(&mut transaction, &registry, 2, 1).await; + let result = rollback_schema_on_connection(&mut transaction, &registry, 3, 2).await; finish_schema_transaction(transaction, result) .await - .expect("rollback to version 1"); + .expect("rollback to version 2"); assert_eq!( - inspect_with_registry(&pool, &registry, 2) + inspect_with_registry(&pool, &registry, 3) .await - .expect("managed version 1"), - RadrootsOutboxSchemaStatus::Managed { version: 1 } + .expect("managed version 2"), + RadrootsOutboxSchemaStatus::Managed { version: 2 } ); let future_objects: i64 = sqlx::query_scalar( "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'outbox_future'", @@ -1218,7 +1227,7 @@ mod tests { .begin_with("BEGIN EXCLUSIVE") .await .expect("ahead transaction"); - let result = rollback_schema_on_connection(&mut transaction, &registry, 2, 2).await; + let result = rollback_schema_on_connection(&mut transaction, &registry, 3, 3).await; assert!(matches!( result, Err(RadrootsOutboxError::RollbackAhead { .. }) @@ -1233,7 +1242,7 @@ mod tests { .begin_with("BEGIN EXCLUSIVE") .await .expect("unmanaged transaction"); - let result = rollback_schema_on_connection(&mut transaction, &registry, 2, 1).await; + let result = rollback_schema_on_connection(&mut transaction, &registry, 3, 2).await; assert!(matches!( result, Err(RadrootsOutboxError::RollbackUnmanaged) @@ -1254,27 +1263,27 @@ mod tests { .execute(&pool) .await .expect("caller state"); - let mut registry = synthetic_v2_registry(&pool).await; + let mut registry = synthetic_v3_registry(&pool).await; let mut connection = pool.acquire().await.expect("catalog connection"); let before_catalog = read_catalog_bounded(&mut connection, &registry) .await .expect("before catalog"); let before_fingerprint = catalog_fingerprint(&governed_catalog(&before_catalog, &registry)); - let before_history = read_history_bounded(&mut connection, 2) + let before_history = read_history_bounded(&mut connection, 3) .await .expect("before history"); drop(connection); - registry[1].owned_object_names = &["outbox_expected"]; - registry[1].owned_table_names = &["outbox_expected"]; - let error = migrate_outbox_schema_with_registry(&pool, &registry, 1, 2) + registry[2].owned_object_names = &["outbox_expected"]; + registry[2].owned_table_names = &["outbox_expected"]; + let error = migrate_outbox_schema_with_registry(&pool, &registry, 1, 3) .await .expect_err("post-UP catalog delta must fail"); assert!(matches!( error, RadrootsOutboxError::MigrationCatalogDeltaMismatch { - version: 2, + version: 3, direction: "up", .. } @@ -1285,7 +1294,7 @@ mod tests { .await .expect("after catalog"); let after_fingerprint = catalog_fingerprint(&governed_catalog(&after_catalog, &registry)); - let after_history = read_history_bounded(&mut connection, 2) + let after_history = read_history_bounded(&mut connection, 3) .await .expect("after history"); drop(connection); @@ -1295,7 +1304,7 @@ mod tests { inspect_outbox_schema_status(&pool) .await .expect("restored managed schema"), - RadrootsOutboxSchemaStatus::Managed { version: 1 } + RadrootsOutboxSchemaStatus::Managed { version: 2 } ); let future_objects: i64 = sqlx::query_scalar( "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'outbox_future'", @@ -1460,7 +1469,7 @@ mod tests { .fetch_one(&pool) .await .expect("ledger rows"); - assert_eq!(rows, 1); + assert_eq!(rows, i64::from(RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT)); } } @@ -1470,7 +1479,7 @@ mod tests { migrate_outbox_schema(&newer).await.expect("migration"); sqlx::query( "INSERT INTO radroots_outbox_schema_migrations(version, name, up_sha256, down_sha256, schema_sha256) - VALUES (2, 'future', ?, ?, ?)", + VALUES (3, 'future', ?, ?, ?)", ) .bind("a".repeat(64)) .bind("b".repeat(64)) @@ -1481,8 +1490,8 @@ mod tests { assert!(matches!( inspect_outbox_schema_status(&newer).await, Err(RadrootsOutboxError::SchemaTooNew { - current: 1, - database: 2 + current: 2, + database: 3 }) )); @@ -1494,7 +1503,7 @@ mod tests { .expect("unknown governed object"); assert!(matches!( inspect_outbox_schema_status(&overflow).await, - Err(RadrootsOutboxError::GovernedCatalogCapacityExceeded { max: 14 }) + Err(RadrootsOutboxError::GovernedCatalogCapacityExceeded { max: 22 }) )); } @@ -1520,11 +1529,15 @@ mod tests { )); assert!(matches!( validate_history_against_registry( - &[row(1, &OUTBOX_MIGRATIONS[0]), row(2, &OUTBOX_MIGRATIONS[0])], + &[ + row(1, &OUTBOX_MIGRATIONS[0]), + row(2, &OUTBOX_MIGRATIONS[1]), + row(3, &OUTBOX_MIGRATIONS[0]), + ], OUTBOX_MIGRATIONS, 3, ), - Err(RadrootsOutboxError::UnknownMigration { version: 2 }) + Err(RadrootsOutboxError::UnknownMigration { version: 3 }) )); assert!(matches!( diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -48,7 +48,7 @@ use std::sync::Arc; #[derive(Clone)] pub struct RadrootsOutbox { - pool: SqlitePool, + pub(crate) pool: SqlitePool, _file_lease: Option<Arc<OutboxFileLease>>, } diff --git a/crates/outbox/tests/fixtures/phase1_publication.v1.json b/crates/outbox/tests/fixtures/phase1_publication.v1.json @@ -0,0 +1,34 @@ +{ + "schema_version": 1, + "contract_id": "radroots_outbox.phase1_publication.v1", + "executor": { + "id": "radroots_outbox.phase1_publication.v1.result_vector_executor.v1", + "path": "crates/outbox/tests/phase1_publication_v1_result_vector.rs", + "test": "phase1_publication_v1_result_vector" + }, + "identity_vector": { + "fixture": "update", + "target_uri": "wss://relay.example/", + "required_target_count": 1, + "artifact_digest": "9ad318496bd4a710fbc7f3e0f6d5a01de808d352db91ca56a12e710904784cc5", + "readiness_digest": "bf9c1ff2a2c26d62d45b9f5a1eea727ce8f8d7ee92e749364900c12ff1546be5", + "endpoint_fingerprint": "59bf3d289a2344bdc61c426484877df2f122d3936115f45e68e812d2189f21d2", + "target_policy_digest": "e7577fbe841b140993f9cf317e1f562d926b64eb19625220a7c0bfaa8bce4fbb", + "operation_digest": "01a16b215292e90af4c7a76556a86b54d7910ed01757acc8429bda5f1c5fb72b", + "dispatch_digest": "0e35ef0a2769586300ecb3b4760e6a034a60ee68344126dc875168e3876618a4" + }, + "cases": [ + { "id": "identity_preimages", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "typed_enqueue", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "duplicate_enqueue", "execution": "direct_executor", "expected_outcome": "accepted_idempotent" }, + { "id": "target_count_exact", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "target_count_one_over", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_target_count" }, + { "id": "target_uri_exact", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "target_uri_one_over", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_target_uri_too_large" }, + { "id": "empty_required_policy", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_required_target_count" }, + { "id": "two_worker_claim_race", "execution": "direct_executor", "expected_outcome": "one_winner" }, + { "id": "expired_lease_reclaim", "execution": "direct_executor", "expected_outcome": "accepted" }, + { "id": "stale_claim_rejected", "execution": "direct_executor", "expected_outcome": "rejected", "expected_error": "phase1_publication_claim_invalid" }, + { "id": "migration_rollback_reopen", "execution": "direct_executor", "expected_outcome": "accepted" } + ] +} diff --git a/crates/outbox/tests/phase1_publication_v1_result_vector.rs b/crates/outbox/tests/phase1_publication_v1_result_vector.rs @@ -0,0 +1,251 @@ +#![forbid(unsafe_code)] +#![cfg(feature = "sqlite")] + +use radroots_event_codec::wire::publication::allowlist::allow_phase1_publication_canonical_json; +use radroots_event_codec::wire::publication::{ + RadrootsPhase1MediaReadyPublicationArtifact, bind_phase1_publication_media_readiness, +}; +use radroots_outbox::{ + RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT, RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES, + RadrootsOutbox, RadrootsOutboxRollbackConfirmation, RadrootsPhase1PublicationEnqueueStatus, + RadrootsPhase1PublicationError, RadrootsPhase1PublicationTargetPolicy, +}; +use serde::Deserialize; +use serde_json::Value; +use std::collections::BTreeSet; + +const VECTOR: &[u8] = include_bytes!("fixtures/phase1_publication.v1.json"); +const ARTIFACT_VECTOR: &[u8] = + include_bytes!("../../event_codec/tests/fixtures/phase1_publication_artifact.v1.json"); + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Vector { + schema_version: u32, + contract_id: String, + executor: Executor, + identity_vector: IdentityVector, + cases: Vec<Case>, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Executor { + id: String, + path: String, + test: String, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct IdentityVector { + fixture: String, + target_uri: String, + required_target_count: usize, + artifact_digest: String, + readiness_digest: String, + endpoint_fingerprint: String, + target_policy_digest: String, + operation_digest: String, + dispatch_digest: String, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Case { + id: String, + execution: String, + expected_outcome: String, + expected_error: Option<String>, +} + +fn ready_fixture(fixture: &str) -> RadrootsPhase1MediaReadyPublicationArtifact { + let root: Value = serde_json::from_slice(ARTIFACT_VECTOR).unwrap(); + let canonical = root["vectors"] + .as_array() + .unwrap() + .iter() + .find(|vector| { + vector["input"]["fixture"].as_str() == Some(fixture) + && vector["kind"] + .as_str() + .is_some_and(|kind| kind.ends_with(".valid")) + }) + .and_then(|vector| vector["expected"]["canonical_json"].as_str()) + .unwrap(); + let allowlisted = allow_phase1_publication_canonical_json(canonical.as_bytes()).unwrap(); + bind_phase1_publication_media_readiness(allowlisted, Vec::new()).unwrap() +} + +#[tokio::test] +async fn phase1_publication_v1_result_vector() { + let vector: Vector = serde_json::from_slice(VECTOR).unwrap(); + assert_eq!(vector.schema_version, 1); + assert_eq!(vector.contract_id, "radroots_outbox.phase1_publication.v1"); + assert_eq!( + vector.executor.id, + "radroots_outbox.phase1_publication.v1.result_vector_executor.v1" + ); + assert_eq!( + vector.executor.path, + "crates/outbox/tests/phase1_publication_v1_result_vector.rs" + ); + assert_eq!(vector.executor.test, "phase1_publication_v1_result_vector"); + assert_eq!( + vector + .cases + .iter() + .map(|case| case.id.as_str()) + .collect::<BTreeSet<_>>(), + BTreeSet::from([ + "duplicate_enqueue", + "empty_required_policy", + "expired_lease_reclaim", + "identity_preimages", + "migration_rollback_reopen", + "stale_claim_rejected", + "target_count_exact", + "target_count_one_over", + "target_uri_exact", + "target_uri_one_over", + "two_worker_claim_race", + "typed_enqueue", + ]) + ); + for case in &vector.cases { + assert_eq!(case.execution, "direct_executor"); + assert!(!case.expected_outcome.is_empty()); + assert_eq!( + case.expected_error.is_some(), + case.expected_outcome == "rejected" + ); + } + + let identity = &vector.identity_vector; + assert_eq!(identity.fixture, "update"); + let ready = ready_fixture(&identity.fixture); + assert_eq!( + ready.artifact().artifact_digest().to_hex(), + identity.artifact_digest + ); + assert_eq!(ready.binding_digest().to_hex(), identity.readiness_digest); + let policy = RadrootsPhase1PublicationTargetPolicy::new( + [identity.target_uri.as_str()], + identity.required_target_count, + ) + .unwrap(); + assert_eq!(hex::encode(policy.digest()), identity.target_policy_digest); + + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let inserted = outbox + .enqueue_phase1_publication(&ready, &policy, 1) + .await + .unwrap(); + assert_eq!( + inserted.status(), + RadrootsPhase1PublicationEnqueueStatus::Inserted + ); + assert_eq!( + hex::encode(inserted.record().operation_digest()), + identity.operation_digest + ); + assert_eq!( + hex::encode(inserted.record().targets()[0].endpoint_fingerprint()), + identity.endpoint_fingerprint + ); + assert_eq!( + hex::encode(inserted.record().targets()[0].dispatch_digest()), + identity.dispatch_digest + ); + assert_eq!( + outbox + .enqueue_phase1_publication(&ready, &policy, 2) + .await + .unwrap() + .status(), + RadrootsPhase1PublicationEnqueueStatus::Existing + ); + + let exact_targets = (0..RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT) + .map(|index| format!("wss://relay-{index}.example")) + .collect::<Vec<_>>(); + assert!(RadrootsPhase1PublicationTargetPolicy::new(exact_targets, 16).is_ok()); + let one_over_targets = (0..=RADROOTS_PHASE1_PUBLICATION_TARGET_MAX_COUNT) + .map(|index| format!("wss://relay-{index}.example")) + .collect::<Vec<_>>(); + assert_eq!( + RadrootsPhase1PublicationTargetPolicy::new(one_over_targets, 1) + .unwrap_err() + .code(), + "phase1_publication_target_count" + ); + let prefix = "wss://example.com/"; + let exact_uri = format!( + "{prefix}{}", + "a".repeat(RADROOTS_PHASE1_PUBLICATION_TARGET_URI_MAX_BYTES - prefix.len()) + ); + assert!(RadrootsPhase1PublicationTargetPolicy::new([exact_uri.clone()], 1).is_ok()); + assert_eq!( + RadrootsPhase1PublicationTargetPolicy::new([format!("{exact_uri}a")], 1) + .unwrap_err() + .code(), + "phase1_publication_target_uri_too_large" + ); + assert_eq!( + RadrootsPhase1PublicationTargetPolicy::new(Vec::<String>::new(), 0) + .unwrap_err() + .code(), + "phase1_publication_required_target_count" + ); + + let publication_id = inserted.record().publication_id(); + let first_worker = outbox.clone(); + let second_worker = outbox.clone(); + let (first, second) = tokio::join!( + first_worker.claim_phase1_publication_for_signing(publication_id, 0, 10, 10), + second_worker.claim_phase1_publication_for_signing(publication_id, 0, 10, 10), + ); + let (winner, loser) = match (first, second) { + (Ok(winner), Err(loser)) | (Err(loser), Ok(winner)) => (winner, loser), + _ => panic!("exactly one worker must win the revision CAS"), + }; + assert_eq!(loser.code(), "phase1_publication_revision_conflict"); + let reclaimed = outbox + .claim_phase1_publication_for_signing(publication_id, 1, 20, 10) + .await + .unwrap(); + assert_eq!( + outbox + .renew_phase1_publication_claim(&winner, 21, 10) + .await + .unwrap_err() + .code(), + "phase1_publication_claim_invalid" + ); + outbox + .release_phase1_publication_claim(&reclaimed, 21) + .await + .unwrap(); + + let temp = tempfile::tempdir().unwrap(); + let path = temp.path().join("vector.sqlite"); + let file_outbox = RadrootsOutbox::open_file(&path).await.unwrap(); + let file_record = file_outbox + .enqueue_phase1_publication(&ready, &policy, 1) + .await + .unwrap(); + let file_publication_id = file_record.record().publication_id(); + file_outbox.close().await; + RadrootsOutbox::rollback_file_schema_offline( + &path, + 1, + RadrootsOutboxRollbackConfirmation::acknowledge_data_loss(), + ) + .await + .unwrap(); + let reopened = RadrootsOutbox::open_file(&path).await.unwrap(); + assert!(matches!( + reopened.load_phase1_publication(file_publication_id).await, + Err(RadrootsPhase1PublicationError::PublicationNotFound { .. }) + )); +} diff --git a/tools/xtask/src/contract.rs b/tools/xtask/src/contract.rs @@ -10,6 +10,7 @@ mod deletion_authority; mod food_availability_projection; mod nip09_reconciliation; mod outbox_migration; +mod outbox_phase1_publication; #[allow(dead_code)] mod phase1_publication_allowlist; mod phase1_publication_media_readiness; @@ -38,6 +39,9 @@ pub(crate) use nip09_reconciliation::{ pub(crate) use outbox_migration::{ validate_outbox_migration_manifest, write_outbox_migration_manifest, }; +pub(crate) use outbox_phase1_publication::{ + validate_outbox_phase1_publication_manifest, write_outbox_phase1_publication_manifest, +}; pub(crate) use phase1_publication_media_readiness::{ validate_immutable_blossom_publication_readiness_predecessor, validate_immutable_phase1_publication_allowlist_predecessor, @@ -87,6 +91,7 @@ pub(crate) fn validate_artifact_contracts(workspace_root: &Path) -> Result<(), S validate_phase1_publication_media_readiness_manifest(workspace_root)?; validate_blossom_raster_decoder_security_manifest(workspace_root)?; validate_outbox_migration_manifest(workspace_root)?; + validate_outbox_phase1_publication_manifest(workspace_root)?; validate_knowledge_contract_manifest(workspace_root) } @@ -150,12 +155,15 @@ const RAW_SOURCE_REBUILD_CONFORMANCE_VECTOR_RELATIVE: &str = "contracts/conformance/vectors/event_store/raw_source_rebuild.v1.json"; const OUTBOX_MIGRATION_CONFORMANCE_VECTOR_RELATIVE: &str = "contracts/conformance/vectors/outbox/migration_authority.v1.json"; -const SPECIALIZED_CONFORMANCE_VECTOR_RELATIVES: [&str; 5] = [ +const OUTBOX_PHASE1_PUBLICATION_CONFORMANCE_VECTOR_RELATIVE: &str = + "contracts/conformance/vectors/outbox/phase1_publication.v1.json"; +const SPECIALIZED_CONFORMANCE_VECTOR_RELATIVES: [&str; 6] = [ NIP09_RECONCILIATION_CONFORMANCE_VECTOR_RELATIVE, FOOD_AVAILABILITY_PROJECTION_CONFORMANCE_VECTOR_RELATIVE, SOURCE_MAINTENANCE_CONFORMANCE_VECTOR_RELATIVE, RAW_SOURCE_REBUILD_CONFORMANCE_VECTOR_RELATIVE, OUTBOX_MIGRATION_CONFORMANCE_VECTOR_RELATIVE, + OUTBOX_PHASE1_PUBLICATION_CONFORMANCE_VECTOR_RELATIVE, ]; const KNOWLEDGE_MANIFEST_RELATIVE: &str = "contracts/knowledge/knowledge_event_contract_manifest.v2.json"; @@ -177,7 +185,7 @@ const REPLICA_CONTRACT_RELATIVE: &str = "contracts/replica.toml"; const REPLICA_CONTRACT_NAME: &str = "radroots_replica_contract"; const REPLICA_TRANSFER_CONSTANT: &str = "RADROOTS_REPLICA_TRANSFER_VERSION"; const REPLICA_TRANSFER_VERSION: u32 = 2; -const CONFORMANCE_VECTOR_MIRRORS: [(&str, &str); 30] = [ +const CONFORMANCE_VECTOR_MIRRORS: [(&str, &str); 31] = [ ( "contracts/conformance/vectors/blossom/bud11_claims.v1.json", "crates/blossom/tests/fixtures/bud11_claims.v1.json", @@ -203,6 +211,10 @@ const CONFORMANCE_VECTOR_MIRRORS: [(&str, &str); 30] = [ "crates/outbox/tests/fixtures/migration_authority.v1.json", ), ( + OUTBOX_PHASE1_PUBLICATION_CONFORMANCE_VECTOR_RELATIVE, + "crates/outbox/tests/fixtures/phase1_publication.v1.json", + ), + ( "contracts/conformance/vectors/blossom/bud11_nostr_adapter.v1.json", "crates/nostr/tests/fixtures/bud11_nostr_adapter.v1.json", ), diff --git a/tools/xtask/src/contract/outbox_migration.rs b/tools/xtask/src/contract/outbox_migration.rs @@ -809,6 +809,11 @@ fn generated_runtime_registry(registry: &Registry) -> Result<String, String> { "pub(crate) const OUTBOX_MIGRATIONS: &[OutboxMigration] = &[OUTBOX_MIGRATION_{:04}];\n", registry.migrations[0].version )); + } else if registry.migrations.len() == 2 { + generated.push_str(&format!( + "pub(crate) const OUTBOX_MIGRATIONS: &[OutboxMigration] =\n &[OUTBOX_MIGRATION_{:04}, OUTBOX_MIGRATION_{:04}];\n", + registry.migrations[0].version, registry.migrations[1].version + )); } else { generated.push_str("pub(crate) const OUTBOX_MIGRATIONS: &[OutboxMigration] = &[\n"); for migration in &registry.migrations { @@ -1128,20 +1133,29 @@ mod tests { #[test] fn outbox_registry_accepts_an_appended_successor_and_rejects_a_gap() { let mut registry = load_and_validate_registry(&workspace_root()).expect("live registry"); + let successor_index = registry.migrations.len(); + let successor_version = registry + .migrations + .last() + .expect("non-empty live registry") + .version + + 1; let mut successor = registry.migrations[0].clone(); - successor.version = 2; + successor.version = successor_version; successor.name = "future".to_owned(); - successor.up_path = "crates/outbox/migrations/0002_future.up.sql".to_owned(); - successor.down_path = "crates/outbox/migrations/0002_future.down.sql".to_owned(); + successor.up_path = + format!("crates/outbox/migrations/{successor_version:04}_future.up.sql"); + successor.down_path = + format!("crates/outbox/migrations/{successor_version:04}_future.down.sql"); successor.owned_objects = vec!["outbox_future".to_owned()]; successor.owned_tables = vec!["outbox_future".to_owned()]; registry.migrations.push(successor); validate_registry_shape(&registry).expect("contiguous successor"); - registry.migrations[1].version = 3; + registry.migrations[successor_index].version = successor_version + 1; assert!( validate_registry_shape(&registry) .expect_err("gap") - .contains("expected 2") + .contains(&format!("expected {successor_version}")) ); } diff --git a/tools/xtask/src/contract/outbox_phase1_publication.rs b/tools/xtask/src/contract/outbox_phase1_publication.rs @@ -0,0 +1,638 @@ +use super::artifact_bundle::{ + GeneratedArtifact, read_regular_file, with_artifact_bundle_transaction, +}; +use serde::Deserialize; +use serde_json::{Value, json}; +use sha2::{Digest, Sha256}; +use std::collections::BTreeSet; +use std::path::Path; + +const CONTRACT_ID: &str = "radroots_outbox.phase1_publication.v1"; +const WRITE_COMMAND: &str = "cargo xtask contract outbox-phase1-publication-manifest --write"; +const DESCRIPTOR_RELATIVE: &str = "crates/outbox/contracts/phase1_publication_v1.descriptor.json"; +const MANIFEST_RELATIVE: &str = "crates/outbox/contracts/phase1_publication_v1.manifest.json"; +const MANIFEST_SCHEMA_RELATIVE: &str = + "crates/outbox/contracts/phase1_publication_v1.manifest.schema.json"; +const MANIFEST_SHA256_RELATIVE: &str = + "crates/outbox/contracts/phase1_publication_v1.manifest.sha256"; +const VECTOR_RELATIVE: &str = "contracts/conformance/vectors/outbox/phase1_publication.v1.json"; +const VECTOR_MIRROR_RELATIVE: &str = "crates/outbox/tests/fixtures/phase1_publication.v1.json"; +const VECTOR_EXECUTOR_RELATIVE: &str = "crates/outbox/tests/phase1_publication_v1_result_vector.rs"; +const VECTOR_EXECUTOR_ID: &str = "radroots_outbox.phase1_publication.v1.result_vector_executor.v1"; +const VECTOR_EXECUTOR_TEST: &str = "phase1_publication_v1_result_vector"; +const RELEASE_RELATIVE: &str = "contracts/releases/1.0.0-alpha.1.toml"; +const CHANGELOG_RELATIVE: &str = "CHANGELOG.md"; +const RELEASE_CHANGE_ID: &str = "outbox-phase1-publication-state"; +const MIGRATION_REGISTRY_RELATIVE: &str = "crates/outbox/contracts/migration_registry.v1.json"; +const MIGRATION_UP_RELATIVE: &str = "crates/outbox/migrations/0002_phase1_publication.up.sql"; +const MIGRATION_DOWN_RELATIVE: &str = "crates/outbox/migrations/0002_phase1_publication.down.sql"; +const MIGRATION_UP_SHA256: &str = + "84f0c9897cff8d002961cb6ad9dee53edcf28853d1407483519b00bdbf029308"; +const MIGRATION_DOWN_SHA256: &str = + "57a5a00ca4257973097acf7f5cc64494dd0bc73fcfa00af2e6c8bb1f61823928"; +const MIGRATION_SCHEMA_SHA256: &str = + "a56af9ba400fd51c97d48886fbb3f3733adb97458d7109fa8989c1b7e0c8bcaf"; + +const SOURCE_FILES: &[(&str, &str)] = &[ + ("outbox_package_manifest", "crates/outbox/Cargo.toml"), + ("outbox_public_surface", "crates/outbox/src/lib.rs"), + ( + "phase1_publication_runtime", + "crates/outbox/src/phase1_publication.rs", + ), + ("outbox_schema_runtime", "crates/outbox/src/schema.rs"), + ( + "migration_registry_runtime", + "crates/outbox/src/generated/outbox_migration_registry.rs", + ), + ("migration_registry_source", MIGRATION_REGISTRY_RELATIVE), + ("vector_executor", VECTOR_EXECUTOR_RELATIVE), + ( + "contract_governance", + "tools/xtask/src/contract/outbox_phase1_publication.rs", + ), + ("contract_dispatch", "tools/xtask/src/contract.rs"), + ("xtask_dispatch", "tools/xtask/src/main.rs"), + ("release_record", RELEASE_RELATIVE), + ("release_notes", CHANGELOG_RELATIVE), +]; + +const EVENT_STATES: &[&str] = &[ + "ready", + "claimed-for-signing", + "signed-ready", + "dispatching", + "published", + "failed-retryable", + "failed-terminal", + "quarantined", + "cancelled", +]; +const TARGET_STATES: &[&str] = &[ + "pending", + "in-flight", + "accepted-observation-pending", + "accepted-observed", + "failed-retryable", + "failed-terminal", + "uncertain", + "cancelled", +]; +const CASE_IDS: &[&str] = &[ + "duplicate_enqueue", + "empty_required_policy", + "expired_lease_reclaim", + "identity_preimages", + "migration_rollback_reopen", + "stale_claim_rejected", + "target_count_exact", + "target_count_one_over", + "target_uri_exact", + "target_uri_one_over", + "two_worker_claim_race", + "typed_enqueue", +]; +const ERROR_CODES: &[&str] = &[ + "phase1_publication_artifact_invalid", + "phase1_publication_claim_invalid", + "phase1_publication_diagnostic_too_large", + "phase1_publication_entropy_unavailable", + "phase1_publication_idempotency_conflict", + "phase1_publication_integer_range", + "phase1_publication_lease_invalid", + "phase1_publication_not_found", + "phase1_publication_readiness_invalid", + "phase1_publication_required_target_count", + "phase1_publication_revision_conflict", + "phase1_publication_signed_event_invalid", + "phase1_publication_signed_event_mismatch", + "phase1_publication_sqlite", + "phase1_publication_state_conflict", + "phase1_publication_stored_authority_invalid", + "phase1_publication_stored_digest_invalid", + "phase1_publication_stored_state_invalid", + "phase1_publication_stored_value_too_large", + "phase1_publication_target_count", + "phase1_publication_target_duplicate", + "phase1_publication_target_not_found", + "phase1_publication_target_uri_invalid", + "phase1_publication_target_uri_too_large", + "phase1_publication_time_invalid", +]; + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Descriptor { + schema_version: u32, + contract_id: String, + migration: Migration, + resource_limits: ResourceLimits, + operation_identity: Identity, + dispatch_identity: Identity, + event_states: Vec<String>, + target_states: Vec<String>, + transitions: Vec<Transition>, + stable_errors: Vec<String>, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Migration { + version: u32, + name: String, + up_path: String, + up_sha256: String, + down_path: String, + down_sha256: String, + schema_sha256: String, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct ResourceLimits { + target_count: usize, + target_uri_bytes: usize, + diagnostic_bytes: usize, + claim_lease_millis: i64, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Identity { + algorithm: String, + domain: String, + domain_terminator_hex: String, + preimage: Vec<String>, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Transition { + id: String, + scope: String, + from: String, + to: String, + revision_cas: bool, + lease_predicate: String, + durable_side_effect: String, + retry_class: String, + repair_edge: bool, + terminal_destination: bool, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Vector { + schema_version: u32, + contract_id: String, + executor: Executor, + identity_vector: Value, + cases: Vec<Case>, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Executor { + id: String, + path: String, + test: String, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct Case { + id: String, + execution: String, + expected_outcome: String, + expected_error: Option<String>, +} + +pub(crate) fn write_outbox_phase1_publication_manifest( + workspace_root: &Path, +) -> Result<(), String> { + with_artifact_bundle_transaction(workspace_root, |transaction| { + transaction.write(expected_artifacts(workspace_root)?)?; + validate_under_lock(workspace_root) + }) +} + +pub(crate) fn validate_outbox_phase1_publication_manifest( + workspace_root: &Path, +) -> Result<(), String> { + with_artifact_bundle_transaction(workspace_root, |_| validate_under_lock(workspace_root)) +} + +fn validate_under_lock(workspace_root: &Path) -> Result<(), String> { + for artifact in expected_artifacts(workspace_root)? { + let actual = read_regular_file(workspace_root, artifact.relative)?; + if actual != artifact.contents { + return Err(format!( + "generated Phase 1 publication artifact {} is stale; run `{WRITE_COMMAND}`", + artifact.relative + )); + } + } + let manifest_bytes = read_regular_file(workspace_root, MANIFEST_RELATIVE)?; + let manifest: Value = serde_json::from_slice(&manifest_bytes) + .map_err(|error| format!("parse {MANIFEST_RELATIVE}: {error}"))?; + let schema_bytes = read_regular_file(workspace_root, MANIFEST_SCHEMA_RELATIVE)?; + let schema: Value = serde_json::from_slice(&schema_bytes) + .map_err(|error| format!("parse {MANIFEST_SCHEMA_RELATIVE}: {error}"))?; + let validator = jsonschema::validator_for(&schema) + .map_err(|error| format!("compile {MANIFEST_SCHEMA_RELATIVE}: {error}"))?; + let errors = validator + .iter_errors(&manifest) + .map(|error| error.to_string()) + .collect::<Vec<_>>(); + if !errors.is_empty() { + return Err(format!( + "{MANIFEST_RELATIVE} violates its schema: {}", + errors.join("; ") + )); + } + let sidecar = read_regular_file(workspace_root, MANIFEST_SHA256_RELATIVE)?; + if sidecar != format!("{}\n", sha256_hex(&manifest_bytes)).as_bytes() { + return Err(format!( + "{MANIFEST_SHA256_RELATIVE} must authenticate exact manifest bytes" + )); + } + Ok(()) +} + +fn expected_artifacts(workspace_root: &Path) -> Result<Vec<GeneratedArtifact>, String> { + let descriptor = load_descriptor(workspace_root)?; + validate_vector(workspace_root)?; + validate_release(workspace_root)?; + let schema_bytes = canonical_json_bytes(&manifest_schema())?; + let manifest = expected_manifest(workspace_root, &descriptor, &schema_bytes)?; + let manifest_bytes = canonical_json_bytes(&manifest)?; + let vector = read_regular_file(workspace_root, VECTOR_RELATIVE)?; + Ok(vec![ + GeneratedArtifact { + relative: MANIFEST_RELATIVE, + contents: manifest_bytes.clone(), + }, + GeneratedArtifact { + relative: MANIFEST_SCHEMA_RELATIVE, + contents: schema_bytes, + }, + GeneratedArtifact { + relative: MANIFEST_SHA256_RELATIVE, + contents: format!("{}\n", sha256_hex(&manifest_bytes)).into_bytes(), + }, + GeneratedArtifact { + relative: VECTOR_MIRROR_RELATIVE, + contents: vector, + }, + ]) +} + +fn load_descriptor(workspace_root: &Path) -> Result<Descriptor, String> { + let bytes = read_regular_file(workspace_root, DESCRIPTOR_RELATIVE)?; + let descriptor: Descriptor = serde_json::from_slice(&bytes) + .map_err(|error| format!("parse {DESCRIPTOR_RELATIVE}: {error}"))?; + if descriptor.schema_version != 1 || descriptor.contract_id != CONTRACT_ID { + return Err("Phase 1 publication descriptor identity is invalid".to_owned()); + } + if descriptor.migration.version != 2 + || descriptor.migration.name != "phase1_publication" + || descriptor.migration.up_path != MIGRATION_UP_RELATIVE + || descriptor.migration.down_path != MIGRATION_DOWN_RELATIVE + || descriptor.migration.up_sha256 != MIGRATION_UP_SHA256 + || descriptor.migration.down_sha256 != MIGRATION_DOWN_SHA256 + || descriptor.migration.schema_sha256 != MIGRATION_SCHEMA_SHA256 + || sha256_hex(&read_regular_file(workspace_root, MIGRATION_UP_RELATIVE)?) + != MIGRATION_UP_SHA256 + || sha256_hex(&read_regular_file(workspace_root, MIGRATION_DOWN_RELATIVE)?) + != MIGRATION_DOWN_SHA256 + { + return Err("Phase 1 publication migration authority is invalid".to_owned()); + } + if descriptor.resource_limits.target_count != 16 + || descriptor.resource_limits.target_uri_bytes != 2_048 + || descriptor.resource_limits.diagnostic_bytes != 4_096 + || descriptor.resource_limits.claim_lease_millis != 300_000 + { + return Err("Phase 1 publication resource limits are invalid".to_owned()); + } + validate_identity( + &descriptor.operation_identity, + "radroots.phase1.publication-operation.v1", + &[ + "artifact_digest_32", + "media_readiness_binding_digest_32", + "expected_author_32", + "target_policy_digest_32", + ], + )?; + validate_identity( + &descriptor.dispatch_identity, + "radroots.phase1.relay-dispatch.v1", + &[ + "event_id_32", + "target_policy_digest_32", + "endpoint_fingerprint_32", + ], + )?; + if descriptor + .event_states + .iter() + .map(String::as_str) + .ne(EVENT_STATES.iter().copied()) + || descriptor + .target_states + .iter() + .map(String::as_str) + .ne(TARGET_STATES.iter().copied()) + { + return Err("Phase 1 publication state inventory is invalid".to_owned()); + } + validate_transitions(&descriptor)?; + if descriptor + .stable_errors + .iter() + .map(String::as_str) + .ne(ERROR_CODES.iter().copied()) + { + return Err("Phase 1 publication stable-error inventory is invalid".to_owned()); + } + let registry: Value = serde_json::from_slice(&read_regular_file( + workspace_root, + MIGRATION_REGISTRY_RELATIVE, + )?) + .map_err(|error| format!("parse {MIGRATION_REGISTRY_RELATIVE}: {error}"))?; + let migration = registry["migrations"] + .as_array() + .and_then(|migrations| migrations.iter().find(|entry| entry["version"] == 2)) + .ok_or_else(|| "migration registry does not contain Phase 1 version 2".to_owned())?; + if migration["name"] != descriptor.migration.name + || migration["up_sha256"] != descriptor.migration.up_sha256 + || migration["down_sha256"] != descriptor.migration.down_sha256 + || migration["schema_sha256"] != descriptor.migration.schema_sha256 + { + return Err("migration registry and Phase 1 descriptor disagree".to_owned()); + } + Ok(descriptor) +} + +fn validate_identity(identity: &Identity, domain: &str, preimage: &[&str]) -> Result<(), String> { + if identity.algorithm != "sha256_raw_fixed_width_v1" + || identity.domain != domain + || identity.domain_terminator_hex != "00" + || identity + .preimage + .iter() + .map(String::as_str) + .ne(preimage.iter().copied()) + { + return Err(format!("Phase 1 identity `{domain}` is invalid")); + } + Ok(()) +} + +fn validate_transitions(descriptor: &Descriptor) -> Result<(), String> { + if descriptor.transitions.len() != 25 { + return Err("Phase 1 publication transition inventory must contain 25 entries".to_owned()); + } + let event_states = descriptor.event_states.iter().collect::<BTreeSet<_>>(); + let target_states = descriptor.target_states.iter().collect::<BTreeSet<_>>(); + let mut ids = BTreeSet::new(); + for transition in &descriptor.transitions { + let states = match transition.scope.as_str() { + "event" => &event_states, + "target" => &target_states, + _ => return Err(format!("invalid transition scope `{}`", transition.scope)), + }; + if !ids.insert(transition.id.as_str()) + || !states.contains(&transition.from) + || !states.contains(&transition.to) + || !transition.revision_cas + || transition.lease_predicate.is_empty() + || transition.durable_side_effect.is_empty() + || !matches!( + transition.retry_class.as_str(), + "none" | "retryable" | "repair" | "terminal" + ) + || transition.repair_edge != (transition.retry_class == "repair") + || transition.terminal_destination + != matches!( + transition.to.as_str(), + "published" + | "failed-terminal" + | "quarantined" + | "cancelled" + | "accepted-observed" + ) + { + return Err(format!("invalid transition `{}`", transition.id)); + } + } + Ok(()) +} + +fn validate_vector(workspace_root: &Path) -> Result<(), String> { + let bytes = read_regular_file(workspace_root, VECTOR_RELATIVE)?; + let vector: Vector = serde_json::from_slice(&bytes) + .map_err(|error| format!("parse {VECTOR_RELATIVE}: {error}"))?; + if vector.schema_version != 1 + || vector.contract_id != CONTRACT_ID + || vector.executor.id != VECTOR_EXECUTOR_ID + || vector.executor.path != VECTOR_EXECUTOR_RELATIVE + || vector.executor.test != VECTOR_EXECUTOR_TEST + || !vector.identity_vector.is_object() + { + return Err("Phase 1 publication vector identity is invalid".to_owned()); + } + let expected = CASE_IDS.iter().copied().collect::<BTreeSet<_>>(); + let mut actual = BTreeSet::new(); + for case in vector.cases { + if case.execution != "direct_executor" + || case.expected_outcome.is_empty() + || case.expected_error.is_some() != (case.expected_outcome == "rejected") + || !actual.insert(case.id) + { + return Err("Phase 1 publication vector case is invalid".to_owned()); + } + } + if actual.iter().map(String::as_str).collect::<BTreeSet<_>>() != expected { + return Err("Phase 1 publication vector case inventory is incomplete".to_owned()); + } + Ok(()) +} + +fn validate_release(workspace_root: &Path) -> Result<(), String> { + let source = read_regular_file(workspace_root, RELEASE_RELATIVE)?; + let release: toml::Value = toml::from_str( + std::str::from_utf8(&source) + .map_err(|error| format!("decode {RELEASE_RELATIVE}: {error}"))?, + ) + .map_err(|error| format!("parse {RELEASE_RELATIVE}: {error}"))?; + let changes = release["changes"] + .as_array() + .ok_or_else(|| format!("{RELEASE_RELATIVE} has no changes"))?; + if changes + .iter() + .filter(|change| change["id"].as_str() == Some(RELEASE_CHANGE_ID)) + .count() + != 1 + { + return Err(format!( + "{RELEASE_RELATIVE} must declare one `{RELEASE_CHANGE_ID}` change" + )); + } + let changelog = read_regular_file(workspace_root, CHANGELOG_RELATIVE)?; + let changelog = std::str::from_utf8(&changelog) + .map_err(|error| format!("decode {CHANGELOG_RELATIVE}: {error}"))?; + let marker = format!("<!-- release-change: {RELEASE_CHANGE_ID} -->"); + if changelog.matches(&marker).count() != 1 { + return Err(format!("{CHANGELOG_RELATIVE} must contain one `{marker}`")); + } + Ok(()) +} + +fn expected_manifest( + workspace_root: &Path, + descriptor: &Descriptor, + schema_bytes: &[u8], +) -> Result<Value, String> { + let sources = SOURCE_FILES + .iter() + .map(|(role, relative)| { + Ok(json!({ + "role": role, + "file": descriptor_for_file(workspace_root, relative)?, + })) + }) + .collect::<Result<Vec<_>, String>>()?; + Ok(json!({ + "schema_version": 1, + "contract_id": CONTRACT_ID, + "manifest_schema": descriptor_for_bytes(MANIFEST_SCHEMA_RELATIVE, schema_bytes), + "descriptor": descriptor_for_file(workspace_root, DESCRIPTOR_RELATIVE)?, + "migration": { + "version": descriptor.migration.version, + "schema_sha256": descriptor.migration.schema_sha256, + "up": descriptor_for_file(workspace_root, MIGRATION_UP_RELATIVE)?, + "down": descriptor_for_file(workspace_root, MIGRATION_DOWN_RELATIVE)?, + }, + "state_machine": { + "event_state_count": descriptor.event_states.len(), + "target_state_count": descriptor.target_states.len(), + "transition_count": descriptor.transitions.len(), + "stable_error_count": descriptor.stable_errors.len(), + }, + "result_vector": { + "canonical": descriptor_for_file(workspace_root, VECTOR_RELATIVE)?, + "mirror_path": VECTOR_MIRROR_RELATIVE, + "executor": descriptor_for_file(workspace_root, VECTOR_EXECUTOR_RELATIVE)?, + "executor_id": VECTOR_EXECUTOR_ID, + "executor_test": VECTOR_EXECUTOR_TEST, + }, + "source_files": sources, + "release": { + "change_id": RELEASE_CHANGE_ID, + "record": RELEASE_RELATIVE, + "changelog": CHANGELOG_RELATIVE, + }, + })) +} + +fn manifest_schema() -> Value { + json!({ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://radroots.org/contracts/outbox/phase1_publication_v1.manifest.schema.json", + "type": "object", + "additionalProperties": false, + "required": ["schema_version", "contract_id", "manifest_schema", "descriptor", "migration", "state_machine", "result_vector", "source_files", "release"], + "properties": { + "schema_version": { "const": 1 }, + "contract_id": { "const": CONTRACT_ID }, + "manifest_schema": { "$ref": "#/$defs/file" }, + "descriptor": { "$ref": "#/$defs/file" }, + "migration": { + "type": "object", "additionalProperties": false, + "required": ["version", "schema_sha256", "up", "down"], + "properties": { + "version": { "const": 2 }, + "schema_sha256": { "$ref": "#/$defs/sha256" }, + "up": { "$ref": "#/$defs/file" }, + "down": { "$ref": "#/$defs/file" } + } + }, + "state_machine": { + "type": "object", "additionalProperties": false, + "required": ["event_state_count", "target_state_count", "transition_count", "stable_error_count"], + "properties": { + "event_state_count": { "const": 9 }, + "target_state_count": { "const": 8 }, + "transition_count": { "const": 25 }, + "stable_error_count": { "const": 25 } + } + }, + "result_vector": { + "type": "object", "additionalProperties": false, + "required": ["canonical", "mirror_path", "executor", "executor_id", "executor_test"], + "properties": { + "canonical": { "$ref": "#/$defs/file" }, + "mirror_path": { "const": VECTOR_MIRROR_RELATIVE }, + "executor": { "$ref": "#/$defs/file" }, + "executor_id": { "const": VECTOR_EXECUTOR_ID }, + "executor_test": { "const": VECTOR_EXECUTOR_TEST } + } + }, + "source_files": { + "type": "array", "minItems": SOURCE_FILES.len(), "maxItems": SOURCE_FILES.len(), + "items": { + "type": "object", "additionalProperties": false, + "required": ["role", "file"], + "properties": { "role": { "type": "string", "minLength": 1 }, "file": { "$ref": "#/$defs/file" } } + } + }, + "release": { + "type": "object", "additionalProperties": false, + "required": ["change_id", "record", "changelog"], + "properties": { + "change_id": { "const": RELEASE_CHANGE_ID }, + "record": { "const": RELEASE_RELATIVE }, + "changelog": { "const": CHANGELOG_RELATIVE } + } + } + }, + "$defs": { + "sha256": { "type": "string", "pattern": "^[0-9a-f]{64}$" }, + "file": { + "type": "object", "additionalProperties": false, + "required": ["path", "byte_length", "sha256"], + "properties": { + "path": { "type": "string", "minLength": 1 }, + "byte_length": { "type": "integer", "minimum": 1 }, + "sha256": { "$ref": "#/$defs/sha256" } + } + } + } + }) +} + +fn descriptor_for_file(workspace_root: &Path, relative: &str) -> Result<Value, String> { + let bytes = read_regular_file(workspace_root, relative)?; + Ok(descriptor_for_bytes(relative, &bytes)) +} + +fn descriptor_for_bytes(relative: &str, bytes: &[u8]) -> Value { + json!({ + "path": relative, + "byte_length": bytes.len(), + "sha256": sha256_hex(bytes), + }) +} + +fn canonical_json_bytes(value: &Value) -> Result<Vec<u8>, String> { + let mut bytes = serde_json::to_vec_pretty(value) + .map_err(|error| format!("serialize Phase 1 publication artifact: {error}"))?; + bytes.push(b'\n'); + Ok(bytes) +} + +fn sha256_hex(bytes: &[u8]) -> String { + hex::encode(Sha256::digest(bytes)) +} diff --git a/tools/xtask/src/main.rs b/tools/xtask/src/main.rs @@ -27,6 +27,7 @@ fn usage() { eprintln!(" cargo xtask contract blossom-publication-readiness-manifest [--write]"); eprintln!(" cargo xtask contract blossom-raster-decoder-security-manifest [--write]"); eprintln!(" cargo xtask contract outbox-migration-manifest [--write]"); + eprintln!(" cargo xtask contract outbox-phase1-publication-manifest [--write]"); eprintln!(" cargo xtask contract phase1-publication-media-readiness-manifest [--write]"); eprintln!(" cargo xtask contract release-provenance-schema [--write]"); eprintln!(" cargo xtask contract knowledge-manifest [--write]"); @@ -227,6 +228,16 @@ fn run_contract(args: &[String]) -> Result<(), String> { Err("outbox-migration-manifest accepts no arguments or exactly --write".to_string()) } }, + Some("outbox-phase1-publication-manifest") => match &args[1..] { + [] => contract::validate_outbox_phase1_publication_manifest(&workspace_root()), + [flag] if flag == "--write" => { + contract::write_outbox_phase1_publication_manifest(&workspace_root()) + } + _ => Err( + "outbox-phase1-publication-manifest accepts no arguments or exactly --write" + .to_string(), + ), + }, Some("phase1-publication-media-readiness-manifest") => match &args[1..] { [] => contract::validate_phase1_publication_media_readiness_manifest(&workspace_root()), [flag] if flag == "--write" => { @@ -402,6 +413,12 @@ mod tests { ]) .expect_err("invalid outbox migration manifest mode"); assert!(invalid_outbox_migration.contains("exactly --write")); + let invalid_outbox_phase1 = run_contract(&[ + "outbox-phase1-publication-manifest".to_string(), + "--invalid".to_string(), + ]) + .expect_err("invalid outbox Phase 1 publication manifest mode"); + assert!(invalid_outbox_phase1.contains("exactly --write")); let invalid_release_provenance_schema = run_contract(&[ "release-provenance-schema".to_string(), "--invalid".to_string(), @@ -527,6 +544,8 @@ mod tests { .expect("contract Blossom raster decoder security manifest"); run_contract(&["outbox-migration-manifest".to_string()]) .expect("contract outbox migration manifest"); + run_contract(&["outbox-phase1-publication-manifest".to_string()]) + .expect("contract outbox Phase 1 publication manifest"); run_contract(&["phase1-publication-media-readiness-manifest".to_string()]) .expect("contract Phase 1 publication media-readiness manifest"); run_contract(&["release-provenance-schema".to_string()])