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 }