commit 9ede1be7a1ef6e59c0e3ab49aad7a568c41fae91
parent 186f08618b22f96b45dc3e54609cfa15cb86a330
Author: triesap <tyson@radroots.org>
Date: Sat, 22 Aug 2026 03:48:56 +0000
state: recover delivery and render NIP-05
Diffstat:
14 files changed, 1553 insertions(+), 185 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -190,6 +190,13 @@
- Treat stored signed bytes as the sole publication authority. Retry, crash
recovery, and reopen must submit the identical bytes and digest without
deserializing, rebuilding, re-signing, or changing targets.
+- Create delivery jobs only inside the atomic signed-response or discovery
+ transaction. Recovery is a bounded, cursor-driven repository operation with
+ injected time/jitter evidence; final supervised startup looping remains a
+ later runtime owner.
+- Render NIP-05 only from verified committed discovery state through an
+ explicit desired/current offline operation. The service library must not
+ silently host the document or acquire network authority while rendering it.
- Distinguish submitted, delivered, failed, and unknown outcomes. Lost
acknowledgement never becomes proof of failure or delivery.
- Bound every queue, pool, request, response, event, tag set, retry schedule,
diff --git a/README b/README
@@ -102,6 +102,18 @@ 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.
+Step 149 closes the repository-owned delivery recovery boundary without
+claiming the final daemon task graph. New jobs can be created only by the
+atomic signed-response or discovery commit and are rejected at the configured
+active-outbox ceiling. Attempt resolution persists caller-injected,
+attempt-cap-bounded full jitter; restart recovery scans fixed 128-job cursor
+batches, revalidates exact signed source bytes and bounded attempt histories,
+recovers only expired leases, and idempotently promotes proven desired
+discovery. It performs no relay I/O and creates no task, clock, or entropy
+authority. Explicit desired/current offline NIP-05 export re-reads verified
+state, works through read-only inspection, and returns canonical compact
+`names` plus NIP-46 discovery JSON without hosting it.
+
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
@@ -207,6 +207,7 @@ pub const fn myc::MycDeliverySourceKind::as_str(self) -> &'static str
pub enum myc::MycDeliveryStateErrorKind
pub myc::MycDeliveryStateErrorKind::InvalidPolicy
pub myc::MycDeliveryStateErrorKind::InvalidRelayId
+pub myc::MycDeliveryStateErrorKind::InvalidRetryJitter
pub myc::MycDeliveryStateErrorKind::InvalidTime
impl myc::MycDeliveryStateErrorKind
pub const fn myc::MycDeliveryStateErrorKind::code(self) -> &'static str
@@ -278,6 +279,9 @@ pub myc::MycLocalSignerTransportErrorKind::Transport
pub myc::MycLocalSignerTransportErrorKind::UnsupportedPlatform
impl myc::MycLocalSignerTransportErrorKind
pub const fn myc::MycLocalSignerTransportErrorKind::code(self) -> &'static str
+pub enum myc::MycNip05ExportSelection
+pub myc::MycNip05ExportSelection::Current
+pub myc::MycNip05ExportSelection::Desired
pub enum myc::MycNip46AdmissionErrorKind
pub myc::MycNip46AdmissionErrorKind::EmptyEvent
pub myc::MycNip46AdmissionErrorKind::EmptyPlaintext
@@ -790,17 +794,32 @@ pub const fn myc::MycDeliveryJobRecord::targets(&self) -> &[myc::MycDeliveryTarg
pub const fn myc::MycDeliveryJobRecord::updated_at(&self) -> myc::MycDeliveryTimeUnixMs
impl core::fmt::Debug for myc::MycDeliveryJobRecord
pub fn myc::MycDeliveryJobRecord::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
-pub struct myc::MycDeliveryJobRequest
-impl myc::MycDeliveryJobRequest
-pub const fn myc::MycDeliveryJobRequest::signer_response(myc::MycSignerOperationId, myc::MycDeliveryArtifactDigest, myc::MycDeliveryTimeUnixMs) -> Self
-impl core::fmt::Debug for myc::MycDeliveryJobRequest
-pub fn myc::MycDeliveryJobRequest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycDeliveryRecoveryEntropy(_)
+impl myc::MycDeliveryRecoveryEntropy
+pub const fn myc::MycDeliveryRecoveryEntropy::from_injected_entropy([u8; 32]) -> Self
+impl core::fmt::Debug for myc::MycDeliveryRecoveryEntropy
+pub fn myc::MycDeliveryRecoveryEntropy::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycDeliveryRecoveryReport
+impl myc::MycDeliveryRecoveryReport
+pub const fn myc::MycDeliveryRecoveryReport::active_targets(self) -> u32
+pub const fn myc::MycDeliveryRecoveryReport::examined_jobs(self) -> u32
+pub const fn myc::MycDeliveryRecoveryReport::finalized_jobs(self) -> u32
+pub const fn myc::MycDeliveryRecoveryReport::promoted_discovery_generations(self) -> u32
+pub const fn myc::MycDeliveryRecoveryReport::ready_targets(self) -> u32
+pub const fn myc::MycDeliveryRecoveryReport::recovered_expired_attempts(self) -> u32
+pub const fn myc::MycDeliveryRecoveryReport::scheduled_targets(self) -> u32
+impl core::fmt::Debug for myc::MycDeliveryRecoveryReport
+pub fn myc::MycDeliveryRecoveryReport::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct myc::MycDeliveryRelayId(_)
impl myc::MycDeliveryRelayId
pub fn myc::MycDeliveryRelayId::as_str(&self) -> &str
pub fn myc::MycDeliveryRelayId::new(&str) -> core::result::Result<Self, myc::MycDeliveryStateError>
impl core::fmt::Debug for myc::MycDeliveryRelayId
pub fn myc::MycDeliveryRelayId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycDeliveryRetryJitter(_)
+impl myc::MycDeliveryRetryJitter
+pub const fn myc::MycDeliveryRetryJitter::get(self) -> u64
+pub fn myc::MycDeliveryRetryJitter::new(u64) -> core::result::Result<Self, myc::MycDeliveryStateError>
pub struct myc::MycDeliveryStateError
impl myc::MycDeliveryStateError
pub const fn myc::MycDeliveryStateError::code(self) -> &'static str
@@ -952,6 +971,19 @@ impl myc::MycLocalSignerUntrustedResponse
pub fn myc::MycLocalSignerUntrustedResponse::verify(self, &myc::MycProviderBinding, &myc::MycProviderOperation, myc::MycProviderResponseObservedAtUnixMs) -> core::result::Result<myc::MycVerifiedProviderResponse, myc::MycProviderVerificationError>
impl core::fmt::Debug for myc::MycLocalSignerUntrustedResponse
pub fn myc::MycLocalSignerUntrustedResponse::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycNip05Document
+impl myc::MycNip05Document
+pub fn myc::MycNip05Document::bytes(&self) -> &[u8]
+pub const fn myc::MycNip05Document::digest(&self) -> myc::MycNip05DocumentDigest
+pub fn myc::MycNip05Document::domain(&self) -> &str
+pub const fn myc::MycNip05Document::selection(&self) -> myc::MycNip05ExportSelection
+impl core::fmt::Debug for myc::MycNip05Document
+pub fn myc::MycNip05Document::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct myc::MycNip05DocumentDigest(_)
+impl myc::MycNip05DocumentDigest
+pub const fn myc::MycNip05DocumentDigest::as_bytes(&self) -> &[u8; 32]
+impl core::fmt::Debug for myc::MycNip05DocumentDigest
+pub fn myc::MycNip05DocumentDigest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct myc::MycNip05ProjectionDigest(_)
impl myc::MycNip05ProjectionDigest
pub const fn myc::MycNip05ProjectionDigest::as_bytes(&self) -> &[u8; 32]
@@ -1414,23 +1446,25 @@ impl myc::MycStateRepository<'_>
pub async fn myc::MycStateRepository<'_>::admit_signer_request(&self, &myc::MycSignerRequest) -> core::result::Result<myc::MycSignerRequestAdmission, myc::MycStateRepositoryError>
impl myc::MycStateRepository<'_>
pub async fn myc::MycStateRepository<'_>::claim_delivery_target(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId, myc::MycDeliveryAttemptNonce, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDeliveryClaim, myc::MycStateRepositoryError>
-pub async fn myc::MycStateRepository<'_>::create_delivery_job(&self, &myc::MycDeliveryJobRequest) -> core::result::Result<myc::MycDeliveryJobAdmission, myc::MycStateRepositoryError>
pub async fn myc::MycStateRepository<'_>::mark_delivery_attempt_submitted(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId, myc::MycDeliveryAttemptId, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDeliveryAttemptRecord, myc::MycStateRepositoryError>
pub async fn myc::MycStateRepository<'_>::read_delivery_attempts(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId) -> core::result::Result<alloc::boxed::Box<[myc::MycDeliveryAttemptRecord]>, myc::MycStateRepositoryError>
pub async fn myc::MycStateRepository<'_>::read_delivery_job(&self, myc::MycDeliveryJobId) -> core::result::Result<core::option::Option<myc::MycDeliveryJobRecord>, myc::MycStateRepositoryError>
-pub async fn myc::MycStateRepository<'_>::record_delivery_attempt_outcome(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId, myc::MycDeliveryAttemptId, myc::MycDeliveryAttemptOutcome, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDeliveryJobRecord, myc::MycStateRepositoryError>
-pub async fn myc::MycStateRepository<'_>::recover_expired_delivery_lease(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId, myc::MycDeliveryAttemptId, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDeliveryJobRecord, myc::MycStateRepositoryError>
+pub async fn myc::MycStateRepository<'_>::record_delivery_attempt_outcome(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId, myc::MycDeliveryAttemptId, myc::MycDeliveryAttemptOutcome, myc::MycDeliveryRetryJitter, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDeliveryJobRecord, myc::MycStateRepositoryError>
+pub async fn myc::MycStateRepository<'_>::recover_expired_delivery_lease(&self, myc::MycDeliveryJobId, &myc::MycDeliveryRelayId, myc::MycDeliveryAttemptId, myc::MycDeliveryRetryJitter, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDeliveryJobRecord, myc::MycStateRepositoryError>
impl myc::MycStateRepository<'_>
pub async fn myc::MycStateRepository<'_>::commit_discovery_desired_state(&self, &myc::MycDiscoveryCommitRequest) -> core::result::Result<myc::MycDiscoveryCommitAdmission, myc::MycStateRepositoryError>
pub async fn myc::MycStateRepository<'_>::promote_delivered_discovery_state(&self, myc::MycDeliveryJobId, myc::MycDeliveryTimeUnixMs) -> core::result::Result<myc::MycDiscoveryPublicationState, myc::MycStateRepositoryError>
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>
+pub async fn myc::MycStateRepository<'_>::render_offline_nip05(&self, myc::MycNip05ExportSelection) -> core::result::Result<myc::MycNip05Document, myc::MycStateRepositoryError>
impl myc::MycStateRepository<'_>
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>
+impl myc::MycStateRepository<'_>
+pub async fn myc::MycStateRepository<'_>::recover_delivery_state(&self, myc::MycDeliveryTimeUnixMs, myc::MycDeliveryRecoveryEntropy) -> core::result::Result<myc::MycDeliveryRecoveryReport, myc::MycStateRepositoryError>
impl<'host> myc::MycStateRepository<'host>
pub async fn myc::MycStateRepository<'host>::verify_binding(&self) -> core::result::Result<(), myc::MycStateRepositoryError>
impl core::fmt::Debug for myc::MycStateRepository<'_>
@@ -1502,7 +1536,9 @@ pub const myc::MYC_CONFIG_SCHEMA: &str
pub const myc::MYC_CONFIG_SCHEMA_VERSION: u32
pub const myc::MYC_CONNECTION_PERMISSION_MAX_COUNT: usize
pub const myc::MYC_DELIVERY_ATTEMPT_MAX_COUNT: u32
+pub const myc::MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT: usize
pub const myc::MYC_DELIVERY_RELAY_ID_MAX_BYTES: usize
+pub const myc::MYC_DELIVERY_RETRY_JITTER_MAX_MS: u64
pub const myc::MYC_DELIVERY_TARGET_MAX_COUNT: usize
pub const myc::MYC_DISCOVERY_DOCUMENT_MAX_BYTES: usize
pub const myc::MYC_ENCRYPTED_IDENTITY_BACKUP_INCLUDED: bool
@@ -1511,6 +1547,7 @@ pub const myc::MYC_ENCRYPTED_IDENTITY_ENVELOPE_MAX_BYTES: usize
pub const myc::MYC_LOCAL_SIGNER_ENDPOINT: &str
pub const myc::MYC_LOCAL_SIGNER_TRANSPORT_CONTRACT_VERSION: u32
pub const myc::MYC_MIGRATION_CATALOG_SHA256: [u8; 32]
+pub const myc::MYC_NIP05_DOCUMENT_MAX_BYTES: usize
pub const myc::MYC_NIP05_PROJECTION_MAX_BYTES: usize
pub const myc::MYC_NIP46_CANONICAL_REQUEST_MAX_BYTES: usize
pub const myc::MYC_NIP46_EVENT_ID_MAX_BYTES: usize
diff --git a/contracts/services_hardening/delivery_recovery_export.v1.json b/contracts/services_hardening/delivery_recovery_export.v1.json
@@ -0,0 +1,93 @@
+{
+ "schema": "radroots.myc.delivery-recovery-export.v1",
+ "contract_version": 1,
+ "step": 149,
+ "schema_version": 9,
+ "delivery_admission": {
+ "production_authority": "atomic_response_or_discovery_commit_only",
+ "configured_active_job_ceiling": "/resource_limits/queues/outbox",
+ "exact_replay_at_ceiling": "admitted",
+ "new_job_at_ceiling": "fail_closed",
+ "partial_delivery_job_constructor": "absent"
+ },
+ "attempt_resolution": {
+ "outcomes": ["delivered", "relay_rejected", "transport_failed", "unknown_acknowledgement"],
+ "unknown_is_failure": false,
+ "retry_jitter": "caller_injected_full_jitter_milliseconds",
+ "jitter_minimum_ms": 0,
+ "jitter_maximum_ms": 300000,
+ "relationship_check": "not_greater_than_exponential_attempt_cap",
+ "persisted_result": "exact_next_attempt_at_unix_ms",
+ "payload_or_target_mutation": false
+ },
+ "restart_recovery": {
+ "batch_maximum_jobs": 128,
+ "continuation": "internal_opaque_job_identity_cursor",
+ "caller_contract": "pause_admission_while_one_public_call_completes_all_internal_batches",
+ "wall_time": "caller_injected_positive_unix_milliseconds",
+ "entropy": "caller_injected_one_use_seed",
+ "transaction": "one_bounded_service_sqlite_transaction_per_batch",
+ "expired_pre_submit": "retryable_failure",
+ "expired_submitted": "unknown_acknowledgement",
+ "already_proven_terminal": "idempotent_finalization",
+ "delivered_desired_discovery": "idempotent_current_promotion",
+ "relay_io": false,
+ "task_spawn": false
+ },
+ "startup_invariants": [
+ "configured_active_job_ceiling",
+ "no_orphan_signed_response",
+ "no_orphan_discovery_document",
+ "no_orphan_target",
+ "no_orphan_attempt",
+ "exact_source_digest_binding",
+ "signature_verified_canonical_active_response_bytes",
+ "signature_verified_canonical_active_discovery_bytes",
+ "consecutive_bounded_attempt_history",
+ "active_target_matches_latest_attempt",
+ "immutable_target_membership"
+ ],
+ "offline_nip05": {
+ "selection": ["desired", "current"],
+ "current_requires_proven_delivery": true,
+ "state_authority": "verified_committed_projection",
+ "read_only_inspection_supported": true,
+ "hosted_by_daemon": false,
+ "wire_order": ["names", "nip46"],
+ "root_name": "_",
+ "nip46_fields": ["relays", "nostrconnect_url"],
+ "encoding": "compact_canonical_utf8_json",
+ "maximum_bytes": 524288,
+ "digest": "sha256_exact_bytes",
+ "debug": "redacted"
+ },
+ "qualification": [
+ "queue_ceiling_and_exact_replay",
+ "jitter_zero_exact_and_just_over_relationship",
+ "submitted_lease_becomes_unknown_after_restart",
+ "recovery_reopen_and_idempotent_replay",
+ "delivered_desired_restart_promotion",
+ "desired_and_current_export_selection",
+ "exact_nip05_nip46_wire_bytes_and_digest",
+ "read_only_offline_export",
+ "source_free_redacted_debug_and_errors"
+ ],
+ "deferrals": {
+ "admin_routes": 150,
+ "rcld_promotion": 151,
+ "final_supervised_startup_loop": 158,
+ "nix": "deferred_and_unclaimed",
+ "oci": "deferred_and_unclaimed"
+ },
+ "forbidden": [
+ "standalone_signer_delivery_job_creation",
+ "response_reconstruction",
+ "target_set_mutation",
+ "ambient_clock",
+ "ambient_entropy",
+ "relay_io",
+ "network_hosting",
+ "task_spawn",
+ "alternate_sqlite_authority"
+ ]
+}
diff --git a/src/lib.rs b/src/lib.rs
@@ -28,6 +28,7 @@ mod state_governance;
mod state_host;
mod state_maintenance;
mod state_metadata;
+mod state_recovery;
mod state_repository;
mod state_request;
mod state_response;
@@ -141,19 +142,21 @@ pub use state_connection::{
MycConnectionStateErrorKind, MycConnectionStatus, MycConnectionTimeUnixMs,
};
pub use state_delivery::{
- MYC_DELIVERY_ATTEMPT_MAX_COUNT, MYC_DELIVERY_RELAY_ID_MAX_BYTES, MYC_DELIVERY_TARGET_MAX_COUNT,
- MycDeliveryArtifactDigest, MycDeliveryAttemptId, MycDeliveryAttemptNonce,
- MycDeliveryAttemptOutcome, MycDeliveryAttemptRecord, MycDeliveryAttemptStatus,
- MycDeliveryClaim, MycDeliveryJobAdmission, MycDeliveryJobId, MycDeliveryJobRecord,
- MycDeliveryJobRequest, MycDeliveryJobStatus, MycDeliveryPolicyMode, MycDeliveryRelayId,
- MycDeliverySourceKind, MycDeliveryStateError, MycDeliveryStateErrorKind,
- MycDeliveryTargetRecord, MycDeliveryTargetStatus, MycDeliveryTimeUnixMs,
+ MYC_DELIVERY_ATTEMPT_MAX_COUNT, MYC_DELIVERY_RELAY_ID_MAX_BYTES,
+ MYC_DELIVERY_RETRY_JITTER_MAX_MS, MYC_DELIVERY_TARGET_MAX_COUNT, MycDeliveryArtifactDigest,
+ MycDeliveryAttemptId, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome,
+ MycDeliveryAttemptRecord, MycDeliveryAttemptStatus, MycDeliveryClaim, MycDeliveryJobAdmission,
+ MycDeliveryJobId, MycDeliveryJobRecord, MycDeliveryJobStatus, MycDeliveryPolicyMode,
+ MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliverySourceKind, MycDeliveryStateError,
+ MycDeliveryStateErrorKind, MycDeliveryTargetRecord, MycDeliveryTargetStatus,
+ MycDeliveryTimeUnixMs,
};
pub use state_discovery::{
- MYC_DISCOVERY_DOCUMENT_MAX_BYTES, MYC_NIP05_PROJECTION_MAX_BYTES, MycDiscoveryCommitAdmission,
- MycDiscoveryCommitRecord, MycDiscoveryCommitRequest, MycDiscoveryDocumentDigest,
- MycDiscoveryDocumentRecord, MycDiscoveryGenerationId, MycDiscoveryPublicationState,
- MycDiscoveryStateError, MycDiscoveryStateErrorKind, MycNip05ProjectionDigest,
+ MYC_DISCOVERY_DOCUMENT_MAX_BYTES, MYC_NIP05_DOCUMENT_MAX_BYTES, MYC_NIP05_PROJECTION_MAX_BYTES,
+ MycDiscoveryCommitAdmission, MycDiscoveryCommitRecord, MycDiscoveryCommitRequest,
+ MycDiscoveryDocumentDigest, MycDiscoveryDocumentRecord, MycDiscoveryGenerationId,
+ MycDiscoveryPublicationState, MycDiscoveryStateError, MycDiscoveryStateErrorKind,
+ MycNip05Document, MycNip05DocumentDigest, MycNip05ExportSelection, MycNip05ProjectionDigest,
};
pub use state_governance::{
MYC_AUDIT_PAGE_MAX_ITEMS, MYC_AUDIT_RETENTION_MAX_MS, MYC_COMPACTION_MAX_ROWS,
@@ -177,6 +180,9 @@ pub use state_metadata::{
MycExpectedIdentities, MycExpectedPublicIdentity, MycNormalizedConfigDigest, MycStateMetadata,
MycStateMetadataError, MycStateMetadataErrorKind, MycStatePolicyVersions,
};
+pub use state_recovery::{
+ MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT, MycDeliveryRecoveryEntropy, MycDeliveryRecoveryReport,
+};
pub use state_repository::{
MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind,
};
diff --git a/src/state_delivery.rs b/src/state_delivery.rs
@@ -21,13 +21,14 @@ pub const MYC_DELIVERY_TARGET_MAX_COUNT: usize = 32;
pub const MYC_DELIVERY_ATTEMPT_MAX_COUNT: u32 = 32;
/// Maximum UTF-8 byte length of a canonical delivery relay identifier.
pub const MYC_DELIVERY_RELAY_ID_MAX_BYTES: usize = 64;
+/// Maximum caller-injected full-jitter delay for one retry.
+pub const MYC_DELIVERY_RETRY_JITTER_MAX_MS: u64 = 300_000;
const JOB_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_job.v1\0";
const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_attempt.v1\0";
-const READ_SIGNER_OPERATION_SQL: &str = r#"SELECT COUNT(*) AS row_count
-FROM nip46_requests
-WHERE operation_id = ?"#;
+const READ_ACTIVE_JOB_COUNT_SQL: &str =
+ "SELECT COUNT(*) AS row_count FROM delivery_jobs WHERE status IN ('pending', 'active')";
const READ_JOB_SQL: &str = r#"SELECT
CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32
@@ -202,6 +203,7 @@ pub enum MycDeliveryStateErrorKind {
InvalidTime,
InvalidRelayId,
InvalidPolicy,
+ InvalidRetryJitter,
}
impl MycDeliveryStateErrorKind {
@@ -212,6 +214,7 @@ impl MycDeliveryStateErrorKind {
Self::InvalidTime => "delivery_time_invalid",
Self::InvalidRelayId => "delivery_relay_id_invalid",
Self::InvalidPolicy => "delivery_policy_invalid",
+ Self::InvalidRetryJitter => "delivery_retry_jitter_invalid",
}
}
}
@@ -246,6 +249,7 @@ impl fmt::Display for MycDeliveryStateError {
MycDeliveryStateErrorKind::InvalidTime => "delivery time is invalid",
MycDeliveryStateErrorKind::InvalidRelayId => "delivery relay identity is invalid",
MycDeliveryStateErrorKind::InvalidPolicy => "delivery policy is invalid",
+ MycDeliveryStateErrorKind::InvalidRetryJitter => "delivery retry jitter is invalid",
})
}
}
@@ -261,6 +265,31 @@ impl fmt::Debug for MycDeliveryStateError {
impl Error for MycDeliveryStateError {}
+/// Caller-injected full-jitter delay, later relationship-checked against a job.
+///
+/// The runtime entropy adapter owns generation. State code accepts only this
+/// bounded value and persists the exact resulting retry schedule.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
+pub struct MycDeliveryRetryJitter(u64);
+
+impl MycDeliveryRetryJitter {
+ /// Validates one injected full-jitter delay.
+ pub fn new(value_ms: u64) -> Result<Self, MycDeliveryStateError> {
+ if value_ms > MYC_DELIVERY_RETRY_JITTER_MAX_MS {
+ return Err(MycDeliveryStateError::new(
+ MycDeliveryStateErrorKind::InvalidRetryJitter,
+ ));
+ }
+ Ok(Self(value_ms))
+ }
+
+ /// Returns the exact injected delay in milliseconds.
+ #[must_use]
+ pub const fn get(self) -> u64 {
+ self.0
+ }
+}
+
macro_rules! redacted_id {
($name:ident, $debug:literal) => {
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
@@ -428,6 +457,7 @@ pub(crate) struct MycDeliveryPolicies {
initial_backoff_ms: u64,
maximum_backoff_ms: u64,
attempt_deadline_ms: u64,
+ outbox_maximum: usize,
targets: Box<[MycDeliveryTargetPolicy]>,
}
@@ -440,6 +470,7 @@ impl MycDeliveryPolicies {
initial_backoff_ms: u64,
maximum_backoff_ms: u64,
attempt_deadline_ms: u64,
+ outbox_maximum: usize,
mut targets: Vec<(MycDeliveryRelayId, bool)>,
) -> Result<Self, MycDeliveryStateError> {
targets.sort_by(|left, right| left.0.cmp(&right.0));
@@ -462,7 +493,8 @@ impl MycDeliveryPolicies {
&& initial_backoff_ms <= maximum_backoff_ms
&& maximum_backoff_ms <= 300_000
&& attempt_deadline_ms != 0
- && attempt_deadline_ms <= 30_000;
+ && attempt_deadline_ms <= 30_000
+ && (1..=65_536).contains(&outbox_maximum);
if !valid {
return Err(MycDeliveryStateError::new(
MycDeliveryStateErrorKind::InvalidPolicy,
@@ -475,12 +507,17 @@ impl MycDeliveryPolicies {
initial_backoff_ms,
maximum_backoff_ms,
attempt_deadline_ms,
+ outbox_maximum,
targets: targets
.into_iter()
.map(|(relay_id, required)| MycDeliveryTargetPolicy { relay_id, required })
.collect(),
})
}
+
+ pub(crate) const fn outbox_maximum(&self) -> usize {
+ self.outbox_maximum
+ }
}
impl fmt::Debug for MycDeliveryPolicies {
@@ -559,43 +596,6 @@ impl fmt::Debug for MycDeliverySource {
}
}
-/// Immutable signer-response delivery job input.
-pub struct MycDeliveryJobRequest {
- operation_id: MycSignerOperationId,
- artifact_digest: MycDeliveryArtifactDigest,
- created_at: MycDeliveryTimeUnixMs,
-}
-
-impl MycDeliveryJobRequest {
- /// Binds one admitted signer operation to one already-verified artifact digest.
- #[must_use]
- pub const fn signer_response(
- operation_id: MycSignerOperationId,
- artifact_digest: MycDeliveryArtifactDigest,
- created_at: MycDeliveryTimeUnixMs,
- ) -> Self {
- Self {
- operation_id,
- artifact_digest,
- created_at,
- }
- }
-
- fn owned(&self) -> Self {
- Self {
- operation_id: self.operation_id,
- artifact_digest: self.artifact_digest,
- created_at: self.created_at,
- }
- }
-}
-
-impl fmt::Debug for MycDeliveryJobRequest {
- fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
- formatter.write_str("MycDeliveryJobRequest([redacted])")
- }
-}
-
/// Durable job state.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum MycDeliveryJobStatus {
@@ -755,6 +755,10 @@ pub struct MycDeliveryJobRecord {
}
impl MycDeliveryJobRecord {
+ pub(crate) const fn source_id(&self) -> &[u8; 32] {
+ &self.source.id
+ }
+
#[must_use]
pub const fn id(&self) -> MycDeliveryJobId {
self.id
@@ -996,34 +1000,6 @@ impl fmt::Debug for MycDeliveryClaim {
}
impl MycStateRepository<'_> {
- /// Atomically creates or exactly replays one config-bound immutable job and target set.
- pub async fn create_delivery_job(
- &self,
- request: &MycDeliveryJobRequest,
- ) -> Result<MycDeliveryJobAdmission, MycStateRepositoryError> {
- let request = request.owned();
- let source = MycDeliverySource::signer_response(request.operation_id);
- let policy = self.expected().delivery_policies().clone();
- let expected = PersistedMetadata::from(self.expected());
- self.host()
- .transaction(move |transaction| {
- Box::pin(async move {
- verify_metadata(transaction, &expected).await?;
- require_signer_operation(transaction, request.operation_id).await?;
- create_job(
- transaction,
- source,
- request.artifact_digest,
- request.created_at,
- &policy,
- )
- .await
- })
- })
- .await
- .map_err(map_transaction_error)
- }
-
/// Claims one eligible target under a bounded expiring attempt lease.
pub async fn claim_delivery_target(
&self,
@@ -1073,6 +1049,7 @@ impl MycStateRepository<'_> {
relay_id: &MycDeliveryRelayId,
attempt_id: MycDeliveryAttemptId,
outcome: MycDeliveryAttemptOutcome,
+ retry_jitter: MycDeliveryRetryJitter,
observed_at: MycDeliveryTimeUnixMs,
) -> Result<MycDeliveryJobRecord, MycStateRepositoryError> {
let relay_id = relay_id.clone();
@@ -1087,6 +1064,7 @@ impl MycStateRepository<'_> {
&relay_id,
attempt_id,
outcome,
+ retry_jitter,
observed_at,
)
.await
@@ -1102,6 +1080,7 @@ impl MycStateRepository<'_> {
job_id: MycDeliveryJobId,
relay_id: &MycDeliveryRelayId,
attempt_id: MycDeliveryAttemptId,
+ retry_jitter: MycDeliveryRetryJitter,
observed_at: MycDeliveryTimeUnixMs,
) -> Result<MycDeliveryJobRecord, MycStateRepositoryError> {
let relay_id = relay_id.clone();
@@ -1110,7 +1089,15 @@ impl MycStateRepository<'_> {
.transaction(move |transaction| {
Box::pin(async move {
verify_metadata(transaction, &expected).await?;
- recover_expired(transaction, job_id, &relay_id, attempt_id, observed_at).await
+ recover_expired(
+ transaction,
+ job_id,
+ &relay_id,
+ attempt_id,
+ retry_jitter,
+ observed_at,
+ )
+ .await
})
})
.await
@@ -1188,6 +1175,18 @@ pub(crate) async fn create_job(
.then_some(MycDeliveryJobAdmission::ExactReplay(existing))
.ok_or(DeliveryOperationError::Binding);
}
+ let active_jobs = sqlx::query(READ_ACTIVE_JOB_COUNT_SQL)
+ .fetch_one(&mut *transaction)
+ .await
+ .map_err(|_| DeliveryOperationError::Storage)?
+ .try_get::<i64, _>("row_count")
+ .map_err(|_| DeliveryOperationError::Binding)?;
+ if usize::try_from(active_jobs)
+ .ok()
+ .is_none_or(|count| count >= policy.outbox_maximum)
+ {
+ return Err(DeliveryOperationError::Binding);
+ }
let job_id = derive_job_id(source, artifact_digest);
let result = sqlx::query(INSERT_JOB_SQL)
.bind(job_id.as_bytes().as_slice())
@@ -1371,6 +1370,7 @@ async fn record_outcome(
relay_id: &MycDeliveryRelayId,
attempt_id: MycDeliveryAttemptId,
outcome: MycDeliveryAttemptOutcome,
+ retry_jitter: MycDeliveryRetryJitter,
observed_at: MycDeliveryTimeUnixMs,
) -> Result<MycDeliveryJobRecord, DeliveryOperationError> {
let job = read_job(transaction, job_id)
@@ -1420,17 +1420,19 @@ async fn record_outcome(
required_prior,
outcome.status(),
outcome.reason(),
+ retry_jitter,
observed_at,
)
.await?;
finalize_job_if_terminal(transaction, job_id, observed_at).await
}
-async fn recover_expired(
+pub(crate) async fn recover_expired(
transaction: &mut ServiceSqliteTransaction<'_>,
job_id: MycDeliveryJobId,
relay_id: &MycDeliveryRelayId,
attempt_id: MycDeliveryAttemptId,
+ retry_jitter: MycDeliveryRetryJitter,
observed_at: MycDeliveryTimeUnixMs,
) -> Result<MycDeliveryJobRecord, DeliveryOperationError> {
let job = read_job(transaction, job_id)
@@ -1471,6 +1473,7 @@ async fn recover_expired(
attempt.status,
terminal,
reason,
+ retry_jitter,
observed_at,
)
.await?;
@@ -1486,6 +1489,7 @@ async fn resolve_attempt_and_target(
prior: MycDeliveryAttemptStatus,
terminal: MycDeliveryAttemptStatus,
reason: &'static str,
+ retry_jitter: MycDeliveryRetryJitter,
observed_at: MycDeliveryTimeUnixMs,
) -> Result<(), DeliveryOperationError> {
let result = sqlx::query(RESOLVE_ATTEMPT_SQL)
@@ -1501,16 +1505,29 @@ async fn resolve_attempt_and_target(
.map_err(|_| DeliveryOperationError::Storage)?;
require_one(result.rows_affected())?;
let attempts_remaining = target.attempt_count < job.max_attempts;
+ let schedules_retry = attempts_remaining
+ && matches!(
+ terminal,
+ MycDeliveryAttemptStatus::Failed | MycDeliveryAttemptStatus::Unknown
+ );
+ if !schedules_retry && retry_jitter.get() != 0 {
+ return Err(DeliveryOperationError::Binding);
+ }
let (target_status, next_attempt) = match terminal {
MycDeliveryAttemptStatus::Delivered => (MycDeliveryTargetStatus::Delivered, None),
MycDeliveryAttemptStatus::Failed if attempts_remaining => (
MycDeliveryTargetStatus::Retryable,
- Some(next_attempt_time(job, attempt.number, observed_at)?),
+ Some(next_attempt_time(
+ job,
+ attempt.number,
+ retry_jitter,
+ observed_at,
+ )?),
),
MycDeliveryAttemptStatus::Unknown => (
MycDeliveryTargetStatus::Unknown,
attempts_remaining
- .then(|| next_attempt_time(job, attempt.number, observed_at))
+ .then(|| next_attempt_time(job, attempt.number, retry_jitter, observed_at))
.transpose()?,
),
MycDeliveryAttemptStatus::Failed => (MycDeliveryTargetStatus::Exhausted, None),
@@ -1536,7 +1553,7 @@ async fn resolve_attempt_and_target(
require_one(result.rows_affected())
}
-async fn finalize_job_if_terminal(
+pub(crate) async fn finalize_job_if_terminal(
transaction: &mut ServiceSqliteTransaction<'_>,
job_id: MycDeliveryJobId,
observed_at: MycDeliveryTimeUnixMs,
@@ -1605,38 +1622,25 @@ async fn finalize_job_if_terminal(
fn next_attempt_time(
job: &MycDeliveryJobRecord,
attempt_number: u32,
+ jitter: MycDeliveryRetryJitter,
observed_at: MycDeliveryTimeUnixMs,
) -> Result<MycDeliveryTimeUnixMs, DeliveryOperationError> {
let exponent = attempt_number.saturating_sub(1).min(31);
let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX);
- let delay = job
+ let maximum_delay = job
.initial_backoff_ms
.saturating_mul(factor)
.min(job.maximum_backoff_ms);
+ if jitter.get() > maximum_delay {
+ return Err(DeliveryOperationError::Binding);
+ }
let value = observed_at
.get()
- .checked_add(delay)
+ .checked_add(jitter.get())
.ok_or(DeliveryOperationError::Binding)?;
MycDeliveryTimeUnixMs::new(value).map_err(|_| DeliveryOperationError::Binding)
}
-async fn require_signer_operation(
- transaction: &mut ServiceSqliteTransaction<'_>,
- operation_id: MycSignerOperationId,
-) -> Result<(), DeliveryOperationError> {
- let row = sqlx::query(READ_SIGNER_OPERATION_SQL)
- .bind(operation_id.as_bytes().as_slice())
- .fetch_one(&mut *transaction)
- .await
- .map_err(|_| DeliveryOperationError::Storage)?;
- let count = row
- .try_get::<i64, _>("row_count")
- .map_err(|_| DeliveryOperationError::Binding)?;
- (count == 1)
- .then_some(())
- .ok_or(DeliveryOperationError::Binding)
-}
-
async fn read_job_by_source(
transaction: &mut ServiceSqliteTransaction<'_>,
source: MycDeliverySource,
@@ -1803,7 +1807,7 @@ fn parse_target(
})
}
-async fn read_attempts(
+pub(crate) async fn read_attempts(
transaction: &mut ServiceSqliteTransaction<'_>,
job_id: MycDeliveryJobId,
target_index: u32,
diff --git a/src/state_discovery.rs b/src/state_discovery.rs
@@ -1,12 +1,15 @@
//! Durable discovery desired/current state and exact committed publication bytes.
use core::fmt;
-use std::{collections::BTreeMap, error::Error};
+use std::{
+ collections::{BTreeMap, BTreeSet},
+ error::Error,
+};
use radroots_service_sqlite::{
ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
};
-use serde::Serialize;
+use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use sqlx::Row;
@@ -28,6 +31,8 @@ use radroots_nostr::event::{
pub const MYC_DISCOVERY_DOCUMENT_MAX_BYTES: usize = 524_288;
/// Maximum deterministic NIP-05 projection input bytes.
pub const MYC_NIP05_PROJECTION_MAX_BYTES: usize = 524_288;
+/// Maximum deterministic offline NIP-05 document bytes.
+pub const MYC_NIP05_DOCUMENT_MAX_BYTES: usize = 524_288;
const NIP46_RPC_KIND: u32 = 24_133;
const NIP89_HANDLER_KIND: u16 = 31_990;
@@ -223,6 +228,60 @@ redacted_id!(
MycNip05ProjectionDigest,
"MycNip05ProjectionDigest([redacted])"
);
+redacted_id!(MycNip05DocumentDigest, "MycNip05DocumentDigest([redacted])");
+
+/// Explicit discovery generation selected for an offline NIP-05 export.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum MycNip05ExportSelection {
+ Desired,
+ Current,
+}
+
+/// Canonical offline NIP-05/NIP-46 document derived from verified state.
+#[derive(Clone, PartialEq, Eq)]
+pub struct MycNip05Document {
+ selection: MycNip05ExportSelection,
+ domain: Box<str>,
+ digest: MycNip05DocumentDigest,
+ bytes: Box<[u8]>,
+}
+
+impl MycNip05Document {
+ #[must_use]
+ pub const fn selection(&self) -> MycNip05ExportSelection {
+ self.selection
+ }
+
+ /// Returns the configured domain associated with this explicit export.
+ #[must_use]
+ pub fn domain(&self) -> &str {
+ &self.domain
+ }
+
+ /// Returns the exact compact UTF-8 JSON bytes to export.
+ #[must_use]
+ pub fn bytes(&self) -> &[u8] {
+ &self.bytes
+ }
+
+ /// Returns the SHA-256 identity of the exact export bytes.
+ #[must_use]
+ pub const fn digest(&self) -> MycNip05DocumentDigest {
+ self.digest
+ }
+}
+
+impl fmt::Debug for MycNip05Document {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("MycNip05Document")
+ .field("selection", &self.selection)
+ .field("domain", &"[redacted]")
+ .field("bytes", &"[redacted]")
+ .field("digest", &"[redacted]")
+ .finish()
+ }
+}
#[derive(Clone, PartialEq, Eq)]
pub(crate) struct MycDiscoveryPolicies {
@@ -713,6 +772,40 @@ impl MycStateRepository<'_> {
.await
.map_err(map_transaction_error)
}
+
+ /// Renders an explicit desired or proven-current NIP-05 document offline.
+ ///
+ /// This performs only verified state reads. It does not host, write, or
+ /// publish the returned bytes and is valid through an inspection host.
+ pub async fn render_offline_nip05(
+ &self,
+ selection: MycNip05ExportSelection,
+ ) -> Result<MycNip05Document, MycStateRepositoryError> {
+ let expected = PersistedMetadata::from(self.expected());
+ let policy = self.expected().discovery_policies().cloned();
+ self.host()
+ .transaction(move |transaction| {
+ Box::pin(async move {
+ verify_metadata(transaction, &expected).await?;
+ let policy = policy.as_ref().ok_or(DiscoveryOperationError::Binding)?;
+ let state = read_state(transaction)
+ .await?
+ .ok_or(DiscoveryOperationError::Binding)?;
+ let generation = match selection {
+ MycNip05ExportSelection::Desired => state.desired_generation_id,
+ MycNip05ExportSelection::Current => state
+ .current_generation_id
+ .ok_or(DiscoveryOperationError::Binding)?,
+ };
+ let document = read_document(transaction, generation, policy)
+ .await?
+ .ok_or(DiscoveryOperationError::Binding)?;
+ render_nip05(selection, &document)
+ })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
@@ -926,6 +1019,40 @@ async fn read_document(
.transpose()
}
+pub(crate) async fn verify_document_for_delivery_job(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ job: &MycDeliveryJobRecord,
+ policy: Option<&MycDiscoveryPolicies>,
+) -> Result<(), DeliveryOperationError> {
+ if job.source_kind() != crate::state_delivery::MycDeliverySourceKind::DiscoveryHandler {
+ return Err(DeliveryOperationError::Binding);
+ }
+ let generation = discovery_generation_for_job(transaction, job.id())
+ .await
+ .map_err(|error| match error {
+ DiscoveryOperationError::Binding => DeliveryOperationError::Binding,
+ DiscoveryOperationError::Storage => DeliveryOperationError::Storage,
+ })?
+ .ok_or(DeliveryOperationError::Binding)?;
+ if generation.as_bytes() != job.source_id() {
+ return Err(DeliveryOperationError::Binding);
+ }
+ let document = read_document(
+ transaction,
+ generation,
+ policy.ok_or(DeliveryOperationError::Binding)?,
+ )
+ .await
+ .map_err(|error| match error {
+ DiscoveryOperationError::Binding => DeliveryOperationError::Binding,
+ DiscoveryOperationError::Storage => DeliveryOperationError::Storage,
+ })?
+ .ok_or(DeliveryOperationError::Binding)?;
+ (document.event_digest() == job.artifact_digest())
+ .then_some(())
+ .ok_or(DeliveryOperationError::Binding)
+}
+
fn parse_document(
row: &sqlx::sqlite::SqliteRow,
policy: &MycDiscoveryPolicies,
@@ -1065,6 +1192,30 @@ async fn promote_current(
.ok_or(DiscoveryOperationError::Binding)
}
+pub(crate) async fn promote_current_if_desired(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ job_id: MycDeliveryJobId,
+ observed_at: MycDeliveryTimeUnixMs,
+) -> Result<bool, DeliveryOperationError> {
+ let Some(state) = read_state(transaction).await.map_err(|error| match error {
+ DiscoveryOperationError::Binding => DeliveryOperationError::Binding,
+ DiscoveryOperationError::Storage => DeliveryOperationError::Storage,
+ })?
+ else {
+ return Err(DeliveryOperationError::Binding);
+ };
+ if state.desired_job_id != job_id {
+ return Ok(false);
+ }
+ promote_current(transaction, job_id, observed_at)
+ .await
+ .map_err(|error| match error {
+ DiscoveryOperationError::Binding => DeliveryOperationError::Binding,
+ DiscoveryOperationError::Storage => DeliveryOperationError::Storage,
+ })?;
+ Ok(true)
+}
+
fn exact_document(
record: &MycDiscoveryDocumentRecord,
request: &MycDiscoveryCommitRequest,
@@ -1129,6 +1280,104 @@ struct Nip05Projection<'a> {
nostrconnect_url: Option<&'a str>,
}
+#[derive(Deserialize, Serialize)]
+#[serde(deny_unknown_fields)]
+struct Nip05ProjectionInput {
+ schema: Box<str>,
+ schema_version: u32,
+ domain: Box<str>,
+ name: Box<str>,
+ public_key: Box<str>,
+ relays: Box<[Box<str>]>,
+ #[serde(skip_serializing_if = "Option::is_none")]
+ nostrconnect_url: Option<Box<str>>,
+}
+
+#[derive(Serialize)]
+struct Nip05Names<'a> {
+ #[serde(rename = "_")]
+ root: &'a str,
+}
+
+#[derive(Serialize)]
+struct Nip46Discovery<'a> {
+ relays: &'a [Box<str>],
+ #[serde(skip_serializing_if = "Option::is_none")]
+ nostrconnect_url: Option<&'a str>,
+}
+
+#[derive(Serialize)]
+struct Nip05Output<'a> {
+ names: Nip05Names<'a>,
+ nip46: Nip46Discovery<'a>,
+}
+
+fn render_nip05(
+ selection: MycNip05ExportSelection,
+ document: &MycDiscoveryDocumentRecord,
+) -> Result<MycNip05Document, DiscoveryOperationError> {
+ let input: Nip05ProjectionInput = serde_json::from_slice(document.nip05_projection_bytes())
+ .map_err(|_| DiscoveryOperationError::Binding)?;
+ let canonical = serde_json::to_vec(&input).map_err(|_| DiscoveryOperationError::Binding)?;
+ let valid_public_key = input.public_key.len() == 64
+ && input
+ .public_key
+ .as_bytes()
+ .iter()
+ .all(u8::is_ascii_hexdigit)
+ && !input
+ .public_key
+ .as_bytes()
+ .iter()
+ .any(u8::is_ascii_uppercase);
+ let unique_relays = input
+ .relays
+ .iter()
+ .map(Box::as_ref)
+ .collect::<BTreeSet<_>>();
+ let valid_relays = (1..=32).contains(&input.relays.len())
+ && input
+ .relays
+ .iter()
+ .all(|relay| RadrootsNostrRelayUrl::parse(relay).is_ok())
+ && unique_relays.len() == input.relays.len();
+ let valid_url = input
+ .nostrconnect_url
+ .as_deref()
+ .is_none_or(|url| nostr::Url::parse(url).is_ok());
+ if canonical.as_slice() != document.nip05_projection_bytes()
+ || input.schema.as_ref() != "radroots.myc.nip05-projection-input.v1"
+ || input.schema_version != 1
+ || input.domain.is_empty()
+ || input.domain.len() > 253
+ || input.name.as_ref() != "_"
+ || !valid_public_key
+ || !valid_relays
+ || !valid_url
+ {
+ return Err(DiscoveryOperationError::Binding);
+ }
+ let bytes = serde_json::to_vec(&Nip05Output {
+ names: Nip05Names {
+ root: &input.public_key,
+ },
+ nip46: Nip46Discovery {
+ relays: &input.relays,
+ nostrconnect_url: input.nostrconnect_url.as_deref(),
+ },
+ })
+ .map_err(|_| DiscoveryOperationError::Binding)?;
+ if bytes.is_empty() || bytes.len() > MYC_NIP05_DOCUMENT_MAX_BYTES {
+ return Err(DiscoveryOperationError::Binding);
+ }
+ Ok(MycNip05Document {
+ selection,
+ domain: input.domain,
+ digest: MycNip05DocumentDigest(Sha256::digest(&bytes).into()),
+ bytes: bytes.into_boxed_slice(),
+ })
+}
+
fn projection_bytes(policy: &MycDiscoveryPolicies) -> Result<Vec<u8>, MycDiscoveryStateError> {
serde_json::to_vec(&Nip05Projection {
schema: "radroots.myc.nip05-projection-input.v1",
diff --git a/src/state_metadata.rs b/src/state_metadata.rs
@@ -339,6 +339,10 @@ impl MycStateMetadata {
&self.delivery
}
+ pub(crate) const fn outbox_maximum(&self) -> usize {
+ self.delivery.outbox_maximum()
+ }
+
pub(crate) const fn discovery_policies(&self) -> Option<&MycDiscoveryPolicies> {
self.discovery.as_ref()
}
@@ -370,6 +374,7 @@ impl fmt::Debug for MycStateMetadata {
.field("identities", &self.identities)
.field("governance", &"[redacted]")
.field("delivery", &"[redacted]")
+ .field("outbox_maximum", &self.delivery.outbox_maximum())
.field("discovery", &self.discovery.as_ref().map(|_| "[redacted]"))
.field("policy_versions", &self.policy_versions)
.field("paths", &"[redacted]")
@@ -571,6 +576,7 @@ fn delivery_policies(
integer("/transport/publish_retry/initial_backoff_ms")?,
integer("/transport/publish_retry/maximum_backoff_ms")?,
integer("/transport/publish_retry/attempt_deadline_ms")?,
+ usize::try_from(integer("/resource_limits/queues/outbox")?).map_err(|_| invalid())?,
targets,
)
.map_err(|_| invalid())
diff --git a/src/state_recovery.rs b/src/state_recovery.rs
@@ -0,0 +1,565 @@
+//! Bounded restart recovery for committed exact-byte delivery work.
+
+use core::fmt;
+
+use radroots_service_sqlite::{
+ ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
+};
+use sha2::{Digest, Sha256};
+use sqlx::Row;
+
+use crate::state_delivery::{
+ DeliveryOperationError, MycDeliveryAttemptId, MycDeliveryAttemptStatus, MycDeliveryJobId,
+ MycDeliveryJobRecord, MycDeliveryJobStatus, MycDeliveryRetryJitter, MycDeliverySourceKind,
+ MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, finalize_job_if_terminal, read_attempts,
+ read_job, recover_expired,
+};
+use crate::state_discovery::{
+ MycDiscoveryPolicies, promote_current_if_desired, verify_document_for_delivery_job,
+};
+use crate::state_repository::{
+ MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata,
+ RepositoryOperationError, require_expected_metadata,
+};
+use crate::state_response::verify_response_for_delivery_job;
+
+/// Maximum delivery jobs examined in one restart-recovery transaction.
+pub const MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT: usize = 128;
+
+const RECOVERY_JITTER_DOMAIN: &[u8] = b"radroots.myc.delivery_recovery_jitter.v1\0";
+
+const READ_INVARIANTS_SQL: &str = r#"SELECT
+ (SELECT COUNT(*) FROM delivery_jobs j
+ WHERE (j.source_kind = 'signer_response' AND NOT EXISTS (
+ SELECT 1 FROM nip46_signed_responses r
+ WHERE r.operation_id = j.source_id
+ AND r.response_sha256 = j.artifact_sha256
+ ))
+ OR (j.source_kind = 'discovery_handler' AND NOT EXISTS (
+ SELECT 1 FROM discovery_documents d
+ WHERE d.generation_id = j.source_id
+ AND d.event_sha256 = j.artifact_sha256
+ ))) AS invalid_sources,
+ (SELECT COUNT(*) FROM nip46_signed_responses r
+ WHERE NOT EXISTS (
+ SELECT 1 FROM delivery_jobs j
+ WHERE j.source_kind = 'signer_response' AND j.source_id = r.operation_id
+ )) AS orphan_responses,
+ (SELECT COUNT(*) FROM discovery_documents d
+ WHERE NOT EXISTS (
+ SELECT 1 FROM delivery_jobs j
+ WHERE j.source_kind = 'discovery_handler' AND j.source_id = d.generation_id
+ )) AS orphan_documents,
+ (SELECT COUNT(*) FROM delivery_targets t
+ WHERE NOT EXISTS (SELECT 1 FROM delivery_jobs j WHERE j.job_id = t.job_id))
+ AS orphan_targets,
+ (SELECT COUNT(*) FROM delivery_attempts a
+ WHERE NOT EXISTS (
+ SELECT 1 FROM delivery_targets t
+ WHERE t.job_id = a.job_id AND t.target_index = a.target_index
+ )) AS orphan_attempts,
+ (SELECT COUNT(*) FROM delivery_jobs WHERE status IN ('pending', 'active')) AS active_jobs"#;
+
+const READ_FIRST_JOB_IDS_SQL: &str = r#"SELECT
+ CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32
+ THEN job_id ELSE NULL END AS job_id
+FROM delivery_jobs
+WHERE status IN ('pending', 'active')
+ OR (status = 'delivered' AND source_kind = 'discovery_handler' AND EXISTS (
+ SELECT 1 FROM discovery_publication_state s
+ WHERE s.singleton = 1
+ AND s.desired_job_id = delivery_jobs.job_id
+ AND (s.current_job_id IS NULL OR s.current_job_id != s.desired_job_id)
+ ))
+ORDER BY job_id
+LIMIT 129"#;
+
+const READ_NEXT_JOB_IDS_SQL: &str = r#"SELECT
+ CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32
+ THEN job_id ELSE NULL END AS job_id
+FROM delivery_jobs
+WHERE (status IN ('pending', 'active')
+ OR (status = 'delivered' AND source_kind = 'discovery_handler' AND EXISTS (
+ SELECT 1 FROM discovery_publication_state s
+ WHERE s.singleton = 1
+ AND s.desired_job_id = delivery_jobs.job_id
+ AND (s.current_job_id IS NULL OR s.current_job_id != s.desired_job_id)
+ ))) AND job_id > ?
+ORDER BY job_id
+LIMIT 129"#;
+
+/// One-use entropy supplied by the runtime recovery boundary.
+pub struct MycDeliveryRecoveryEntropy([u8; 32]);
+
+impl MycDeliveryRecoveryEntropy {
+ /// Wraps exact entropy from the caller's injected source.
+ #[must_use]
+ pub const fn from_injected_entropy(bytes: [u8; 32]) -> Self {
+ Self(bytes)
+ }
+}
+
+impl fmt::Debug for MycDeliveryRecoveryEntropy {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("MycDeliveryRecoveryEntropy([redacted])")
+ }
+}
+
+#[derive(Clone, Copy, PartialEq, Eq)]
+struct RecoveryCursor([u8; 32]);
+
+/// Bounded, identity-free result of one restart-recovery transaction.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct MycDeliveryRecoveryReport {
+ examined_jobs: u32,
+ recovered_expired_attempts: u32,
+ finalized_jobs: u32,
+ promoted_discovery_generations: u32,
+ ready_targets: u32,
+ scheduled_targets: u32,
+ active_targets: u32,
+ continuation: Option<RecoveryCursor>,
+}
+
+impl MycDeliveryRecoveryReport {
+ #[must_use]
+ pub const fn examined_jobs(self) -> u32 {
+ self.examined_jobs
+ }
+
+ #[must_use]
+ pub const fn recovered_expired_attempts(self) -> u32 {
+ self.recovered_expired_attempts
+ }
+
+ #[must_use]
+ pub const fn finalized_jobs(self) -> u32 {
+ self.finalized_jobs
+ }
+
+ #[must_use]
+ pub const fn promoted_discovery_generations(self) -> u32 {
+ self.promoted_discovery_generations
+ }
+
+ #[must_use]
+ pub const fn ready_targets(self) -> u32 {
+ self.ready_targets
+ }
+
+ #[must_use]
+ pub const fn scheduled_targets(self) -> u32 {
+ self.scheduled_targets
+ }
+
+ #[must_use]
+ pub const fn active_targets(self) -> u32 {
+ self.active_targets
+ }
+
+ fn merge(&mut self, batch: Self) -> Result<(), RecoveryOperationError> {
+ self.examined_jobs = add(self.examined_jobs, batch.examined_jobs)?;
+ self.recovered_expired_attempts = add(
+ self.recovered_expired_attempts,
+ batch.recovered_expired_attempts,
+ )?;
+ self.finalized_jobs = add(self.finalized_jobs, batch.finalized_jobs)?;
+ self.promoted_discovery_generations = add(
+ self.promoted_discovery_generations,
+ batch.promoted_discovery_generations,
+ )?;
+ self.ready_targets = add(self.ready_targets, batch.ready_targets)?;
+ self.scheduled_targets = add(self.scheduled_targets, batch.scheduled_targets)?;
+ self.active_targets = add(self.active_targets, batch.active_targets)?;
+ self.continuation = batch.continuation;
+ Ok(())
+ }
+
+ const fn empty() -> Self {
+ Self {
+ examined_jobs: 0,
+ recovered_expired_attempts: 0,
+ finalized_jobs: 0,
+ promoted_discovery_generations: 0,
+ ready_targets: 0,
+ scheduled_targets: 0,
+ active_targets: 0,
+ continuation: None,
+ }
+ }
+}
+
+impl fmt::Debug for MycDeliveryRecoveryReport {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("MycDeliveryRecoveryReport")
+ .field("examined_jobs", &self.examined_jobs)
+ .field(
+ "recovered_expired_attempts",
+ &self.recovered_expired_attempts,
+ )
+ .field("finalized_jobs", &self.finalized_jobs)
+ .field(
+ "promoted_discovery_generations",
+ &self.promoted_discovery_generations,
+ )
+ .field("ready_targets", &self.ready_targets)
+ .field("scheduled_targets", &self.scheduled_targets)
+ .field("active_targets", &self.active_targets)
+ .field("has_continuation", &self.continuation.is_some())
+ .finish()
+ }
+}
+
+impl MycStateRepository<'_> {
+ /// Recovers startup delivery state in fixed bounded transactions without relay I/O.
+ ///
+ /// The later runtime owner must pause admission while invoking this method.
+ /// Pagination remains sealed inside the repository, time and entropy are
+ /// injected, and exact payload bytes and target membership never change.
+ pub async fn recover_delivery_state(
+ &self,
+ observed_at: MycDeliveryTimeUnixMs,
+ entropy: MycDeliveryRecoveryEntropy,
+ ) -> Result<MycDeliveryRecoveryReport, MycStateRepositoryError> {
+ let outbox_maximum = self.expected().outbox_maximum();
+ let discovery_policy = self.expected().discovery_policies().cloned();
+ let mut cursor = None;
+ let mut aggregate = MycDeliveryRecoveryReport::empty();
+ loop {
+ let expected = PersistedMetadata::from(self.expected());
+ let discovery_policy = discovery_policy.clone();
+ let entropy_bytes = entropy.0;
+ let batch = self
+ .host()
+ .transaction(move |transaction| {
+ Box::pin(async move {
+ require_expected_metadata(transaction, &expected)
+ .await
+ .map_err(RecoveryOperationError::from)?;
+ recover_batch(
+ transaction,
+ observed_at,
+ &entropy_bytes,
+ cursor,
+ outbox_maximum,
+ discovery_policy.as_ref(),
+ )
+ .await
+ })
+ })
+ .await
+ .map_err(map_transaction_error)?;
+ let next = batch.continuation;
+ aggregate
+ .merge(batch)
+ .map_err(|_| MycStateRepositoryError::new(MycStateRepositoryErrorKind::Binding))?;
+ let Some(next) = next else {
+ aggregate.continuation = None;
+ return Ok(aggregate);
+ };
+ cursor = Some(next);
+ }
+ }
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum RecoveryOperationError {
+ Binding,
+ Storage,
+}
+
+impl From<RepositoryOperationError> for RecoveryOperationError {
+ fn from(error: RepositoryOperationError) -> Self {
+ match error {
+ RepositoryOperationError::Binding => Self::Binding,
+ RepositoryOperationError::Storage => Self::Storage,
+ }
+ }
+}
+
+impl From<DeliveryOperationError> for RecoveryOperationError {
+ fn from(error: DeliveryOperationError) -> Self {
+ match error {
+ DeliveryOperationError::Binding => Self::Binding,
+ DeliveryOperationError::Storage => Self::Storage,
+ }
+ }
+}
+
+async fn recover_batch(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ observed_at: MycDeliveryTimeUnixMs,
+ entropy: &[u8; 32],
+ cursor: Option<RecoveryCursor>,
+ outbox_maximum: usize,
+ discovery_policy: Option<&MycDiscoveryPolicies>,
+) -> Result<MycDeliveryRecoveryReport, RecoveryOperationError> {
+ verify_global_invariants(transaction, outbox_maximum).await?;
+ let mut job_ids = read_job_ids(transaction, cursor).await?;
+ let has_more = job_ids.len() > MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT;
+ if has_more {
+ job_ids.truncate(MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT);
+ }
+ let continuation = has_more
+ .then(|| job_ids.last().copied().map(RecoveryCursor))
+ .flatten();
+ let mut report = MycDeliveryRecoveryReport {
+ examined_jobs: 0,
+ recovered_expired_attempts: 0,
+ finalized_jobs: 0,
+ promoted_discovery_generations: 0,
+ ready_targets: 0,
+ scheduled_targets: 0,
+ active_targets: 0,
+ continuation,
+ };
+ for job_id in job_ids {
+ let before = read_job(transaction, MycDeliveryJobId::from_persisted(job_id))
+ .await?
+ .ok_or(RecoveryOperationError::Binding)?;
+ match before.source_kind() {
+ MycDeliverySourceKind::SignerResponse => {
+ verify_response_for_delivery_job(transaction, before.id()).await?;
+ }
+ MycDeliverySourceKind::DiscoveryHandler => {
+ verify_document_for_delivery_job(transaction, &before, discovery_policy).await?;
+ }
+ }
+ validate_attempt_histories(transaction, &before).await?;
+ let mut current = before.clone();
+ for target in before.targets() {
+ let Some(attempt_id) = target.active_attempt_id() else {
+ continue;
+ };
+ let attempts = read_attempts(transaction, before.id(), target.index()).await?;
+ let attempt = attempts
+ .last()
+ .filter(|attempt| attempt.id() == attempt_id)
+ .ok_or(RecoveryOperationError::Binding)?;
+ if observed_at > attempt.lease_expires_at() {
+ let jitter = recovery_jitter(¤t, attempt.id(), attempt.number(), entropy)?;
+ current = recover_expired(
+ transaction,
+ before.id(),
+ target.relay_id(),
+ attempt.id(),
+ jitter,
+ observed_at,
+ )
+ .await?;
+ report.recovered_expired_attempts = increment(report.recovered_expired_attempts)?;
+ }
+ }
+ current = finalize_job_if_terminal(transaction, before.id(), observed_at).await?;
+ if !is_terminal(before.status()) && is_terminal(current.status()) {
+ report.finalized_jobs = increment(report.finalized_jobs)?;
+ }
+ if current.status() == MycDeliveryJobStatus::Delivered
+ && current.source_kind() == MycDeliverySourceKind::DiscoveryHandler
+ && promote_current_if_desired(transaction, current.id(), observed_at).await?
+ {
+ report.promoted_discovery_generations =
+ increment(report.promoted_discovery_generations)?;
+ }
+ let verified = read_job(transaction, current.id())
+ .await?
+ .ok_or(RecoveryOperationError::Binding)?;
+ validate_attempt_histories(transaction, &verified).await?;
+ classify_targets(&verified, observed_at, &mut report)?;
+ report.examined_jobs = increment(report.examined_jobs)?;
+ }
+ Ok(report)
+}
+
+async fn verify_global_invariants(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ outbox_maximum: usize,
+) -> Result<(), RecoveryOperationError> {
+ let rows = sqlx::query(READ_INVARIANTS_SQL)
+ .fetch_all(&mut *transaction)
+ .await
+ .map_err(|_| RecoveryOperationError::Storage)?;
+ if rows.len() != 1 {
+ return Err(RecoveryOperationError::Binding);
+ }
+ let row = &rows[0];
+ for column in [
+ "invalid_sources",
+ "orphan_responses",
+ "orphan_documents",
+ "orphan_targets",
+ "orphan_attempts",
+ ] {
+ if row
+ .try_get::<i64, _>(column)
+ .map_err(|_| RecoveryOperationError::Binding)?
+ != 0
+ {
+ return Err(RecoveryOperationError::Binding);
+ }
+ }
+ let active_jobs = row
+ .try_get::<i64, _>("active_jobs")
+ .map_err(|_| RecoveryOperationError::Binding)?;
+ usize::try_from(active_jobs)
+ .ok()
+ .filter(|count| *count <= outbox_maximum)
+ .map(|_| ())
+ .ok_or(RecoveryOperationError::Binding)
+}
+
+async fn read_job_ids(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ cursor: Option<RecoveryCursor>,
+) -> Result<Vec<[u8; 32]>, RecoveryOperationError> {
+ let query = match cursor {
+ Some(cursor) => sqlx::query(READ_NEXT_JOB_IDS_SQL).bind(cursor.0.as_slice()),
+ None => sqlx::query(READ_FIRST_JOB_IDS_SQL),
+ };
+ let rows = query
+ .fetch_all(&mut *transaction)
+ .await
+ .map_err(|_| RecoveryOperationError::Storage)?;
+ if rows.len() > MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT + 1 {
+ return Err(RecoveryOperationError::Binding);
+ }
+ rows.iter()
+ .map(|row| {
+ let bytes = row
+ .try_get::<Option<Vec<u8>>, _>("job_id")
+ .map_err(|_| RecoveryOperationError::Binding)?
+ .ok_or(RecoveryOperationError::Binding)?;
+ bytes
+ .try_into()
+ .map_err(|_| RecoveryOperationError::Binding)
+ })
+ .collect()
+}
+
+async fn validate_attempt_histories(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ job: &MycDeliveryJobRecord,
+) -> Result<(), RecoveryOperationError> {
+ for target in job.targets() {
+ let attempts = read_attempts(transaction, job.id(), target.index()).await?;
+ if usize::try_from(target.attempt_count()) != Ok(attempts.len())
+ || attempts
+ .iter()
+ .enumerate()
+ .any(|(index, attempt)| usize::try_from(attempt.number()) != Ok(index + 1))
+ {
+ return Err(RecoveryOperationError::Binding);
+ }
+ let active = matches!(
+ target.status(),
+ MycDeliveryTargetStatus::Leased | MycDeliveryTargetStatus::Submitted
+ );
+ if active {
+ let attempt = attempts.last().ok_or(RecoveryOperationError::Binding)?;
+ let expected_status = match target.status() {
+ MycDeliveryTargetStatus::Leased => MycDeliveryAttemptStatus::Leased,
+ MycDeliveryTargetStatus::Submitted => MycDeliveryAttemptStatus::Submitted,
+ _ => return Err(RecoveryOperationError::Binding),
+ };
+ if target.active_attempt_id() != Some(attempt.id())
+ || attempt.status() != expected_status
+ {
+ return Err(RecoveryOperationError::Binding);
+ }
+ } else if attempts.last().is_some_and(|attempt| {
+ matches!(
+ attempt.status(),
+ MycDeliveryAttemptStatus::Leased | MycDeliveryAttemptStatus::Submitted
+ )
+ }) {
+ return Err(RecoveryOperationError::Binding);
+ }
+ }
+ Ok(())
+}
+
+fn recovery_jitter(
+ job: &MycDeliveryJobRecord,
+ attempt_id: MycDeliveryAttemptId,
+ attempt_number: u32,
+ entropy: &[u8; 32],
+) -> Result<MycDeliveryRetryJitter, RecoveryOperationError> {
+ let exponent = attempt_number.saturating_sub(1).min(31);
+ let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX);
+ let maximum = job
+ .initial_backoff_ms()
+ .saturating_mul(factor)
+ .min(job.maximum_backoff_ms());
+ let mut hasher = Sha256::new();
+ hasher.update(RECOVERY_JITTER_DOMAIN);
+ hasher.update(entropy);
+ hasher.update(attempt_id.as_bytes());
+ let digest: [u8; 32] = hasher.finalize().into();
+ let raw = u64::from_be_bytes(digest[..8].try_into().expect("fixed SHA-256 prefix"));
+ let value = raw
+ % maximum
+ .checked_add(1)
+ .ok_or(RecoveryOperationError::Binding)?;
+ MycDeliveryRetryJitter::new(value).map_err(|_| RecoveryOperationError::Binding)
+}
+
+fn classify_targets(
+ job: &MycDeliveryJobRecord,
+ observed_at: MycDeliveryTimeUnixMs,
+ report: &mut MycDeliveryRecoveryReport,
+) -> Result<(), RecoveryOperationError> {
+ for target in job.targets() {
+ match target.status() {
+ MycDeliveryTargetStatus::Pending => {
+ report.ready_targets = increment(report.ready_targets)?;
+ }
+ MycDeliveryTargetStatus::Retryable | MycDeliveryTargetStatus::Unknown => {
+ if target
+ .next_attempt_at()
+ .is_some_and(|next| next > observed_at)
+ {
+ report.scheduled_targets = increment(report.scheduled_targets)?;
+ } else if target.next_attempt_at().is_some() {
+ report.ready_targets = increment(report.ready_targets)?;
+ }
+ }
+ MycDeliveryTargetStatus::Leased | MycDeliveryTargetStatus::Submitted => {
+ report.active_targets = increment(report.active_targets)?;
+ }
+ MycDeliveryTargetStatus::Delivered | MycDeliveryTargetStatus::Exhausted => {}
+ }
+ }
+ Ok(())
+}
+
+const fn is_terminal(status: MycDeliveryJobStatus) -> bool {
+ matches!(
+ status,
+ MycDeliveryJobStatus::Delivered
+ | MycDeliveryJobStatus::Failed
+ | MycDeliveryJobStatus::Unknown
+ )
+}
+
+fn increment(value: u32) -> Result<u32, RecoveryOperationError> {
+ value.checked_add(1).ok_or(RecoveryOperationError::Binding)
+}
+
+fn add(left: u32, right: u32) -> Result<u32, RecoveryOperationError> {
+ left.checked_add(right)
+ .ok_or(RecoveryOperationError::Binding)
+}
+
+fn map_transaction_error(
+ error: ServiceSqliteTransactionError<RecoveryOperationError>,
+) -> MycStateRepositoryError {
+ if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
+ return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown);
+ }
+ let kind = match error.operation_error() {
+ Some(RecoveryOperationError::Binding) => MycStateRepositoryErrorKind::Binding,
+ Some(RecoveryOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction,
+ };
+ MycStateRepositoryError::new(kind)
+}
diff --git a/src/state_response.rs b/src/state_response.rs
@@ -544,6 +544,20 @@ async fn read_response_by_job(
read_response(transaction, READ_RESPONSE_BY_JOB_SQL, job_id.as_bytes()).await
}
+pub(crate) async fn verify_response_for_delivery_job(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ job_id: MycDeliveryJobId,
+) -> Result<(), DeliveryOperationError> {
+ read_response_by_job(transaction, job_id)
+ .await
+ .map_err(|error| match error {
+ AtomicOperationError::Binding => DeliveryOperationError::Binding,
+ AtomicOperationError::Storage => DeliveryOperationError::Storage,
+ })?
+ .map(|_| ())
+ .ok_or(DeliveryOperationError::Binding)
+}
+
async fn read_response(
transaction: &mut ServiceSqliteTransaction<'_>,
sql: &'static str,
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -12,6 +12,8 @@ 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 DELIVERY_RECOVERY: &str = include_str!("../src/state_recovery.rs");
+const DISCOVERY_STATE: &str = include_str!("../src/state_discovery.rs");
const NIP46_VERIFICATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/nip46_verification.v1.json");
const NIP46_REPLAY_CONTRACT: &str =
@@ -26,6 +28,8 @@ 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 DELIVERY_RECOVERY_EXPORT_CONTRACT: &str =
+ include_str!("../contracts/services_hardening/delivery_recovery_export.v1.json");
const SOURCES: &[&str] = &[
include_str!("../src/cli_v1.rs"),
include_str!("../src/config_v1.rs"),
@@ -51,6 +55,7 @@ const SOURCES: &[&str] = &[
include_str!("../src/state_maintenance.rs"),
include_str!("../src/state_metadata.rs"),
include_str!("../src/state_repository.rs"),
+ include_str!("../src/state_recovery.rs"),
include_str!("../src/state_request.rs"),
include_str!("../src/state_response.rs"),
];
@@ -89,6 +94,7 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() {
"state_maintenance",
"state_metadata",
"state_repository",
+ "state_recovery",
"state_request",
"state_response",
])
@@ -148,6 +154,12 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() {
"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<'_>::recover_delivery_state",
+ "pub async fn myc::MycStateRepository<'_>::render_offline_nip05",
+ "pub struct myc::MycDeliveryRecoveryEntropy",
+ "pub struct myc::MycDeliveryRecoveryReport",
+ "pub struct myc::MycNip05Document",
+ "pub enum myc::MycNip05ExportSelection",
"pub async fn myc::MycStateRepository<'_>::read_connection_decision",
"pub fn myc::admit_myc_nip46_event",
"pub fn myc::admit_myc_nip46_request",
@@ -192,6 +204,7 @@ fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() {
"state_maintenance",
"state_metadata",
"state_repository",
+ "state_recovery",
"state_request",
"state_response",
] {
@@ -551,3 +564,83 @@ fn step142_verification_remains_pure_and_defers_later_authority() {
);
}
}
+
+#[test]
+fn step149_recovery_and_offline_export_are_bounded_and_non_networked() {
+ let contract: serde_json::Value =
+ serde_json::from_str(DELIVERY_RECOVERY_EXPORT_CONTRACT).expect("Step 149 contract");
+ assert_eq!(
+ contract["schema"],
+ "radroots.myc.delivery-recovery-export.v1"
+ );
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(contract["step"], 149);
+ assert_eq!(contract["schema_version"], 9);
+ for required in [
+ "atomic_response_or_discovery_commit_only",
+ "caller_injected_full_jitter_milliseconds",
+ "internal_opaque_job_identity_cursor",
+ "one_bounded_service_sqlite_transaction_per_batch",
+ "signature_verified_canonical_active_response_bytes",
+ "signature_verified_canonical_active_discovery_bytes",
+ "delivered_desired_restart_promotion",
+ "verified_committed_projection",
+ "compact_canonical_utf8_json",
+ "standalone_signer_delivery_job_creation",
+ "final_supervised_startup_loop",
+ ] {
+ assert!(
+ DELIVERY_RECOVERY_EXPORT_CONTRACT.contains(required),
+ "Step 149 contract is missing `{required}`"
+ );
+ }
+ for required in [
+ "MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT: usize = 128",
+ "READ_INVARIANTS_SQL",
+ "recover_delivery_state",
+ "verify_response_for_delivery_job",
+ "verify_document_for_delivery_job",
+ "recover_expired",
+ "promote_current_if_desired",
+ ] {
+ assert!(
+ DELIVERY_RECOVERY.contains(required),
+ "Step 149 recovery is missing `{required}`"
+ );
+ }
+ for required in [
+ "render_offline_nip05",
+ "MycNip05ExportSelection::Desired",
+ "MycNip05ExportSelection::Current",
+ "struct Nip05Output",
+ "names: Nip05Names",
+ "nip46: Nip46Discovery",
+ ] {
+ assert!(
+ DISCOVERY_STATE.contains(required),
+ "Step 149 offline export is missing `{required}`"
+ );
+ }
+ for removed in ["MycDeliveryJobRequest", "create_delivery_job"] {
+ assert!(
+ !PUBLIC_API.contains(removed),
+ "partial delivery authority remains public: `{removed}`"
+ );
+ }
+ for forbidden in [
+ "reqwest",
+ "RelayPool",
+ ".publish(",
+ "tokio::spawn",
+ "std::time::SystemTime",
+ "Timestamp::now",
+ "rand::",
+ "getrandom",
+ "rusqlite",
+ ] {
+ assert!(
+ !DELIVERY_RECOVERY.contains(forbidden),
+ "Step 149 recovery gained forbidden authority `{forbidden}`"
+ );
+ }
+}
diff --git a/tests/services_hardening_delivery_state.rs b/tests/services_hardening_delivery_state.rs
@@ -4,17 +4,22 @@
use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path};
use myc::{
- MYC_DELIVERY_RELAY_ID_MAX_BYTES, MYC_STATE_SCHEMA_VERSION, MycConfigProfile,
- MycDeliveryArtifactDigest, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome,
- MycDeliveryAttemptStatus, MycDeliveryClaim, MycDeliveryJobAdmission, MycDeliveryJobRequest,
- MycDeliveryJobStatus, MycDeliveryPolicyMode, MycDeliveryRelayId, MycDeliveryStateErrorKind,
- MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, MycNip46ClientPublicKey, MycNip46EventId,
- MycNip46RequestId, MycRequestReceivedAtUnixMs, MycSignerOperationNonce, MycSignerRequest,
- MycSignerRequestDigest, MycSignerRequestMethod, MycStateMetadata, MycStateRepositoryErrorKind,
+ MYC_DELIVERY_RELAY_ID_MAX_BYTES, MYC_DELIVERY_RETRY_JITTER_MAX_MS, MYC_STATE_SCHEMA_VERSION,
+ MycConfigProfile, MycDeliveryArtifactDigest, MycDeliveryAttemptNonce,
+ MycDeliveryAttemptOutcome, MycDeliveryAttemptStatus, MycDeliveryClaim, MycDeliveryJobStatus,
+ MycDeliveryPolicyMode, MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliveryStateErrorKind,
+ MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, MycDiscoveryCommitAdmission,
+ MycDiscoveryCommitRequest, MycStateMetadata, MycStateRepositoryErrorKind,
RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, initialize_myc_state,
open_myc_state_read_write, parse_myc_cli_v1_from, parse_myc_config_v1,
resolve_myc_runtime_context,
};
+use nostr::{Keys, SecretKey};
+use radroots_nostr::event::{
+ ApplicationHandlerSpec as RadrootsNostrApplicationHandlerSpec,
+ Metadata as RadrootsNostrMetadata, Timestamp as RadrootsNostrTimestamp,
+ build_application_handler as radroots_nostr_build_application_handler_event,
+};
use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
use radroots_storage::event::SourceGeneration;
use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions};
@@ -24,7 +29,7 @@ const CONFIG_EXAMPLE: &[u8] =
const DELIVERY_SOURCE: &str = include_str!("../src/state_delivery.rs");
const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs");
const LIB_SOURCE: &str = include_str!("../src/lib.rs");
-const CLIENT_PUBLIC_KEY: &str = "2222222222222222222222222222222222222222222222222222222222222222";
+const DISCOVERY_SECRET: &str = "3333333333333333333333333333333333333333333333333333333333333333";
fn runtime(root: &Path) -> myc::MycRuntimeContext {
let root = root.to_str().expect("UTF-8 temporary root");
@@ -64,6 +69,58 @@ fn metadata(runtime: &myc::MycRuntimeContext, source: &[u8]) -> MycStateMetadata
.expect("metadata")
}
+fn discovery_keys() -> Keys {
+ Keys::new(SecretKey::parse(DISCOVERY_SECRET).expect("discovery secret"))
+}
+
+fn config_source(source: &[u8]) -> Vec<u8> {
+ String::from_utf8(source.to_vec())
+ .expect("UTF-8 configuration")
+ .replace(
+ "3333333333333333333333333333333333333333333333333333333333333333",
+ &discovery_keys().public_key().to_hex(),
+ )
+ .into_bytes()
+}
+
+fn nostrconnect_url() -> String {
+ let mut query = url::form_urlencoded::Serializer::new(String::new());
+ query.append_pair("relay", "wss://relay-primary.example.test/");
+ query.append_pair("relay", "wss://relay-secondary.example.test/");
+ let bunker = format!(
+ "bunker://{}?{}",
+ "2222222222222222222222222222222222222222222222222222222222222222",
+ query.finish()
+ );
+ let encoded: String = url::form_urlencoded::byte_serialize(bunker.as_bytes()).collect();
+ format!("https://myc.example.test/connect?uri={encoded}")
+}
+
+fn signed_handler_event(created_at: u64) -> Vec<u8> {
+ let metadata = RadrootsNostrMetadata {
+ name: Some("myc".to_owned()),
+ display_name: Some("Radroots Myc".to_owned()),
+ about: Some("NIP-46 signer".to_owned()),
+ website: Some("https://myc.example.test/".to_owned()),
+ picture: Some("https://myc.example.test/myc.png".to_owned()),
+ ..RadrootsNostrMetadata::default()
+ };
+ let spec = RadrootsNostrApplicationHandlerSpec::new(vec![24_133])
+ .with_identifier("myc")
+ .with_relays(vec![
+ "wss://relay-primary.example.test/".to_owned(),
+ "wss://relay-secondary.example.test/".to_owned(),
+ ])
+ .with_nostr_connect_url(nostrconnect_url())
+ .with_metadata(metadata);
+ let event = radroots_nostr_build_application_handler_event(&spec)
+ .expect("typed handler event")
+ .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at))
+ .sign_with_keys(&discovery_keys())
+ .expect("signed event");
+ serde_json::to_vec(&event).expect("canonical event bytes")
+}
+
fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) {
let applied_at = MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time");
let build = MigrationBuildIdentity::new(
@@ -83,23 +140,14 @@ fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentit
(applied_at, build)
}
-fn signer_request(request_id: &str, nonce: u8, received_at: u64) -> MycSignerRequest {
- let canonical = format!(r#"{{"id":"{request_id}","method":"ping","params":[]}}"#);
- MycSignerRequest::new(
- MycNip46ClientPublicKey::new(CLIENT_PUBLIC_KEY).expect("client identity"),
- MycNip46RequestId::new(request_id).expect("request ID"),
- MycNip46EventId::from_bytes([nonce; 32]),
- MycSignerRequestMethod::Ping,
- MycSignerRequestDigest::for_canonical_request(canonical.as_bytes()).expect("digest"),
- MycSignerOperationNonce::from_injected_entropy([nonce; 32]),
- MycRequestReceivedAtUnixMs::new(received_at).expect("time"),
- )
-}
-
fn time(value: u64) -> MycDeliveryTimeUnixMs {
MycDeliveryTimeUnixMs::new(value).expect("delivery time")
}
+fn jitter(value: u64) -> MycDeliveryRetryJitter {
+ MycDeliveryRetryJitter::new(value).expect("retry jitter")
+}
+
#[test]
fn delivery_inputs_and_diagnostics_are_closed_bounded_and_redacted() {
let maximum = format!("a{}", "1".repeat(MYC_DELIVERY_RELAY_ID_MAX_BYTES - 1));
@@ -129,6 +177,13 @@ fn delivery_inputs_and_diagnostics_are_closed_bounded_and_redacted() {
);
}
assert!(MycDeliveryTimeUnixMs::new(i64::MAX.unsigned_abs()).is_ok());
+ assert!(MycDeliveryRetryJitter::new(MYC_DELIVERY_RETRY_JITTER_MAX_MS).is_ok());
+ assert_eq!(
+ MycDeliveryRetryJitter::new(MYC_DELIVERY_RETRY_JITTER_MAX_MS + 1)
+ .expect_err("oversize retry jitter")
+ .kind(),
+ MycDeliveryStateErrorKind::InvalidRetryJitter
+ );
assert_eq!(
[
MycDeliveryPolicyMode::AtLeastOneRequired,
@@ -155,7 +210,8 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
let directory = tempfile::tempdir().expect("temporary root");
let runtime = runtime(directory.path());
prepare_state_directory(&runtime);
- let metadata = metadata(&runtime, CONFIG_EXAMPLE);
+ let source = config_source(CONFIG_EXAMPLE);
+ let metadata = metadata(&runtime, &source);
let (applied_at, build) = migration_evidence();
initialize_myc_state(&runtime, &metadata, applied_at, &build)
.await
@@ -163,37 +219,37 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build)
.await
.expect("writer");
- let admitted = host
- .repository()
- .admit_signer_request(&signer_request("delivery-01", 0x31, 100))
- .await
- .expect("signer request");
- let request = MycDeliveryJobRequest::signer_response(
- admitted.record().operation_id(),
- MycDeliveryArtifactDigest::from_bytes([0x44; 32]),
- time(110),
- );
+ let request =
+ MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_000), time(110))
+ .expect("discovery request");
let job = host
.repository()
- .create_delivery_job(&request)
+ .commit_discovery_desired_state(&request)
.await
.expect("job");
- assert!(matches!(job, MycDeliveryJobAdmission::Created(_)));
- let job_id = job.record().id();
+ assert!(matches!(job, MycDiscoveryCommitAdmission::Created(_)));
+ let job_id = job.record().job().id();
assert_eq!(
- job.record().policy_mode(),
+ job.record().job().policy_mode(),
MycDeliveryPolicyMode::AllRequired
);
- assert_eq!(job.record().required_acknowledgements(), 2);
- assert_eq!(job.record().max_attempts(), 5);
- assert_eq!(job.record().initial_backoff_ms(), 250);
- assert_eq!(job.record().maximum_backoff_ms(), 30_000);
- assert_eq!(job.record().attempt_deadline_ms(), 15_000);
- assert_eq!(job.record().targets().len(), 2);
- assert_eq!(job.record().targets()[0].relay_id().as_str(), "primary");
- assert_eq!(job.record().targets()[1].relay_id().as_str(), "secondary");
+ assert_eq!(job.record().job().required_acknowledgements(), 2);
+ assert_eq!(job.record().job().max_attempts(), 5);
+ assert_eq!(job.record().job().initial_backoff_ms(), 250);
+ assert_eq!(job.record().job().maximum_backoff_ms(), 30_000);
+ assert_eq!(job.record().job().attempt_deadline_ms(), 15_000);
+ assert_eq!(job.record().job().targets().len(), 2);
+ assert_eq!(
+ job.record().job().targets()[0].relay_id().as_str(),
+ "primary"
+ );
+ assert_eq!(
+ job.record().job().targets()[1].relay_id().as_str(),
+ "secondary"
+ );
assert!(
job.record()
+ .job()
.targets()
.iter()
.all(|target| target.required())
@@ -201,18 +257,19 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
let replay = host
.repository()
- .create_delivery_job(&request)
+ .commit_discovery_desired_state(&request)
.await
.expect("exact replay");
- assert!(matches!(replay, MycDeliveryJobAdmission::ExactReplay(_)));
- let conflicting = MycDeliveryJobRequest::signer_response(
- admitted.record().operation_id(),
- MycDeliveryArtifactDigest::from_bytes([0x45; 32]),
- time(110),
- );
+ assert!(matches!(
+ replay,
+ MycDiscoveryCommitAdmission::ExactReplay(_)
+ ));
+ let conflicting =
+ MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_001), time(110))
+ .expect("conflicting discovery request");
assert_eq!(
host.repository()
- .create_delivery_job(&conflicting)
+ .commit_discovery_desired_state(&conflicting)
.await
.expect_err("conflicting job")
.kind(),
@@ -252,6 +309,7 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
&primary,
claimed.id(),
MycDeliveryAttemptOutcome::Delivered,
+ jitter(0),
time(121),
)
.await
@@ -272,6 +330,7 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
&primary,
claimed.id(),
MycDeliveryAttemptOutcome::Delivered,
+ jitter(0),
time(122),
)
.await
@@ -297,6 +356,21 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
.mark_delivery_attempt_submitted(job_id, &secondary, first.id(), time(124))
.await
.expect("secondary submitted");
+ assert_eq!(
+ host.repository()
+ .record_delivery_attempt_outcome(
+ job_id,
+ &secondary,
+ first.id(),
+ MycDeliveryAttemptOutcome::UnknownAcknowledgement,
+ jitter(251),
+ time(125),
+ )
+ .await
+ .expect_err("jitter exceeds first retry cap")
+ .kind(),
+ MycStateRepositoryErrorKind::Binding
+ );
let unknown = host
.repository()
.record_delivery_attempt_outcome(
@@ -304,6 +378,7 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
&secondary,
first.id(),
MycDeliveryAttemptOutcome::UnknownAcknowledgement,
+ jitter(250),
time(125),
)
.await
@@ -342,7 +417,7 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
};
let retry = host
.repository()
- .recover_expired_delivery_lease(job_id, &secondary, second.id(), time(15_376))
+ .recover_expired_delivery_lease(job_id, &secondary, second.id(), jitter(500), time(15_376))
.await
.expect("expired pre-submit lease");
assert_eq!(
@@ -375,6 +450,7 @@ async fn delivery_jobs_are_config_bound_idempotent_restart_safe_and_unknown_awar
&secondary,
third.id(),
MycDeliveryAttemptOutcome::Delivered,
+ jitter(0),
time(15_878),
)
.await
@@ -414,7 +490,7 @@ async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_e
let directory = tempfile::tempdir().expect("temporary root");
let runtime = runtime(directory.path());
prepare_state_directory(&runtime);
- let source = String::from_utf8(CONFIG_EXAMPLE.to_vec())
+ let source = String::from_utf8(config_source(CONFIG_EXAMPLE))
.expect("UTF-8 config")
.replace("max_attempts = 5", "max_attempts = 1");
let metadata = metadata(&runtime, source.as_bytes());
@@ -425,21 +501,19 @@ async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_e
let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build)
.await
.expect("writer");
- let admitted = host
- .repository()
- .admit_signer_request(&signer_request("delivery-unknown", 0x71, 100))
- .await
- .expect("request");
let job = host
.repository()
- .create_delivery_job(&MycDeliveryJobRequest::signer_response(
- admitted.record().operation_id(),
- MycDeliveryArtifactDigest::from_bytes([0x72; 32]),
- time(110),
- ))
+ .commit_discovery_desired_state(
+ &MycDiscoveryCommitRequest::new(
+ &metadata,
+ &signed_handler_event(1_725_000_010),
+ time(110),
+ )
+ .expect("discovery request"),
+ )
.await
.expect("job");
- let job_id = job.record().id();
+ let job_id = job.record().job().id();
let primary = MycDeliveryRelayId::new("primary").expect("primary");
let attempt = match host
.repository()
@@ -466,6 +540,7 @@ async fn terminal_unknown_is_not_relabelled_as_failure_and_sql_guards_preserve_e
&primary,
attempt.id(),
MycDeliveryAttemptOutcome::UnknownAcknowledgement,
+ jitter(0),
time(122),
)
.await
diff --git a/tests/services_hardening_discovery_state.rs b/tests/services_hardening_discovery_state.rs
@@ -6,11 +6,12 @@ use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path};
use myc::{
MYC_DISCOVERY_DOCUMENT_MAX_BYTES, MYC_STATE_SCHEMA_VERSION, MycConfigProfile,
MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, MycDeliveryClaim, MycDeliveryJobStatus,
- MycDeliveryRelayId, MycDeliverySourceKind, MycDeliveryTimeUnixMs, MycDiscoveryCommitAdmission,
- MycDiscoveryCommitRequest, MycDiscoveryStateErrorKind, MycStateMetadata,
+ MycDeliveryRecoveryEntropy, MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliverySourceKind,
+ MycDeliveryTimeUnixMs, MycDiscoveryCommitAdmission, MycDiscoveryCommitRequest,
+ MycDiscoveryStateErrorKind, MycNip05ExportSelection, MycStateMetadata,
MycStateRepositoryErrorKind, RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform,
- initialize_myc_state, open_myc_state_read_write, parse_myc_cli_v1_from, parse_myc_config_v1,
- resolve_myc_runtime_context,
+ initialize_myc_state, open_myc_state_inspection, open_myc_state_read_write,
+ parse_myc_cli_v1_from, parse_myc_config_v1, resolve_myc_runtime_context,
};
use nostr::{Keys, SecretKey};
use radroots_nostr::event::{
@@ -20,6 +21,7 @@ use radroots_nostr::event::{
};
use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
use radroots_storage::event::SourceGeneration;
+use sha2::{Digest, Sha256};
use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions};
const CONFIG_EXAMPLE: &[u8] =
@@ -198,6 +200,7 @@ async fn deliver_all_required(
&relay,
claimed.id(),
MycDeliveryAttemptOutcome::Delivered,
+ MycDeliveryRetryJitter::new(0).expect("zero retry jitter"),
time(start + offset as u64 * 10 + 2),
)
.await
@@ -319,6 +322,31 @@ async fn discovery_desired_current_and_exact_documents_are_atomic_restart_safe_a
assert_eq!(record.state().current_generation_id(), None);
assert_eq!(record.state().current_job_id(), None);
let job_id = record.job().id();
+ let desired_export = host
+ .repository()
+ .render_offline_nip05(MycNip05ExportSelection::Desired)
+ .await
+ .expect("desired NIP-05 export");
+ assert_eq!(desired_export.selection(), MycNip05ExportSelection::Desired);
+ assert_eq!(desired_export.domain(), "myc.example.test");
+ let expected_export = format!(
+ "{{\"names\":{{\"_\":\"{}\"}},\"nip46\":{{\"relays\":[\"wss://relay-primary.example.test/\",\"wss://relay-secondary.example.test/\"],\"nostrconnect_url\":\"{}\"}}}}",
+ discovery_keys().public_key().to_hex(),
+ nostrconnect_url()
+ );
+ assert_eq!(desired_export.bytes(), expected_export.as_bytes());
+ assert_eq!(
+ desired_export.digest().as_bytes(),
+ &<[u8; 32]>::from(Sha256::digest(expected_export.as_bytes()))
+ );
+ assert_eq!(
+ host.repository()
+ .render_offline_nip05(MycNip05ExportSelection::Current)
+ .await
+ .expect_err("no proven-current generation")
+ .kind(),
+ MycStateRepositoryErrorKind::Binding
+ );
let replay = host
.repository()
@@ -362,6 +390,14 @@ async fn discovery_desired_current_and_exact_documents_are_atomic_restart_safe_a
assert_eq!(promoted.current_job_id(), Some(job_id));
assert_eq!(
host.repository()
+ .render_offline_nip05(MycNip05ExportSelection::Current)
+ .await
+ .expect("current NIP-05 export")
+ .bytes(),
+ expected_export.as_bytes()
+ );
+ assert_eq!(
+ host.repository()
.read_discovery_document_for_job(job_id)
.await
.expect("document read")
@@ -420,6 +456,176 @@ async fn discovery_desired_current_and_exact_documents_are_atomic_restart_safe_a
.expect("committed state");
assert_eq!(state.current_job_id(), Some(next_job_id));
host.close().await.expect("final close");
+
+ let inspection = open_myc_state_inspection(&runtime, &metadata)
+ .await
+ .expect("inspection host");
+ let offline = inspection
+ .repository()
+ .render_offline_nip05(MycNip05ExportSelection::Current)
+ .await
+ .expect("read-only offline export");
+ assert_eq!(offline.bytes(), expected_export.as_bytes());
+ let rendered = format!("{offline:?}");
+ assert!(!rendered.contains("myc.example.test"));
+ assert!(!rendered.contains(&discovery_keys().public_key().to_hex()));
+ inspection.close().await.expect("inspection close");
+}
+
+#[tokio::test]
+async fn restart_recovery_is_bounded_jittered_and_idempotent() {
+ let directory = tempfile::tempdir().expect("temporary root");
+ let runtime = runtime(directory.path());
+ prepare_state_directory(&runtime);
+ let bounded_source = String::from_utf8(config_source(true))
+ .expect("UTF-8 configuration")
+ .replace("outbox = 4096", "outbox = 1");
+ let metadata = state_metadata(&runtime, bounded_source.as_bytes());
+ let (applied_at, build) = migration_evidence();
+ initialize_myc_state(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("initialization");
+ let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("writer");
+ let request =
+ MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_020), time(100))
+ .expect("discovery request");
+ let committed = host
+ .repository()
+ .commit_discovery_desired_state(&request)
+ .await
+ .expect("desired state");
+ let job_id = committed.record().job().id();
+ assert!(matches!(
+ host.repository()
+ .commit_discovery_desired_state(&request)
+ .await
+ .expect("exact replay at queue ceiling"),
+ MycDiscoveryCommitAdmission::ExactReplay(_)
+ ));
+ let saturated =
+ MycDiscoveryCommitRequest::new(&metadata, &signed_handler_event(1_725_000_021), time(101))
+ .expect("second desired state");
+ assert_eq!(
+ host.repository()
+ .commit_discovery_desired_state(&saturated)
+ .await
+ .expect_err("outbox ceiling")
+ .kind(),
+ MycStateRepositoryErrorKind::Binding
+ );
+ let primary = MycDeliveryRelayId::new("primary").expect("primary");
+ let attempt = match host
+ .repository()
+ .claim_delivery_target(
+ job_id,
+ &primary,
+ MycDeliveryAttemptNonce::from_injected_entropy([0xa1; 32]),
+ time(120),
+ )
+ .await
+ .expect("claim")
+ {
+ MycDeliveryClaim::Claimed(attempt) => attempt,
+ other => panic!("unexpected claim: {other:?}"),
+ };
+ host.repository()
+ .mark_delivery_attempt_submitted(job_id, &primary, attempt.id(), time(121))
+ .await
+ .expect("submitted");
+ host.close().await.expect("close before restart");
+
+ let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("reopened writer");
+ let recovered_at = time(15_121);
+ let report = host
+ .repository()
+ .recover_delivery_state(
+ recovered_at,
+ MycDeliveryRecoveryEntropy::from_injected_entropy([0xb2; 32]),
+ )
+ .await
+ .expect("restart recovery");
+ assert_eq!(report.examined_jobs(), 1);
+ assert_eq!(report.recovered_expired_attempts(), 1);
+ assert_eq!(report.finalized_jobs(), 0);
+ assert_eq!(report.promoted_discovery_generations(), 0);
+ assert_eq!(report.active_targets(), 0);
+ assert_eq!(report.ready_targets() + report.scheduled_targets(), 2);
+ let job = host
+ .repository()
+ .read_delivery_job(job_id)
+ .await
+ .expect("job read")
+ .expect("job");
+ let next = job.targets()[0]
+ .next_attempt_at()
+ .expect("persisted jittered retry");
+ assert!(next >= recovered_at);
+ assert!(next.get() <= recovered_at.get() + 250);
+ let replay = host
+ .repository()
+ .recover_delivery_state(
+ recovered_at,
+ MycDeliveryRecoveryEntropy::from_injected_entropy([0xc3; 32]),
+ )
+ .await
+ .expect("idempotent recovery");
+ assert_eq!(replay.recovered_expired_attempts(), 0);
+ assert_eq!(replay.examined_jobs(), 1);
+ assert_eq!(
+ host.repository()
+ .read_delivery_attempts(job_id, &primary)
+ .await
+ .expect("attempt history")
+ .len(),
+ 1
+ );
+ deliver_all_required(&host.repository(), job_id, 16_000).await;
+ assert_eq!(
+ host.repository()
+ .render_offline_nip05(MycNip05ExportSelection::Current)
+ .await
+ .expect_err("delivered desired state is not promoted implicitly")
+ .kind(),
+ MycStateRepositoryErrorKind::Binding
+ );
+ host.close().await.expect("close after delivered job");
+
+ let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("reopened writer after delivered job");
+ let promotion = host
+ .repository()
+ .recover_delivery_state(
+ time(16_100),
+ MycDeliveryRecoveryEntropy::from_injected_entropy([0xd4; 32]),
+ )
+ .await
+ .expect("restart promotion recovery");
+ assert_eq!(promotion.examined_jobs(), 1);
+ assert_eq!(promotion.recovered_expired_attempts(), 0);
+ assert_eq!(promotion.finalized_jobs(), 0);
+ assert_eq!(promotion.promoted_discovery_generations(), 1);
+ assert!(
+ host.repository()
+ .render_offline_nip05(MycNip05ExportSelection::Current)
+ .await
+ .is_ok()
+ );
+ let settled = host
+ .repository()
+ .recover_delivery_state(
+ time(16_101),
+ MycDeliveryRecoveryEntropy::from_injected_entropy([0xe5; 32]),
+ )
+ .await
+ .expect("settled recovery");
+ assert_eq!(settled.examined_jobs(), 0);
+ assert_eq!(settled.promoted_discovery_generations(), 0);
+ host.close().await.expect("final close");
}
#[tokio::test]
diff --git a/tests/services_hardening_legacy_removal.rs b/tests/services_hardening_legacy_removal.rs
@@ -16,6 +16,7 @@ const ACTIVE_STATE_SOURCES: &[&str] = &[
include_str!("../src/state_host.rs"),
include_str!("../src/state_maintenance.rs"),
include_str!("../src/state_metadata.rs"),
+ include_str!("../src/state_recovery.rs"),
include_str!("../src/state_repository.rs"),
include_str!("../src/state_request.rs"),
include_str!("../src/state_response.rs"),
@@ -113,7 +114,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(), 12);
+ assert_eq!(LIB_SOURCE.matches("mod state_").count(), 13);
assert!(!LIB_SOURCE.contains("pub mod state_"));
let active_state = ACTIVE_STATE_SOURCES.join("\n");
for forbidden in [