myc

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

state_response.rs (43269B)


      1 //! Atomic completion, exact signed-response, and initial delivery-state commit.
      2 
      3 use core::fmt;
      4 use std::error::Error;
      5 
      6 use radroots_nostr::event::{Event as RadrootsNostrEvent, Kind as RadrootsNostrKind};
      7 use radroots_service_sqlite::{
      8     ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
      9 };
     10 use sha2::{Digest, Sha256};
     11 use sqlx::Row;
     12 
     13 use crate::state_completion::{
     14     CommitOperationError, MycNip46CommitAdmission, MycNip46CommitRecord, MycNip46CommitRequest,
     15     commit_operation,
     16 };
     17 use crate::state_delivery::{
     18     DeliveryOperationError, MycDeliveryArtifactDigest, MycDeliveryJobAdmission, MycDeliveryJobId,
     19     MycDeliveryJobRecord, MycDeliverySource, MycDeliveryTimeUnixMs, create_job, read_job,
     20 };
     21 use crate::state_repository::{
     22     MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata,
     23     RepositoryOperationError, require_expected_metadata,
     24 };
     25 use crate::{
     26     MYC_PROVIDER_OUTPUT_MAX_BYTES, MycConnectionDecision, MycConnectionDecisionRecord,
     27     MycConnectionId, MycConnectionPolicyGeneration, MycConnectionStatus, MycNip46ClientPublicKey,
     28     MycNip46Work, MycNip46WorkKind, MycProviderCapability, MycProviderOperation, MycProviderRole,
     29     MycSignerOperationId, MycSignerRequestMethod, MycVerifiedProviderResponse,
     30 };
     31 
     32 const NIP46_RPC_KIND: u16 = 24_133;
     33 
     34 const INSERT_RESPONSE_SQL: &str = r#"INSERT INTO nip46_signed_responses (
     35     operation_id, response_provider_operation_id, response_event_id,
     36     response_sha256, response_bytes, authored_at_unix_s, committed_at_unix_ms
     37 ) VALUES (?, ?, ?, ?, ?, ?, ?)"#;
     38 
     39 const INSERT_PENDING_RESPONSE_SQL: &str = r#"INSERT INTO nip46_pending_responses (
     40     operation_id, connection_id, response_kind, response_provider_operation_id,
     41     response_event_id, response_sha256, response_bytes, authored_at_unix_s,
     42     committed_at_unix_ms
     43 ) VALUES (?, ?, 'pending_approval', ?, ?, ?, ?, ?, ?)"#;
     44 
     45 const READ_RESPONSE_BY_OPERATION_SQL: &str = r#"WITH response_authority AS (
     46     SELECT 'terminal' AS authority_kind, operation_id,
     47         response_provider_operation_id, response_event_id, response_sha256,
     48         response_bytes, authored_at_unix_s, committed_at_unix_ms
     49     FROM nip46_signed_responses
     50     UNION ALL
     51     SELECT response_kind AS authority_kind, operation_id,
     52         response_provider_operation_id, response_event_id, response_sha256,
     53         response_bytes, authored_at_unix_s, committed_at_unix_ms
     54     FROM nip46_pending_responses
     55 )
     56 SELECT
     57     CASE WHEN typeof(r.authority_kind) = 'text'
     58             AND length(CAST(r.authority_kind AS BLOB)) <= 16
     59         THEN r.authority_kind ELSE NULL END AS authority_kind,
     60     CASE WHEN typeof(r.operation_id) = 'blob' AND length(r.operation_id) = 32
     61         THEN r.operation_id ELSE NULL END AS operation_id,
     62     CASE WHEN typeof(r.response_provider_operation_id) = 'blob'
     63             AND length(r.response_provider_operation_id) = 32
     64         THEN r.response_provider_operation_id ELSE NULL END AS response_provider_operation_id,
     65     CASE WHEN typeof(r.response_event_id) = 'blob' AND length(r.response_event_id) = 32
     66         THEN r.response_event_id ELSE NULL END AS response_event_id,
     67     CASE WHEN typeof(r.response_sha256) = 'blob' AND length(r.response_sha256) = 32
     68         THEN r.response_sha256 ELSE NULL END AS response_sha256,
     69     CASE WHEN typeof(r.response_bytes) = 'blob'
     70             AND length(r.response_bytes) BETWEEN 1 AND 1048576
     71         THEN r.response_bytes ELSE NULL END AS response_bytes,
     72     r.authored_at_unix_s, r.committed_at_unix_ms,
     73     CASE WHEN typeof(q.client_public_key) = 'text'
     74             AND length(CAST(q.client_public_key AS BLOB)) = 64
     75         THEN q.client_public_key ELSE NULL END AS client_public_key,
     76     CASE WHEN typeof(j.job_id) = 'blob' AND length(j.job_id) = 32
     77         THEN j.job_id ELSE NULL END AS job_id
     78 FROM response_authority r
     79 JOIN nip46_requests q ON q.operation_id = r.operation_id
     80 JOIN delivery_jobs j ON j.source_kind = 'signer_response'
     81     AND j.source_id = r.operation_id
     82 WHERE r.operation_id = ?
     83 LIMIT 2"#;
     84 
     85 const READ_RESPONSE_BY_JOB_SQL: &str = r#"WITH response_authority AS (
     86     SELECT 'terminal' AS authority_kind, operation_id,
     87         response_provider_operation_id, response_event_id, response_sha256,
     88         response_bytes, authored_at_unix_s, committed_at_unix_ms
     89     FROM nip46_signed_responses
     90     UNION ALL
     91     SELECT response_kind AS authority_kind, operation_id,
     92         response_provider_operation_id, response_event_id, response_sha256,
     93         response_bytes, authored_at_unix_s, committed_at_unix_ms
     94     FROM nip46_pending_responses
     95 )
     96 SELECT
     97     CASE WHEN typeof(r.authority_kind) = 'text'
     98             AND length(CAST(r.authority_kind AS BLOB)) <= 16
     99         THEN r.authority_kind ELSE NULL END AS authority_kind,
    100     CASE WHEN typeof(r.operation_id) = 'blob' AND length(r.operation_id) = 32
    101         THEN r.operation_id ELSE NULL END AS operation_id,
    102     CASE WHEN typeof(r.response_provider_operation_id) = 'blob'
    103             AND length(r.response_provider_operation_id) = 32
    104         THEN r.response_provider_operation_id ELSE NULL END AS response_provider_operation_id,
    105     CASE WHEN typeof(r.response_event_id) = 'blob' AND length(r.response_event_id) = 32
    106         THEN r.response_event_id ELSE NULL END AS response_event_id,
    107     CASE WHEN typeof(r.response_sha256) = 'blob' AND length(r.response_sha256) = 32
    108         THEN r.response_sha256 ELSE NULL END AS response_sha256,
    109     CASE WHEN typeof(r.response_bytes) = 'blob'
    110             AND length(r.response_bytes) BETWEEN 1 AND 1048576
    111         THEN r.response_bytes ELSE NULL END AS response_bytes,
    112     r.authored_at_unix_s, r.committed_at_unix_ms,
    113     CASE WHEN typeof(q.client_public_key) = 'text'
    114             AND length(CAST(q.client_public_key AS BLOB)) = 64
    115         THEN q.client_public_key ELSE NULL END AS client_public_key,
    116     CASE WHEN typeof(j.job_id) = 'blob' AND length(j.job_id) = 32
    117         THEN j.job_id ELSE NULL END AS job_id
    118 FROM delivery_jobs j
    119 JOIN response_authority r ON r.operation_id = j.source_id
    120 JOIN nip46_requests q ON q.operation_id = r.operation_id
    121 WHERE j.job_id = ? AND j.source_kind = 'signer_response'
    122 LIMIT 2"#;
    123 
    124 const READ_PENDING_BINDING_SQL: &str = r#"SELECT
    125     CASE WHEN typeof(decision.connection_id) = 'blob'
    126             AND length(decision.connection_id) = 32
    127         THEN decision.connection_id ELSE NULL END AS connection_id,
    128     decision.policy_generation, decision.decided_at_unix_ms,
    129     CASE WHEN typeof(decision.decision) = 'text'
    130             AND length(CAST(decision.decision AS BLOB)) <= 32
    131         THEN decision.decision ELSE NULL END AS decision,
    132     CASE WHEN typeof(decision.reason_code) = 'text'
    133             AND length(CAST(decision.reason_code AS BLOB)) <= 32
    134         THEN decision.reason_code ELSE NULL END AS reason_code,
    135     CASE WHEN typeof(connection.status) = 'text'
    136             AND length(CAST(connection.status AS BLOB)) <= 16
    137         THEN connection.status ELSE NULL END AS connection_status,
    138     CASE WHEN typeof(connection.client_public_key) = 'text'
    139             AND length(CAST(connection.client_public_key AS BLOB)) = 64
    140         THEN connection.client_public_key ELSE NULL END AS client_public_key
    141 FROM nip46_request_decisions AS decision
    142 JOIN connections AS connection ON connection.connection_id = decision.connection_id
    143 JOIN nip46_requests AS request ON request.operation_id = decision.operation_id
    144 WHERE decision.operation_id = ?
    145     AND connection.client_public_key = request.client_public_key
    146     AND connection.policy_generation = decision.policy_generation
    147     AND connection.requested_permissions_sha256 = decision.requested_permissions_sha256
    148     AND request.method = 'connect'
    149 LIMIT 2"#;
    150 
    151 /// Stable construction failure classes for an atomic NIP-46 response commit.
    152 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    153 pub enum MycNip46ResponseCommitErrorKind {
    154     InvalidBinding,
    155     InvalidResponse,
    156     InvalidTime,
    157 }
    158 
    159 impl MycNip46ResponseCommitErrorKind {
    160     /// Returns the stable machine-readable failure code.
    161     #[must_use]
    162     pub const fn code(self) -> &'static str {
    163         match self {
    164             Self::InvalidBinding => "nip46_response_binding_invalid",
    165             Self::InvalidResponse => "nip46_response_invalid",
    166             Self::InvalidTime => "nip46_response_time_invalid",
    167         }
    168     }
    169 }
    170 
    171 /// Source-free, path-free construction failure.
    172 #[derive(Clone, Copy, PartialEq, Eq)]
    173 pub struct MycNip46ResponseCommitError {
    174     kind: MycNip46ResponseCommitErrorKind,
    175 }
    176 
    177 impl MycNip46ResponseCommitError {
    178     const fn new(kind: MycNip46ResponseCommitErrorKind) -> Self {
    179         Self { kind }
    180     }
    181 
    182     #[must_use]
    183     pub const fn kind(self) -> MycNip46ResponseCommitErrorKind {
    184         self.kind
    185     }
    186 
    187     #[must_use]
    188     pub const fn code(self) -> &'static str {
    189         self.kind.code()
    190     }
    191 }
    192 
    193 impl fmt::Display for MycNip46ResponseCommitError {
    194     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    195         formatter.write_str(match self.kind {
    196             MycNip46ResponseCommitErrorKind::InvalidBinding => "NIP-46 response binding is invalid",
    197             MycNip46ResponseCommitErrorKind::InvalidResponse => "NIP-46 signed response is invalid",
    198             MycNip46ResponseCommitErrorKind::InvalidTime => "NIP-46 response time is invalid",
    199         })
    200     }
    201 }
    202 
    203 impl fmt::Debug for MycNip46ResponseCommitError {
    204     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    205         formatter
    206             .debug_struct("MycNip46ResponseCommitError")
    207             .field("kind", &self.kind)
    208             .finish()
    209     }
    210 }
    211 
    212 impl Error for MycNip46ResponseCommitError {}
    213 
    214 /// Exact independently verified response and completion prepared before SQLite.
    215 pub struct MycNip46ResponseCommitRequest {
    216     completion: MycNip46CommitRequest,
    217     response_provider_operation_id: [u8; 32],
    218     response_event_id: [u8; 32],
    219     response_digest: MycDeliveryArtifactDigest,
    220     response_bytes: Box<[u8]>,
    221     authored_at_unix_s: u64,
    222     committed_at: MycDeliveryTimeUnixMs,
    223     #[cfg(test)]
    224     fail_after_completion: bool,
    225     #[cfg(test)]
    226     fail_after_response: bool,
    227 }
    228 
    229 impl MycNip46ResponseCommitRequest {
    230     /// Binds one completion to a signature-verified canonical NIP-46 response event.
    231     pub fn new(
    232         completion: &MycNip46CommitRequest,
    233         response_operation: &MycProviderOperation,
    234         response: &MycVerifiedProviderResponse,
    235         committed_at: MycDeliveryTimeUnixMs,
    236     ) -> Result<Self, MycNip46ResponseCommitError> {
    237         if response_operation.role() != MycProviderRole::Transport
    238             || response_operation.input().capability() != MycProviderCapability::SignEvent
    239             || response.operation_id() != response_operation.operation_id()
    240             || response.correlation_id() != response_operation.correlation_id()
    241             || response.instance() != response_operation.instance()
    242             || response.role() != response_operation.role()
    243             || response.capability() != response_operation.input().capability()
    244             || !response.matches_operation(response_operation)
    245             || response_operation.operation_id().as_bytes()
    246                 == completion.signer_request().operation_id().as_bytes()
    247         {
    248             return Err(Self::error(MycNip46ResponseCommitErrorKind::InvalidBinding));
    249         }
    250         let bytes = response
    251             .signed_event_bytes()
    252             .ok_or_else(|| Self::error(MycNip46ResponseCommitErrorKind::InvalidResponse))?;
    253         let event = validate_response_event(
    254             bytes,
    255             completion.signer_request().client_public_key().as_hex(),
    256             Some(response_operation.expected_identity().as_hex()),
    257         )?;
    258         let authored_at_unix_s = event.created_at.as_secs();
    259         if committed_at.get() < completion.completed_at().get()
    260             || authored_at_unix_s
    261                 .checked_mul(1_000)
    262                 .is_none_or(|authored_ms| authored_ms > committed_at.get())
    263         {
    264             return Err(Self::error(MycNip46ResponseCommitErrorKind::InvalidTime));
    265         }
    266         let response_digest = MycDeliveryArtifactDigest::from_bytes(Sha256::digest(bytes).into());
    267         Ok(Self {
    268             completion: completion.owned(),
    269             response_provider_operation_id: *response_operation.operation_id().as_bytes(),
    270             response_event_id: *event.id.as_bytes(),
    271             response_digest,
    272             response_bytes: Box::from(bytes),
    273             authored_at_unix_s,
    274             committed_at,
    275             #[cfg(test)]
    276             fail_after_completion: false,
    277             #[cfg(test)]
    278             fail_after_response: false,
    279         })
    280     }
    281 
    282     const fn error(kind: MycNip46ResponseCommitErrorKind) -> MycNip46ResponseCommitError {
    283         MycNip46ResponseCommitError::new(kind)
    284     }
    285 
    286     fn owned(&self) -> Self {
    287         Self {
    288             completion: self.completion.owned(),
    289             response_provider_operation_id: self.response_provider_operation_id,
    290             response_event_id: self.response_event_id,
    291             response_digest: self.response_digest,
    292             response_bytes: self.response_bytes.clone(),
    293             authored_at_unix_s: self.authored_at_unix_s,
    294             committed_at: self.committed_at,
    295             #[cfg(test)]
    296             fail_after_completion: self.fail_after_completion,
    297             #[cfg(test)]
    298             fail_after_response: self.fail_after_response,
    299         }
    300     }
    301 
    302     #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    303     pub(crate) fn fail_after_completion_for_test(&self) -> Self {
    304         let mut request = self.owned();
    305         request.fail_after_completion = true;
    306         request
    307     }
    308 
    309     #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    310     pub(crate) fn fail_after_response_for_test(&self) -> Self {
    311         let mut request = self.owned();
    312         request.fail_after_response = true;
    313         request
    314     }
    315 }
    316 
    317 impl fmt::Debug for MycNip46ResponseCommitRequest {
    318     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    319         formatter.write_str("MycNip46ResponseCommitRequest([redacted])")
    320     }
    321 }
    322 
    323 #[derive(Clone)]
    324 pub(crate) struct MycNip46PendingResponseCommitRequest {
    325     operation_id: MycSignerOperationId,
    326     connection_id: MycConnectionId,
    327     policy_generation: MycConnectionPolicyGeneration,
    328     client_public_key: MycNip46ClientPublicKey,
    329     decided_at_unix_ms: u64,
    330     response_provider_operation_id: [u8; 32],
    331     response_event_id: [u8; 32],
    332     response_digest: MycDeliveryArtifactDigest,
    333     response_bytes: Box<[u8]>,
    334     authored_at_unix_s: u64,
    335     committed_at: MycDeliveryTimeUnixMs,
    336     #[cfg(test)]
    337     fail_after_response: bool,
    338 }
    339 
    340 impl MycNip46PendingResponseCommitRequest {
    341     pub(crate) fn new(
    342         work: &MycNip46Work,
    343         decision: &MycConnectionDecisionRecord,
    344         response_operation: &MycProviderOperation,
    345         response: &MycVerifiedProviderResponse,
    346         committed_at: MycDeliveryTimeUnixMs,
    347     ) -> Result<Self, MycNip46ResponseCommitError> {
    348         let connection = decision.connection().ok_or_else(|| {
    349             MycNip46ResponseCommitRequest::error(MycNip46ResponseCommitErrorKind::InvalidBinding)
    350         })?;
    351         if work.kind() != MycNip46WorkKind::Connect
    352             || work.method() != MycSignerRequestMethod::Connect
    353             || decision.operation_id() != work.request_record().operation_id()
    354             || decision.decision() != MycConnectionDecision::PendingApproval
    355             || decision.policy_generation() != connection.policy_generation()
    356             || connection.status() != MycConnectionStatus::Pending
    357             || connection.client_public_key() != work.request_record().client_public_key()
    358             || response_operation.role() != MycProviderRole::Transport
    359             || response_operation.input().capability() != MycProviderCapability::SignEvent
    360             || response.operation_id() != response_operation.operation_id()
    361             || response.correlation_id() != response_operation.correlation_id()
    362             || response.instance() != response_operation.instance()
    363             || response.role() != response_operation.role()
    364             || response.capability() != response_operation.input().capability()
    365             || !response.matches_operation(response_operation)
    366             || response_operation.operation_id().as_bytes()
    367                 == work.request_record().operation_id().as_bytes()
    368         {
    369             return Err(MycNip46ResponseCommitRequest::error(
    370                 MycNip46ResponseCommitErrorKind::InvalidBinding,
    371             ));
    372         }
    373         let bytes = response.signed_event_bytes().ok_or_else(|| {
    374             MycNip46ResponseCommitRequest::error(MycNip46ResponseCommitErrorKind::InvalidResponse)
    375         })?;
    376         let event = validate_response_event(
    377             bytes,
    378             work.request_record().client_public_key().as_hex(),
    379             Some(response_operation.expected_identity().as_hex()),
    380         )?;
    381         let authored_at_unix_s = event.created_at.as_secs();
    382         if committed_at.get() < decision.decided_at().get()
    383             || authored_at_unix_s
    384                 .checked_mul(1_000)
    385                 .is_none_or(|authored_ms| authored_ms > committed_at.get())
    386         {
    387             return Err(MycNip46ResponseCommitRequest::error(
    388                 MycNip46ResponseCommitErrorKind::InvalidTime,
    389             ));
    390         }
    391         Ok(Self {
    392             operation_id: work.request_record().operation_id(),
    393             connection_id: connection.id(),
    394             policy_generation: connection.policy_generation(),
    395             client_public_key: connection.client_public_key().clone(),
    396             decided_at_unix_ms: decision.decided_at().get(),
    397             response_provider_operation_id: *response_operation.operation_id().as_bytes(),
    398             response_event_id: *event.id.as_bytes(),
    399             response_digest: MycDeliveryArtifactDigest::from_bytes(Sha256::digest(bytes).into()),
    400             response_bytes: Box::from(bytes),
    401             authored_at_unix_s,
    402             committed_at,
    403             #[cfg(test)]
    404             fail_after_response: false,
    405         })
    406     }
    407 
    408     #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    409     pub(crate) fn fail_after_response_for_test(&self) -> Self {
    410         let mut request = self.clone();
    411         request.fail_after_response = true;
    412         request
    413     }
    414 }
    415 
    416 impl fmt::Debug for MycNip46PendingResponseCommitRequest {
    417     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    418         formatter.write_str("MycNip46PendingResponseCommitRequest([redacted])")
    419     }
    420 }
    421 
    422 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    423 enum ResponseAuthorityKind {
    424     Terminal,
    425     PendingApproval,
    426 }
    427 
    428 impl ResponseAuthorityKind {
    429     fn parse(value: &str) -> Option<Self> {
    430         match value {
    431             "terminal" => Some(Self::Terminal),
    432             "pending_approval" => Some(Self::PendingApproval),
    433             _ => None,
    434         }
    435     }
    436 }
    437 
    438 /// Immutable exact signed response and its config-bound initial delivery state.
    439 #[derive(Clone, PartialEq, Eq)]
    440 pub struct MycNip46ResponseRecord {
    441     authority_kind: ResponseAuthorityKind,
    442     operation_id: MycSignerOperationId,
    443     response_provider_operation_id: [u8; 32],
    444     response_event_id: [u8; 32],
    445     response_digest: MycDeliveryArtifactDigest,
    446     response_bytes: Box<[u8]>,
    447     authored_at_unix_s: u64,
    448     committed_at: MycDeliveryTimeUnixMs,
    449     delivery_job: MycDeliveryJobRecord,
    450 }
    451 
    452 impl MycNip46ResponseRecord {
    453     #[must_use]
    454     pub const fn operation_id(&self) -> MycSignerOperationId {
    455         self.operation_id
    456     }
    457     #[must_use]
    458     pub const fn response_provider_operation_id(&self) -> &[u8; 32] {
    459         &self.response_provider_operation_id
    460     }
    461     #[must_use]
    462     pub const fn response_event_id(&self) -> &[u8; 32] {
    463         &self.response_event_id
    464     }
    465     #[must_use]
    466     pub const fn response_digest(&self) -> MycDeliveryArtifactDigest {
    467         self.response_digest
    468     }
    469     /// Returns the exact committed bytes that every retry must submit unchanged.
    470     #[must_use]
    471     pub fn signed_response_bytes(&self) -> &[u8] {
    472         &self.response_bytes
    473     }
    474     #[must_use]
    475     pub const fn authored_at_unix_s(&self) -> u64 {
    476         self.authored_at_unix_s
    477     }
    478     #[must_use]
    479     pub const fn committed_at(&self) -> MycDeliveryTimeUnixMs {
    480         self.committed_at
    481     }
    482     #[must_use]
    483     pub const fn delivery_job(&self) -> &MycDeliveryJobRecord {
    484         &self.delivery_job
    485     }
    486 }
    487 
    488 impl fmt::Debug for MycNip46ResponseRecord {
    489     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    490         formatter
    491             .debug_struct("MycNip46ResponseRecord")
    492             .field("delivery_job", &self.delivery_job)
    493             .field("response", &"[redacted]")
    494             .field("identity", &"[redacted]")
    495             .finish()
    496     }
    497 }
    498 
    499 /// Immutable result of the one response-authority transaction.
    500 #[derive(Clone, PartialEq, Eq)]
    501 pub struct MycNip46ResponseCommitRecord {
    502     completion: MycNip46CommitRecord,
    503     response: MycNip46ResponseRecord,
    504 }
    505 
    506 impl MycNip46ResponseCommitRecord {
    507     #[must_use]
    508     pub const fn completion(&self) -> &MycNip46CommitRecord {
    509         &self.completion
    510     }
    511     #[must_use]
    512     pub const fn response(&self) -> &MycNip46ResponseRecord {
    513         &self.response
    514     }
    515 }
    516 
    517 impl fmt::Debug for MycNip46ResponseCommitRecord {
    518     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    519         formatter.write_str("MycNip46ResponseCommitRecord([redacted])")
    520     }
    521 }
    522 
    523 /// New or exact-replayed atomic response authority.
    524 #[derive(Clone, PartialEq, Eq)]
    525 pub enum MycNip46ResponseCommitAdmission {
    526     Committed(MycNip46ResponseCommitRecord),
    527     ExactReplay(MycNip46ResponseCommitRecord),
    528 }
    529 
    530 impl MycNip46ResponseCommitAdmission {
    531     #[must_use]
    532     pub const fn record(&self) -> &MycNip46ResponseCommitRecord {
    533         match self {
    534             Self::Committed(record) | Self::ExactReplay(record) => record,
    535         }
    536     }
    537 }
    538 
    539 impl fmt::Debug for MycNip46ResponseCommitAdmission {
    540     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    541         formatter.write_str(match self {
    542             Self::Committed(_) => "MycNip46ResponseCommitAdmission::Committed([redacted])",
    543             Self::ExactReplay(_) => "MycNip46ResponseCommitAdmission::ExactReplay([redacted])",
    544         })
    545     }
    546 }
    547 
    548 impl MycStateRepository<'_> {
    549     /// Atomically commits completion, exact response bytes, targets, and pending attempt state.
    550     pub async fn commit_nip46_response(
    551         &self,
    552         request: &MycNip46ResponseCommitRequest,
    553     ) -> Result<MycNip46ResponseCommitAdmission, MycStateRepositoryError> {
    554         let expected = PersistedMetadata::from(self.expected());
    555         let policy = self.expected().delivery_policies().clone();
    556         let request = request.owned();
    557         self.host()
    558             .transaction(move |transaction| {
    559                 Box::pin(async move {
    560                     require_expected_metadata(transaction, &expected)
    561                         .await
    562                         .map_err(AtomicOperationError::from)?;
    563                     let completion = commit_operation(transaction, &request.completion)
    564                         .await
    565                         .map_err(AtomicOperationError::from)?;
    566                     match completion {
    567                         MycNip46CommitAdmission::ExactReplay(completion) => {
    568                             let response = read_response_by_operation(
    569                                 transaction,
    570                                 request.completion.signer_request().operation_id(),
    571                             )
    572                             .await?
    573                             .ok_or(AtomicOperationError::Binding)?;
    574                             exact_response(&response, &request)?;
    575                             Ok(MycNip46ResponseCommitAdmission::ExactReplay(
    576                                 MycNip46ResponseCommitRecord {
    577                                     completion,
    578                                     response,
    579                                 },
    580                             ))
    581                         }
    582                         MycNip46CommitAdmission::Committed(completion) => {
    583                             #[cfg(test)]
    584                             if request.fail_after_completion {
    585                                 return Err(AtomicOperationError::Storage);
    586                             }
    587                             insert_response(transaction, &request).await?;
    588                             #[cfg(test)]
    589                             if request.fail_after_response {
    590                                 return Err(AtomicOperationError::Storage);
    591                             }
    592                             let delivery = create_job(
    593                                 transaction,
    594                                 MycDeliverySource::signer_response(
    595                                     request.completion.signer_request().operation_id(),
    596                                 ),
    597                                 request.response_digest,
    598                                 request.committed_at,
    599                                 &policy,
    600                             )
    601                             .await
    602                             .map_err(AtomicOperationError::from)?;
    603                             if !matches!(delivery, MycDeliveryJobAdmission::Created(_)) {
    604                                 return Err(AtomicOperationError::Binding);
    605                             }
    606                             let response = read_response_by_operation(
    607                                 transaction,
    608                                 request.completion.signer_request().operation_id(),
    609                             )
    610                             .await?
    611                             .ok_or(AtomicOperationError::Binding)?;
    612                             exact_response(&response, &request)?;
    613                             Ok(MycNip46ResponseCommitAdmission::Committed(
    614                                 MycNip46ResponseCommitRecord {
    615                                     completion,
    616                                     response,
    617                                 },
    618                             ))
    619                         }
    620                     }
    621                 })
    622             })
    623             .await
    624             .map_err(map_transaction_error)
    625     }
    626 
    627     pub(crate) async fn commit_nip46_pending_response(
    628         &self,
    629         request: &MycNip46PendingResponseCommitRequest,
    630     ) -> Result<MycNip46ResponseRecord, MycStateRepositoryError> {
    631         let expected = PersistedMetadata::from(self.expected());
    632         let policy = self.expected().delivery_policies().clone();
    633         let request = request.clone();
    634         self.host()
    635             .transaction(move |transaction| {
    636                 Box::pin(async move {
    637                     require_expected_metadata(transaction, &expected)
    638                         .await
    639                         .map_err(AtomicOperationError::from)?;
    640                     require_pending_binding(transaction, &request).await?;
    641                     if let Some(response) =
    642                         read_response_by_operation(transaction, request.operation_id).await?
    643                     {
    644                         exact_pending_response(&response, &request)?;
    645                         return Ok(response);
    646                     }
    647                     insert_pending_response(transaction, &request).await?;
    648                     #[cfg(test)]
    649                     if request.fail_after_response {
    650                         return Err(AtomicOperationError::Storage);
    651                     }
    652                     let delivery = create_job(
    653                         transaction,
    654                         MycDeliverySource::signer_response(request.operation_id),
    655                         request.response_digest,
    656                         request.committed_at,
    657                         &policy,
    658                     )
    659                     .await
    660                     .map_err(AtomicOperationError::from)?;
    661                     if !matches!(delivery, MycDeliveryJobAdmission::Created(_)) {
    662                         return Err(AtomicOperationError::Binding);
    663                     }
    664                     let response = read_response_by_operation(transaction, request.operation_id)
    665                         .await?
    666                         .ok_or(AtomicOperationError::Binding)?;
    667                     exact_pending_response(&response, &request)?;
    668                     Ok(response)
    669                 })
    670             })
    671             .await
    672             .map_err(map_transaction_error)
    673     }
    674 
    675     /// Reads the exact committed response bytes for a retained delivery job.
    676     pub async fn read_nip46_response(
    677         &self,
    678         job_id: MycDeliveryJobId,
    679     ) -> Result<Option<MycNip46ResponseRecord>, MycStateRepositoryError> {
    680         let expected = PersistedMetadata::from(self.expected());
    681         self.host()
    682             .transaction(move |transaction| {
    683                 Box::pin(async move {
    684                     require_expected_metadata(transaction, &expected)
    685                         .await
    686                         .map_err(AtomicOperationError::from)?;
    687                     read_response_by_job(transaction, job_id).await
    688                 })
    689             })
    690             .await
    691             .map_err(map_transaction_error)
    692     }
    693 
    694     /// Reads an already committed response by stable signer operation identity.
    695     ///
    696     /// Runtime replay handling uses this lookup before any provider call so an
    697     /// exact completed replay always reuses the originally committed bytes.
    698     pub(crate) async fn read_nip46_response_by_operation(
    699         &self,
    700         operation_id: MycSignerOperationId,
    701     ) -> Result<Option<MycNip46ResponseRecord>, MycStateRepositoryError> {
    702         let expected = PersistedMetadata::from(self.expected());
    703         self.host()
    704             .transaction(move |transaction| {
    705                 Box::pin(async move {
    706                     require_expected_metadata(transaction, &expected)
    707                         .await
    708                         .map_err(AtomicOperationError::from)?;
    709                     read_response_by_operation(transaction, operation_id).await
    710                 })
    711             })
    712             .await
    713             .map_err(map_transaction_error)
    714     }
    715 }
    716 
    717 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    718 enum AtomicOperationError {
    719     Binding,
    720     Storage,
    721 }
    722 
    723 impl From<RepositoryOperationError> for AtomicOperationError {
    724     fn from(error: RepositoryOperationError) -> Self {
    725         match error {
    726             RepositoryOperationError::Binding => Self::Binding,
    727             RepositoryOperationError::Storage => Self::Storage,
    728         }
    729     }
    730 }
    731 
    732 impl From<CommitOperationError> for AtomicOperationError {
    733     fn from(error: CommitOperationError) -> Self {
    734         match error {
    735             CommitOperationError::Binding => Self::Binding,
    736             CommitOperationError::Storage => Self::Storage,
    737         }
    738     }
    739 }
    740 
    741 impl From<DeliveryOperationError> for AtomicOperationError {
    742     fn from(error: DeliveryOperationError) -> Self {
    743         match error {
    744             DeliveryOperationError::Binding => Self::Binding,
    745             DeliveryOperationError::Storage => Self::Storage,
    746         }
    747     }
    748 }
    749 
    750 async fn require_pending_binding(
    751     transaction: &mut ServiceSqliteTransaction<'_>,
    752     request: &MycNip46PendingResponseCommitRequest,
    753 ) -> Result<(), AtomicOperationError> {
    754     let rows = sqlx::query(READ_PENDING_BINDING_SQL)
    755         .bind(request.operation_id.as_bytes().as_slice())
    756         .fetch_all(&mut *transaction)
    757         .await
    758         .map_err(|_| AtomicOperationError::Storage)?;
    759     if rows.len() != 1 {
    760         return Err(AtomicOperationError::Binding);
    761     }
    762     let row = &rows[0];
    763     let connection_id = MycConnectionId::from_bytes(blob32(row, "connection_id")?);
    764     let policy_generation = positive_i64(row, "policy_generation")?;
    765     let decided_at_unix_ms = positive_i64(row, "decided_at_unix_ms")?;
    766     let decision = bounded_text(row, "decision", 32)?;
    767     let reason_code = bounded_text(row, "reason_code", 32)?;
    768     let connection_status = bounded_text(row, "connection_status", 16)?;
    769     let client_public_key = bounded_text(row, "client_public_key", 64)?;
    770     let valid = connection_id == request.connection_id
    771         && policy_generation == request.policy_generation.get()
    772         && decided_at_unix_ms == request.decided_at_unix_ms
    773         && decision == "pending_approval"
    774         && reason_code == "explicit_approval_required"
    775         && connection_status == "pending"
    776         && client_public_key == request.client_public_key.as_hex();
    777     valid.then_some(()).ok_or(AtomicOperationError::Binding)
    778 }
    779 
    780 async fn insert_pending_response(
    781     transaction: &mut ServiceSqliteTransaction<'_>,
    782     request: &MycNip46PendingResponseCommitRequest,
    783 ) -> Result<(), AtomicOperationError> {
    784     let result = sqlx::query(INSERT_PENDING_RESPONSE_SQL)
    785         .bind(request.operation_id.as_bytes().as_slice())
    786         .bind(request.connection_id.as_bytes().as_slice())
    787         .bind(request.response_provider_operation_id.as_slice())
    788         .bind(request.response_event_id.as_slice())
    789         .bind(request.response_digest.as_bytes().as_slice())
    790         .bind(request.response_bytes.as_ref())
    791         .bind(to_i64(request.authored_at_unix_s)?)
    792         .bind(to_i64(request.committed_at.get())?)
    793         .execute(&mut *transaction)
    794         .await
    795         .map_err(|_| AtomicOperationError::Storage)?;
    796     (result.rows_affected() == 1)
    797         .then_some(())
    798         .ok_or(AtomicOperationError::Storage)
    799 }
    800 
    801 async fn insert_response(
    802     transaction: &mut ServiceSqliteTransaction<'_>,
    803     request: &MycNip46ResponseCommitRequest,
    804 ) -> Result<(), AtomicOperationError> {
    805     let result = sqlx::query(INSERT_RESPONSE_SQL)
    806         .bind(
    807             request
    808                 .completion
    809                 .signer_request()
    810                 .operation_id()
    811                 .as_bytes()
    812                 .as_slice(),
    813         )
    814         .bind(request.response_provider_operation_id.as_slice())
    815         .bind(request.response_event_id.as_slice())
    816         .bind(request.response_digest.as_bytes().as_slice())
    817         .bind(request.response_bytes.as_ref())
    818         .bind(to_i64(request.authored_at_unix_s)?)
    819         .bind(to_i64(request.committed_at.get())?)
    820         .execute(&mut *transaction)
    821         .await
    822         .map_err(|_| AtomicOperationError::Storage)?;
    823     (result.rows_affected() == 1)
    824         .then_some(())
    825         .ok_or(AtomicOperationError::Storage)
    826 }
    827 
    828 async fn read_response_by_operation(
    829     transaction: &mut ServiceSqliteTransaction<'_>,
    830     operation_id: MycSignerOperationId,
    831 ) -> Result<Option<MycNip46ResponseRecord>, AtomicOperationError> {
    832     read_response(
    833         transaction,
    834         READ_RESPONSE_BY_OPERATION_SQL,
    835         operation_id.as_bytes(),
    836     )
    837     .await
    838 }
    839 
    840 async fn read_response_by_job(
    841     transaction: &mut ServiceSqliteTransaction<'_>,
    842     job_id: MycDeliveryJobId,
    843 ) -> Result<Option<MycNip46ResponseRecord>, AtomicOperationError> {
    844     read_response(transaction, READ_RESPONSE_BY_JOB_SQL, job_id.as_bytes()).await
    845 }
    846 
    847 pub(crate) async fn verify_response_for_delivery_job(
    848     transaction: &mut ServiceSqliteTransaction<'_>,
    849     job_id: MycDeliveryJobId,
    850 ) -> Result<(), DeliveryOperationError> {
    851     read_response_by_job(transaction, job_id)
    852         .await
    853         .map_err(|error| match error {
    854             AtomicOperationError::Binding => DeliveryOperationError::Binding,
    855             AtomicOperationError::Storage => DeliveryOperationError::Storage,
    856         })?
    857         .map(|_| ())
    858         .ok_or(DeliveryOperationError::Binding)
    859 }
    860 
    861 async fn read_response(
    862     transaction: &mut ServiceSqliteTransaction<'_>,
    863     sql: &'static str,
    864     identity: &[u8; 32],
    865 ) -> Result<Option<MycNip46ResponseRecord>, AtomicOperationError> {
    866     let rows = sqlx::query(sql)
    867         .bind(identity.as_slice())
    868         .fetch_all(&mut *transaction)
    869         .await
    870         .map_err(|_| AtomicOperationError::Storage)?;
    871     if rows.len() > 1 {
    872         return Err(AtomicOperationError::Binding);
    873     }
    874     let Some(row) = rows.first() else {
    875         return Ok(None);
    876     };
    877     let authority_kind = ResponseAuthorityKind::parse(bounded_text(row, "authority_kind", 16)?)
    878         .ok_or(AtomicOperationError::Binding)?;
    879     let operation_id = MycSignerOperationId::from_persisted(blob32(row, "operation_id")?);
    880     let response_provider_operation_id = blob32(row, "response_provider_operation_id")?;
    881     let response_event_id = blob32(row, "response_event_id")?;
    882     let response_digest = MycDeliveryArtifactDigest::from_bytes(blob32(row, "response_sha256")?);
    883     let response_bytes = row
    884         .try_get::<Option<Vec<u8>>, _>("response_bytes")
    885         .map_err(|_| AtomicOperationError::Binding)?
    886         .ok_or(AtomicOperationError::Binding)?
    887         .into_boxed_slice();
    888     let authored_at_unix_s = positive_i64(row, "authored_at_unix_s")?;
    889     let committed_at = MycDeliveryTimeUnixMs::new(positive_i64(row, "committed_at_unix_ms")?)
    890         .map_err(|_| AtomicOperationError::Binding)?;
    891     let client_public_key = row
    892         .try_get::<Option<String>, _>("client_public_key")
    893         .map_err(|_| AtomicOperationError::Binding)?
    894         .ok_or(AtomicOperationError::Binding)?;
    895     let job_id = MycDeliveryJobId::from_persisted(blob32(row, "job_id")?);
    896     let actual_digest: [u8; 32] = Sha256::digest(&response_bytes).into();
    897     let event = validate_response_event(&response_bytes, &client_public_key, None)
    898         .map_err(|_| AtomicOperationError::Binding)?;
    899     if actual_digest != *response_digest.as_bytes()
    900         || event.id.as_bytes() != &response_event_id
    901         || event.created_at.as_secs() != authored_at_unix_s
    902     {
    903         return Err(AtomicOperationError::Binding);
    904     }
    905     let delivery_job = read_job(transaction, job_id)
    906         .await
    907         .map_err(AtomicOperationError::from)?
    908         .ok_or(AtomicOperationError::Binding)?;
    909     if delivery_job.operation_id() != Some(operation_id)
    910         || delivery_job.artifact_digest() != response_digest
    911         || delivery_job.created_at() != committed_at
    912     {
    913         return Err(AtomicOperationError::Binding);
    914     }
    915     Ok(Some(MycNip46ResponseRecord {
    916         authority_kind,
    917         operation_id,
    918         response_provider_operation_id,
    919         response_event_id,
    920         response_digest,
    921         response_bytes,
    922         authored_at_unix_s,
    923         committed_at,
    924         delivery_job,
    925     }))
    926 }
    927 
    928 fn exact_response(
    929     response: &MycNip46ResponseRecord,
    930     request: &MycNip46ResponseCommitRequest,
    931 ) -> Result<(), AtomicOperationError> {
    932     (response.authority_kind == ResponseAuthorityKind::Terminal
    933         && response.operation_id == request.completion.signer_request().operation_id()
    934         && response.response_provider_operation_id == request.response_provider_operation_id
    935         && response.response_event_id == request.response_event_id
    936         && response.response_digest == request.response_digest
    937         && response.response_bytes.as_ref() == request.response_bytes.as_ref()
    938         && response.authored_at_unix_s == request.authored_at_unix_s
    939         && response.committed_at == request.committed_at)
    940         .then_some(())
    941         .ok_or(AtomicOperationError::Binding)
    942 }
    943 
    944 fn exact_pending_response(
    945     response: &MycNip46ResponseRecord,
    946     request: &MycNip46PendingResponseCommitRequest,
    947 ) -> Result<(), AtomicOperationError> {
    948     (response.authority_kind == ResponseAuthorityKind::PendingApproval
    949         && response.operation_id == request.operation_id
    950         && response.response_provider_operation_id == request.response_provider_operation_id
    951         && response.response_event_id == request.response_event_id
    952         && response.response_digest == request.response_digest
    953         && response.response_bytes.as_ref() == request.response_bytes.as_ref()
    954         && response.authored_at_unix_s == request.authored_at_unix_s
    955         && response.committed_at == request.committed_at)
    956         .then_some(())
    957         .ok_or(AtomicOperationError::Binding)
    958 }
    959 
    960 fn validate_response_event(
    961     bytes: &[u8],
    962     client_public_key: &str,
    963     expected_responder: Option<&str>,
    964 ) -> Result<RadrootsNostrEvent, MycNip46ResponseCommitError> {
    965     if bytes.is_empty() || bytes.len() > MYC_PROVIDER_OUTPUT_MAX_BYTES {
    966         return Err(MycNip46ResponseCommitRequest::error(
    967             MycNip46ResponseCommitErrorKind::InvalidResponse,
    968         ));
    969     }
    970     let event: RadrootsNostrEvent = serde_json::from_slice(bytes).map_err(|_| {
    971         MycNip46ResponseCommitRequest::error(MycNip46ResponseCommitErrorKind::InvalidResponse)
    972     })?;
    973     let canonical = serde_json::to_vec(&event).map_err(|_| {
    974         MycNip46ResponseCommitRequest::error(MycNip46ResponseCommitErrorKind::InvalidResponse)
    975     })?;
    976     let expected_recipient = ["p", client_public_key];
    977     let valid_recipient = event.tags.len() == 1
    978         && event.tags.as_slice()[0]
    979             .as_slice()
    980             .iter()
    981             .map(String::as_str)
    982             .eq(expected_recipient);
    983     if canonical != bytes
    984         || event.kind != RadrootsNostrKind::Custom(NIP46_RPC_KIND)
    985         || event.content.is_empty()
    986         || !valid_recipient
    987         || expected_responder.is_some_and(|expected| event.pubkey.to_hex() != expected)
    988         || event.verify().is_err()
    989     {
    990         return Err(MycNip46ResponseCommitRequest::error(
    991             MycNip46ResponseCommitErrorKind::InvalidResponse,
    992         ));
    993     }
    994     Ok(event)
    995 }
    996 
    997 fn blob32(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<[u8; 32], AtomicOperationError> {
    998     row.try_get::<Option<Vec<u8>>, _>(column)
    999         .map_err(|_| AtomicOperationError::Binding)?
   1000         .ok_or(AtomicOperationError::Binding)?
   1001         .try_into()
   1002         .map_err(|_| AtomicOperationError::Binding)
   1003 }
   1004 
   1005 fn bounded_text<'row>(
   1006     row: &'row sqlx::sqlite::SqliteRow,
   1007     column: &str,
   1008     maximum_bytes: usize,
   1009 ) -> Result<&'row str, AtomicOperationError> {
   1010     row.try_get::<Option<&str>, _>(column)
   1011         .map_err(|_| AtomicOperationError::Binding)?
   1012         .filter(|value| !value.is_empty() && value.len() <= maximum_bytes)
   1013         .ok_or(AtomicOperationError::Binding)
   1014 }
   1015 
   1016 fn positive_i64(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<u64, AtomicOperationError> {
   1017     row.try_get::<i64, _>(column)
   1018         .map_err(|_| AtomicOperationError::Binding)
   1019         .and_then(|value| u64::try_from(value).map_err(|_| AtomicOperationError::Binding))
   1020         .and_then(|value| {
   1021             value
   1022                 .checked_sub(1)
   1023                 .map(|_| value)
   1024                 .ok_or(AtomicOperationError::Binding)
   1025         })
   1026 }
   1027 
   1028 fn to_i64(value: u64) -> Result<i64, AtomicOperationError> {
   1029     i64::try_from(value).map_err(|_| AtomicOperationError::Binding)
   1030 }
   1031 
   1032 fn map_transaction_error(
   1033     error: ServiceSqliteTransactionError<AtomicOperationError>,
   1034 ) -> MycStateRepositoryError {
   1035     if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
   1036         return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown);
   1037     }
   1038     let kind = match error.operation_error() {
   1039         Some(AtomicOperationError::Binding) => MycStateRepositoryErrorKind::Binding,
   1040         Some(AtomicOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction,
   1041     };
   1042     MycStateRepositoryError::new(kind)
   1043 }