myc

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

commit d6bc2a722f72f403bd9ba2519401d8e0a874750c
parent 67775a6f61c962e375cb72b7686f821a0d5ab865
Author: triesap <tyson@radroots.org>
Date:   Sun, 23 Aug 2026 11:02:36 +0000

delivery: seal provider and relay execution

- execute both governed provider kinds behind verified private adapters
- persist durable delivery state before exact Nostr submission
- bind dependencies and machine contract to source-lock v2
- qualify cancellation, replay, redaction, and public API boundaries

Diffstat:
MAGENTS.md | 10++++++++++
MCargo.lock | 350++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
MCargo.toml | 3+++
MREADME | 17+++++++++++++++++
Acontracts/services_hardening/provider_delivery.v1.json | 52++++++++++++++++++++++++++++++++++++++++++++++++++++
Mradroots.service.source-lock.v2.toml | 2+-
Asrc/delivery_worker.rs | 552+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/lib.rs | 3+++
Msrc/nip46_wave_080_b.rs | 4++--
Msrc/provider_contract.rs | 2+-
Msrc/provider_envelope.rs | 16++++++++++++++++
Asrc/provider_executor.rs | 502+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/provider_verification.rs | 56++++++++++++++++++++++++++++++++++++++++++++++++--------
Msrc/runtime_supervision.rs | 11+++++++++++
Asrc/transport_nostr_adapter.rs | 405+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/build_policy.rs | 16++++++++++++++++
Mtests/package_boundary.rs | 100+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_legacy_removal.rs | 4+++-
Mtests/services_hardening_native_release.rs | 2+-
Mtests/services_hardening_signer_request_state.rs | 10++++++++--
20 files changed, 2088 insertions(+), 29 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -151,6 +151,16 @@ journal. Database-only mutations must later compose their effect, audit, and completion in one transaction; online backup records Prepared before capture and completes only after the bundle is durable. +- Step 159 unit 12 owns the crate-private provider executor, the exact + source-locked `radroots_transport_nostr` adapter, and the durable delivery + worker. Provider results remain untrusted until independently verified. + Delivery preparation performs no relay I/O; the worker persists Submitted + immediately before execution, maps post-submit cancellation or lost + acknowledgement to UnknownAcknowledgement, and retries only the exact + committed signed bytes. Never hold a SQLite transaction across provider or + relay work, detach protected blocking work, expose the executor/client, or + create one task per relay. Unit 13 alone wires these components into the + fixed runtime graph and startup handshake. - Treat checked-in source, tests, and prototype behavior as implementation evidence, not permission to preserve behavior that the active requirement removes. diff --git a/Cargo.lock b/Cargo.lock @@ -115,6 +115,37 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50" [[package]] +name = "async-utility" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "188f83b9a198af8c336e505611edb00d6d2ac5c694241c5a4f9a12316938cfe9" +dependencies = [ + "futures-util", + "gloo-timers", + "tokio", + "wasm-bindgen-futures", +] + +[[package]] +name = "async-wsocket" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1c92385c7c8b3eb2de1b78aeca225212e4c9a69a78b802832759b108681a5069" +dependencies = [ + "async-utility", + "futures", + "futures-util", + "js-sys", + "tokio", + "tokio-rustls", + "tokio-socks", + "tokio-tungstenite", + "url", + "wasm-bindgen", + "web-sys", +] + +[[package]] name = "atoi" version = "2.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -124,6 +155,12 @@ dependencies = [ ] [[package]] +name = "atomic-destructor" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef49f5882e4b6afaac09ad239a4f8c70a24b8f2b0897edb1f706008efd109cf4" + +[[package]] name = "atomic-waker" version = "1.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -417,7 +454,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0dc92fb57ca44df6db8059111ab3af99a63d5d0f8375d9972e319a379c6bab76" dependencies = [ "generic-array", - "rand_core", + "rand_core 0.6.4", "subtle", "zeroize", ] @@ -429,7 +466,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a" dependencies = [ "generic-array", - "rand_core", + "rand_core 0.6.4", "typenum", ] @@ -497,7 +534,7 @@ dependencies = [ "ff", "generic-array", "group", - "rand_core", + "rand_core 0.6.4", "sec1", "subtle", "zeroize", @@ -525,7 +562,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -562,7 +599,7 @@ version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c0b50bfb653653f9ca9095b427bed08ab8d75a137839d9ad64eb11810d5b6393" dependencies = [ - "rand_core", + "rand_core 0.6.4", "subtle", ] @@ -786,13 +823,25 @@ dependencies = [ ] [[package]] +name = "gloo-timers" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbb143cf96099802033e0d4f4963b19fd2e0b728bcf076cd9cf7f6634f092994" +dependencies = [ + "futures-channel", + "futures-core", + "js-sys", + "wasm-bindgen", +] + +[[package]] name = "group" version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f0f9ef7462f7c099f518d754361858f86d8a07af53ba9af0fe635bbccb151a63" dependencies = [ "ff", - "rand_core", + "rand_core 0.6.4", "subtle", ] @@ -1123,6 +1172,8 @@ version = "0.3.94" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2e04e2ef80ce82e13552136fabeef8a5ed1f985a96805761cbb9a2c34e7664d9" dependencies = [ + "cfg-if", + "futures-util", "once_cell", "wasm-bindgen", ] @@ -1246,6 +1297,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897" [[package]] +name = "lru" +version = "0.16.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7f66e8d5d03f609abc3a39e6f08e4164ebf1447a732906d39eb9b99b7919ef39" + +[[package]] name = "mediatype" version = "0.21.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1285,6 +1342,7 @@ dependencies = [ "hex", "jsonschema", "nostr", + "radroots_event_codec", "radroots_nostr", "radroots_nostr_connect", "radroots_runtime_paths", @@ -1292,6 +1350,8 @@ dependencies = [ "radroots_service_host", "radroots_service_sqlite", "radroots_storage", + "radroots_transport", + "radroots_transport_nostr", "rustix", "serde", "serde_json", @@ -1305,6 +1365,12 @@ dependencies = [ ] [[package]] +name = "negentropy" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "81c353b400a5503efdcf398f11a83fb7aa84f59f5d76fc4bf5bbc1e4f5366caa" + +[[package]] name = "nostr" version = "0.44.7" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1331,6 +1397,59 @@ dependencies = [ ] [[package]] +name = "nostr-database" +version = "0.44.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7462c9d8ae5ef6a28d66a192d399ad2530f1f2130b13186296dbb11bdef5b3d1" +dependencies = [ + "lru", + "nostr", + "tokio", +] + +[[package]] +name = "nostr-gossip" +version = "0.44.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ade30de16869618919c6b5efc8258f47b654a98b51541eb77f85e8ec5e3c83a6" +dependencies = [ + "nostr", +] + +[[package]] +name = "nostr-relay-pool" +version = "0.44.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c85c54d6ca9aae4ae2bf19a7663ba9db5f45f783f1d24aff55f006386b8b99a1" +dependencies = [ + "async-utility", + "async-wsocket", + "atomic-destructor", + "hex", + "lru", + "negentropy", + "nostr", + "nostr-database", + "tokio", + "tracing", +] + +[[package]] +name = "nostr-sdk" +version = "0.44.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "471732576710e779b64f04c55e3f8b5292f865fea228436daf19694f0bf70393" +dependencies = [ + "async-utility", + "nostr", + "nostr-database", + "nostr-gossip", + "nostr-relay-pool", + "tokio", + "tracing", +] + +[[package]] name = "num" version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1468,7 +1587,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "346f04948ba92c43e8469c1ee6736c7563d71012b17d40745260fe106aac2166" dependencies = [ "base64ct", - "rand_core", + "rand_core 0.6.4", "subtle", ] @@ -1768,14 +1887,44 @@ dependencies = [ ] [[package]] +name = "radroots_transport_nostr" +version = "0.1.0-alpha" +source = "git+https://github.com/radrootslabs/lib?rev=7d7b454b4c9ed86569671993bd03ca868b676665#7d7b454b4c9ed86569671993bd03ca868b676665" +dependencies = [ + "async-wsocket", + "futures", + "nostr-relay-pool", + "nostr-sdk", + "radroots_event_codec", + "radroots_nostr", + "radroots_protocol", + "radroots_transport", + "serde_json", + "sha2", + "tokio", + "tokio-tungstenite", + "url", +] + +[[package]] name = "rand" version = "0.8.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" dependencies = [ "libc", - "rand_chacha", - "rand_core", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + +[[package]] +name = "rand" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9ef1d0d795eb7d84685bca4f72f3649f064e6641543d3a8c415898726a57b41" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.5", ] [[package]] @@ -1785,7 +1934,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.5", ] [[package]] @@ -1798,6 +1957,15 @@ dependencies = [ ] [[package]] +name = "rand_core" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" +dependencies = [ + "getrandom 0.3.4", +] + +[[package]] name = "redox_syscall" version = "0.5.18" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1873,6 +2041,20 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dc897dd8d9e8bd1ed8cdad82b5966c3e0ecae09fb1907d58efaa013543185d0a" [[package]] +name = "ring" +version = "0.17.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" +dependencies = [ + "cc", + "cfg-if", + "getrandom 0.2.17", + "libc", + "untrusted", + "windows-sys 0.52.0", +] + +[[package]] name = "rust_decimal" version = "1.41.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1894,7 +2076,41 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", +] + +[[package]] +name = "rustls" +version = "0.23.43" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" +dependencies = [ + "once_cell", + "ring", + "rustls-pki-types", + "rustls-webpki", + "subtle", + "zeroize", +] + +[[package]] +name = "rustls-pki-types" +version = "1.15.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4925028c7eb5d1fcdaf196971378ed9d2c1c4efc7dc5d011256f76c99c0a96" +dependencies = [ + "zeroize", +] + +[[package]] +name = "rustls-webpki" +version = "0.103.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" +dependencies = [ + "ring", + "rustls-pki-types", + "untrusted", ] [[package]] @@ -1949,7 +2165,7 @@ version = "0.29.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9465315bc9d4566e1724f0fffcbcc446268cb522e60f9a27bcded6b19c108113" dependencies = [ - "rand", + "rand 0.8.5", "secp256k1-sys", "serde", ] @@ -2022,6 +2238,17 @@ dependencies = [ ] [[package]] +name = "sha1" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + +[[package]] name = "sha2" version = "0.10.9" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2255,7 +2482,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2350,6 +2577,28 @@ dependencies = [ ] [[package]] +name = "tokio-rustls" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" +dependencies = [ + "rustls", + "tokio", +] + +[[package]] +name = "tokio-socks" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7e2948f60dbe26b35f2c7fb74ac2854c1fddded0fe9d7548fcc674a246f7615" +dependencies = [ + "either", + "futures-util", + "thiserror 1.0.69", + "tokio", +] + +[[package]] name = "tokio-stream" version = "0.1.19" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2361,6 +2610,22 @@ dependencies = [ ] [[package]] +name = "tokio-tungstenite" +version = "0.26.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a9daff607c6d2bf6c16fd681ccb7eecc83e4e2cdc1ca067ffaadfca5de7f084" +dependencies = [ + "futures-util", + "log", + "rustls", + "rustls-pki-types", + "tokio", + "tokio-rustls", + "tungstenite", + "webpki-roots 0.26.11", +] + +[[package]] name = "tokio-util" version = "0.7.19" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2454,6 +2719,25 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] +name = "tungstenite" +version = "0.26.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4793cb5e56680ecbb1d843515b23b6de9a75eb04b66643e256a396d43be33c13" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.5", + "rustls", + "rustls-pki-types", + "sha1", + "thiserror 2.0.18", + "utf-8", +] + +[[package]] name = "typenum" version = "1.19.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2503,6 +2787,12 @@ dependencies = [ ] [[package]] +name = "untrusted" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" + +[[package]] name = "url" version = "2.5.8" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2528,6 +2818,12 @@ dependencies = [ ] [[package]] +name = "utf-8" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" + +[[package]] name = "utf8_iter" version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2615,6 +2911,16 @@ dependencies = [ ] [[package]] +name = "wasm-bindgen-futures" +version = "0.4.67" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "03623de6905b7206edd0a75f69f747f134b7f0a2323392d664448bf2d3c5d87e" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + +[[package]] name = "wasm-bindgen-macro" version = "0.2.117" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2691,6 +2997,24 @@ dependencies = [ ] [[package]] +name = "webpki-roots" +version = "0.26.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "521bc38abb08001b01866da9f51eb7c5d647a19260e00054a8c7fd5f9e57f7a9" +dependencies = [ + "webpki-roots 1.0.9", +] + +[[package]] +name = "webpki-roots" +version = "1.0.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7dcd9d09a39985f5344844e66b0c530a33843579125f23e21e9f0f220850f22a" +dependencies = [ + "rustls-pki-types", +] + +[[package]] name = "winapi" version = "0.3.9" source = "registry+https://github.com/rust-lang/crates.io-index" diff --git a/Cargo.toml b/Cargo.toml @@ -58,11 +58,14 @@ jsonschema = { version = "0.48.1", default-features = false } nostr = { version = "0.44.2", features = ["nip04", "nip44", "nip46", "nip49"] } radroots_nostr = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", features = ["events"] } radroots_nostr_connect = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } +radroots_event_codec = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", features = ["json"] } radroots_runtime_paths = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } radroots_service_host = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } radroots_service_sqlite = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } radroots_secrets = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", features = ["std"] } radroots_storage = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", default-features = false } +radroots_transport = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", default-features = false } +radroots_transport_nostr = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } serde = { version = "1.0", features = ["derive"] } serde_json = { version = "1.0", features = ["raw_value"] } sha2 = "0.10" diff --git a/README b/README @@ -185,6 +185,23 @@ The admin response transport admits at least 8,382 UTF-8 bytes so the maximum retained model plus the maximum safe correlation identity always fits its canonical success envelope. +The Step 159 provider and delivery boundary is sealed inside the crate. Both +governed provider kinds execute only fully bound operations, and every result +is independently verified before it becomes authority. Protected blocking +work is joined, local-signer calls use the hardened Unix-admin client, and no +raw provider client, secret, callback, or dependency-owned error crosses the +public API. + +Relay publication uses only the exact source-locked +`radroots_transport_nostr` adapter. Preparation validates the exact committed, +signature-verified event bytes without network I/O. The delivery worker then +persists Submitted immediately before execution. Accepted, rejected, +transport-failed, and unknown acknowledgements remain distinct; cancellation +or lost acknowledgement after Submitted is durably unknown, and retries never +alter the committed bytes. Provider or relay work never occurs inside a SQLite +transaction. Runtime task-graph wiring and startup handshakes remain the next +ordered Step 159 unit. + Schema v8 adds immutable NIP-46 operation-completion evidence. The Step 147 integration checkpoint binds each durable request to its stable operation and correlation identities, terminal connect authority or exact active session, diff --git a/contracts/services_hardening/provider_delivery.v1.json b/contracts/services_hardening/provider_delivery.v1.json @@ -0,0 +1,52 @@ +{ + "schema": "radroots.myc.provider-delivery.v1", + "contract_version": 1, + "step": 159, + "unit": "myc-provider-delivery", + "provider_executor": { + "visibility": "crate_private_sealed", + "providers": ["encrypted_file", "local_signer"], + "operation_input": "closed_provider_operation", + "encrypted_file_execution": "joined_bounded_blocking_worker", + "local_signer_execution": "hardened_unix_admin_client", + "result_authority": "independently_verified_provider_response", + "raw_secret_exposure": false, + "raw_client_exposure": false, + "detached_blocking_work": false + }, + "relay_adapter": { + "implementation": "radroots_transport_nostr", + "source_locked_revision": "7d7b454b4c9ed86569671993bd03ca868b676665", + "preparation_performs_io": false, + "target_cardinality_per_attempt": 1, + "public_profile": "wss_only", + "repo_local_simulator_profile": "ws_loopback_only", + "raw_transport_error_exposure": false + }, + "durable_delivery": { + "payload_authority": "exact_committed_signature_verified_event_bytes", + "submitted_transition": "immediately_after_prepare_before_execute", + "outcomes": ["delivered", "relay_rejected", "transport_failed", "unknown_acknowledgement"], + "cancellation_before_submission": "proven_transport_failure", + "cancellation_after_submission": "unknown_acknowledgement", + "execution_error_after_submission": "unknown_acknowledgement", + "exact_replay_performs_io": false, + "retry_payload_mutation": false, + "external_io_inside_sqlite_transaction": false + }, + "resource_policy": { + "one_task_per_relay": false, + "unbounded_queue": false, + "runtime_creation": false, + "ambient_time": false, + "ambient_entropy": false + }, + "deferred": [ + "runtime_task_graph_wiring", + "provider_startup_handshake", + "relay_ingress_subscription", + "delivery_outbox_coordinator", + "ordered_shutdown" + ], + "nonclaims": ["step_159_closure", "rcld_promotion", "nix", "oci"] +} diff --git a/radroots.service.source-lock.v2.toml b/radroots.service.source-lock.v2.toml @@ -7,7 +7,7 @@ architecture = "radroots.crates.release.v2" workspace_catalog_sha256 = "deca0c080deae187ff8186c0708903e42f41ea57f77c5f91581e23aa561164a4" version = "0.1.0-alpha" source_archive_sha256 = "b425371c134be96cce46b37f7035d6212f1efe8cff50bef366631ba5632991b0" -cargo_lock_sha256 = "57be61e2ce5f5cf37c960e8aaf223a8716c7cd49922a752b437f0be45fd71a08" +cargo_lock_sha256 = "a49a45076e3fecb49cdf8d98ea1082f62d7025cbe5cd9359df1e7b89e36df6d4" rust_version = "1.97.1" host_feature_profile = "service-host" diff --git a/src/delivery_worker.rs b/src/delivery_worker.rs @@ -0,0 +1,552 @@ +//! Durable delivery orchestration over exact committed event bytes. + +#![allow( + dead_code, + reason = "Step 159 Unit 12 seals the worker before Unit 13 runtime graph wiring" +)] + +use core::fmt; +use std::error::Error; + +use sha2::{Digest as _, Sha256}; + +use crate::transport_nostr_adapter::{ + MycNostrDeliveryAdapter, MycRelayAdapter, MycRelayExecutionOutcome, +}; +use crate::{ + MycConfigDocumentV1, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, MycDeliveryClaim, + MycDeliveryJobId, MycDeliveryJobRecord, MycDeliveryRelayId, MycDeliveryRetryJitter, + MycDeliverySourceKind, MycDeliveryTimeUnixMs, MycStateRepository, MycTaskCancellation, +}; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum MycDeliveryWorkerErrorKind { + Configuration, + State, + Artifact, +} + +pub(crate) struct MycDeliveryWorkerError { + kind: MycDeliveryWorkerErrorKind, +} + +impl MycDeliveryWorkerError { + pub(crate) const fn kind(&self) -> MycDeliveryWorkerErrorKind { + self.kind + } +} + +impl fmt::Debug for MycDeliveryWorkerError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDeliveryWorkerError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for MycDeliveryWorkerError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("Myc delivery worker failed") + } +} + +impl Error for MycDeliveryWorkerError {} + +const fn worker_error(kind: MycDeliveryWorkerErrorKind) -> MycDeliveryWorkerError { + MycDeliveryWorkerError { kind } +} + +#[derive(Clone, PartialEq, Eq)] +pub(crate) enum MycDeliveryWorkerResult { + Completed(MycDeliveryJobRecord), + ExactReplay, + NotReady, + Terminal, +} + +impl fmt::Debug for MycDeliveryWorkerResult { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::Completed(_) => "MycDeliveryWorkerResult::Completed([redacted])", + Self::ExactReplay => "MycDeliveryWorkerResult::ExactReplay", + Self::NotReady => "MycDeliveryWorkerResult::NotReady", + Self::Terminal => "MycDeliveryWorkerResult::Terminal", + }) + } +} + +pub(crate) struct MycDeliveryExecutionEvidence { + pub(crate) claimed_at: MycDeliveryTimeUnixMs, + pub(crate) submitted_at: MycDeliveryTimeUnixMs, + pub(crate) observed_at: MycDeliveryTimeUnixMs, + pub(crate) retry_jitter: MycDeliveryRetryJitter, +} + +impl fmt::Debug for MycDeliveryExecutionEvidence { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycDeliveryExecutionEvidence([sealed])") + } +} + +pub(crate) struct MycDeliveryWorker { + adapter: MycNostrDeliveryAdapter, +} + +impl MycDeliveryWorker { + pub(crate) fn from_configuration( + configuration: &MycConfigDocumentV1, + ) -> Result<Self, MycDeliveryWorkerError> { + let adapter = MycNostrDeliveryAdapter::from_configuration(configuration) + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::Configuration))?; + Ok(Self { adapter }) + } + + pub(crate) async fn run_one( + &self, + repository: &MycStateRepository<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + nonce: MycDeliveryAttemptNonce, + evidence: MycDeliveryExecutionEvidence, + cancellation: &MycTaskCancellation, + ) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { + run_with_adapter( + &self.adapter, + repository, + job_id, + relay_id, + nonce, + evidence, + cancellation, + ) + .await + } +} + +#[allow(clippy::too_many_arguments)] +async fn run_with_adapter<A: MycRelayAdapter>( + adapter: &A, + repository: &MycStateRepository<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + nonce: MycDeliveryAttemptNonce, + evidence: MycDeliveryExecutionEvidence, + cancellation: &MycTaskCancellation, +) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { + if cancellation.is_cancelled() { + return Ok(MycDeliveryWorkerResult::NotReady); + } + let attempt = match repository + .claim_delivery_target(job_id, relay_id, nonce, evidence.claimed_at) + .await + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? + { + MycDeliveryClaim::Claimed(attempt) => attempt, + MycDeliveryClaim::ExactReplay(_) => return Ok(MycDeliveryWorkerResult::ExactReplay), + MycDeliveryClaim::NotReady => return Ok(MycDeliveryWorkerResult::NotReady), + MycDeliveryClaim::Terminal => return Ok(MycDeliveryWorkerResult::Terminal), + }; + let job = repository + .read_delivery_job(job_id) + .await + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? + .ok_or_else(|| worker_error(MycDeliveryWorkerErrorKind::Artifact))?; + let exact_event_bytes = read_exact_event(repository, &job).await?; + let request_id = request_id(job_id, attempt.id()); + let prepared = match adapter.prepare( + relay_id, + request_id, + &exact_event_bytes, + attempt.lease_expires_at().get(), + ) { + Ok(prepared) => prepared, + Err(_) => { + return persist_outcome( + repository, + job_id, + relay_id, + attempt.id(), + MycDeliveryAttemptOutcome::TransportFailed, + evidence, + ) + .await; + } + }; + if cancellation.is_cancelled() { + return persist_outcome( + repository, + job_id, + relay_id, + attempt.id(), + MycDeliveryAttemptOutcome::TransportFailed, + evidence, + ) + .await; + } + repository + .mark_delivery_attempt_submitted(job_id, relay_id, attempt.id(), evidence.submitted_at) + .await + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))?; + let outcome = tokio::select! { + result = adapter.execute(prepared) => result.unwrap_or(MycRelayExecutionOutcome::UnknownAcknowledgement), + () = cancellation.cancelled() => MycRelayExecutionOutcome::UnknownAcknowledgement, + }; + let outcome = match outcome { + MycRelayExecutionOutcome::Accepted => MycDeliveryAttemptOutcome::Delivered, + MycRelayExecutionOutcome::Rejected => MycDeliveryAttemptOutcome::RelayRejected, + MycRelayExecutionOutcome::TransportFailed => MycDeliveryAttemptOutcome::TransportFailed, + MycRelayExecutionOutcome::UnknownAcknowledgement => { + MycDeliveryAttemptOutcome::UnknownAcknowledgement + } + }; + persist_outcome( + repository, + job_id, + relay_id, + attempt.id(), + outcome, + evidence, + ) + .await +} + +async fn read_exact_event( + repository: &MycStateRepository<'_>, + job: &MycDeliveryJobRecord, +) -> Result<Box<[u8]>, MycDeliveryWorkerError> { + let bytes: Box<[u8]> = match job.source_kind() { + MycDeliverySourceKind::SignerResponse => repository + .read_nip46_response(job.id()) + .await + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? + .filter(|record| record.delivery_job().id() == job.id()) + .map(|record| Box::from(record.signed_response_bytes())) + .ok_or_else(|| worker_error(MycDeliveryWorkerErrorKind::Artifact))?, + MycDeliverySourceKind::DiscoveryHandler => repository + .read_discovery_document_for_job(job.id()) + .await + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State))? + .map(|record| Box::from(record.event_bytes())) + .ok_or_else(|| worker_error(MycDeliveryWorkerErrorKind::Artifact))?, + }; + let digest: [u8; 32] = Sha256::digest(&bytes).into(); + if &digest != job.artifact_digest().as_bytes() { + return Err(worker_error(MycDeliveryWorkerErrorKind::Artifact)); + } + Ok(bytes) +} + +async fn persist_outcome( + repository: &MycStateRepository<'_>, + job_id: MycDeliveryJobId, + relay_id: &MycDeliveryRelayId, + attempt_id: crate::MycDeliveryAttemptId, + outcome: MycDeliveryAttemptOutcome, + evidence: MycDeliveryExecutionEvidence, +) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { + repository + .record_delivery_attempt_outcome( + job_id, + relay_id, + attempt_id, + outcome, + evidence.retry_jitter, + evidence.observed_at, + ) + .await + .map(MycDeliveryWorkerResult::Completed) + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)) +} + +fn request_id(job_id: MycDeliveryJobId, attempt_id: crate::MycDeliveryAttemptId) -> String { + let mut request = hex::encode(job_id.as_bytes()); + request.push(':'); + request.push_str(&hex::encode(attempt_id.as_bytes())); + request +} + +#[cfg(all(test, any(target_os = "linux", target_os = "macos")))] +mod tests { + use std::{ + fs, + os::unix::fs::PermissionsExt as _, + sync::{Arc, Mutex}, + }; + + use super::*; + use crate::nip46_wave_080_a::{ + OBSERVED_AT_SECONDS, RECEIVED_AT_MS, configuration, connection_time, metadata, + migration_evidence, runtime, + }; + use crate::nip46_wave_080_b::{active_connection, atomic_response_request}; + use crate::{ + MycDeliveryAttemptStatus, MycDeliveryTargetStatus, MycNip46CommitRequest, + MycTaskCancellation, initialize_myc_state, open_myc_state_read_write, + }; + use tokio::sync::Notify; + + struct FakeAdapter { + entered: Arc<Notify>, + release: Arc<Notify>, + prepared_bytes: Arc<Mutex<Vec<u8>>>, + outcome: MycRelayExecutionOutcome, + } + + impl MycRelayAdapter for FakeAdapter { + type Prepared = (); + + fn prepare( + &self, + _relay_id: &MycDeliveryRelayId, + request_id: String, + exact_event_bytes: &[u8], + deadline_unix_ms: u64, + ) -> Result<Self::Prepared, crate::transport_nostr_adapter::MycRelayAdapterError> { + assert_eq!(request_id.len(), 129); + assert!(deadline_unix_ms > RECEIVED_AT_MS); + *self.prepared_bytes.lock().expect("prepared bytes") = exact_event_bytes.to_vec(); + Ok(()) + } + + fn execute<'a>( + &'a self, + (): Self::Prepared, + ) -> crate::transport_nostr_adapter::RelayExecutionFuture<'a> { + Box::pin(async move { + self.entered.notify_one(); + self.release.notified().await; + Ok(self.outcome) + }) + } + } + + struct NoIoAdapter; + + impl MycRelayAdapter for NoIoAdapter { + type Prepared = (); + + fn prepare( + &self, + _relay_id: &MycDeliveryRelayId, + _request_id: String, + _exact_event_bytes: &[u8], + _deadline_unix_ms: u64, + ) -> Result<Self::Prepared, crate::transport_nostr_adapter::MycRelayAdapterError> { + panic!("an exact replay must not prepare a second relay operation") + } + + fn execute<'a>( + &'a self, + (): Self::Prepared, + ) -> crate::transport_nostr_adapter::RelayExecutionFuture<'a> { + panic!("an exact replay must not execute a second relay operation") + } + } + + async fn committed_response_host( + root: &std::path::Path, + ) -> (crate::MycStateHost, MycDeliveryJobId, Vec<u8>) { + let runtime = runtime(root); + fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); + fs::set_permissions( + runtime.context().paths().state(), + fs::Permissions::from_mode(0o700), + ) + .expect("state mode"); + let metadata = metadata(&runtime); + let config = configuration(); + let (applied_at, build) = migration_evidence(); + initialize_myc_state(&runtime, &metadata, applied_at, &build) + .await + .expect("initialization"); + let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("state host"); + let repository = host.repository(); + let (work, _active, decision) = active_connection(&repository, &config).await; + let completion = MycNip46CommitRequest::new( + &work, + Some(&decision), + None, + connection_time(RECEIVED_AT_MS + 2_001), + ) + .expect("completion"); + let (response, exact_bytes) = atomic_response_request(&config, &completion); + let committed = repository + .commit_nip46_response(&response) + .await + .expect("committed response"); + ( + host, + committed.record().response().delivery_job().id(), + exact_bytes, + ) + } + + fn evidence() -> MycDeliveryExecutionEvidence { + MycDeliveryExecutionEvidence { + claimed_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_000).expect("claim time"), + submitted_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_001) + .expect("submitted time"), + observed_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_002).expect("observed time"), + retry_jitter: MycDeliveryRetryJitter::new(0).expect("jitter"), + } + } + + #[tokio::test] + async fn exact_bytes_are_submitted_only_after_durable_submitted_state() { + let root = tempfile::tempdir().expect("temporary root"); + let (host, job_id, exact_bytes) = committed_response_host(root.path()).await; + let repository = host.repository(); + let relay = MycDeliveryRelayId::new("primary").expect("relay"); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let prepared_bytes = Arc::new(Mutex::new(Vec::new())); + let adapter = FakeAdapter { + entered: Arc::clone(&entered), + release: Arc::clone(&release), + prepared_bytes: Arc::clone(&prepared_bytes), + outcome: MycRelayExecutionOutcome::Accepted, + }; + let (cancellation, _cancel) = MycTaskCancellation::test_pair(); + let run = run_with_adapter( + &adapter, + &repository, + job_id, + &relay, + MycDeliveryAttemptNonce::from_injected_entropy([0xa1; 32]), + evidence(), + &cancellation, + ); + let observe = async { + entered.notified().await; + let attempts = repository + .read_delivery_attempts(job_id, &relay) + .await + .expect("submitted attempt"); + assert_eq!(attempts.len(), 1); + assert_eq!(attempts[0].status(), MycDeliveryAttemptStatus::Submitted); + release.notify_one(); + }; + let (result, ()) = tokio::join!(run, observe); + let MycDeliveryWorkerResult::Completed(job) = result.expect("delivery") else { + panic!("delivery must complete"); + }; + assert_eq!( + job.targets()[0].status(), + MycDeliveryTargetStatus::Delivered + ); + assert_eq!( + prepared_bytes.lock().expect("prepared bytes").as_slice(), + exact_bytes + ); + host.close().await.expect("close"); + } + + #[tokio::test] + async fn cancellation_after_submission_is_durably_unknown() { + let root = tempfile::tempdir().expect("temporary root"); + let (host, job_id, _) = committed_response_host(root.path()).await; + let repository = host.repository(); + let relay = MycDeliveryRelayId::new("primary").expect("relay"); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let adapter = FakeAdapter { + entered: Arc::clone(&entered), + release, + prepared_bytes: Arc::new(Mutex::new(Vec::new())), + outcome: MycRelayExecutionOutcome::Accepted, + }; + let (cancellation, cancel) = MycTaskCancellation::test_pair(); + let run = run_with_adapter( + &adapter, + &repository, + job_id, + &relay, + MycDeliveryAttemptNonce::from_injected_entropy([0xa2; 32]), + evidence(), + &cancellation, + ); + let cancel_after_submit = async { + entered.notified().await; + cancel.cancel(); + }; + let (result, ()) = tokio::join!(run, cancel_after_submit); + let MycDeliveryWorkerResult::Completed(job) = result.expect("unknown delivery") else { + panic!("delivery must resolve"); + }; + assert_eq!(job.targets()[0].status(), MycDeliveryTargetStatus::Unknown); + let attempts = repository + .read_delivery_attempts(job_id, &relay) + .await + .expect("attempts"); + assert_eq!(attempts[0].status(), MycDeliveryAttemptStatus::Unknown); + host.close().await.expect("close"); + } + + #[tokio::test] + async fn concurrent_exact_replay_performs_no_second_external_operation() { + let root = tempfile::tempdir().expect("temporary root"); + let (host, job_id, _) = committed_response_host(root.path()).await; + let repository = host.repository(); + let relay = MycDeliveryRelayId::new("primary").expect("relay"); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let adapter = FakeAdapter { + entered: Arc::clone(&entered), + release: Arc::clone(&release), + prepared_bytes: Arc::new(Mutex::new(Vec::new())), + outcome: MycRelayExecutionOutcome::Accepted, + }; + let nonce = MycDeliveryAttemptNonce::from_injected_entropy([0xa3; 32]); + let replay_nonce = MycDeliveryAttemptNonce::from_injected_entropy([0xa3; 32]); + let (cancellation, _cancel) = MycTaskCancellation::test_pair(); + let run = run_with_adapter( + &adapter, + &repository, + job_id, + &relay, + nonce, + evidence(), + &cancellation, + ); + let replay = async { + entered.notified().await; + let result = run_with_adapter( + &NoIoAdapter, + &repository, + job_id, + &relay, + replay_nonce, + evidence(), + &cancellation, + ) + .await + .expect("exact replay"); + assert_eq!(result, MycDeliveryWorkerResult::ExactReplay); + release.notify_one(); + }; + let (result, ()) = tokio::join!(run, replay); + assert!(matches!( + result.expect("delivery"), + MycDeliveryWorkerResult::Completed(_) + )); + host.close().await.expect("close"); + } + + #[test] + fn worker_errors_are_source_free_and_redacted() { + for kind in [ + MycDeliveryWorkerErrorKind::Configuration, + MycDeliveryWorkerErrorKind::State, + MycDeliveryWorkerErrorKind::Artifact, + ] { + let error = worker_error(kind); + assert_eq!(error.kind(), kind); + assert!(Error::source(&error).is_none()); + assert!(!format!("{error} {error:?}").contains("secret")); + } + assert_eq!(OBSERVED_AT_SECONDS, 1_725_000_000); + } +} diff --git a/src/lib.rs b/src/lib.rs @@ -7,6 +7,7 @@ mod cli_v1; mod config_v1; #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] mod control_plane_wave_090_a; +mod delivery_worker; mod diagnostics_v1; mod doctor_v1; mod nip46_admission; @@ -22,6 +23,7 @@ mod operations_v1; mod provider_contract; mod provider_credential; mod provider_envelope; +mod provider_executor; mod provider_local_signer; mod provider_verification; mod runtime_context; @@ -44,6 +46,7 @@ mod state_repository; mod state_request; mod state_response; mod status_v1; +mod transport_nostr_adapter; #[cfg(any(target_os = "linux", target_os = "macos"))] pub use admin_v1::{ diff --git a/src/nip46_wave_080_b.rs b/src/nip46_wave_080_b.rs @@ -25,7 +25,7 @@ use super::nip46_wave_080_a::{ prepared_request, runtime, unsigned_sign_event, untrusted_response, }; -async fn active_connection( +pub(crate) async fn active_connection( repository: &crate::MycStateRepository<'_>, config: &crate::MycConfigDocumentV1, ) -> ( @@ -92,7 +92,7 @@ async fn active_connection( (work, active, decision) } -fn atomic_response_request( +pub(crate) fn atomic_response_request( config: &crate::MycConfigDocumentV1, completion: &MycNip46CommitRequest, ) -> (MycNip46ResponseCommitRequest, Vec<u8>) { diff --git a/src/provider_contract.rs b/src/provider_contract.rs @@ -85,7 +85,7 @@ pub enum MycProviderCapability { } impl MycProviderCapability { - const ALL: [Self; 7] = [ + pub(crate) const ALL: [Self; 7] = [ Self::Describe, Self::PublicIdentity, Self::SignEvent, diff --git a/src/provider_envelope.rs b/src/provider_envelope.rs @@ -248,6 +248,22 @@ impl MycDecryptedIdentity { pub fn public_identity(&self) -> &MycProviderPublicIdentity { &self.public_identity } + + pub(crate) fn secret_bytes(&self) -> &[u8; IDENTITY_SECRET_BYTES] { + &self.secret + } + + #[cfg(test)] + pub(crate) fn from_test_secret(secret: [u8; IDENTITY_SECRET_BYTES]) -> Self { + let secret_key = SecretKey::from_slice(&secret).expect("test secret must be valid"); + let public_identity = + MycProviderPublicIdentity::new(&Keys::new(secret_key).public_key().to_hex()) + .expect("test public identity must be valid"); + Self { + secret: Zeroizing::new(secret), + public_identity, + } + } } impl fmt::Debug for MycDecryptedIdentity { diff --git a/src/provider_executor.rs b/src/provider_executor.rs @@ -0,0 +1,502 @@ +//! Sealed execution over the two governed signer-provider implementations. + +#![allow( + dead_code, + reason = "Step 159 Unit 12 seals the executor before Unit 13 runtime graph wiring" +)] + +use core::fmt; +use std::{error::Error, sync::Arc}; + +use nostr::{ + JsonUtil as _, Keys, PublicKey, SecretKey, UnsignedEvent, + nips::{nip04, nip44}, +}; + +use crate::provider_local_signer::{ProtectedWireHex, WireCapability, WireProviderResult}; +use crate::provider_verification::verify_encrypted_provider_response; +use crate::{ + MYC_LOCAL_SIGNER_TRANSPORT_CONTRACT_VERSION, MYC_PROVIDER_INPUT_MAX_BYTES, MycConfigDocumentV1, + MycDecryptedIdentity, MycLocalSignerClient, MycProviderBinding, MycProviderCapability, + MycProviderKind, MycProviderOperation, MycProviderResponseObservedAtUnixMs, MycProviderRole, + MycRuntimeContext, MycTaskCancellation, MycVerifiedProviderResponse, + open_myc_encrypted_identity, resolve_myc_wrapping_credential, +}; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum MycProviderExecutionErrorKind { + Binding, + Open, + Operation, + Cancelled, + Transport, + Verification, + UnsupportedPlatform, +} + +pub(crate) struct MycProviderExecutionError { + kind: MycProviderExecutionErrorKind, +} + +impl MycProviderExecutionError { + pub(crate) const fn kind(&self) -> MycProviderExecutionErrorKind { + self.kind + } +} + +impl fmt::Debug for MycProviderExecutionError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycProviderExecutionError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for MycProviderExecutionError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("Myc provider execution failed") + } +} + +impl Error for MycProviderExecutionError {} + +const fn execution_error(kind: MycProviderExecutionErrorKind) -> MycProviderExecutionError { + MycProviderExecutionError { kind } +} + +enum ExecutableProvider { + EncryptedFile { + binding: MycProviderBinding, + identity: Arc<MycDecryptedIdentity>, + }, + LocalSigner { + binding: MycProviderBinding, + client: Box<MycLocalSignerClient>, + }, +} + +impl ExecutableProvider { + const fn role(&self) -> MycProviderRole { + match self { + Self::EncryptedFile { binding, .. } | Self::LocalSigner { binding, .. } => { + binding.role() + } + } + } +} + +pub(crate) struct MycProviderExecutor { + providers: Box<[ExecutableProvider]>, +} + +impl MycProviderExecutor { + pub(crate) async fn open( + runtime: &MycRuntimeContext, + configuration: &MycConfigDocumentV1, + cancellation: &MycTaskCancellation, + ) -> Result<Self, MycProviderExecutionError> { + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + let _ = (runtime, configuration, cancellation); + return Err(execution_error( + MycProviderExecutionErrorKind::UnsupportedPlatform, + )); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + let mut providers = + Vec::with_capacity(configuration.provider_contract().bindings().len()); + for binding in configuration.provider_contract().bindings() { + if cancellation.is_cancelled() { + return Err(execution_error(MycProviderExecutionErrorKind::Cancelled)); + } + match binding.kind() { + MycProviderKind::EncryptedFile => { + let runtime = runtime.clone(); + let binding = binding.clone(); + let worker_binding = binding.clone(); + let mut worker = tokio::task::spawn_blocking(move || { + let credential = + resolve_myc_wrapping_credential(&runtime, &worker_binding) + .map_err(|_| ())?; + open_myc_encrypted_identity(&worker_binding, &credential) + .map_err(|_| ()) + }); + let identity = tokio::select! { + result = &mut worker => result + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Open))? + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Open))?, + () = cancellation.cancelled() => { + let _ = worker.await; + return Err(execution_error(MycProviderExecutionErrorKind::Cancelled)); + } + }; + providers.push(ExecutableProvider::EncryptedFile { + binding, + identity: Arc::new(identity), + }); + } + MycProviderKind::LocalSigner => { + let client = MycLocalSignerClient::new(binding) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Open))?; + providers.push(ExecutableProvider::LocalSigner { + binding: binding.clone(), + client: Box::new(client), + }); + } + } + } + if providers.len() != configuration.provider_contract().bindings().len() { + return Err(execution_error(MycProviderExecutionErrorKind::Binding)); + } + Ok(Self { + providers: providers.into_boxed_slice(), + }) + } + } + + pub(crate) async fn execute( + &self, + operation: MycProviderOperation, + observed_at: MycProviderResponseObservedAtUnixMs, + cancellation: &MycTaskCancellation, + ) -> Result<MycVerifiedProviderResponse, MycProviderExecutionError> { + let provider = self + .providers + .iter() + .find(|provider| provider.role() == operation.role()) + .ok_or_else(|| execution_error(MycProviderExecutionErrorKind::Binding))?; + if cancellation.is_cancelled() { + return Err(execution_error(MycProviderExecutionErrorKind::Cancelled)); + } + match provider { + ExecutableProvider::EncryptedFile { binding, identity } => { + let binding = binding.clone(); + let identity = Arc::clone(identity); + let mut worker = tokio::task::spawn_blocking(move || { + let result = execute_encrypted(&identity, &operation)?; + Ok::<_, MycProviderExecutionError>((operation, result)) + }); + tokio::select! { + joined = &mut worker => { + let (operation, result) = joined + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))??; + verify_encrypted_provider_response(&binding, &operation, observed_at, result) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Verification)) + } + () = cancellation.cancelled() => { + // A blocking cryptographic call cannot be abandoned. Join it before + // returning cancellation so no protected operation is detached. + // The result is deliberately discarded and never becomes domain authority. + let _ = worker.await; + Err(execution_error(MycProviderExecutionErrorKind::Cancelled)) + } + } + } + ExecutableProvider::LocalSigner { binding, client } => { + tokio::select! { + result = client.execute(&operation) => { + let response = result.map_err(|_| execution_error(MycProviderExecutionErrorKind::Transport))?; + response + .verify(binding, &operation, observed_at) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Verification)) + } + () = cancellation.cancelled() => { + Err(execution_error(MycProviderExecutionErrorKind::Cancelled)) + } + } + } + } + } + + pub(crate) fn contains_role(&self, role: MycProviderRole) -> bool { + self.providers + .iter() + .any(|provider| provider.role() == role) + } +} + +impl fmt::Debug for MycProviderExecutor { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycProviderExecutor") + .field("provider_count", &self.providers.len()) + .finish() + } +} + +fn execute_encrypted( + identity: &MycDecryptedIdentity, + operation: &MycProviderOperation, +) -> Result<WireProviderResult, MycProviderExecutionError> { + if operation.provider() != MycProviderKind::EncryptedFile + || operation.expected_identity() != identity.public_identity() + { + return Err(execution_error(MycProviderExecutionErrorKind::Binding)); + } + let secret = SecretKey::from_slice(identity.secret_bytes()) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + let keys = Keys::new(secret); + let input = operation.input(); + let peer = || { + input + .peer() + .ok_or_else(|| execution_error(MycProviderExecutionErrorKind::Binding)) + .and_then(|identity| { + PublicKey::from_hex(identity.as_hex()) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Binding)) + }) + }; + let bytes = || { + input + .bytes() + .ok_or_else(|| execution_error(MycProviderExecutionErrorKind::Binding)) + }; + match input.capability() { + MycProviderCapability::Describe => Ok(WireProviderResult::Describe { + public_identity: identity.public_identity().as_hex().to_owned(), + protocol_version: MYC_LOCAL_SIGNER_TRANSPORT_CONTRACT_VERSION, + capabilities: MycProviderCapability::ALL + .into_iter() + .filter(|capability| capability_allowed_for_role(operation.role(), *capability)) + .map(WireCapability::from) + .collect(), + maximum_request_bytes: MYC_PROVIDER_INPUT_MAX_BYTES as u64, + }), + MycProviderCapability::PublicIdentity => Ok(WireProviderResult::PublicIdentity { + public_identity: identity.public_identity().as_hex().to_owned(), + }), + MycProviderCapability::SignEvent => { + let unsigned: UnsignedEvent = serde_json::from_slice(bytes()?) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + let signed = unsigned + .sign_with_keys(&keys) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + Ok(WireProviderResult::SignEvent { + payload_hex: ProtectedWireHex::from_bytes(signed.as_json().as_bytes()), + }) + } + MycProviderCapability::Nip04Encrypt => { + let peer = peer()?; + let payload = nip04::encrypt(keys.secret_key(), &peer, bytes()?) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + Ok(WireProviderResult::Nip04Encrypt { + peer: peer.to_hex(), + payload_hex: ProtectedWireHex::from_bytes(payload.as_bytes()), + }) + } + MycProviderCapability::Nip04Decrypt => { + let peer = peer()?; + let payload = core::str::from_utf8(bytes()?) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + let plaintext = nip04::decrypt(keys.secret_key(), &peer, payload) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + Ok(WireProviderResult::Nip04Decrypt { + peer: peer.to_hex(), + payload_hex: ProtectedWireHex::from_bytes(plaintext.as_bytes()), + }) + } + MycProviderCapability::Nip44Encrypt => { + let peer = peer()?; + let payload = nip44::encrypt(keys.secret_key(), &peer, bytes()?, nip44::Version::V2) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + Ok(WireProviderResult::Nip44Encrypt { + peer: peer.to_hex(), + version: 2, + payload_hex: ProtectedWireHex::from_bytes(payload.as_bytes()), + }) + } + MycProviderCapability::Nip44Decrypt => { + let peer = peer()?; + let payload = core::str::from_utf8(bytes()?) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + let plaintext = nip44::decrypt(keys.secret_key(), &peer, payload) + .map_err(|_| execution_error(MycProviderExecutionErrorKind::Operation))?; + Ok(WireProviderResult::Nip44Decrypt { + peer: peer.to_hex(), + version: 2, + payload_hex: ProtectedWireHex::from_bytes(plaintext.as_bytes()), + }) + } + } +} + +const fn capability_allowed_for_role( + role: MycProviderRole, + capability: MycProviderCapability, +) -> bool { + match role { + MycProviderRole::Transport => !matches!(capability, MycProviderCapability::SignEvent), + MycProviderRole::User => true, + MycProviderRole::Discovery => matches!( + capability, + MycProviderCapability::Describe + | MycProviderCapability::PublicIdentity + | MycProviderCapability::SignEvent + ), + } +} + +#[cfg(test)] +mod tests { + use std::error::Error as _; + + use nostr::{JsonUtil as _, Kind, Tag, Timestamp}; + + use super::*; + use crate::{ + MycConfigProfile, MycProviderCorrelationId, MycProviderDeadlineUnixMs, + MycProviderNip44Version, MycProviderOperationId, MycProviderOperationInput, + parse_myc_config_v1, + }; + + const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); + + fn keys(seed: u8) -> Keys { + Keys::parse(&format!("{seed:02x}{}", "00".repeat(31))).expect("test keys") + } + + fn secret(seed: u8) -> [u8; 32] { + let mut secret = [0_u8; 32]; + secret[0] = seed; + secret + } + + fn configuration() -> MycConfigDocumentV1 { + let source = CONFIG + .replacen( + "4444444444444444444444444444444444444444444444444444444444444444", + &keys(2).public_key().to_hex(), + 1, + ) + .replacen( + "3333333333333333333333333333333333333333333333333333333333333333", + &keys(4).public_key().to_hex(), + 1, + ); + parse_myc_config_v1(source.as_bytes(), MycConfigProfile::RepoLocal).expect("configuration") + } + + fn operation( + binding: &MycProviderBinding, + seed: u8, + input: MycProviderOperationInput, + ) -> MycProviderOperation { + MycProviderOperation::new( + binding, + MycProviderOperationId::from_bytes([seed; 32]), + MycProviderCorrelationId::from_bytes([seed.wrapping_add(1); 32]), + MycProviderDeadlineUnixMs::new(2_000_000_000_000).expect("deadline"), + input, + ) + .expect("operation") + } + + #[test] + fn encrypted_executor_produces_independently_verified_results() { + let configuration = configuration(); + let binding = configuration + .provider_contract() + .binding(MycProviderRole::Transport) + .expect("transport binding"); + let identity = MycDecryptedIdentity::from_test_secret(secret(2)); + let peer = MycDecryptedIdentity::from_test_secret(secret(5)); + let observed = + MycProviderResponseObservedAtUnixMs::new(1_999_999_999_999).expect("observed time"); + + for (seed, input) in [ + (1, MycProviderOperationInput::describe()), + (2, MycProviderOperationInput::public_identity()), + ( + 3, + MycProviderOperationInput::nip04_encrypt( + peer.public_identity().clone(), + b"nip04 protected", + ) + .expect("NIP-04 input"), + ), + ( + 4, + MycProviderOperationInput::nip44_encrypt( + peer.public_identity().clone(), + MycProviderNip44Version::V2, + b"nip44 protected", + ) + .expect("NIP-44 input"), + ), + ] { + let operation = operation(binding, seed, input); + let result = execute_encrypted(&identity, &operation).expect("encrypted result"); + let verified = + verify_encrypted_provider_response(binding, &operation, observed, result) + .expect("independent verification"); + assert!(verified.matches_operation(&operation)); + } + + let discovery = configuration + .provider_contract() + .binding(MycProviderRole::Discovery) + .expect("discovery binding"); + let discovery_identity = MycDecryptedIdentity::from_test_secret(secret(4)); + let unsigned = nostr::UnsignedEvent::new( + discovery_identity + .public_identity() + .as_hex() + .parse() + .expect("public key"), + Timestamp::from_secs(1_725_000_000), + Kind::Custom(31_990), + Vec::<Tag>::new(), + "{}", + ); + let operation = operation( + discovery, + 5, + MycProviderOperationInput::sign_event(unsigned.as_json().as_bytes()) + .expect("sign input"), + ); + let result = execute_encrypted(&discovery_identity, &operation).expect("signature"); + let verified = verify_encrypted_provider_response(discovery, &operation, observed, result) + .expect("verified signature"); + assert!(verified.matches_operation(&operation)); + } + + #[test] + fn capability_matrix_and_diagnostics_are_closed() { + for capability in MycProviderCapability::ALL { + assert_eq!( + capability_allowed_for_role(MycProviderRole::Transport, capability), + !matches!(capability, MycProviderCapability::SignEvent) + ); + assert!(capability_allowed_for_role( + MycProviderRole::User, + capability + )); + assert_eq!( + capability_allowed_for_role(MycProviderRole::Discovery, capability), + matches!( + capability, + MycProviderCapability::Describe + | MycProviderCapability::PublicIdentity + | MycProviderCapability::SignEvent + ) + ); + } + for kind in [ + MycProviderExecutionErrorKind::Binding, + MycProviderExecutionErrorKind::Open, + MycProviderExecutionErrorKind::Operation, + MycProviderExecutionErrorKind::Cancelled, + MycProviderExecutionErrorKind::Transport, + MycProviderExecutionErrorKind::Verification, + MycProviderExecutionErrorKind::UnsupportedPlatform, + ] { + let error = execution_error(kind); + assert_eq!(error.kind(), kind); + assert!(error.source().is_none()); + assert!(!format!("{error} {error:?}").contains("protected")); + } + } +} diff --git a/src/provider_verification.rs b/src/provider_verification.rs @@ -14,10 +14,10 @@ use crate::provider_local_signer::{ WireProviderResult, WireRole, }; use crate::{ - MYC_PROVIDER_OUTPUT_MAX_BYTES, MycProviderBinding, MycProviderCapability, - MycProviderCapabilitySet, MycProviderCorrelationId, MycProviderInstanceId, MycProviderKind, - MycProviderNip44Version, MycProviderOperation, MycProviderOperationId, - MycProviderPublicIdentity, MycProviderRole, + MYC_PROVIDER_INPUT_MAX_BYTES, MYC_PROVIDER_OUTPUT_MAX_BYTES, MycProviderBinding, + MycProviderCapability, MycProviderCapabilitySet, MycProviderCorrelationId, + MycProviderInstanceId, MycProviderKind, MycProviderNip44Version, MycProviderOperation, + MycProviderOperationId, MycProviderPublicIdentity, MycProviderRole, }; const NIP44_V2_VERSION: u8 = 2; @@ -308,6 +308,42 @@ impl MycLocalSignerUntrustedResponse { } } +pub(crate) fn verify_encrypted_provider_response( + binding: &MycProviderBinding, + operation: &MycProviderOperation, + observed_at: MycProviderResponseObservedAtUnixMs, + result: WireProviderResult, +) -> Result<MycVerifiedProviderResponse, MycProviderVerificationError> { + if binding.kind() != MycProviderKind::EncryptedFile + || operation.provider() != MycProviderKind::EncryptedFile + || operation.role() != binding.role() + || operation.instance() != binding.instance() + || operation.expected_identity() != binding.expected_identity() + || !binding + .required_capabilities() + .contains(operation.input().capability()) + { + return Err(verification_error( + MycProviderVerificationErrorKind::InvalidBinding, + )); + } + if observed_at.get() > operation.deadline().get() { + return Err(verification_error( + MycProviderVerificationErrorKind::LateResponse, + )); + } + let result = verify_result(binding, operation, result)?; + Ok(MycVerifiedProviderResponse { + operation_id: operation.operation_id(), + correlation_id: operation.correlation_id(), + instance: operation.instance(), + role: operation.role(), + capability: operation.input().capability(), + operation_binding: operation.binding_digest(), + result, + }) +} + fn verify_response( binding: &MycProviderBinding, operation: &MycProviderOperation, @@ -555,10 +591,14 @@ fn verify_describe( MycProviderVerificationErrorKind::Capability, )); } - let limits = binding - .local_signer_limits() - .ok_or_else(|| verification_error(MycProviderVerificationErrorKind::InvalidBinding))?; - if maximum_request_bytes != limits.request_max_bytes() { + let expected_maximum = match binding.kind() { + MycProviderKind::EncryptedFile => MYC_PROVIDER_INPUT_MAX_BYTES as u64, + MycProviderKind::LocalSigner => binding + .local_signer_limits() + .ok_or_else(|| verification_error(MycProviderVerificationErrorKind::InvalidBinding))? + .request_max_bytes(), + }; + if maximum_request_bytes != expected_maximum { return Err(verification_error(MycProviderVerificationErrorKind::Size)); } Ok(VerifiedProviderResult::Describe { diff --git a/src/runtime_supervision.rs b/src/runtime_supervision.rs @@ -42,6 +42,17 @@ impl MycTaskCancellation { pub async fn cancelled(&self) { self.inner.cancelled().await; } + + #[cfg(test)] + pub(crate) fn test_pair() -> (Self, CancellationToken) { + let token = CancellationToken::new(); + ( + Self { + inner: token.clone(), + }, + token, + ) + } } impl fmt::Debug for MycTaskCancellation { diff --git a/src/transport_nostr_adapter.rs b/src/transport_nostr_adapter.rs @@ -0,0 +1,405 @@ +//! Exact source-locked Nostr delivery adapter owned by the Myc runtime. + +#![allow( + dead_code, + reason = "Step 159 Unit 12 seals the adapter before Unit 13 runtime graph wiring" +)] + +use core::{fmt, future::Future, pin::Pin}; +use std::{collections::BTreeMap, error::Error}; + +use radroots_event_codec::Codec; +use radroots_transport::{ + Target, TargetSet, + outcome::DeliveryOutcomeKind, + policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, + sink::{DeliveryPayload, DeliveryRequest, DeliveryTargetReceipt}, +}; +use radroots_transport_nostr::{ + Config, NostrTransport, PreparedDelivery, RelayAccess, RelayEndpoint, RelayProfile, + RelayProfileKind, RelayUrlPolicy, +}; + +use crate::{MycConfigDocumentV1, MycConfigProfile, MycDeliveryRelayId}; + +pub(crate) type RelayExecutionFuture<'a> = Pin< + Box<dyn Future<Output = Result<MycRelayExecutionOutcome, MycRelayAdapterError>> + Send + 'a>, +>; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum MycRelayAdapterErrorKind { + Configuration, + Target, + Payload, + Preparation, + Execution, +} + +pub(crate) struct MycRelayAdapterError { + kind: MycRelayAdapterErrorKind, +} + +impl MycRelayAdapterError { + pub(crate) const fn kind(&self) -> MycRelayAdapterErrorKind { + self.kind + } +} + +impl fmt::Debug for MycRelayAdapterError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycRelayAdapterError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for MycRelayAdapterError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("Myc relay adapter failed") + } +} + +impl Error for MycRelayAdapterError {} + +const fn adapter_error(kind: MycRelayAdapterErrorKind) -> MycRelayAdapterError { + MycRelayAdapterError { kind } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum MycRelayExecutionOutcome { + Accepted, + Rejected, + TransportFailed, + UnknownAcknowledgement, +} + +pub(crate) struct MycPreparedRelayDelivery { + transport: NostrTransport, + prepared: PreparedDelivery, +} + +impl fmt::Debug for MycPreparedRelayDelivery { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycPreparedRelayDelivery([redacted])") + } +} + +pub(crate) trait MycRelayAdapter: Send + Sync { + type Prepared: Send; + + fn prepare( + &self, + relay_id: &MycDeliveryRelayId, + request_id: String, + exact_event_bytes: &[u8], + deadline_unix_ms: u64, + ) -> Result<Self::Prepared, MycRelayAdapterError>; + + fn execute<'a>(&'a self, prepared: Self::Prepared) -> RelayExecutionFuture<'a>; +} + +pub(crate) struct MycNostrDeliveryAdapter { + targets: BTreeMap<MycDeliveryRelayId, RelayTarget>, +} + +struct RelayTarget { + transport: NostrTransport, + target: Target, +} + +impl MycNostrDeliveryAdapter { + pub(crate) fn from_configuration( + configuration: &MycConfigDocumentV1, + ) -> Result<Self, MycRelayAdapterError> { + let relays = configuration + .normalized() + .pointer("/relays") + .and_then(serde_json::Value::as_array) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let connect_timeout = + configuration_integer(configuration, "/transport/connect_deadline_ms")?; + let request_timeout = configuration_integer( + configuration, + "/transport/publish_retry/attempt_deadline_ms", + )?; + let mut public = Vec::new(); + let mut local = Vec::new(); + let mut definitions = Vec::with_capacity(relays.len()); + for relay in relays { + let id = relay + .pointer("/id") + .and_then(serde_json::Value::as_str) + .and_then(|value| MycDeliveryRelayId::new(value).ok()) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let url = relay + .pointer("/url") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let read = relay + .pointer("/read") + .and_then(serde_json::Value::as_bool) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let write = relay + .pointer("/write") + .and_then(serde_json::Value::as_bool) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let access = if write { + RelayAccess::ReadWrite + } else if read { + RelayAccess::ReadOnly + } else { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + }; + let (kind, policy) = if url.starts_with("wss://") { + (RelayProfileKind::Public, RelayUrlPolicy::Public) + } else if configuration.profile() == MycConfigProfile::RepoLocal + && url.starts_with("ws://") + { + (RelayProfileKind::Simulator, RelayUrlPolicy::Local) + } else { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + }; + let endpoint = RelayEndpoint::new(url, policy, access) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + match kind { + RelayProfileKind::Public => public.push(endpoint), + RelayProfileKind::Simulator => local.push(endpoint), + RelayProfileKind::Device => { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + _ => return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)), + } + definitions.push((id, url.to_owned(), kind)); + } + let public_transport = build_transport( + RelayProfileKind::Public, + public, + connect_timeout, + request_timeout, + )?; + let local_transport = build_transport( + RelayProfileKind::Simulator, + local, + connect_timeout, + request_timeout, + )?; + let mut targets = BTreeMap::new(); + for (id, url, kind) in definitions { + let transport = match kind { + RelayProfileKind::Public => public_transport.clone(), + RelayProfileKind::Simulator => local_transport.clone(), + RelayProfileKind::Device => None, + _ => None, + } + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let target = Target::nostr_relay(url) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; + if targets + .insert(id, RelayTarget { transport, target }) + .is_some() + { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + } + if targets.is_empty() { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + Ok(Self { targets }) + } +} + +impl MycRelayAdapter for MycNostrDeliveryAdapter { + type Prepared = MycPreparedRelayDelivery; + + fn prepare( + &self, + relay_id: &MycDeliveryRelayId, + request_id: String, + exact_event_bytes: &[u8], + deadline_unix_ms: u64, + ) -> Result<Self::Prepared, MycRelayAdapterError> { + let binding = self + .targets + .get(relay_id) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Target))?; + let raw = core::str::from_utf8(exact_event_bytes) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Payload))?; + let signed = Codec::decode_signed_event(raw) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Payload))?; + if signed.raw_json().as_bytes() != exact_event_bytes { + return Err(adapter_error(MycRelayAdapterErrorKind::Payload)); + } + let targets = TargetSet::new(vec![binding.target.clone()]) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; + let request = DeliveryRequest::new( + request_id, + DeliveryPayload::new(signed), + targets, + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + deadline_unix_ms, + ) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Preparation))?; + let prepared = binding + .transport + .prepare_delivery(request) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Preparation))?; + Ok(MycPreparedRelayDelivery { + transport: binding.transport.clone(), + prepared, + }) + } + + fn execute<'a>(&'a self, prepared: Self::Prepared) -> RelayExecutionFuture<'a> { + Box::pin(async move { + let MycPreparedRelayDelivery { + transport, + prepared, + } = prepared; + match transport.execute_prepared_delivery(prepared).await { + Err(_) => Ok(MycRelayExecutionOutcome::UnknownAcknowledgement), + Ok(receipt) => classify_receipt(receipt.target_receipts()), + } + }) + } +} + +fn classify_receipt( + receipts: &[DeliveryTargetReceipt], +) -> Result<MycRelayExecutionOutcome, MycRelayAdapterError> { + let [receipt] = receipts else { + return Err(adapter_error(MycRelayAdapterErrorKind::Execution)); + }; + Ok(match receipt.outcome().kind() { + DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered => { + MycRelayExecutionOutcome::Accepted + } + DeliveryOutcomeKind::Rejected => MycRelayExecutionOutcome::Rejected, + DeliveryOutcomeKind::Unavailable | DeliveryOutcomeKind::Failed => { + MycRelayExecutionOutcome::TransportFailed + } + }) +} + +fn build_transport( + kind: RelayProfileKind, + endpoints: Vec<RelayEndpoint>, + connect_timeout: u64, + request_timeout: u64, +) -> Result<Option<NostrTransport>, MycRelayAdapterError> { + if endpoints.is_empty() { + return Ok(None); + } + let maximum_connections = endpoints.len().min(8); + let profile = RelayProfile::explicit(kind, endpoints) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let config = Config::from_profile(profile) + .with_timeouts(connect_timeout, request_timeout, connect_timeout) + .and_then(|config| config.with_max_connections(maximum_connections)) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + Ok(Some(NostrTransport::new(config))) +} + +fn configuration_integer( + configuration: &MycConfigDocumentV1, + pointer: &str, +) -> Result<u64, MycRelayAdapterError> { + configuration + .normalized() + .pointer(pointer) + .and_then(serde_json::Value::as_u64) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration)) +} + +#[cfg(test)] +mod tests { + use std::error::Error as _; + + use nostr::{EventBuilder, JsonUtil as _, Keys}; + use radroots_transport::{ + outcome::{DeliveryOutcome, Retryability}, + sink::DeliveryTargetReceipt, + }; + + use super::*; + use crate::{MycConfigProfile, parse_myc_config_v1}; + + const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); + + fn receipt(outcome: DeliveryOutcome) -> DeliveryTargetReceipt { + DeliveryTargetReceipt::attempted( + Target::nostr_relay("wss://relay.example.test/").expect("target"), + outcome, + ) + } + + #[test] + fn receipt_classification_preserves_all_four_durable_outcomes() { + for (outcome, expected) in [ + ( + DeliveryOutcome::accepted(), + MycRelayExecutionOutcome::Accepted, + ), + ( + DeliveryOutcome::delivered(), + MycRelayExecutionOutcome::Accepted, + ), + ( + DeliveryOutcome::rejected(), + MycRelayExecutionOutcome::Rejected, + ), + ( + DeliveryOutcome::unavailable(), + MycRelayExecutionOutcome::TransportFailed, + ), + ( + DeliveryOutcome::failed(Retryability::Retryable).expect("failed"), + MycRelayExecutionOutcome::TransportFailed, + ), + ] { + assert_eq!(classify_receipt(&[receipt(outcome)]).unwrap(), expected); + } + assert!(classify_receipt(&[]).is_err()); + assert!( + classify_receipt(&[ + receipt(DeliveryOutcome::accepted()), + receipt(DeliveryOutcome::accepted()), + ]) + .is_err() + ); + } + + #[test] + fn exact_config_builds_without_network_io_and_errors_are_redacted() { + let configuration = parse_myc_config_v1(CONFIG.as_bytes(), MycConfigProfile::RepoLocal) + .expect("configuration"); + let adapter = MycNostrDeliveryAdapter::from_configuration(&configuration) + .expect("offline adapter construction"); + assert_eq!(adapter.targets.len(), 2); + let signing_keys = Keys::parse(&format!("02{}", "00".repeat(31))).expect("test keys"); + let event = EventBuilder::text_note("exact committed payload") + .sign_with_keys(&signing_keys) + .expect("signed event"); + let exact = event.as_json(); + let relay = MycDeliveryRelayId::new("primary").expect("relay"); + let prepared = adapter + .prepare(&relay, "offline-prepare".to_owned(), exact.as_bytes(), 1) + .expect("offline prepared delivery"); + assert_eq!( + format!("{prepared:?}"), + "MycPreparedRelayDelivery([redacted])" + ); + for kind in [ + MycRelayAdapterErrorKind::Configuration, + MycRelayAdapterErrorKind::Target, + MycRelayAdapterErrorKind::Payload, + MycRelayAdapterErrorKind::Preparation, + MycRelayAdapterErrorKind::Execution, + ] { + let error = adapter_error(kind); + assert_eq!(error.kind(), kind); + assert!(error.source().is_none()); + assert!(!format!("{error} {error:?}").contains("relay.example")); + } + } +} diff --git a/tests/build_policy.rs b/tests/build_policy.rs @@ -56,6 +56,22 @@ fn shared_storage_evidence_is_exactly_pinned_to_the_source_locked_lib() { } #[test] +fn delivery_dependencies_are_exactly_source_locked() { + for dependency in [ + "radroots_event_codec", + "radroots_transport", + "radroots_transport_nostr", + ] { + assert!( + MANIFEST.contains(&format!( + "{dependency} = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"7d7b454b4c9ed86569671993bd03ca868b676665\", version = \"=0.1.0-alpha\"" + )), + "{dependency} is not pinned to the exact source lock" + ); + } +} + +#[test] fn release_acceptance_checks_both_feature_profiles() { assert!( RELEASE_ACCEPTANCE.contains("cargo check --locked --all-targets --no-default-features\n") diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -14,6 +14,11 @@ const NIP46_WAVE_080_A: &str = include_str!("../src/nip46_wave_080_a.rs"); const NIP46_COMPLETION: &str = include_str!("../src/state_completion.rs"); const NIP46_RESPONSE: &str = include_str!("../src/state_response.rs"); const DELIVERY_RECOVERY: &str = include_str!("../src/state_recovery.rs"); +const DELIVERY_WORKER: &str = include_str!("../src/delivery_worker.rs"); +const PROVIDER_EXECUTOR: &str = include_str!("../src/provider_executor.rs"); +const TRANSPORT_NOSTR_ADAPTER: &str = include_str!("../src/transport_nostr_adapter.rs"); +const PROVIDER_DELIVERY_CONTRACT: &str = + include_str!("../contracts/services_hardening/provider_delivery.v1.json"); const DOCTOR_V1: &str = include_str!("../src/doctor_v1.rs"); const CONTROL_PLANE_WAVE_090_A: &str = include_str!("../src/control_plane_wave_090_a.rs"); const CONTROL_PLANE_WAVE_090_A_CONTRACT: &str = @@ -52,6 +57,7 @@ const SOURCES: &[&str] = &[ include_str!("../src/cli_v1.rs"), include_str!("../src/config_v1.rs"), include_str!("../src/control_plane_wave_090_a.rs"), + include_str!("../src/delivery_worker.rs"), include_str!("../src/doctor_v1.rs"), include_str!("../src/diagnostics_v1.rs"), include_str!("../src/nip46_admission.rs"), @@ -63,12 +69,14 @@ const SOURCES: &[&str] = &[ include_str!("../src/provider_contract.rs"), include_str!("../src/provider_credential.rs"), include_str!("../src/provider_envelope.rs"), + include_str!("../src/provider_executor.rs"), include_str!("../src/provider_local_signer.rs"), include_str!("../src/provider_verification.rs"), include_str!("../src/runtime_context.rs"), include_str!("../src/runtime_foundation.rs"), include_str!("../src/runtime_supervision.rs"), include_str!("../src/status_v1.rs"), + include_str!("../src/transport_nostr_adapter.rs"), include_str!("../src/state_catalog.rs"), include_str!("../src/state_admin.rs"), include_str!("../src/state_completion.rs"), @@ -98,6 +106,7 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() { "cli_v1", "config_v1", "control_plane_wave_090_a", + "delivery_worker", "doctor_v1", "diagnostics_v1", "nip46_admission", @@ -111,12 +120,14 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() { "provider_contract", "provider_credential", "provider_envelope", + "provider_executor", "provider_local_signer", "provider_verification", "runtime_context", "runtime_foundation", "runtime_supervision", "status_v1", + "transport_nostr_adapter", "state_catalog", "state_admin", "state_completion", @@ -152,6 +163,9 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() { "caps a\nreplayed response model at 8,192 bytes", "admits at least 8,382 UTF-8 bytes", "The journal stores no request body, path,\ncorrelation ID, credential, bundle path, or secret", + "The Step 159 provider and delivery boundary is sealed inside the crate", + "persists Submitted immediately before execution", + "Runtime task-graph wiring and startup handshakes remain the next\nordered Step 159 unit", ] { assert!(README.contains(required), "README is missing `{required}`"); } @@ -284,6 +298,7 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() { "cli_v1", "config_v1", "control_plane_wave_090_a", + "delivery_worker", "doctor_v1", "diagnostics_v1", "nip46_admission", @@ -297,12 +312,14 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() { "provider_contract", "provider_credential", "provider_envelope", + "provider_executor", "provider_local_signer", "provider_verification", "runtime_context", "runtime_foundation", "runtime_supervision", "status_v1", + "transport_nostr_adapter", "state_catalog", "state_admin", "state_completion", @@ -1122,3 +1139,86 @@ fn step149_recovery_and_offline_export_are_bounded_and_non_networked() { ); } } + +#[test] +fn step159_provider_delivery_is_sealed_exact_and_durability_ordered() { + let contract: serde_json::Value = serde_json::from_str(PROVIDER_DELIVERY_CONTRACT) + .expect("Step 159 provider-delivery contract"); + assert_eq!(contract["schema"], "radroots.myc.provider-delivery.v1"); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["step"], 159); + assert_eq!(contract["unit"], "myc-provider-delivery"); + assert_eq!( + contract["relay_adapter"]["implementation"], + "radroots_transport_nostr" + ); + assert_eq!( + contract["durable_delivery"]["submitted_transition"], + "immediately_after_prepare_before_execute" + ); + for required in [ + "spawn_blocking", + "let _ = worker.await", + "verify_encrypted_provider_response", + "MycLocalSignerClient", + ] { + assert!( + PROVIDER_EXECUTOR.contains(required), + "provider executor is missing `{required}`" + ); + } + for required in [ + "radroots_transport_nostr", + "prepare_delivery(request)", + "execute_prepared_delivery(prepared).await", + "DeliveryOutcomeKind::Accepted", + "DeliveryOutcomeKind::Rejected", + "DeliveryOutcomeKind::Unavailable", + ] { + assert!( + TRANSPORT_NOSTR_ADAPTER.contains(required), + "transport adapter is missing `{required}`" + ); + } + for required in [ + "claim_delivery_target", + "read_nip46_response", + "read_discovery_document_for_job", + "mark_delivery_attempt_submitted", + "adapter.execute(prepared)", + "MycDeliveryAttemptOutcome::UnknownAcknowledgement", + "record_delivery_attempt_outcome", + ] { + assert!( + DELIVERY_WORKER.contains(required), + "delivery worker is missing `{required}`" + ); + } + for forbidden in [ + "pub struct myc::MycProviderExecutor", + "pub struct myc::MycDeliveryWorker", + "pub struct myc::MycNostrDeliveryAdapter", + "radroots_transport_nostr::", + "radroots_transport::", + ] { + assert!( + !PUBLIC_API.contains(forbidden), + "sealed delivery authority escaped: `{forbidden}`" + ); + } + for source in [PROVIDER_EXECUTOR, TRANSPORT_NOSTR_ADAPTER, DELIVERY_WORKER] { + for forbidden in [ + "std::time::SystemTime", + "Timestamp::now", + "getrandom", + "rand::", + "tokio::runtime::Runtime", + "tokio::spawn(", + ] { + assert!( + !source.contains(forbidden), + "Step 159 provider-delivery gained `{forbidden}`" + ); + } + } +} diff --git a/tests/services_hardening_legacy_removal.rs b/tests/services_hardening_legacy_removal.rs @@ -9,8 +9,10 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MAIN_SOURCE: &str = include_str!("../src/main.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); const ACTIVE_STATE_SOURCES: &[&str] = &[ + include_str!("../src/state_admin.rs"), include_str!("../src/state_catalog.rs"), include_str!("../src/state_completion.rs"), + include_str!("../src/state_config.rs"), include_str!("../src/state_connection.rs"), include_str!("../src/state_delivery.rs"), include_str!("../src/state_discovery.rs"), @@ -120,7 +122,7 @@ fn prototype_environment_and_cli_sources_are_absent() { #[test] fn active_state_tree_has_one_shared_database_and_no_legacy_backend() { - assert_eq!(LIB_SOURCE.matches("mod state_").count(), 13); + assert_eq!(LIB_SOURCE.matches("mod state_").count(), 15); assert!(!LIB_SOURCE.contains("pub mod state_")); let active_state = ACTIVE_STATE_SOURCES.join("\n"); for forbidden in [ diff --git a/tests/services_hardening_native_release.rs b/tests/services_hardening_native_release.rs @@ -149,7 +149,7 @@ fn every_radroots_dependency_is_exactly_source_locked() { .iter() .filter(|(name, _)| name.starts_with("radroots_")) .collect::<Vec<_>>(); - assert_eq!(radroots.len(), 7); + assert_eq!(radroots.len(), 10); for (name, dependency) in radroots { let dependency = dependency.as_table().expect("detailed dependency"); assert_eq!( diff --git a/tests/services_hardening_signer_request_state.rs b/tests/services_hardening_signer_request_state.rs @@ -458,7 +458,7 @@ async fn concurrent_identical_admission_creates_one_request_and_bounded_replay_e } #[tokio::test] -async fn exact_schema_v3_state_advances_to_v9_before_request_admission() { +async fn exact_schema_v3_state_advances_to_v11_before_request_admission() { let directory = tempfile::tempdir().expect("temporary root"); let runtime = runtime(directory.path()); prepare_state_directory(&runtime); @@ -513,6 +513,12 @@ async fn exact_schema_v3_state_advances_to_v9_before_request_admission() { "DROP TRIGGER connection_permissions_no_update", "DROP TRIGGER connections_no_delete", "DROP TRIGGER connections_guard_update", + "DROP TRIGGER myc_admin_operations_guard_update", + "DROP TABLE myc_admin_operations", + "DROP TRIGGER myc_config_bindings_no_delete", + "DROP TRIGGER myc_config_bindings_no_update", + "DROP TRIGGER myc_config_bindings_guard_insert", + "DROP TABLE myc_config_bindings", "DROP TRIGGER nip46_signed_responses_no_delete", "DROP TRIGGER nip46_signed_responses_no_update", "DROP TABLE nip46_signed_responses", @@ -535,7 +541,7 @@ async fn exact_schema_v3_state_advances_to_v9_before_request_admission() { "DROP TABLE connections", "UPDATE radroots_service_metadata SET state_schema_version = 3 WHERE singleton = 1", "UPDATE myc_state_metadata SET state_contract_version = 3 WHERE singleton = 1", - "DELETE FROM schema_migrations WHERE version IN (4, 5, 6, 7, 8, 9)", + "DELETE FROM schema_migrations WHERE version IN (4, 5, 6, 7, 8, 9, 10, 11)", ] { sqlx::query(sql) .execute(&mut connection)