commit 186f08618b22f96b45dc3e54609cfa15cb86a330
parent d156f91c18f5c8a7a76570a06470c826c88019ed
Author: triesap <tyson@radroots.org>
Date: Sat, 22 Aug 2026 03:12:09 +0000
state: commit NIP-46 responses atomically
Diffstat:
23 files changed, 1295 insertions(+), 69 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -95,6 +95,11 @@
source archive, and the generated service lock. Native target metadata does
not qualify an artifact; Nix, OCI, signing, tags, publication, and deployment
remain deferred.
+- Step 148 closes the production NIP-46 response authority in
+ `contracts/services_hardening/nip46_response_commit.v1.json`. Production
+ code may commit a completion only through the atomic exact-response method;
+ a retained completion without its immutable response and initial delivery
+ state is inconsistent evidence and must never be repaired implicitly.
- 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.toml b/Cargo.toml
@@ -17,7 +17,7 @@ resolver = "3"
service = "myc"
host_feature_profile = "service-host"
config_contract_version = 1
-state_contract_version = 8
+state_contract_version = 9
admin_contract_version = 1
status_contract_version = 1
provider_contract_version = 1
diff --git a/README b/README
@@ -75,10 +75,10 @@ without changing the canonical common artifact inventory.
Create-new initialization reserves the shared schema-v1 metadata and migration
ledger, retains exclusive writer authority, applies the exact Myc schema-v2
-through schema-v8 migrations, binds the normalized configuration, expected
+through schema-v9 migrations, binds the normalized configuration, expected
identity roles, and policy versions through a sealed typed repository, and
explicitly closes the host before reporting success. Existing writable open can
-resume any exact v1 through v7 prefix; read-only inspection requires the current
+resume any exact v1 through v8 prefix; read-only inspection requires the current
catalog and exact immutable Myc binding.
Schema v8 adds immutable NIP-46 operation-completion evidence. The Step 147
@@ -93,6 +93,15 @@ must compose it with the outer signed response bytes, immutable relay targets,
and initial outbox state in one transaction before RCLD-RSHR-080 can be
promoted to `master`.
+Schema v9 closes that production boundary. `commit_nip46_response` admits only
+an independently signature-verified, canonical kind-24133 response bound to the
+original client and exact provider signing operation. One SQLite transaction
+commits the Step 147 completion, exact response bytes and SHA-256 identity,
+immutable configured relay targets, pending delivery job, and zero-attempt
+target state. Exact replay returns only the retained response bytes; partial
+legacy completion state and mismatched replay fail closed. Provider execution
+and relay I/O remain outside the transaction.
+
The public Myc repository exposes no raw pool, connection, transaction-control
handle, path, or SQL. Binding and request-admission mutations execute only
inside the shared `ServiceSqliteTransaction` runner. Provider and relay work
diff --git a/contracts/api_baselines/myc.txt b/contracts/api_baselines/myc.txt
@@ -320,6 +320,19 @@ pub myc::MycNip46ReplayDisposition::ConflictingRequestReuse
pub myc::MycNip46ReplayDisposition::DistinctRequest
pub myc::MycNip46ReplayDisposition::DuplicateEvent
pub myc::MycNip46ReplayDisposition::ExactRequestReplay
+pub enum myc::MycNip46ResponseCommitAdmission
+pub myc::MycNip46ResponseCommitAdmission::Committed(myc::MycNip46ResponseCommitRecord)
+pub myc::MycNip46ResponseCommitAdmission::ExactReplay(myc::MycNip46ResponseCommitRecord)
+impl myc::MycNip46ResponseCommitAdmission
+pub const fn myc::MycNip46ResponseCommitAdmission::record(&self) -> &myc::MycNip46ResponseCommitRecord
+impl core::fmt::Debug for myc::MycNip46ResponseCommitAdmission
+pub fn myc::MycNip46ResponseCommitAdmission::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub enum myc::MycNip46ResponseCommitErrorKind
+pub myc::MycNip46ResponseCommitErrorKind::InvalidBinding
+pub myc::MycNip46ResponseCommitErrorKind::InvalidResponse
+pub myc::MycNip46ResponseCommitErrorKind::InvalidTime
+impl myc::MycNip46ResponseCommitErrorKind
+pub const fn myc::MycNip46ResponseCommitErrorKind::code(self) -> &'static str
pub enum myc::MycNip46SessionEffect
pub myc::MycNip46SessionEffect::ConnectionAdmitted
pub myc::MycNip46SessionEffect::ConnectionRevoked
@@ -1036,6 +1049,38 @@ pub fn myc::MycNip46RequestId::as_str(&self) -> &str
pub fn myc::MycNip46RequestId::new(&str) -> core::result::Result<Self, myc::MycSignerRequestError>
impl core::fmt::Debug for myc::MycNip46RequestId
pub fn myc::MycNip46RequestId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycNip46ResponseCommitError
+impl myc::MycNip46ResponseCommitError
+pub const fn myc::MycNip46ResponseCommitError::code(self) -> &'static str
+pub const fn myc::MycNip46ResponseCommitError::kind(self) -> myc::MycNip46ResponseCommitErrorKind
+impl core::error::Error for myc::MycNip46ResponseCommitError
+impl core::fmt::Debug for myc::MycNip46ResponseCommitError
+pub fn myc::MycNip46ResponseCommitError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl core::fmt::Display for myc::MycNip46ResponseCommitError
+pub fn myc::MycNip46ResponseCommitError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycNip46ResponseCommitRecord
+impl myc::MycNip46ResponseCommitRecord
+pub const fn myc::MycNip46ResponseCommitRecord::completion(&self) -> &myc::MycNip46CommitRecord
+pub const fn myc::MycNip46ResponseCommitRecord::response(&self) -> &myc::MycNip46ResponseRecord
+impl core::fmt::Debug for myc::MycNip46ResponseCommitRecord
+pub fn myc::MycNip46ResponseCommitRecord::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycNip46ResponseCommitRequest
+impl myc::MycNip46ResponseCommitRequest
+pub fn myc::MycNip46ResponseCommitRequest::new(&myc::MycNip46CommitRequest, &myc::MycProviderOperation, &myc::MycVerifiedProviderResponse, myc::MycDeliveryTimeUnixMs) -> core::result::Result<Self, myc::MycNip46ResponseCommitError>
+impl core::fmt::Debug for myc::MycNip46ResponseCommitRequest
+pub fn myc::MycNip46ResponseCommitRequest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycNip46ResponseRecord
+impl myc::MycNip46ResponseRecord
+pub const fn myc::MycNip46ResponseRecord::authored_at_unix_s(&self) -> u64
+pub const fn myc::MycNip46ResponseRecord::committed_at(&self) -> myc::MycDeliveryTimeUnixMs
+pub const fn myc::MycNip46ResponseRecord::delivery_job(&self) -> &myc::MycDeliveryJobRecord
+pub const fn myc::MycNip46ResponseRecord::operation_id(&self) -> myc::MycSignerOperationId
+pub const fn myc::MycNip46ResponseRecord::response_digest(&self) -> myc::MycDeliveryArtifactDigest
+pub const fn myc::MycNip46ResponseRecord::response_event_id(&self) -> &[u8; 32]
+pub const fn myc::MycNip46ResponseRecord::response_provider_operation_id(&self) -> &[u8; 32]
+pub fn myc::MycNip46ResponseRecord::signed_response_bytes(&self) -> &[u8]
+impl core::fmt::Debug for myc::MycNip46ResponseRecord
+pub fn myc::MycNip46ResponseRecord::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct myc::MycNip46VerificationError
impl myc::MycNip46VerificationError
pub const fn myc::MycNip46VerificationError::kind(self) -> myc::MycNip46VerificationErrorKind
@@ -1381,7 +1426,8 @@ pub async fn myc::MycStateRepository<'_>::promote_delivered_discovery_state(&sel
pub async fn myc::MycStateRepository<'_>::read_discovery_document_for_job(&self, myc::MycDeliveryJobId) -> core::result::Result<core::option::Option<myc::MycDiscoveryDocumentRecord>, myc::MycStateRepositoryError>
pub async fn myc::MycStateRepository<'_>::read_discovery_publication_state(&self) -> core::result::Result<core::option::Option<myc::MycDiscoveryPublicationState>, myc::MycStateRepositoryError>
impl myc::MycStateRepository<'_>
-pub async fn myc::MycStateRepository<'_>::commit_nip46_operation(&self, &myc::MycNip46CommitRequest) -> core::result::Result<myc::MycNip46CommitAdmission, myc::MycStateRepositoryError>
+pub async fn myc::MycStateRepository<'_>::commit_nip46_response(&self, &myc::MycNip46ResponseCommitRequest) -> core::result::Result<myc::MycNip46ResponseCommitAdmission, myc::MycStateRepositoryError>
+pub async fn myc::MycStateRepository<'_>::read_nip46_response(&self, myc::MycDeliveryJobId) -> core::result::Result<core::option::Option<myc::MycNip46ResponseRecord>, myc::MycStateRepositoryError>
impl myc::MycStateRepository<'_>
pub async fn myc::MycStateRepository<'_>::compact_governance_evidence(&self, myc::MycConnectionTimeUnixMs, myc::MycGovernanceCompactionPolicy, myc::MycAuditCorrelationId) -> core::result::Result<myc::MycGovernanceCompactionOutcome, myc::MycStateRepositoryError>
pub async fn myc::MycStateRepository<'_>::read_audit_page(&self, myc::MycAuditPageLimit, core::option::Option<u64>, core::option::Option<u64>) -> core::result::Result<myc::MycAuditPage, myc::MycStateRepositoryError>
@@ -1513,6 +1559,9 @@ pub const myc::MYC_STATE_SCHEMA_VERSION_7_SHA256: [u8; 32]
pub const myc::MYC_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256: [u8; 32]
pub const myc::MYC_STATE_SCHEMA_VERSION_8_OBJECT_COUNT: u32
pub const myc::MYC_STATE_SCHEMA_VERSION_8_SHA256: [u8; 32]
+pub const myc::MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256: [u8; 32]
+pub const myc::MYC_STATE_SCHEMA_VERSION_9_OBJECT_COUNT: u32
+pub const myc::MYC_STATE_SCHEMA_VERSION_9_SHA256: [u8; 32]
pub const myc::MYC_WRAPPING_CREDENTIAL_ARTIFACT_BYTES: usize
pub const myc::MYC_WRAPPING_CREDENTIAL_CONTRACT_VERSION: u32
pub fn myc::admit_myc_nip46_event(myc::MycNip46AdmissionLimits, &[u8]) -> core::result::Result<myc::MycBoundedNip46Event, myc::MycNip46AdmissionError>
diff --git a/contracts/services_hardening/native_release.v1.json b/contracts/services_hardening/native_release.v1.json
@@ -32,7 +32,7 @@
},
"contract_versions": {
"config": 1,
- "state": 8,
+ "state": 9,
"admin": 1,
"status": 1,
"provider": 1
diff --git a/contracts/services_hardening/nip46_completion.v1.json b/contracts/services_hardening/nip46_completion.v1.json
@@ -30,7 +30,8 @@
},
"production_composition": {
"step147_checkpoint": "integration_only_not_promotable",
- "step147_component": "must_be_composed_before_master_promotion",
+ "step147_component": "composed_by_step148_atomic_response_commit",
+ "composition_status": "satisfied_on_rcld_080_integration_branch",
"commit_owner": 148,
"required_same_transaction_members": [
"step147_completion",
diff --git a/contracts/services_hardening/nip46_response_commit.v1.json b/contracts/services_hardening/nip46_response_commit.v1.json
@@ -0,0 +1,72 @@
+{
+ "schema": "radroots.myc.nip46-response-commit.v1",
+ "contract_version": 1,
+ "step": 148,
+ "schema_version": 9,
+ "authority": {
+ "completion": "sealed_step147_completion_component",
+ "response": "independently_signature_verified_canonical_kind_24133_event",
+ "response_signer": "exact_bound_user_provider_operation",
+ "recipient": "exact_original_nip46_client_public_key",
+ "targets_and_retry_policy": "normalized_configuration_bound_in_state_metadata",
+ "time": "caller_injected_positive_unix_milliseconds",
+ "transaction": "existing_service_sqlite_transaction_runner"
+ },
+ "atomic_commit": [
+ "request_decision_and_session_effect",
+ "immutable_completion_evidence",
+ "exact_signed_response_bytes_sha256_and_event_id",
+ "immutable_relay_target_set",
+ "pending_delivery_job",
+ "zero_attempt_target_state"
+ ],
+ "replay": {
+ "source": "committed_response_bytes_only",
+ "same_bound_commit": "exact_replay",
+ "mismatched_commit": "fail_closed",
+ "completion_without_response_or_job": "fail_closed_without_repair",
+ "response_or_job_without_complete_composition": "fail_closed",
+ "failed_transaction": "no_completion_response_job_target_or_session_effect"
+ },
+ "response_event": {
+ "kind": 24133,
+ "canonical_bytes": true,
+ "signature_verified": true,
+ "single_recipient_tag": "original_client_public_key",
+ "maximum_bytes": 1048576,
+ "debug": "redacted"
+ },
+ "delivery": {
+ "job_state": "pending",
+ "target_state": "pending",
+ "attempt_count": 0,
+ "target_count_maximum": 32,
+ "retry_bytes": "exact_committed_response_bytes",
+ "relay_io": "after_local_commit_only"
+ },
+ "qualification": [
+ "completion_edge_failure_rolls_back_all_effects",
+ "response_edge_failure_rolls_back_all_effects",
+ "exact_replay_returns_original_bytes",
+ "targets_are_config_bound_and_zero_attempt",
+ "partial_step147_state_is_not_repaired",
+ "public_errors_and_debug_are_source_free_and_redacted"
+ ],
+ "deferrals": {
+ "relay_submission_and_acknowledgement_execution": 149,
+ "crash_recovery_and_full_wave_matrix": 150,
+ "rcld_promotion": 151,
+ "nix": "deferred_and_unclaimed",
+ "oci": "deferred_and_unclaimed"
+ },
+ "forbidden": [
+ "standalone_production_completion_commit",
+ "provider_execution_inside_transaction",
+ "relay_io_inside_transaction",
+ "response_reconstruction_on_retry",
+ "alternate_sqlite_authority",
+ "ambient_clock",
+ "ambient_entropy",
+ "task_spawn"
+ ]
+}
diff --git a/radroots.service.source-lock.v1.toml b/radroots.service.source-lock.v1.toml
@@ -14,7 +14,7 @@ host_feature_profile = "service-host"
[contract_versions]
config = 1
-state = 8
+state = 9
admin = 1
status = 1
provider = 1
diff --git a/src/lib.rs b/src/lib.rs
@@ -30,6 +30,7 @@ mod state_maintenance;
mod state_metadata;
mod state_repository;
mod state_request;
+mod state_response;
pub use cli_v1::{
MycBootstrapProfileV1, MycCliInvocationV1, MycCliV1Error, MycCliV1ErrorKind, MycCommandV1,
@@ -120,8 +121,9 @@ pub use state_catalog::{
MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT,
MYC_STATE_SCHEMA_VERSION_7_SHA256, MYC_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256,
MYC_STATE_SCHEMA_VERSION_8_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_8_SHA256,
- MycStateCatalogError, MycStateCatalogErrorKind, myc_migration_catalog, myc_schema_catalog,
- validate_myc_state_catalogs,
+ MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_9_OBJECT_COUNT,
+ MYC_STATE_SCHEMA_VERSION_9_SHA256, MycStateCatalogError, MycStateCatalogErrorKind,
+ myc_migration_catalog, myc_schema_catalog, validate_myc_state_catalogs,
};
pub use state_completion::{
MycNip46CommitAdmission, MycNip46CommitError, MycNip46CommitErrorKind, MycNip46CommitRecord,
@@ -185,3 +187,7 @@ pub use state_request::{
MycSignerRequestAdmission, MycSignerRequestDigest, MycSignerRequestError,
MycSignerRequestErrorKind, MycSignerRequestMethod, MycSignerRequestRecord,
};
+pub use state_response::{
+ MycNip46ResponseCommitAdmission, MycNip46ResponseCommitError, MycNip46ResponseCommitErrorKind,
+ MycNip46ResponseCommitRecord, MycNip46ResponseCommitRequest, MycNip46ResponseRecord,
+};
diff --git a/src/nip46_wave_080_b.rs b/src/nip46_wave_080_b.rs
@@ -2,7 +2,7 @@
use std::{error::Error as _, fs, os::unix::fs::PermissionsExt};
-use nostr::UnsignedEvent as NostrUnsignedEvent;
+use nostr::{JsonUtil as _, Kind, Tag, Timestamp, UnsignedEvent as NostrUnsignedEvent};
use radroots_nostr_connect::message::Request;
use sha2::{Digest, Sha256};
use sqlx::{ConnectOptions as _, Connection as _, Row as _, sqlite::SqliteConnectOptions};
@@ -10,8 +10,10 @@ use sqlx::{ConnectOptions as _, Connection as _, Row as _, sqlite::SqliteConnect
use crate::{
MycAuditCorrelationId, MycConnectionAdmissionPolicy, MycConnectionNonce,
MycConnectionOperatorDecision, MycConnectionPolicyGeneration, MycLocalSignerUntrustedResponse,
- MycNip46CommitAdmission, MycNip46CommitRequest, MycNip46SessionEffect,
- MycProviderDeadlineUnixMs, MycProviderResponseObservedAtUnixMs, MycProviderRole,
+ MycNip46CommitAdmission, MycNip46CommitRequest, MycNip46ResponseCommitAdmission,
+ MycNip46ResponseCommitRequest, MycNip46SessionEffect, MycProviderCorrelationId,
+ MycProviderDeadlineUnixMs, MycProviderOperation, MycProviderOperationId,
+ MycProviderOperationInput, MycProviderResponseObservedAtUnixMs, MycProviderRole,
MycRateRelayId, MycSignerRequestAdmission, MycSignerRequestMethod, MycStateRepositoryErrorKind,
initialize_myc_state, open_myc_state_read_write, prepare_myc_nip46_work,
provider_local_signer::{ProtectedWireHex, WireProviderResult},
@@ -90,8 +92,61 @@ async fn active_connection(
(work, active, decision)
}
+fn atomic_response_request(
+ config: &crate::MycConfigDocumentV1,
+ completion: &MycNip46CommitRequest,
+) -> (MycNip46ResponseCommitRequest, Vec<u8>) {
+ let unsigned = NostrUnsignedEvent::new(
+ keys(3).public_key(),
+ Timestamp::from_secs(OBSERVED_AT_SECONDS + 3),
+ Kind::Custom(24_133),
+ vec![Tag::public_key(keys(10).public_key())],
+ "encrypted-response",
+ );
+ let operation = MycProviderOperation::new(
+ config
+ .provider_contract()
+ .binding(MycProviderRole::User)
+ .expect("user binding"),
+ MycProviderOperationId::from_bytes([0x91; 32]),
+ MycProviderCorrelationId::from_bytes([0x92; 32]),
+ MycProviderDeadlineUnixMs::new(PROVIDER_DEADLINE_MS).expect("provider deadline"),
+ MycProviderOperationInput::sign_event(unsigned.as_json().as_bytes())
+ .expect("response signing input"),
+ )
+ .expect("response operation");
+ let signed = unsigned.sign_with_keys(&keys(3)).expect("signed response");
+ let bytes = serde_json::to_vec(&signed).expect("canonical response");
+ let response: MycLocalSignerUntrustedResponse = untrusted_response(
+ &operation,
+ hex::encode(operation.correlation_id().as_bytes()),
+ WireProviderResult::SignEvent {
+ payload_hex: ProtectedWireHex::from_bytes(&bytes),
+ },
+ );
+ let verified = response
+ .verify(
+ config
+ .provider_contract()
+ .binding(MycProviderRole::User)
+ .expect("user binding"),
+ &operation,
+ MycProviderResponseObservedAtUnixMs::new(RECEIVED_AT_MS + 3_001)
+ .expect("response time"),
+ )
+ .expect("verified response");
+ let request = MycNip46ResponseCommitRequest::new(
+ completion,
+ &operation,
+ &verified,
+ crate::MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 3_003).expect("commit time"),
+ )
+ .expect("atomic response request");
+ (request, bytes)
+}
+
#[tokio::test]
-async fn verified_signed_artifact_commits_exactly_once_without_outbox_state() {
+async fn verified_signed_response_commits_completion_and_outbox_exactly_once() {
let directory = tempfile::tempdir().expect("temporary root");
let runtime = runtime(directory.path());
fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
@@ -163,6 +218,15 @@ async fn verified_signed_artifact_commits_exactly_once_without_outbox_state() {
.session_effect(),
MycNip46SessionEffect::ConnectionAdmitted
);
+ let (partial_response, _) = atomic_response_request(&config, &connect_commit);
+ assert_eq!(
+ repository
+ .commit_nip46_response(&partial_response)
+ .await
+ .expect_err("completion-only legacy state is not repaired")
+ .kind(),
+ MycStateRepositoryErrorKind::Binding
+ );
let prepared = prepared_request(
&config,
@@ -223,28 +287,75 @@ async fn verified_signed_artifact_commits_exactly_once_without_outbox_state() {
assert!(!request_debug.contains(&hex::encode(operation.operation_id().as_bytes())));
assert!(!request_debug.contains(&hex::encode(operation.correlation_id().as_bytes())));
assert!(!request_debug.contains(String::from_utf8_lossy(&signed_bytes).as_ref()));
+ let (atomic_request, response_bytes) = atomic_response_request(&config, &request);
+ assert_eq!(
+ repository
+ .commit_nip46_response(&atomic_request.fail_after_completion_for_test())
+ .await
+ .expect_err("completion-edge failure rolls back")
+ .kind(),
+ MycStateRepositoryErrorKind::Transaction
+ );
+ assert_eq!(
+ repository
+ .commit_nip46_response(&atomic_request.fail_after_response_for_test())
+ .await
+ .expect_err("response-edge failure rolls back")
+ .kind(),
+ MycStateRepositoryErrorKind::Transaction
+ );
let committed = repository
- .commit_nip46_operation(&request)
+ .commit_nip46_response(&atomic_request)
.await
- .expect("completion commit");
- assert!(matches!(committed, MycNip46CommitAdmission::Committed(_)));
+ .expect("atomic response commit");
+ assert!(matches!(
+ committed,
+ MycNip46ResponseCommitAdmission::Committed(_)
+ ));
assert_eq!(
- committed.record().method(),
+ committed.record().completion().method(),
MycSignerRequestMethod::SignEvent
);
assert_eq!(
- committed.record().artifact_sha256(),
+ committed.record().completion().artifact_sha256(),
Some(&<[u8; 32]>::from(Sha256::digest(&signed_bytes)))
);
+ assert_eq!(
+ committed.record().response().signed_response_bytes(),
+ response_bytes
+ );
+ assert!(
+ committed
+ .record()
+ .response()
+ .delivery_job()
+ .targets()
+ .iter()
+ .all(|target| target.attempt_count() == 0
+ && target.status() == crate::MycDeliveryTargetStatus::Pending)
+ );
let record_debug = format!("{:?}", committed.record());
- assert!(!record_debug.contains(&hex::encode(committed.record().operation_id().as_bytes())));
- assert!(!record_debug.contains(&hex::encode(committed.record().correlation_id().as_bytes())));
+ assert!(!record_debug.contains(&hex::encode(
+ committed.record().completion().operation_id().as_bytes()
+ )));
+ assert!(!record_debug.contains(&hex::encode(
+ committed.record().completion().correlation_id().as_bytes()
+ )));
assert!(!record_debug.contains(&hex::encode(Sha256::digest(&signed_bytes))));
let replay = repository
- .commit_nip46_operation(&request)
+ .commit_nip46_response(&atomic_request)
+ .await
+ .expect("exact atomic replay");
+ assert!(matches!(
+ replay,
+ MycNip46ResponseCommitAdmission::ExactReplay(_)
+ ));
+ let retry = repository
+ .read_nip46_response(committed.record().response().delivery_job().id())
.await
- .expect("exact completion replay");
- assert!(matches!(replay, MycNip46CommitAdmission::ExactReplay(_)));
+ .expect("response retry read")
+ .expect("retained response");
+ assert_eq!(retry.signed_response_bytes(), response_bytes);
host.close().await.expect("host close");
let options = SqliteConnectOptions::new()
@@ -271,7 +382,14 @@ async fn verified_signed_artifact_commits_exactly_once_without_outbox_state() {
.fetch_one(&mut connection)
.await
.expect("delivery job count"),
- 0
+ 1
+ );
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_signed_responses")
+ .fetch_one(&mut connection)
+ .await
+ .expect("response count"),
+ 1
);
connection.close().await.expect("inspection close");
}
diff --git a/src/state_catalog.rs b/src/state_catalog.rs
@@ -12,7 +12,7 @@ use radroots_service_sqlite::{
pub const MYC_STATE_BASE_SCHEMA_VERSION: u32 = 1;
/// The newest governed Myc state schema understood by this binary.
-pub const MYC_STATE_SCHEMA_VERSION: u32 = 8;
+pub const MYC_STATE_SCHEMA_VERSION: u32 = 9;
/// The shared metadata and migration-ledger objects present at schema v1.
pub const MYC_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6;
@@ -38,6 +38,9 @@ pub const MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT: u32 = 53;
/// The shared objects plus immutable NIP-46 operation completion evidence.
pub const MYC_STATE_SCHEMA_VERSION_8_OBJECT_COUNT: u32 = 56;
+/// The shared objects plus immutable exact NIP-46 response authority.
+pub const MYC_STATE_SCHEMA_VERSION_9_OBJECT_COUNT: u32 = 59;
+
/// SHA-256 identity of the exact schema-v1 object snapshot.
pub const MYC_STATE_SCHEMA_VERSION_1_SHA256: [u8; 32] = [
0x94, 0xdc, 0x66, 0xfb, 0xca, 0x60, 0x16, 0x79, 0x61, 0x5c, 0x05, 0x52, 0x29, 0xdc, 0x0d, 0xb6,
@@ -52,14 +55,14 @@ pub const MYC_STATE_SCHEMA_VERSION_2_SHA256: [u8; 32] = [
/// SHA-256 identity of the ordered Myc migration catalog.
pub const MYC_MIGRATION_CATALOG_SHA256: [u8; 32] = [
- 0xa2, 0xdd, 0xd3, 0x20, 0xe2, 0xf9, 0x6d, 0x08, 0xc8, 0x17, 0x7d, 0xa1, 0x17, 0x5d, 0x0e, 0xd3,
- 0x29, 0x83, 0x16, 0x8d, 0xc1, 0x6b, 0x3d, 0x32, 0x6f, 0xd4, 0x3d, 0x9d, 0x65, 0x5d, 0x8e, 0x82,
+ 0x86, 0x04, 0xc5, 0x1e, 0xa6, 0xf1, 0x6f, 0x0a, 0xa3, 0x7a, 0x63, 0xee, 0xe0, 0x9c, 0x60, 0x0a,
+ 0x0b, 0x70, 0x31, 0xef, 0xbc, 0x8e, 0x5e, 0x0e, 0x41, 0xb5, 0xf0, 0xf5, 0x3c, 0x48, 0xe7, 0x0b,
];
/// SHA-256 identity of the schema catalog bound to the migration catalog.
pub const MYC_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [
- 0x49, 0x62, 0x28, 0xfc, 0xd4, 0xc2, 0xa5, 0x83, 0xd7, 0xff, 0xe2, 0x11, 0x3e, 0x0a, 0x7a, 0xdd,
- 0x68, 0xbc, 0xf7, 0xb7, 0x02, 0x18, 0xce, 0x5b, 0x94, 0x6e, 0x2f, 0x92, 0xf3, 0x0e, 0xda, 0xc7,
+ 0x63, 0x2c, 0x8d, 0x52, 0x16, 0xd2, 0xa5, 0x41, 0xfd, 0x28, 0xae, 0x3a, 0x2c, 0xae, 0x80, 0x69,
+ 0xc5, 0x8b, 0xef, 0x11, 0xfe, 0x81, 0xd2, 0xf0, 0x43, 0x99, 0xa7, 0x7c, 0xf8, 0x35, 0x9e, 0xa7,
];
/// SHA-256 identity of the schema-v2 migration content.
@@ -140,6 +143,18 @@ pub const MYC_STATE_SCHEMA_VERSION_8_SHA256: [u8; 32] = [
0x80, 0x96, 0x65, 0xc8, 0x9f, 0xde, 0x92, 0x5d, 0x87, 0x23, 0x65, 0x90, 0x9d, 0x28, 0xf0, 0x6e,
];
+/// SHA-256 identity of the schema-v9 atomic response migration.
+pub const MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256: [u8; 32] = [
+ 0xfd, 0xb0, 0x39, 0xd4, 0x72, 0xcd, 0x62, 0xe7, 0xda, 0x46, 0xc4, 0x0a, 0x97, 0x86, 0xc3, 0xa9,
+ 0xa7, 0xec, 0x55, 0xf2, 0xae, 0x7d, 0xb6, 0x89, 0xd0, 0xf8, 0x55, 0xfb, 0xbf, 0x8a, 0x95, 0x66,
+];
+
+/// SHA-256 identity of the schema-v9 object snapshot.
+pub const MYC_STATE_SCHEMA_VERSION_9_SHA256: [u8; 32] = [
+ 0xee, 0xe7, 0x6f, 0x4f, 0xf0, 0xbd, 0x2d, 0xc2, 0xc0, 0x61, 0xae, 0x38, 0x4e, 0x00, 0xde, 0x16,
+ 0xc5, 0xf5, 0xef, 0xc4, 0xb6, 0xed, 0xd0, 0xac, 0x7f, 0x73, 0xca, 0xd8, 0x3a, 0x99, 0x1f, 0xf7,
+];
+
/// SHA-256 identity of the Myc metadata table definition.
const MYC_STATE_METADATA_TABLE_SHA256: [u8; 32] = [
0x16, 0x17, 0x46, 0xa2, 0x26, 0x42, 0x46, 0x2f, 0x2b, 0xdb, 0x08, 0x5b, 0xae, 0xde, 0xb2, 0x3b,
@@ -1645,6 +1660,62 @@ const CREATE_NIP46_OPERATION_COMPLETION_MIGRATION_SQL: &str = concat!(
nip46_operation_commits_no_delete_sql!(),
);
+macro_rules! nip46_signed_responses_table_sql {
+ () => {
+ r#"CREATE TABLE nip46_signed_responses (
+ operation_id BLOB NOT NULL PRIMARY KEY CHECK (length(operation_id) = 32)
+ REFERENCES nip46_operation_commits(operation_id),
+ response_provider_operation_id BLOB NOT NULL UNIQUE
+ CHECK (length(response_provider_operation_id) = 32),
+ response_event_id BLOB NOT NULL UNIQUE CHECK (length(response_event_id) = 32),
+ response_sha256 BLOB NOT NULL CHECK (length(response_sha256) = 32),
+ response_bytes BLOB NOT NULL
+ CHECK (length(response_bytes) BETWEEN 1 AND 1048576),
+ authored_at_unix_s INTEGER NOT NULL
+ CHECK (authored_at_unix_s BETWEEN 1 AND 9223372036854775807),
+ committed_at_unix_ms INTEGER NOT NULL
+ CHECK (committed_at_unix_ms BETWEEN 1 AND 9223372036854775807)
+) STRICT"#
+ };
+}
+
+macro_rules! nip46_signed_responses_no_update_sql {
+ () => {
+ r#"CREATE TRIGGER nip46_signed_responses_no_update
+BEFORE UPDATE ON nip46_signed_responses
+BEGIN
+ SELECT RAISE(ABORT, 'NIP-46 signed response is immutable');
+END"#
+ };
+}
+
+macro_rules! nip46_signed_responses_no_delete_sql {
+ () => {
+ r#"CREATE TRIGGER nip46_signed_responses_no_delete
+BEFORE DELETE ON nip46_signed_responses
+BEGIN
+ SELECT RAISE(ABORT, 'NIP-46 signed response is retained');
+END"#
+ };
+}
+
+const CREATE_NIP46_SIGNED_RESPONSES_TABLE_SQL: &str = nip46_signed_responses_table_sql!();
+const CREATE_NIP46_SIGNED_RESPONSES_NO_UPDATE_SQL: &str = nip46_signed_responses_no_update_sql!();
+const CREATE_NIP46_SIGNED_RESPONSES_NO_DELETE_SQL: &str = nip46_signed_responses_no_delete_sql!();
+
+const CREATE_NIP46_ATOMIC_RESPONSE_MIGRATION_SQL: &str = concat!(
+ "DROP TRIGGER myc_state_metadata_no_update;\n",
+ "UPDATE myc_state_metadata SET state_contract_version = CASE ",
+ "WHEN state_contract_version = 8 THEN 9 ELSE 0 END WHERE singleton = 1;\n",
+ myc_state_metadata_no_update_sql!(),
+ ";\n",
+ nip46_signed_responses_table_sql!(),
+ ";\n",
+ nip46_signed_responses_no_update_sql!(),
+ ";\n",
+ nip46_signed_responses_no_delete_sql!(),
+);
+
const CONNECTIONS_TABLE_SHA256: [u8; 32] = [
0x72, 0xd5, 0xd8, 0xba, 0x24, 0x68, 0x9c, 0x93, 0x34, 0xb3, 0x8f, 0xbf, 0x64, 0x21, 0xe1, 0x65,
0xfd, 0xc3, 0x80, 0x46, 0xf1, 0x3f, 0x56, 0x49, 0x3a, 0xef, 0xd7, 0x42, 0xc4, 0xe6, 0x49, 0x85,
@@ -1853,6 +1924,18 @@ const NIP46_OPERATION_COMMITS_NO_DELETE_SHA256: [u8; 32] = [
0x99, 0x67, 0xd8, 0x61, 0xd2, 0xa7, 0x1e, 0x1a, 0x74, 0x3f, 0xb5, 0xe9, 0x7b, 0x96, 0xe0, 0x00,
0xa9, 0xaa, 0x38, 0xe4, 0xb6, 0x76, 0x47, 0xa0, 0x87, 0x0b, 0x48, 0x9b, 0x55, 0x20, 0xa9, 0xa3,
];
+const NIP46_SIGNED_RESPONSES_TABLE_SHA256: [u8; 32] = [
+ 0x12, 0xf2, 0x9a, 0x4b, 0x24, 0x62, 0x5a, 0xf4, 0xcb, 0x0e, 0xa1, 0x16, 0x9d, 0x64, 0x9c, 0x1f,
+ 0xe6, 0x4a, 0xe7, 0x10, 0x14, 0x52, 0xbc, 0x43, 0x9c, 0xdd, 0x86, 0x23, 0x47, 0x60, 0x96, 0x19,
+];
+const NIP46_SIGNED_RESPONSES_NO_UPDATE_SHA256: [u8; 32] = [
+ 0xee, 0xd8, 0xb0, 0x82, 0xc7, 0x46, 0x37, 0x53, 0x5a, 0x77, 0x77, 0x95, 0xb5, 0x29, 0x58, 0x41,
+ 0x5a, 0x3f, 0xab, 0xc3, 0xa6, 0x7a, 0xcf, 0x44, 0xef, 0xe0, 0xc6, 0x82, 0xeb, 0xcb, 0x2d, 0x05,
+];
+const NIP46_SIGNED_RESPONSES_NO_DELETE_SHA256: [u8; 32] = [
+ 0x8e, 0x6a, 0x36, 0xe5, 0x89, 0x50, 0x81, 0x90, 0x5f, 0x1c, 0x32, 0x19, 0x25, 0xec, 0x37, 0x07,
+ 0x1f, 0x42, 0xd9, 0x1a, 0x73, 0xd3, 0x89, 0x56, 0x3d, 0x2b, 0x4b, 0x3e, 0x9d, 0x3a, 0xf3, 0xa4,
+];
/// Stable classes for invalid embedded Myc catalog definitions.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
@@ -1974,6 +2057,13 @@ pub fn myc_migration_catalog() -> Result<MigrationCatalog, MycStateCatalogError>
MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256),
)
.map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?;
+ let response = MigrationDescriptor::sql(
+ 9,
+ "create_nip46_atomic_response",
+ CREATE_NIP46_ATOMIC_RESPONSE_MIGRATION_SQL,
+ MigrationChecksum::from_bytes(MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256),
+ )
+ .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?;
let catalog = MigrationCatalog::new([
metadata,
requests,
@@ -1982,10 +2072,11 @@ pub fn myc_migration_catalog() -> Result<MigrationCatalog, MycStateCatalogError>
delivery,
discovery,
completion,
+ response,
])
.map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::MigrationCatalog))?;
if catalog.current_version() != MYC_STATE_SCHEMA_VERSION
- || catalog.descriptors().len() != 7
+ || catalog.descriptors().len() != 8
|| catalog.digest().as_bytes() != &MYC_MIGRATION_CATALOG_SHA256
{
return Err(MycStateCatalogError::new(
@@ -2046,6 +2137,12 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> {
SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_8_SHA256),
)
.map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?;
+ let version_nine = SchemaVersionCatalog::new(
+ 9,
+ myc_state_response_objects()?,
+ SchemaDigest::from_bytes(MYC_STATE_SCHEMA_VERSION_9_SHA256),
+ )
+ .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?;
let catalog = SchemaCatalog::new(
&migrations,
[
@@ -2057,6 +2154,7 @@ pub fn myc_schema_catalog() -> Result<SchemaCatalog, MycStateCatalogError> {
version_six,
version_seven,
version_eight,
+ version_nine,
],
)
.map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))?;
@@ -2580,6 +2678,38 @@ fn myc_state_completion_objects() -> Result<Vec<SchemaObject>, MycStateCatalogEr
Ok(objects)
}
+fn myc_state_response_objects() -> Result<Vec<SchemaObject>, MycStateCatalogError> {
+ let mut objects = myc_state_completion_objects()?;
+ let object = |kind, name, table, sql, digest| {
+ SchemaObject::new(kind, name, table, sql, SchemaDigest::from_bytes(digest))
+ .map_err(|_| MycStateCatalogError::new(MycStateCatalogErrorKind::SchemaCatalog))
+ };
+ objects.extend([
+ object(
+ SchemaObjectKind::Table,
+ "nip46_signed_responses",
+ "nip46_signed_responses",
+ CREATE_NIP46_SIGNED_RESPONSES_TABLE_SQL,
+ NIP46_SIGNED_RESPONSES_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "nip46_signed_responses_no_update",
+ "nip46_signed_responses",
+ CREATE_NIP46_SIGNED_RESPONSES_NO_UPDATE_SQL,
+ NIP46_SIGNED_RESPONSES_NO_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "nip46_signed_responses_no_delete",
+ "nip46_signed_responses",
+ CREATE_NIP46_SIGNED_RESPONSES_NO_DELETE_SQL,
+ NIP46_SIGNED_RESPONSES_NO_DELETE_SHA256,
+ )?,
+ ]);
+ Ok(objects)
+}
+
/// Independently validates exact catalog versions, counts, and digests.
pub fn validate_myc_state_catalogs(
migrations: &MigrationCatalog,
@@ -2588,7 +2718,7 @@ pub fn validate_myc_state_catalogs(
let versions = schema.versions();
let descriptors = migrations.descriptors();
let valid = migrations.current_version() == MYC_STATE_SCHEMA_VERSION
- && descriptors.len() == 7
+ && descriptors.len() == 8
&& descriptors[0].target_version() == 2
&& descriptors[0].name().as_str() == "create_myc_state_metadata"
&& descriptors[0].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256
@@ -2610,9 +2740,12 @@ pub fn validate_myc_state_catalogs(
&& descriptors[6].target_version() == 8
&& descriptors[6].name().as_str() == "create_nip46_operation_completion"
&& descriptors[6].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256
+ && descriptors[7].target_version() == 9
+ && descriptors[7].name().as_str() == "create_nip46_atomic_response"
+ && descriptors[7].checksum().as_bytes() == &MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256
&& migrations.digest().as_bytes() == &MYC_MIGRATION_CATALOG_SHA256
&& schema.migration_catalog_digest() == migrations.digest()
- && versions.len() == 8
+ && versions.len() == 9
&& versions[0].version() == MYC_STATE_BASE_SCHEMA_VERSION
&& versions[0].object_count() == MYC_STATE_SCHEMA_VERSION_1_OBJECT_COUNT
&& versions[0].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_1_SHA256
@@ -2637,6 +2770,9 @@ pub fn validate_myc_state_catalogs(
&& versions[7].version() == 8
&& versions[7].object_count() == MYC_STATE_SCHEMA_VERSION_8_OBJECT_COUNT
&& versions[7].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_8_SHA256
+ && versions[8].version() == 9
+ && versions[8].object_count() == MYC_STATE_SCHEMA_VERSION_9_OBJECT_COUNT
+ && versions[8].digest().as_bytes() == &MYC_STATE_SCHEMA_VERSION_9_SHA256
&& schema.digest().as_bytes() == &MYC_STATE_SCHEMA_CATALOG_SHA256;
if valid {
Ok(())
diff --git a/src/state_completion.rs b/src/state_completion.rs
@@ -3,12 +3,13 @@
use core::fmt;
use std::error::Error;
-use radroots_service_sqlite::{
- ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
-};
+use radroots_service_sqlite::ServiceSqliteTransaction;
+#[cfg(test)]
+use radroots_service_sqlite::{ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind};
use sha2::{Digest, Sha256};
use sqlx::Row;
+#[cfg(test)]
use crate::state_repository::{
MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata,
RepositoryOperationError, require_expected_metadata,
@@ -409,6 +410,7 @@ impl fmt::Debug for MycNip46CommitAdmission {
}
}
+#[cfg(test)]
impl MycStateRepository<'_> {
/// Atomically records the Step 147 completion component.
///
@@ -416,7 +418,7 @@ impl MycStateRepository<'_> {
/// response commit. Step 148 must compose this component with the outer
/// signed response, immutable relay targets, and initial outbox state in
/// the same transaction before RCLD-RSHR-080 may be promoted to `master`.
- pub async fn commit_nip46_operation(
+ pub(crate) async fn commit_nip46_operation(
&self,
request: &MycNip46CommitRequest,
) -> Result<MycNip46CommitAdmission, MycStateRepositoryError> {
@@ -437,7 +439,15 @@ impl MycStateRepository<'_> {
}
impl MycNip46CommitRequest {
- fn owned(&self) -> Self {
+ pub(crate) const fn signer_request(&self) -> &MycSignerRequestRecord {
+ &self.request
+ }
+
+ pub(crate) const fn completed_at(&self) -> MycConnectionTimeUnixMs {
+ self.completed_at
+ }
+
+ pub(crate) fn owned(&self) -> Self {
Self {
request: self.request.clone(),
connection_id: self.connection_id,
@@ -456,12 +466,12 @@ impl MycNip46CommitRequest {
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
-enum CommitOperationError {
+pub(crate) enum CommitOperationError {
Binding,
Storage,
}
-async fn commit_operation(
+pub(crate) async fn commit_operation(
transaction: &mut ServiceSqliteTransaction<'_>,
request: &MycNip46CommitRequest,
) -> Result<MycNip46CommitAdmission, CommitOperationError> {
@@ -804,6 +814,7 @@ fn require_one(rows: u64) -> Result<(), CommitOperationError> {
.ok_or(CommitOperationError::Storage)
}
+#[cfg(test)]
const fn map_repository_error(error: RepositoryOperationError) -> CommitOperationError {
match error {
RepositoryOperationError::Binding => CommitOperationError::Binding,
@@ -811,6 +822,7 @@ const fn map_repository_error(error: RepositoryOperationError) -> CommitOperatio
}
}
+#[cfg(test)]
fn map_transaction_error(
error: ServiceSqliteTransactionError<CommitOperationError>,
) -> MycStateRepositoryError {
diff --git a/src/state_delivery.rs b/src/state_delivery.rs
@@ -1655,7 +1655,7 @@ async fn read_job_by_source(
}
}
-async fn read_job(
+pub(crate) async fn read_job(
transaction: &mut ServiceSqliteTransaction<'_>,
job_id: MycDeliveryJobId,
) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> {
diff --git a/src/state_host.rs b/src/state_host.rs
@@ -377,20 +377,21 @@ fn require_migration_build(
fn exact_initialization_outcome(outcome: MigrationApplicationOutcome) -> bool {
outcome.initial_version() == MYC_STATE_BASE_SCHEMA_VERSION
&& outcome.final_version() == MYC_STATE_SCHEMA_VERSION
- && outcome.applied_count() == 7
+ && outcome.applied_count() == 8
}
fn exact_existing_outcome(outcome: MigrationApplicationOutcome) -> bool {
outcome.final_version() == MYC_STATE_SCHEMA_VERSION
&& matches!(
(outcome.initial_version(), outcome.applied_count()),
- (MYC_STATE_BASE_SCHEMA_VERSION, 7)
- | (2, 6)
- | (3, 5)
- | (4, 4)
- | (5, 3)
- | (6, 2)
- | (7, 1)
+ (MYC_STATE_BASE_SCHEMA_VERSION, 8)
+ | (2, 7)
+ | (3, 6)
+ | (4, 5)
+ | (5, 4)
+ | (6, 3)
+ | (7, 2)
+ | (8, 1)
| (MYC_STATE_SCHEMA_VERSION, 0)
)
}
diff --git a/src/state_request.rs b/src/state_request.rs
@@ -512,6 +512,10 @@ pub struct MycSignerRequestRecord {
}
impl MycSignerRequestRecord {
+ pub(crate) const fn client_public_key(&self) -> &MycNip46ClientPublicKey {
+ &self.client_public_key
+ }
+
#[must_use]
/// Returns the stable logical operation identity.
pub const fn operation_id(&self) -> MycSignerOperationId {
diff --git a/src/state_response.rs b/src/state_response.rs
@@ -0,0 +1,698 @@
+//! Atomic completion, exact signed-response, and initial delivery-state commit.
+
+use core::fmt;
+use std::error::Error;
+
+use radroots_nostr::event::{Event as RadrootsNostrEvent, Kind as RadrootsNostrKind};
+use radroots_service_sqlite::{
+ ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
+};
+use sha2::{Digest, Sha256};
+use sqlx::Row;
+
+use crate::state_completion::{
+ CommitOperationError, MycNip46CommitAdmission, MycNip46CommitRecord, MycNip46CommitRequest,
+ commit_operation,
+};
+use crate::state_delivery::{
+ DeliveryOperationError, MycDeliveryArtifactDigest, MycDeliveryJobAdmission, MycDeliveryJobId,
+ MycDeliveryJobRecord, MycDeliverySource, MycDeliveryTimeUnixMs, create_job, read_job,
+};
+use crate::state_repository::{
+ MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata,
+ RepositoryOperationError, require_expected_metadata,
+};
+use crate::{
+ MYC_PROVIDER_OUTPUT_MAX_BYTES, MycProviderCapability, MycProviderOperation, MycProviderRole,
+ MycSignerOperationId, MycVerifiedProviderResponse,
+};
+
+const NIP46_RPC_KIND: u16 = 24_133;
+
+const INSERT_RESPONSE_SQL: &str = r#"INSERT INTO nip46_signed_responses (
+ operation_id, response_provider_operation_id, response_event_id,
+ response_sha256, response_bytes, authored_at_unix_s, committed_at_unix_ms
+) VALUES (?, ?, ?, ?, ?, ?, ?)"#;
+
+const READ_RESPONSE_BY_OPERATION_SQL: &str = r#"SELECT
+ CASE WHEN typeof(r.operation_id) = 'blob' AND length(r.operation_id) = 32
+ THEN r.operation_id ELSE NULL END AS operation_id,
+ CASE WHEN typeof(r.response_provider_operation_id) = 'blob'
+ AND length(r.response_provider_operation_id) = 32
+ THEN r.response_provider_operation_id ELSE NULL END AS response_provider_operation_id,
+ CASE WHEN typeof(r.response_event_id) = 'blob' AND length(r.response_event_id) = 32
+ THEN r.response_event_id ELSE NULL END AS response_event_id,
+ CASE WHEN typeof(r.response_sha256) = 'blob' AND length(r.response_sha256) = 32
+ THEN r.response_sha256 ELSE NULL END AS response_sha256,
+ CASE WHEN typeof(r.response_bytes) = 'blob'
+ AND length(r.response_bytes) BETWEEN 1 AND 1048576
+ THEN r.response_bytes ELSE NULL END AS response_bytes,
+ r.authored_at_unix_s, r.committed_at_unix_ms,
+ CASE WHEN typeof(q.client_public_key) = 'text'
+ AND length(CAST(q.client_public_key AS BLOB)) = 64
+ THEN q.client_public_key ELSE NULL END AS client_public_key,
+ CASE WHEN typeof(j.job_id) = 'blob' AND length(j.job_id) = 32
+ THEN j.job_id ELSE NULL END AS job_id
+FROM nip46_signed_responses r
+JOIN nip46_requests q ON q.operation_id = r.operation_id
+JOIN delivery_jobs j ON j.source_kind = 'signer_response'
+ AND j.source_id = r.operation_id
+WHERE r.operation_id = ?
+LIMIT 2"#;
+
+const READ_RESPONSE_BY_JOB_SQL: &str = r#"SELECT
+ CASE WHEN typeof(r.operation_id) = 'blob' AND length(r.operation_id) = 32
+ THEN r.operation_id ELSE NULL END AS operation_id,
+ CASE WHEN typeof(r.response_provider_operation_id) = 'blob'
+ AND length(r.response_provider_operation_id) = 32
+ THEN r.response_provider_operation_id ELSE NULL END AS response_provider_operation_id,
+ CASE WHEN typeof(r.response_event_id) = 'blob' AND length(r.response_event_id) = 32
+ THEN r.response_event_id ELSE NULL END AS response_event_id,
+ CASE WHEN typeof(r.response_sha256) = 'blob' AND length(r.response_sha256) = 32
+ THEN r.response_sha256 ELSE NULL END AS response_sha256,
+ CASE WHEN typeof(r.response_bytes) = 'blob'
+ AND length(r.response_bytes) BETWEEN 1 AND 1048576
+ THEN r.response_bytes ELSE NULL END AS response_bytes,
+ r.authored_at_unix_s, r.committed_at_unix_ms,
+ CASE WHEN typeof(q.client_public_key) = 'text'
+ AND length(CAST(q.client_public_key AS BLOB)) = 64
+ THEN q.client_public_key ELSE NULL END AS client_public_key,
+ CASE WHEN typeof(j.job_id) = 'blob' AND length(j.job_id) = 32
+ THEN j.job_id ELSE NULL END AS job_id
+FROM delivery_jobs j
+JOIN nip46_signed_responses r ON r.operation_id = j.source_id
+JOIN nip46_requests q ON q.operation_id = r.operation_id
+WHERE j.job_id = ? AND j.source_kind = 'signer_response'
+LIMIT 2"#;
+
+/// Stable construction failure classes for an atomic NIP-46 response commit.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum MycNip46ResponseCommitErrorKind {
+ InvalidBinding,
+ InvalidResponse,
+ InvalidTime,
+}
+
+impl MycNip46ResponseCommitErrorKind {
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::InvalidBinding => "nip46_response_binding_invalid",
+ Self::InvalidResponse => "nip46_response_invalid",
+ Self::InvalidTime => "nip46_response_time_invalid",
+ }
+ }
+}
+
+/// Source-free, path-free construction failure.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct MycNip46ResponseCommitError {
+ kind: MycNip46ResponseCommitErrorKind,
+}
+
+impl MycNip46ResponseCommitError {
+ const fn new(kind: MycNip46ResponseCommitErrorKind) -> Self {
+ Self { kind }
+ }
+
+ #[must_use]
+ pub const fn kind(self) -> MycNip46ResponseCommitErrorKind {
+ self.kind
+ }
+
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ self.kind.code()
+ }
+}
+
+impl fmt::Display for MycNip46ResponseCommitError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ MycNip46ResponseCommitErrorKind::InvalidBinding => "NIP-46 response binding is invalid",
+ MycNip46ResponseCommitErrorKind::InvalidResponse => "NIP-46 signed response is invalid",
+ MycNip46ResponseCommitErrorKind::InvalidTime => "NIP-46 response time is invalid",
+ })
+ }
+}
+
+impl fmt::Debug for MycNip46ResponseCommitError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("MycNip46ResponseCommitError")
+ .field("kind", &self.kind)
+ .finish()
+ }
+}
+
+impl Error for MycNip46ResponseCommitError {}
+
+/// Exact independently verified response and completion prepared before SQLite.
+pub struct MycNip46ResponseCommitRequest {
+ completion: MycNip46CommitRequest,
+ response_provider_operation_id: [u8; 32],
+ response_event_id: [u8; 32],
+ response_digest: MycDeliveryArtifactDigest,
+ response_bytes: Box<[u8]>,
+ authored_at_unix_s: u64,
+ committed_at: MycDeliveryTimeUnixMs,
+ #[cfg(test)]
+ fail_after_completion: bool,
+ #[cfg(test)]
+ fail_after_response: bool,
+}
+
+impl MycNip46ResponseCommitRequest {
+ /// Binds one completion to a signature-verified canonical NIP-46 response event.
+ pub fn new(
+ completion: &MycNip46CommitRequest,
+ response_operation: &MycProviderOperation,
+ response: &MycVerifiedProviderResponse,
+ committed_at: MycDeliveryTimeUnixMs,
+ ) -> Result<Self, MycNip46ResponseCommitError> {
+ if response_operation.role() != MycProviderRole::User
+ || response_operation.input().capability() != MycProviderCapability::SignEvent
+ || response.operation_id() != response_operation.operation_id()
+ || response.correlation_id() != response_operation.correlation_id()
+ || response.instance() != response_operation.instance()
+ || response.role() != response_operation.role()
+ || response.capability() != response_operation.input().capability()
+ || !response.matches_operation(response_operation)
+ || response_operation.operation_id().as_bytes()
+ == completion.signer_request().operation_id().as_bytes()
+ {
+ return Err(Self::error(MycNip46ResponseCommitErrorKind::InvalidBinding));
+ }
+ let bytes = response
+ .signed_event_bytes()
+ .ok_or_else(|| Self::error(MycNip46ResponseCommitErrorKind::InvalidResponse))?;
+ let event = validate_response_event(
+ bytes,
+ completion.signer_request().client_public_key().as_hex(),
+ Some(response_operation.expected_identity().as_hex()),
+ )?;
+ let authored_at_unix_s = event.created_at.as_secs();
+ if committed_at.get() < completion.completed_at().get()
+ || authored_at_unix_s
+ .checked_mul(1_000)
+ .is_none_or(|authored_ms| authored_ms > committed_at.get())
+ {
+ return Err(Self::error(MycNip46ResponseCommitErrorKind::InvalidTime));
+ }
+ let response_digest = MycDeliveryArtifactDigest::from_bytes(Sha256::digest(bytes).into());
+ Ok(Self {
+ completion: completion.owned(),
+ response_provider_operation_id: *response_operation.operation_id().as_bytes(),
+ response_event_id: *event.id.as_bytes(),
+ response_digest,
+ response_bytes: Box::from(bytes),
+ authored_at_unix_s,
+ committed_at,
+ #[cfg(test)]
+ fail_after_completion: false,
+ #[cfg(test)]
+ fail_after_response: false,
+ })
+ }
+
+ const fn error(kind: MycNip46ResponseCommitErrorKind) -> MycNip46ResponseCommitError {
+ MycNip46ResponseCommitError::new(kind)
+ }
+
+ fn owned(&self) -> Self {
+ Self {
+ completion: self.completion.owned(),
+ response_provider_operation_id: self.response_provider_operation_id,
+ response_event_id: self.response_event_id,
+ response_digest: self.response_digest,
+ response_bytes: self.response_bytes.clone(),
+ authored_at_unix_s: self.authored_at_unix_s,
+ committed_at: self.committed_at,
+ #[cfg(test)]
+ fail_after_completion: self.fail_after_completion,
+ #[cfg(test)]
+ fail_after_response: self.fail_after_response,
+ }
+ }
+
+ #[cfg(test)]
+ pub(crate) fn fail_after_completion_for_test(&self) -> Self {
+ let mut request = self.owned();
+ request.fail_after_completion = true;
+ request
+ }
+
+ #[cfg(test)]
+ pub(crate) fn fail_after_response_for_test(&self) -> Self {
+ let mut request = self.owned();
+ request.fail_after_response = true;
+ request
+ }
+}
+
+impl fmt::Debug for MycNip46ResponseCommitRequest {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("MycNip46ResponseCommitRequest([redacted])")
+ }
+}
+
+/// Immutable exact signed response and its config-bound initial delivery state.
+#[derive(Clone, PartialEq, Eq)]
+pub struct MycNip46ResponseRecord {
+ operation_id: MycSignerOperationId,
+ response_provider_operation_id: [u8; 32],
+ response_event_id: [u8; 32],
+ response_digest: MycDeliveryArtifactDigest,
+ response_bytes: Box<[u8]>,
+ authored_at_unix_s: u64,
+ committed_at: MycDeliveryTimeUnixMs,
+ delivery_job: MycDeliveryJobRecord,
+}
+
+impl MycNip46ResponseRecord {
+ #[must_use]
+ pub const fn operation_id(&self) -> MycSignerOperationId {
+ self.operation_id
+ }
+ #[must_use]
+ pub const fn response_provider_operation_id(&self) -> &[u8; 32] {
+ &self.response_provider_operation_id
+ }
+ #[must_use]
+ pub const fn response_event_id(&self) -> &[u8; 32] {
+ &self.response_event_id
+ }
+ #[must_use]
+ pub const fn response_digest(&self) -> MycDeliveryArtifactDigest {
+ self.response_digest
+ }
+ /// Returns the exact committed bytes that every retry must submit unchanged.
+ #[must_use]
+ pub fn signed_response_bytes(&self) -> &[u8] {
+ &self.response_bytes
+ }
+ #[must_use]
+ pub const fn authored_at_unix_s(&self) -> u64 {
+ self.authored_at_unix_s
+ }
+ #[must_use]
+ pub const fn committed_at(&self) -> MycDeliveryTimeUnixMs {
+ self.committed_at
+ }
+ #[must_use]
+ pub const fn delivery_job(&self) -> &MycDeliveryJobRecord {
+ &self.delivery_job
+ }
+}
+
+impl fmt::Debug for MycNip46ResponseRecord {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("MycNip46ResponseRecord")
+ .field("delivery_job", &self.delivery_job)
+ .field("response", &"[redacted]")
+ .field("identity", &"[redacted]")
+ .finish()
+ }
+}
+
+/// Immutable result of the one response-authority transaction.
+#[derive(Clone, PartialEq, Eq)]
+pub struct MycNip46ResponseCommitRecord {
+ completion: MycNip46CommitRecord,
+ response: MycNip46ResponseRecord,
+}
+
+impl MycNip46ResponseCommitRecord {
+ #[must_use]
+ pub const fn completion(&self) -> &MycNip46CommitRecord {
+ &self.completion
+ }
+ #[must_use]
+ pub const fn response(&self) -> &MycNip46ResponseRecord {
+ &self.response
+ }
+}
+
+impl fmt::Debug for MycNip46ResponseCommitRecord {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("MycNip46ResponseCommitRecord([redacted])")
+ }
+}
+
+/// New or exact-replayed atomic response authority.
+#[derive(Clone, PartialEq, Eq)]
+pub enum MycNip46ResponseCommitAdmission {
+ Committed(MycNip46ResponseCommitRecord),
+ ExactReplay(MycNip46ResponseCommitRecord),
+}
+
+impl MycNip46ResponseCommitAdmission {
+ #[must_use]
+ pub const fn record(&self) -> &MycNip46ResponseCommitRecord {
+ match self {
+ Self::Committed(record) | Self::ExactReplay(record) => record,
+ }
+ }
+}
+
+impl fmt::Debug for MycNip46ResponseCommitAdmission {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self {
+ Self::Committed(_) => "MycNip46ResponseCommitAdmission::Committed([redacted])",
+ Self::ExactReplay(_) => "MycNip46ResponseCommitAdmission::ExactReplay([redacted])",
+ })
+ }
+}
+
+impl MycStateRepository<'_> {
+ /// Atomically commits completion, exact response bytes, targets, and pending attempt state.
+ pub async fn commit_nip46_response(
+ &self,
+ request: &MycNip46ResponseCommitRequest,
+ ) -> Result<MycNip46ResponseCommitAdmission, MycStateRepositoryError> {
+ let expected = PersistedMetadata::from(self.expected());
+ let policy = self.expected().delivery_policies().clone();
+ let request = request.owned();
+ self.host()
+ .transaction(move |transaction| {
+ Box::pin(async move {
+ require_expected_metadata(transaction, &expected)
+ .await
+ .map_err(AtomicOperationError::from)?;
+ let completion = commit_operation(transaction, &request.completion)
+ .await
+ .map_err(AtomicOperationError::from)?;
+ match completion {
+ MycNip46CommitAdmission::ExactReplay(completion) => {
+ let response = read_response_by_operation(
+ transaction,
+ request.completion.signer_request().operation_id(),
+ )
+ .await?
+ .ok_or(AtomicOperationError::Binding)?;
+ exact_response(&response, &request)?;
+ Ok(MycNip46ResponseCommitAdmission::ExactReplay(
+ MycNip46ResponseCommitRecord {
+ completion,
+ response,
+ },
+ ))
+ }
+ MycNip46CommitAdmission::Committed(completion) => {
+ #[cfg(test)]
+ if request.fail_after_completion {
+ return Err(AtomicOperationError::Storage);
+ }
+ insert_response(transaction, &request).await?;
+ #[cfg(test)]
+ if request.fail_after_response {
+ return Err(AtomicOperationError::Storage);
+ }
+ let delivery = create_job(
+ transaction,
+ MycDeliverySource::signer_response(
+ request.completion.signer_request().operation_id(),
+ ),
+ request.response_digest,
+ request.committed_at,
+ &policy,
+ )
+ .await
+ .map_err(AtomicOperationError::from)?;
+ if !matches!(delivery, MycDeliveryJobAdmission::Created(_)) {
+ return Err(AtomicOperationError::Binding);
+ }
+ let response = read_response_by_operation(
+ transaction,
+ request.completion.signer_request().operation_id(),
+ )
+ .await?
+ .ok_or(AtomicOperationError::Binding)?;
+ exact_response(&response, &request)?;
+ Ok(MycNip46ResponseCommitAdmission::Committed(
+ MycNip46ResponseCommitRecord {
+ completion,
+ response,
+ },
+ ))
+ }
+ }
+ })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+
+ /// Reads the exact committed response bytes for a retained delivery job.
+ pub async fn read_nip46_response(
+ &self,
+ job_id: MycDeliveryJobId,
+ ) -> Result<Option<MycNip46ResponseRecord>, MycStateRepositoryError> {
+ let expected = PersistedMetadata::from(self.expected());
+ self.host()
+ .transaction(move |transaction| {
+ Box::pin(async move {
+ require_expected_metadata(transaction, &expected)
+ .await
+ .map_err(AtomicOperationError::from)?;
+ read_response_by_job(transaction, job_id).await
+ })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum AtomicOperationError {
+ Binding,
+ Storage,
+}
+
+impl From<RepositoryOperationError> for AtomicOperationError {
+ fn from(error: RepositoryOperationError) -> Self {
+ match error {
+ RepositoryOperationError::Binding => Self::Binding,
+ RepositoryOperationError::Storage => Self::Storage,
+ }
+ }
+}
+
+impl From<CommitOperationError> for AtomicOperationError {
+ fn from(error: CommitOperationError) -> Self {
+ match error {
+ CommitOperationError::Binding => Self::Binding,
+ CommitOperationError::Storage => Self::Storage,
+ }
+ }
+}
+
+impl From<DeliveryOperationError> for AtomicOperationError {
+ fn from(error: DeliveryOperationError) -> Self {
+ match error {
+ DeliveryOperationError::Binding => Self::Binding,
+ DeliveryOperationError::Storage => Self::Storage,
+ }
+ }
+}
+
+async fn insert_response(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ request: &MycNip46ResponseCommitRequest,
+) -> Result<(), AtomicOperationError> {
+ let result = sqlx::query(INSERT_RESPONSE_SQL)
+ .bind(
+ request
+ .completion
+ .signer_request()
+ .operation_id()
+ .as_bytes()
+ .as_slice(),
+ )
+ .bind(request.response_provider_operation_id.as_slice())
+ .bind(request.response_event_id.as_slice())
+ .bind(request.response_digest.as_bytes().as_slice())
+ .bind(request.response_bytes.as_ref())
+ .bind(to_i64(request.authored_at_unix_s)?)
+ .bind(to_i64(request.committed_at.get())?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| AtomicOperationError::Storage)?;
+ (result.rows_affected() == 1)
+ .then_some(())
+ .ok_or(AtomicOperationError::Storage)
+}
+
+async fn read_response_by_operation(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ operation_id: MycSignerOperationId,
+) -> Result<Option<MycNip46ResponseRecord>, AtomicOperationError> {
+ read_response(
+ transaction,
+ READ_RESPONSE_BY_OPERATION_SQL,
+ operation_id.as_bytes(),
+ )
+ .await
+}
+
+async fn read_response_by_job(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ job_id: MycDeliveryJobId,
+) -> Result<Option<MycNip46ResponseRecord>, AtomicOperationError> {
+ read_response(transaction, READ_RESPONSE_BY_JOB_SQL, job_id.as_bytes()).await
+}
+
+async fn read_response(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ sql: &'static str,
+ identity: &[u8; 32],
+) -> Result<Option<MycNip46ResponseRecord>, AtomicOperationError> {
+ let rows = sqlx::query(sql)
+ .bind(identity.as_slice())
+ .fetch_all(&mut *transaction)
+ .await
+ .map_err(|_| AtomicOperationError::Storage)?;
+ if rows.len() > 1 {
+ return Err(AtomicOperationError::Binding);
+ }
+ let Some(row) = rows.first() else {
+ return Ok(None);
+ };
+ let operation_id = MycSignerOperationId::from_persisted(blob32(row, "operation_id")?);
+ let response_provider_operation_id = blob32(row, "response_provider_operation_id")?;
+ let response_event_id = blob32(row, "response_event_id")?;
+ let response_digest = MycDeliveryArtifactDigest::from_bytes(blob32(row, "response_sha256")?);
+ let response_bytes = row
+ .try_get::<Option<Vec<u8>>, _>("response_bytes")
+ .map_err(|_| AtomicOperationError::Binding)?
+ .ok_or(AtomicOperationError::Binding)?
+ .into_boxed_slice();
+ let authored_at_unix_s = positive_i64(row, "authored_at_unix_s")?;
+ let committed_at = MycDeliveryTimeUnixMs::new(positive_i64(row, "committed_at_unix_ms")?)
+ .map_err(|_| AtomicOperationError::Binding)?;
+ let client_public_key = row
+ .try_get::<Option<String>, _>("client_public_key")
+ .map_err(|_| AtomicOperationError::Binding)?
+ .ok_or(AtomicOperationError::Binding)?;
+ let job_id = MycDeliveryJobId::from_persisted(blob32(row, "job_id")?);
+ let actual_digest: [u8; 32] = Sha256::digest(&response_bytes).into();
+ let event = validate_response_event(&response_bytes, &client_public_key, None)
+ .map_err(|_| AtomicOperationError::Binding)?;
+ if actual_digest != *response_digest.as_bytes()
+ || event.id.as_bytes() != &response_event_id
+ || event.created_at.as_secs() != authored_at_unix_s
+ {
+ return Err(AtomicOperationError::Binding);
+ }
+ let delivery_job = read_job(transaction, job_id)
+ .await
+ .map_err(AtomicOperationError::from)?
+ .ok_or(AtomicOperationError::Binding)?;
+ if delivery_job.operation_id() != Some(operation_id)
+ || delivery_job.artifact_digest() != response_digest
+ || delivery_job.created_at() != committed_at
+ {
+ return Err(AtomicOperationError::Binding);
+ }
+ Ok(Some(MycNip46ResponseRecord {
+ operation_id,
+ response_provider_operation_id,
+ response_event_id,
+ response_digest,
+ response_bytes,
+ authored_at_unix_s,
+ committed_at,
+ delivery_job,
+ }))
+}
+
+fn exact_response(
+ response: &MycNip46ResponseRecord,
+ request: &MycNip46ResponseCommitRequest,
+) -> Result<(), AtomicOperationError> {
+ (response.operation_id == request.completion.signer_request().operation_id()
+ && response.response_provider_operation_id == request.response_provider_operation_id
+ && response.response_event_id == request.response_event_id
+ && response.response_digest == request.response_digest
+ && response.response_bytes.as_ref() == request.response_bytes.as_ref()
+ && response.authored_at_unix_s == request.authored_at_unix_s
+ && response.committed_at == request.committed_at)
+ .then_some(())
+ .ok_or(AtomicOperationError::Binding)
+}
+
+fn validate_response_event(
+ bytes: &[u8],
+ client_public_key: &str,
+ expected_responder: Option<&str>,
+) -> Result<RadrootsNostrEvent, MycNip46ResponseCommitError> {
+ if bytes.is_empty() || bytes.len() > MYC_PROVIDER_OUTPUT_MAX_BYTES {
+ return Err(MycNip46ResponseCommitRequest::error(
+ MycNip46ResponseCommitErrorKind::InvalidResponse,
+ ));
+ }
+ let event: RadrootsNostrEvent = serde_json::from_slice(bytes).map_err(|_| {
+ MycNip46ResponseCommitRequest::error(MycNip46ResponseCommitErrorKind::InvalidResponse)
+ })?;
+ let canonical = serde_json::to_vec(&event).map_err(|_| {
+ MycNip46ResponseCommitRequest::error(MycNip46ResponseCommitErrorKind::InvalidResponse)
+ })?;
+ let expected_recipient = ["p", client_public_key];
+ let valid_recipient = event.tags.len() == 1
+ && event.tags.as_slice()[0]
+ .as_slice()
+ .iter()
+ .map(String::as_str)
+ .eq(expected_recipient);
+ if canonical != bytes
+ || event.kind != RadrootsNostrKind::Custom(NIP46_RPC_KIND)
+ || event.content.is_empty()
+ || !valid_recipient
+ || expected_responder.is_some_and(|expected| event.pubkey.to_hex() != expected)
+ || event.verify().is_err()
+ {
+ return Err(MycNip46ResponseCommitRequest::error(
+ MycNip46ResponseCommitErrorKind::InvalidResponse,
+ ));
+ }
+ Ok(event)
+}
+
+fn blob32(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<[u8; 32], AtomicOperationError> {
+ row.try_get::<Option<Vec<u8>>, _>(column)
+ .map_err(|_| AtomicOperationError::Binding)?
+ .ok_or(AtomicOperationError::Binding)?
+ .try_into()
+ .map_err(|_| AtomicOperationError::Binding)
+}
+
+fn positive_i64(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<u64, AtomicOperationError> {
+ row.try_get::<i64, _>(column)
+ .map_err(|_| AtomicOperationError::Binding)
+ .and_then(|value| u64::try_from(value).map_err(|_| AtomicOperationError::Binding))
+ .and_then(|value| {
+ value
+ .checked_sub(1)
+ .map(|_| value)
+ .ok_or(AtomicOperationError::Binding)
+ })
+}
+
+fn to_i64(value: u64) -> Result<i64, AtomicOperationError> {
+ i64::try_from(value).map_err(|_| AtomicOperationError::Binding)
+}
+
+fn map_transaction_error(
+ error: ServiceSqliteTransactionError<AtomicOperationError>,
+) -> MycStateRepositoryError {
+ if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
+ return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown);
+ }
+ let kind = match error.operation_error() {
+ Some(AtomicOperationError::Binding) => MycStateRepositoryErrorKind::Binding,
+ Some(AtomicOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction,
+ };
+ MycStateRepositoryError::new(kind)
+}
diff --git a/tests/build_policy.rs b/tests/build_policy.rs
@@ -83,6 +83,6 @@ fn source_lock_binds_the_current_cargo_lock() {
"source_archive_sha256 = \"975474804e6358b9228981add0a23181dbdd1afddf5ae12579c82220876bc379\""
));
assert!(SOURCE_LOCK.ends_with(
- "[contract_versions]\nconfig = 1\nstate = 7\nadmin = 1\nstatus = 1\nprovider = 1\n"
+ "[contract_versions]\nconfig = 1\nstate = 9\nadmin = 1\nstatus = 1\nprovider = 1\n"
));
}
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -11,6 +11,7 @@ const NIP46_REPLAY: &str = include_str!("../src/nip46_replay.rs");
const NIP46_WORK: &str = include_str!("../src/nip46_work.rs");
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 NIP46_VERIFICATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/nip46_verification.v1.json");
const NIP46_REPLAY_CONTRACT: &str =
@@ -23,6 +24,8 @@ const NIP46_WAVE_080_A_CONTRACT: &str =
include_str!("../contracts/services_hardening/nip46_wave_080_a.v1.json");
const NIP46_COMPLETION_CONTRACT: &str =
include_str!("../contracts/services_hardening/nip46_completion.v1.json");
+const NIP46_RESPONSE_CONTRACT: &str =
+ include_str!("../contracts/services_hardening/nip46_response_commit.v1.json");
const SOURCES: &[&str] = &[
include_str!("../src/cli_v1.rs"),
include_str!("../src/config_v1.rs"),
@@ -49,6 +52,7 @@ const SOURCES: &[&str] = &[
include_str!("../src/state_metadata.rs"),
include_str!("../src/state_repository.rs"),
include_str!("../src/state_request.rs"),
+ include_str!("../src/state_response.rs"),
];
#[test]
@@ -86,6 +90,7 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() {
"state_metadata",
"state_repository",
"state_request",
+ "state_response",
])
);
assert!(!ROOT.contains("pub mod "));
@@ -136,7 +141,13 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() {
"pub enum myc::MycNip46CommitAdmission",
"pub enum myc::MycNip46SessionEffect",
"pub enum myc::MycNip46CommitErrorKind",
- "pub async fn myc::MycStateRepository<'_>::commit_nip46_operation",
+ "pub struct myc::MycNip46ResponseCommitRequest",
+ "pub struct myc::MycNip46ResponseRecord",
+ "pub struct myc::MycNip46ResponseCommitRecord",
+ "pub enum myc::MycNip46ResponseCommitAdmission",
+ "pub enum myc::MycNip46ResponseCommitErrorKind",
+ "pub async fn myc::MycStateRepository<'_>::commit_nip46_response",
+ "pub async fn myc::MycStateRepository<'_>::read_nip46_response",
"pub async fn myc::MycStateRepository<'_>::read_connection_decision",
"pub fn myc::admit_myc_nip46_event",
"pub fn myc::admit_myc_nip46_request",
@@ -182,6 +193,7 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() {
"state_metadata",
"state_repository",
"state_request",
+ "state_response",
] {
assert!(
!PUBLIC_API.contains(&format!("pub mod myc::{module}")),
@@ -215,6 +227,64 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() {
}
#[test]
+fn step148_response_commit_is_one_atomic_exact_byte_authority() {
+ let contract: serde_json::Value =
+ serde_json::from_str(NIP46_RESPONSE_CONTRACT).expect("Step 148 contract");
+ assert_eq!(contract["schema"], "radroots.myc.nip46-response-commit.v1");
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(contract["step"], 148);
+ assert_eq!(contract["schema_version"], 9);
+ for required in [
+ "sealed_step147_completion_component",
+ "independently_signature_verified_canonical_kind_24133_event",
+ "exact_signed_response_bytes_sha256_and_event_id",
+ "zero_attempt_target_state",
+ "committed_response_bytes_only",
+ "completion_without_response_or_job",
+ "fail_closed_without_repair",
+ "relay_io_inside_transaction",
+ "response_reconstruction_on_retry",
+ ] {
+ assert!(
+ NIP46_RESPONSE_CONTRACT.contains(required),
+ "Step 148 contract is missing `{required}`"
+ );
+ }
+ for required in [
+ "ServiceSqliteTransaction",
+ "commit_operation(transaction",
+ "nip46_signed_responses",
+ "create_job(",
+ "read_response_by_operation",
+ "signed_response_bytes",
+ "fail_after_completion_for_test",
+ "fail_after_response_for_test",
+ ] {
+ assert!(
+ NIP46_RESPONSE.contains(required),
+ "Step 148 implementation is missing `{required}`"
+ );
+ }
+ assert!(!PUBLIC_API.contains("commit_nip46_operation"));
+ for forbidden in [
+ "RelayPool",
+ ".publish(",
+ "tokio::spawn",
+ "std::time::SystemTime",
+ "Timestamp::now",
+ "rand::",
+ "getrandom",
+ "SqlitePool",
+ "rusqlite",
+ ] {
+ assert!(
+ !NIP46_RESPONSE.contains(forbidden),
+ "Step 148 gained forbidden authority `{forbidden}`"
+ );
+ }
+}
+
+#[test]
fn step147_completion_is_atomic_redacted_and_defers_delivery_authority() {
let contract: serde_json::Value =
serde_json::from_str(NIP46_COMPLETION_CONTRACT).expect("Step 147 contract");
@@ -231,7 +301,8 @@ fn step147_completion_is_atomic_redacted_and_defers_delivery_authority() {
"protected_provider_output\": \"not_persisted",
"failed_transaction\": \"no_session_or_completion_effect",
"step147_checkpoint\": \"integration_only_not_promotable",
- "step147_component\": \"must_be_composed_before_master_promotion",
+ "step147_component\": \"composed_by_step148_atomic_response_commit",
+ "composition_status\": \"satisfied_on_rcld_080_integration_branch",
"commit_owner\": 148",
"promotion_owner\": 151",
"outer_signed_nip46_response\": 148",
@@ -444,7 +515,7 @@ fn public_errors_remain_crate_owned_redacted_and_source_free() {
.lines()
.filter(|line| line.starts_with("pub struct myc::") && line.ends_with("Error"))
.count();
- assert_eq!(public_error_count, 23);
+ assert_eq!(public_error_count, 24);
assert!(!PUBLIC_API.contains("pub struct myc::MycRuntimeFoundation {"));
assert!(!PUBLIC_API.contains("pub struct myc::MycStateHost {"));
}
diff --git a/tests/services_hardening_legacy_removal.rs b/tests/services_hardening_legacy_removal.rs
@@ -8,6 +8,7 @@ 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_catalog.rs"),
+ include_str!("../src/state_completion.rs"),
include_str!("../src/state_connection.rs"),
include_str!("../src/state_delivery.rs"),
include_str!("../src/state_discovery.rs"),
@@ -17,6 +18,7 @@ const ACTIVE_STATE_SOURCES: &[&str] = &[
include_str!("../src/state_metadata.rs"),
include_str!("../src/state_repository.rs"),
include_str!("../src/state_request.rs"),
+ include_str!("../src/state_response.rs"),
];
#[test]
@@ -111,7 +113,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(), 10);
+ assert_eq!(LIB_SOURCE.matches("mod state_").count(), 12);
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
@@ -52,7 +52,7 @@ fn native_release_contract_and_manifest_metadata_are_exact() {
},
"contract_versions": {
"config": 1,
- "state": 8,
+ "state": 9,
"admin": 1,
"status": 1,
"provider": 1
@@ -112,7 +112,7 @@ fn native_release_contract_and_manifest_metadata_are_exact() {
service = "myc"
host_feature_profile = "service-host"
config_contract_version = 1
- state_contract_version = 8
+ state_contract_version = 9
admin_contract_version = 1
status_contract_version = 1
provider_contract_version = 1
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_v8_before_request_admission() {
+async fn exact_schema_v3_state_advances_to_v9_before_request_admission() {
let directory = tempfile::tempdir().expect("temporary root");
let runtime = runtime(directory.path());
prepare_state_directory(&runtime);
@@ -513,6 +513,9 @@ async fn exact_schema_v3_state_advances_to_v8_before_request_admission() {
"DROP TRIGGER connection_permissions_no_update",
"DROP TRIGGER connections_no_delete",
"DROP TRIGGER connections_guard_update",
+ "DROP TRIGGER nip46_signed_responses_no_delete",
+ "DROP TRIGGER nip46_signed_responses_no_update",
+ "DROP TABLE nip46_signed_responses",
"DROP TRIGGER nip46_operation_commits_no_delete",
"DROP TRIGGER nip46_operation_commits_no_update",
"DROP TABLE nip46_operation_commits",
@@ -532,7 +535,7 @@ async fn exact_schema_v3_state_advances_to_v8_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)",
+ "DELETE FROM schema_migrations WHERE version IN (4, 5, 6, 7, 8, 9)",
] {
sqlx::query(sql)
.execute(&mut connection)
diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs
@@ -16,8 +16,9 @@ use myc::{
MYC_STATE_SCHEMA_VERSION_7_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_7_OBJECT_COUNT,
MYC_STATE_SCHEMA_VERSION_7_SHA256, MYC_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256,
MYC_STATE_SCHEMA_VERSION_8_OBJECT_COUNT, MYC_STATE_SCHEMA_VERSION_8_SHA256,
- MycStateCatalogErrorKind, myc_migration_catalog, myc_schema_catalog,
- validate_myc_state_catalogs,
+ MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256, MYC_STATE_SCHEMA_VERSION_9_OBJECT_COUNT,
+ MYC_STATE_SCHEMA_VERSION_9_SHA256, MycStateCatalogErrorKind, myc_migration_catalog,
+ myc_schema_catalog, validate_myc_state_catalogs,
};
use radroots_service_sqlite::{
MigrationCatalog, MigrationChecksum, MigrationDescriptor, SchemaCatalog, SchemaDigest,
@@ -29,13 +30,13 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs");
const MANIFEST: &str = include_str!("../Cargo.toml");
#[test]
-fn schema_v1_through_v8_and_all_migrations_have_exact_literal_identities() {
+fn schema_v1_through_v9_and_all_migrations_have_exact_literal_identities() {
let migrations = myc_migration_catalog().expect("Myc migration catalog");
let schema = myc_schema_catalog().expect("Myc schema catalog");
assert_eq!(MYC_STATE_BASE_SCHEMA_VERSION, 1);
- assert_eq!(MYC_STATE_SCHEMA_VERSION, 8);
- assert_eq!(migrations.descriptors().len(), 7);
+ assert_eq!(MYC_STATE_SCHEMA_VERSION, 9);
+ assert_eq!(migrations.descriptors().len(), 8);
let metadata = &migrations.descriptors()[0];
assert_eq!(metadata.target_version(), 2);
assert_eq!(metadata.name().as_str(), "create_myc_state_metadata");
@@ -94,13 +95,20 @@ fn schema_v1_through_v8_and_all_migrations_have_exact_literal_identities() {
completion.checksum().as_bytes(),
&MYC_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256
);
- assert_eq!(migrations.current_version(), 8);
+ let response = &migrations.descriptors()[7];
+ assert_eq!(response.target_version(), 9);
+ assert_eq!(response.name().as_str(), "create_nip46_atomic_response");
+ assert_eq!(
+ response.checksum().as_bytes(),
+ &MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256
+ );
+ assert_eq!(migrations.current_version(), 9);
assert_eq!(
migrations.digest().as_bytes(),
&MYC_MIGRATION_CATALOG_SHA256
);
- assert_eq!(schema.versions().len(), 8);
+ assert_eq!(schema.versions().len(), 9);
assert_eq!(schema.versions()[0].version(), 1);
assert_eq!(
schema.versions()[0].object_count(),
@@ -180,6 +188,16 @@ fn schema_v1_through_v8_and_all_migrations_have_exact_literal_identities() {
schema.versions()[7].digest().as_bytes(),
&MYC_STATE_SCHEMA_VERSION_8_SHA256
);
+ assert_eq!(schema.versions()[8].version(), 9);
+ assert_eq!(
+ schema.versions()[8].object_count(),
+ MYC_STATE_SCHEMA_VERSION_9_OBJECT_COUNT
+ );
+ assert_eq!(schema.versions()[8].object_count(), 59);
+ assert_eq!(
+ schema.versions()[8].digest().as_bytes(),
+ &MYC_STATE_SCHEMA_VERSION_9_SHA256
+ );
assert_eq!(schema.digest().as_bytes(), &MYC_STATE_SCHEMA_CATALOG_SHA256);
assert_eq!(schema.migration_catalog_digest(), migrations.digest());
validate_myc_state_catalogs(&migrations, &schema).expect("exact catalogs");
@@ -190,7 +208,7 @@ fn schema_v1_through_v8_and_all_migrations_have_exact_literal_identities() {
);
assert_eq!(
hex::encode(MYC_MIGRATION_CATALOG_SHA256),
- "a2ddd320e2f96d08c8177da1175d0ed32983168dc16b3d326fd43d9d655d8e82"
+ "8604c51ea6f16f0aa37a63eee09c600a0b7031efbc8e5e0e41b5f0f53c48e70b"
);
assert_eq!(
hex::encode(MYC_STATE_SCHEMA_VERSION_1_SHA256),
@@ -202,7 +220,7 @@ fn schema_v1_through_v8_and_all_migrations_have_exact_literal_identities() {
);
assert_eq!(
hex::encode(MYC_STATE_SCHEMA_CATALOG_SHA256),
- "496228fcd4c2a583d7ffe2113e0a7add68bcf7b70218ce5b946e2f92f30edac7"
+ "632c8d5216d2a541fd28ae3a2cae8069c58bef11fe81d2f04399a77cf8359ea7"
);
assert_eq!(
hex::encode(MYC_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256),
@@ -252,6 +270,14 @@ fn schema_v1_through_v8_and_all_migrations_have_exact_literal_identities() {
hex::encode(MYC_STATE_SCHEMA_VERSION_8_SHA256),
"50158cb093ed70b3d5783762b18d21b7809665c89fde925d872365909d28f06e"
);
+ assert_eq!(
+ hex::encode(MYC_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256),
+ "fdb039d472cd62e7da46c40a9786c3a9a7ec55f2ae7db689d0f855fbbf8a9566"
+ );
+ assert_eq!(
+ hex::encode(MYC_STATE_SCHEMA_VERSION_9_SHA256),
+ "eee76f4ff0bd2dc2c061ae384e00de16c5f5efc4b6edd0ac7f73cad83a991ff7"
+ );
}
#[test]
@@ -307,8 +333,11 @@ fn independent_validator_rejects_migration_or_schema_drift() {
let v7 = SchemaVersionCatalog::new(7, [object.clone()], v7_digest).expect("schema v7");
let v8_digest =
SchemaVersionCatalog::computed_digest(8, [object.clone()]).expect("schema-v8 digest");
- let v8 = SchemaVersionCatalog::new(8, [object], v8_digest).expect("schema v8");
- let schema = SchemaCatalog::new(&expected_migrations, [v1, v2, v3, v4, v5, v6, v7, v8])
+ let v8 = SchemaVersionCatalog::new(8, [object.clone()], v8_digest).expect("schema v8");
+ let v9_digest =
+ SchemaVersionCatalog::computed_digest(9, [object.clone()]).expect("schema-v9 digest");
+ let v9 = SchemaVersionCatalog::new(9, [object], v9_digest).expect("schema v9");
+ let schema = SchemaCatalog::new(&expected_migrations, [v1, v2, v3, v4, v5, v6, v7, v8, v9])
.expect("drift schema catalog");
assert_eq!(
validate_myc_state_catalogs(&expected_migrations, &schema)
diff --git a/tests/services_hardening_state_repository.rs b/tests/services_hardening_state_repository.rs
@@ -120,7 +120,7 @@ async fn initialization_migrates_and_binds_exact_metadata_before_inspection() {
.fetch_all(&mut connection)
.await
.expect("migration rows");
- assert_eq!(migrations.len(), 6);
+ assert_eq!(migrations.len(), 8);
assert_eq!(migrations[0].get::<i64, _>(0), 2);
assert_eq!(
migrations[0].get::<String, _>(1),
@@ -151,6 +151,16 @@ async fn initialization_migrates_and_binds_exact_metadata_before_inspection() {
migrations[5].get::<String, _>(1),
"create_discovery_desired_state"
);
+ assert_eq!(migrations[6].get::<i64, _>(0), 8);
+ assert_eq!(
+ migrations[6].get::<String, _>(1),
+ "create_nip46_operation_completion"
+ );
+ assert_eq!(migrations[7].get::<i64, _>(0), 9);
+ assert_eq!(
+ migrations[7].get::<String, _>(1),
+ "create_nip46_atomic_response"
+ );
let binding = sqlx::query(
"SELECT normalized_config_sha256, transport_public_key, user_public_key, \
discovery_public_key, config_contract_version, state_contract_version, \