state_completion.rs (32415B)
1 //! Atomic durable completion of one admitted NIP-46 operation. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use radroots_service_sqlite::ServiceSqliteTransaction; 7 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 8 use radroots_service_sqlite::{ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind}; 9 use sha2::{Digest, Sha256}; 10 use sqlx::Row; 11 12 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 13 use crate::state_repository::{ 14 MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, 15 RepositoryOperationError, require_expected_metadata, 16 }; 17 use crate::{ 18 MYC_PROVIDER_OUTPUT_MAX_BYTES, MycConnectionDecision, MycConnectionDecisionRecord, 19 MycConnectionId, MycConnectionStatus, MycConnectionTimeUnixMs, MycNip46Work, MycNip46WorkKind, 20 MycSignerCorrelationId, MycSignerOperationId, MycSignerRequestMethod, MycSignerRequestRecord, 21 MycVerifiedProviderResponse, 22 }; 23 24 const READ_REQUEST_SQL: &str = r#"SELECT 25 CASE WHEN typeof(correlation_id) = 'blob' AND length(correlation_id) = 32 26 THEN correlation_id ELSE NULL END AS correlation_id, 27 CASE WHEN typeof(method) = 'text' AND length(CAST(method AS BLOB)) BETWEEN 1 AND 32 28 THEN method ELSE NULL END AS method, 29 received_at_unix_ms 30 FROM nip46_requests WHERE operation_id = ? LIMIT 2"#; 31 32 const READ_CONNECTION_SQL: &str = r#"SELECT 33 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 34 THEN status ELSE NULL END AS status, 35 policy_generation, updated_at_unix_ms 36 FROM connections WHERE connection_id = ? LIMIT 2"#; 37 38 const READ_CONNECT_DECISION_SQL: &str = r#"SELECT 39 CASE WHEN typeof(decision) = 'text' AND length(CAST(decision AS BLOB)) <= 32 40 THEN decision ELSE NULL END AS decision, 41 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 42 THEN connection_id ELSE NULL END AS connection_id, 43 typeof(connection_id) AS connection_id_type 44 FROM nip46_request_decisions WHERE operation_id = ? LIMIT 2"#; 45 46 const READ_COMMIT_SQL: &str = r#"SELECT 47 CASE WHEN typeof(correlation_id) = 'blob' AND length(correlation_id) = 32 48 THEN correlation_id ELSE NULL END AS correlation_id, 49 CASE WHEN typeof(method) = 'text' AND length(CAST(method AS BLOB)) BETWEEN 1 AND 32 50 THEN method ELSE NULL END AS method, 51 CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 52 THEN connection_id ELSE NULL END AS connection_id, 53 typeof(connection_id) AS connection_id_type, 54 CASE WHEN typeof(session_effect) = 'text' AND length(CAST(session_effect AS BLOB)) <= 32 55 THEN session_effect ELSE NULL END AS session_effect, 56 CASE WHEN typeof(provider_operation_id) = 'blob' AND length(provider_operation_id) = 32 57 THEN provider_operation_id ELSE NULL END AS provider_operation_id, 58 typeof(provider_operation_id) AS provider_operation_id_type, 59 CASE WHEN typeof(provider_artifact_kind) = 'text' 60 AND length(CAST(provider_artifact_kind AS BLOB)) <= 16 61 THEN provider_artifact_kind ELSE NULL END AS provider_artifact_kind, 62 CASE WHEN typeof(provider_artifact_sha256) = 'blob' 63 AND length(provider_artifact_sha256) = 32 64 THEN provider_artifact_sha256 ELSE NULL END AS provider_artifact_sha256, 65 typeof(provider_artifact_sha256) AS provider_artifact_sha256_type, 66 CASE WHEN typeof(provider_artifact) = 'blob' 67 AND length(provider_artifact) BETWEEN 1 AND 1048576 68 THEN provider_artifact ELSE NULL END AS provider_artifact, 69 typeof(provider_artifact) AS provider_artifact_type, 70 CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 71 THEN reason_code ELSE NULL END AS reason_code, 72 CASE WHEN typeof(outcome) = 'text' AND length(CAST(outcome AS BLOB)) <= 16 73 THEN outcome ELSE NULL END AS outcome, 74 completed_at_unix_ms 75 FROM nip46_operation_commits WHERE operation_id = ? LIMIT 2"#; 76 77 const INSERT_COMMIT_SQL: &str = r#"INSERT INTO nip46_operation_commits ( 78 operation_id, correlation_id, method, connection_id, session_effect, 79 provider_operation_id, provider_artifact_kind, provider_artifact_sha256, 80 provider_artifact, outcome, reason_code, completed_at_unix_ms 81 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'succeeded', ?, ?)"#; 82 83 const REVOKE_CONNECTION_SQL: &str = r#"UPDATE connections 84 SET status = 'expired', updated_at_unix_ms = ?, authorized_until_unix_ms = NULL 85 WHERE connection_id = ? AND status = 'active' AND policy_generation = ?"#; 86 87 /// Closed session mutation committed with a NIP-46 operation decision. 88 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 89 pub enum MycNip46SessionEffect { 90 None, 91 ConnectionAdmitted, 92 ConnectionRevoked, 93 } 94 95 impl MycNip46SessionEffect { 96 const fn as_str(self) -> &'static str { 97 match self { 98 Self::None => "none", 99 Self::ConnectionAdmitted => "connection_admitted", 100 Self::ConnectionRevoked => "connection_revoked", 101 } 102 } 103 104 fn parse(value: &str) -> Option<Self> { 105 match value { 106 "none" => Some(Self::None), 107 "connection_admitted" => Some(Self::ConnectionAdmitted), 108 "connection_revoked" => Some(Self::ConnectionRevoked), 109 _ => None, 110 } 111 } 112 } 113 114 /// Stable construction failure classes for a NIP-46 completion commit. 115 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 116 pub enum MycNip46CommitErrorKind { 117 InvalidBinding, 118 InvalidTime, 119 } 120 121 impl MycNip46CommitErrorKind { 122 /// Returns the stable machine-readable error code. 123 #[must_use] 124 pub const fn code(self) -> &'static str { 125 match self { 126 Self::InvalidBinding => "nip46_commit_binding_invalid", 127 Self::InvalidTime => "nip46_commit_time_invalid", 128 } 129 } 130 } 131 132 /// Source-free redacted construction failure. 133 #[derive(Clone, Copy, PartialEq, Eq)] 134 pub struct MycNip46CommitError { 135 kind: MycNip46CommitErrorKind, 136 } 137 138 impl MycNip46CommitError { 139 const fn new(kind: MycNip46CommitErrorKind) -> Self { 140 Self { kind } 141 } 142 143 #[must_use] 144 pub const fn kind(self) -> MycNip46CommitErrorKind { 145 self.kind 146 } 147 148 #[must_use] 149 pub const fn code(self) -> &'static str { 150 self.kind.code() 151 } 152 } 153 154 impl fmt::Display for MycNip46CommitError { 155 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 156 formatter.write_str(match self.kind { 157 MycNip46CommitErrorKind::InvalidBinding => "NIP-46 completion binding is invalid", 158 MycNip46CommitErrorKind::InvalidTime => "NIP-46 completion time is invalid", 159 }) 160 } 161 } 162 163 impl fmt::Debug for MycNip46CommitError { 164 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 165 formatter 166 .debug_struct("MycNip46CommitError") 167 .field("kind", &self.kind) 168 .finish() 169 } 170 } 171 172 impl Error for MycNip46CommitError {} 173 174 /// One sealed, independently bound NIP-46 completion request. 175 pub struct MycNip46CommitRequest { 176 request: MycSignerRequestRecord, 177 connection_id: Option<MycConnectionId>, 178 connection_generation: Option<u64>, 179 connect_decision: Option<MycConnectionDecision>, 180 session_effect: MycNip46SessionEffect, 181 provider_operation_id: Option<[u8; 32]>, 182 artifact: Option<Box<[u8]>>, 183 artifact_sha256: Option<[u8; 32]>, 184 reason_code: &'static str, 185 completed_at: MycConnectionTimeUnixMs, 186 #[cfg(test)] 187 fail_after_session_effect: bool, 188 } 189 190 impl MycNip46CommitRequest { 191 /// Binds already-admitted work and any independently verified provider output. 192 /// 193 /// A terminal connect decision is required for connect work. Provider work 194 /// requires its exact verified result. Local work accepts neither. Only an 195 /// independently verified signed-event result is retained as an artifact; 196 /// protected provider output is deliberately not persisted here. 197 pub fn new( 198 work: &MycNip46Work, 199 connect_decision: Option<&MycConnectionDecisionRecord>, 200 provider_response: Option<&MycVerifiedProviderResponse>, 201 completed_at: MycConnectionTimeUnixMs, 202 ) -> Result<Self, MycNip46CommitError> { 203 if completed_at.get() < work.request_record().received_at().get() 204 || work 205 .connection() 206 .is_some_and(|connection| completed_at < connection.updated_at()) 207 { 208 return Err(MycNip46CommitError::new( 209 MycNip46CommitErrorKind::InvalidTime, 210 )); 211 } 212 213 let mut connection_id = work.connection().map(|connection| connection.id()); 214 let mut connection_generation = work 215 .connection() 216 .map(|connection| connection.policy_generation().get()); 217 let mut terminal_connect_decision = None; 218 let (session_effect, reason_code) = match work.kind() { 219 MycNip46WorkKind::Connect => { 220 if provider_response.is_some() { 221 return Err(MycNip46CommitError::new( 222 MycNip46CommitErrorKind::InvalidBinding, 223 )); 224 } 225 let decision = connect_decision.ok_or_else(|| { 226 MycNip46CommitError::new(MycNip46CommitErrorKind::InvalidBinding) 227 })?; 228 if decision.operation_id() != work.request_record().operation_id() 229 || !matches!( 230 decision.decision(), 231 MycConnectionDecision::Allowed | MycConnectionDecision::Denied 232 ) 233 { 234 return Err(MycNip46CommitError::new( 235 MycNip46CommitErrorKind::InvalidBinding, 236 )); 237 } 238 if completed_at < decision.decided_at() { 239 return Err(MycNip46CommitError::new( 240 MycNip46CommitErrorKind::InvalidTime, 241 )); 242 } 243 connection_id = decision.connection().map(|connection| connection.id()); 244 connection_generation = decision 245 .connection() 246 .map(|connection| connection.policy_generation().get()); 247 terminal_connect_decision = Some(decision.decision()); 248 match decision.decision() { 249 MycConnectionDecision::Allowed => { 250 if decision.connection().is_none_or(|connection| { 251 connection.status() != MycConnectionStatus::Active 252 }) { 253 return Err(MycNip46CommitError::new( 254 MycNip46CommitErrorKind::InvalidBinding, 255 )); 256 } 257 ( 258 MycNip46SessionEffect::ConnectionAdmitted, 259 "connection_admitted", 260 ) 261 } 262 MycConnectionDecision::Denied => { 263 if decision.connection().is_some() { 264 return Err(MycNip46CommitError::new( 265 MycNip46CommitErrorKind::InvalidBinding, 266 )); 267 } 268 (MycNip46SessionEffect::None, "connection_denied") 269 } 270 MycConnectionDecision::PendingApproval | MycConnectionDecision::Challenged => { 271 unreachable!("nonterminal decisions rejected above") 272 } 273 } 274 } 275 MycNip46WorkKind::Local => { 276 if connect_decision.is_some() || provider_response.is_some() { 277 return Err(MycNip46CommitError::new( 278 MycNip46CommitErrorKind::InvalidBinding, 279 )); 280 } 281 if work.method() == MycSignerRequestMethod::Logout { 282 (MycNip46SessionEffect::ConnectionRevoked, "session_revoked") 283 } else { 284 (MycNip46SessionEffect::None, "completed") 285 } 286 } 287 MycNip46WorkKind::Provider => { 288 if connect_decision.is_some() { 289 return Err(MycNip46CommitError::new( 290 MycNip46CommitErrorKind::InvalidBinding, 291 )); 292 } 293 (MycNip46SessionEffect::None, "completed") 294 } 295 }; 296 297 let provider = bind_provider_result(work, provider_response)?; 298 Ok(Self { 299 request: work.request_record().clone(), 300 connection_id, 301 connection_generation, 302 connect_decision: terminal_connect_decision, 303 session_effect, 304 provider_operation_id: provider.operation_id, 305 artifact: provider.artifact, 306 artifact_sha256: provider.artifact_sha256, 307 reason_code, 308 completed_at, 309 #[cfg(test)] 310 fail_after_session_effect: false, 311 }) 312 } 313 314 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 315 pub(crate) fn fail_after_session_effect_for_test(mut self) -> Self { 316 self.fail_after_session_effect = true; 317 self 318 } 319 } 320 321 impl fmt::Debug for MycNip46CommitRequest { 322 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 323 formatter 324 .debug_struct("MycNip46CommitRequest") 325 .field("method", &self.request.method()) 326 .field("session_effect", &self.session_effect) 327 .field("artifact", &self.artifact.as_ref().map(|_| "[redacted]")) 328 .field("identity", &"[redacted]") 329 .finish() 330 } 331 } 332 333 /// Immutable result of the atomic local completion transaction. 334 #[derive(Clone, PartialEq, Eq)] 335 pub struct MycNip46CommitRecord { 336 operation_id: MycSignerOperationId, 337 correlation_id: MycSignerCorrelationId, 338 method: MycSignerRequestMethod, 339 connection_id: Option<MycConnectionId>, 340 session_effect: MycNip46SessionEffect, 341 artifact_sha256: Option<[u8; 32]>, 342 completed_at: MycConnectionTimeUnixMs, 343 } 344 345 impl MycNip46CommitRecord { 346 #[must_use] 347 pub const fn operation_id(&self) -> MycSignerOperationId { 348 self.operation_id 349 } 350 #[must_use] 351 pub const fn correlation_id(&self) -> MycSignerCorrelationId { 352 self.correlation_id 353 } 354 #[must_use] 355 pub const fn method(&self) -> MycSignerRequestMethod { 356 self.method 357 } 358 #[must_use] 359 pub const fn connection_id(&self) -> Option<MycConnectionId> { 360 self.connection_id 361 } 362 #[must_use] 363 pub const fn session_effect(&self) -> MycNip46SessionEffect { 364 self.session_effect 365 } 366 #[must_use] 367 pub const fn artifact_sha256(&self) -> Option<&[u8; 32]> { 368 self.artifact_sha256.as_ref() 369 } 370 #[must_use] 371 pub const fn completed_at(&self) -> MycConnectionTimeUnixMs { 372 self.completed_at 373 } 374 } 375 376 impl fmt::Debug for MycNip46CommitRecord { 377 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 378 formatter 379 .debug_struct("MycNip46CommitRecord") 380 .field("method", &self.method) 381 .field("session_effect", &self.session_effect) 382 .field("artifact", &self.artifact_sha256.map(|_| "[redacted]")) 383 .field("identity", &"[redacted]") 384 .finish() 385 } 386 } 387 388 /// New or exact-replayed completion result. 389 #[derive(Clone, PartialEq, Eq)] 390 pub enum MycNip46CommitAdmission { 391 Committed(MycNip46CommitRecord), 392 ExactReplay(MycNip46CommitRecord), 393 } 394 395 impl MycNip46CommitAdmission { 396 #[must_use] 397 pub const fn record(&self) -> &MycNip46CommitRecord { 398 match self { 399 Self::Committed(record) | Self::ExactReplay(record) => record, 400 } 401 } 402 } 403 404 impl fmt::Debug for MycNip46CommitAdmission { 405 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 406 formatter.write_str(match self { 407 Self::Committed(_) => "MycNip46CommitAdmission::Committed([redacted])", 408 Self::ExactReplay(_) => "MycNip46CommitAdmission::ExactReplay([redacted])", 409 }) 410 } 411 } 412 413 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 414 impl MycStateRepository<'_> { 415 /// Atomically records the Step 147 completion component. 416 /// 417 /// This integration-branch checkpoint is not a complete production 418 /// response commit. Step 148 must compose this component with the outer 419 /// signed response, immutable relay targets, and initial outbox state in 420 /// the same transaction before RCLD-RSHR-080 may be promoted to `master`. 421 pub(crate) async fn commit_nip46_operation( 422 &self, 423 request: &MycNip46CommitRequest, 424 ) -> Result<MycNip46CommitAdmission, MycStateRepositoryError> { 425 let expected = PersistedMetadata::from(self.expected()); 426 let request = request.owned(); 427 self.host() 428 .transaction(move |transaction| { 429 Box::pin(async move { 430 require_expected_metadata(transaction, &expected) 431 .await 432 .map_err(map_repository_error)?; 433 commit_operation(transaction, &request).await 434 }) 435 }) 436 .await 437 .map_err(map_transaction_error) 438 } 439 } 440 441 impl MycNip46CommitRequest { 442 pub(crate) const fn signer_request(&self) -> &MycSignerRequestRecord { 443 &self.request 444 } 445 446 pub(crate) const fn completed_at(&self) -> MycConnectionTimeUnixMs { 447 self.completed_at 448 } 449 450 pub(crate) fn owned(&self) -> Self { 451 Self { 452 request: self.request.clone(), 453 connection_id: self.connection_id, 454 connection_generation: self.connection_generation, 455 connect_decision: self.connect_decision, 456 session_effect: self.session_effect, 457 provider_operation_id: self.provider_operation_id, 458 artifact: self.artifact.clone(), 459 artifact_sha256: self.artifact_sha256, 460 reason_code: self.reason_code, 461 completed_at: self.completed_at, 462 #[cfg(test)] 463 fail_after_session_effect: self.fail_after_session_effect, 464 } 465 } 466 } 467 468 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 469 pub(crate) enum CommitOperationError { 470 Binding, 471 Storage, 472 } 473 474 pub(crate) async fn commit_operation( 475 transaction: &mut ServiceSqliteTransaction<'_>, 476 request: &MycNip46CommitRequest, 477 ) -> Result<MycNip46CommitAdmission, CommitOperationError> { 478 verify_request(transaction, request).await?; 479 if let Some(existing) = read_commit(transaction, request).await? { 480 return Ok(MycNip46CommitAdmission::ExactReplay(existing)); 481 } 482 verify_connect_decision(transaction, request).await?; 483 verify_or_apply_session_effect(transaction, request).await?; 484 #[cfg(test)] 485 if request.fail_after_session_effect { 486 return Err(CommitOperationError::Storage); 487 } 488 let artifact_kind = if request.artifact.is_some() { 489 "signed_event" 490 } else { 491 "none" 492 }; 493 let result = sqlx::query(INSERT_COMMIT_SQL) 494 .bind(request.request.operation_id().as_bytes().as_slice()) 495 .bind(request.request.correlation_id().as_bytes().as_slice()) 496 .bind(request.request.method().as_str()) 497 .bind(request.connection_id.map(|id| id.as_bytes().to_vec())) 498 .bind(request.session_effect.as_str()) 499 .bind(request.provider_operation_id.map(|id| id.to_vec())) 500 .bind(artifact_kind) 501 .bind(request.artifact_sha256.map(|digest| digest.to_vec())) 502 .bind(request.artifact.as_deref()) 503 .bind(request.reason_code) 504 .bind(to_i64(request.completed_at.get())?) 505 .execute(&mut *transaction) 506 .await 507 .map_err(|_| CommitOperationError::Storage)?; 508 require_one(result.rows_affected())?; 509 read_commit(transaction, request) 510 .await? 511 .map(MycNip46CommitAdmission::Committed) 512 .ok_or(CommitOperationError::Binding) 513 } 514 515 async fn verify_request( 516 transaction: &mut ServiceSqliteTransaction<'_>, 517 request: &MycNip46CommitRequest, 518 ) -> Result<(), CommitOperationError> { 519 let rows = sqlx::query(READ_REQUEST_SQL) 520 .bind(request.request.operation_id().as_bytes().as_slice()) 521 .fetch_all(&mut *transaction) 522 .await 523 .map_err(|_| CommitOperationError::Storage)?; 524 let [row] = rows.as_slice() else { 525 return Err(CommitOperationError::Binding); 526 }; 527 let correlation = exact_bytes(row, "correlation_id")?; 528 let method = bounded_text(row, "method")?; 529 let received_at = positive_i64(row, "received_at_unix_ms")?; 530 (correlation == *request.request.correlation_id().as_bytes() 531 && method == request.request.method().as_str() 532 && received_at == request.request.received_at().get() 533 && request.completed_at.get() >= received_at) 534 .then_some(()) 535 .ok_or(CommitOperationError::Binding) 536 } 537 538 async fn verify_connect_decision( 539 transaction: &mut ServiceSqliteTransaction<'_>, 540 request: &MycNip46CommitRequest, 541 ) -> Result<(), CommitOperationError> { 542 let Some(expected_decision) = request.connect_decision else { 543 return Ok(()); 544 }; 545 let rows = sqlx::query(READ_CONNECT_DECISION_SQL) 546 .bind(request.request.operation_id().as_bytes().as_slice()) 547 .fetch_all(&mut *transaction) 548 .await 549 .map_err(|_| CommitOperationError::Storage)?; 550 let [row] = rows.as_slice() else { 551 return Err(CommitOperationError::Binding); 552 }; 553 let decision = bounded_text(row, "decision")?; 554 let connection = optional_exact_bytes(row, "connection_id", "connection_id_type")?; 555 (decision 556 == match expected_decision { 557 MycConnectionDecision::Allowed => "allowed", 558 MycConnectionDecision::Denied => "denied", 559 MycConnectionDecision::PendingApproval | MycConnectionDecision::Challenged => { 560 return Err(CommitOperationError::Binding); 561 } 562 } 563 && connection == request.connection_id.map(|id| *id.as_bytes())) 564 .then_some(()) 565 .ok_or(CommitOperationError::Binding) 566 } 567 568 async fn verify_or_apply_session_effect( 569 transaction: &mut ServiceSqliteTransaction<'_>, 570 request: &MycNip46CommitRequest, 571 ) -> Result<(), CommitOperationError> { 572 let Some(connection_id) = request.connection_id else { 573 return (request.connection_generation.is_none() 574 && request.session_effect == MycNip46SessionEffect::None) 575 .then_some(()) 576 .ok_or(CommitOperationError::Binding); 577 }; 578 let generation = request 579 .connection_generation 580 .ok_or(CommitOperationError::Binding)?; 581 let rows = sqlx::query(READ_CONNECTION_SQL) 582 .bind(connection_id.as_bytes().as_slice()) 583 .fetch_all(&mut *transaction) 584 .await 585 .map_err(|_| CommitOperationError::Storage)?; 586 let [row] = rows.as_slice() else { 587 return Err(CommitOperationError::Binding); 588 }; 589 let status = bounded_text(row, "status")?; 590 let actual_generation = positive_i64(row, "policy_generation")?; 591 let updated_at = positive_i64(row, "updated_at_unix_ms")?; 592 if actual_generation != generation || request.completed_at.get() < updated_at { 593 return Err(CommitOperationError::Binding); 594 } 595 match request.session_effect { 596 MycNip46SessionEffect::ConnectionRevoked => { 597 if status != "active" { 598 return Err(CommitOperationError::Binding); 599 } 600 let result = sqlx::query(REVOKE_CONNECTION_SQL) 601 .bind(to_i64(request.completed_at.get())?) 602 .bind(connection_id.as_bytes().as_slice()) 603 .bind(to_i64(generation)?) 604 .execute(&mut *transaction) 605 .await 606 .map_err(|_| CommitOperationError::Storage)?; 607 require_one(result.rows_affected()) 608 } 609 MycNip46SessionEffect::None | MycNip46SessionEffect::ConnectionAdmitted => (status 610 == "active") 611 .then_some(()) 612 .ok_or(CommitOperationError::Binding), 613 } 614 } 615 616 async fn read_commit( 617 transaction: &mut ServiceSqliteTransaction<'_>, 618 request: &MycNip46CommitRequest, 619 ) -> Result<Option<MycNip46CommitRecord>, CommitOperationError> { 620 let rows = sqlx::query(READ_COMMIT_SQL) 621 .bind(request.request.operation_id().as_bytes().as_slice()) 622 .fetch_all(&mut *transaction) 623 .await 624 .map_err(|_| CommitOperationError::Storage)?; 625 if rows.len() > 1 { 626 return Err(CommitOperationError::Binding); 627 } 628 let Some(row) = rows.first() else { 629 return Ok(None); 630 }; 631 let correlation = exact_bytes(row, "correlation_id")?; 632 let method = MycSignerRequestMethod::parse(bounded_text(row, "method")?) 633 .ok_or(CommitOperationError::Binding)?; 634 let connection = optional_exact_bytes(row, "connection_id", "connection_id_type")?; 635 let session_effect = MycNip46SessionEffect::parse(bounded_text(row, "session_effect")?) 636 .ok_or(CommitOperationError::Binding)?; 637 let provider_operation = 638 optional_exact_bytes(row, "provider_operation_id", "provider_operation_id_type")?; 639 let artifact_kind = bounded_text(row, "provider_artifact_kind")?; 640 let artifact_digest = optional_exact_bytes( 641 row, 642 "provider_artifact_sha256", 643 "provider_artifact_sha256_type", 644 )?; 645 let artifact = optional_bounded_blob(row, "provider_artifact", "provider_artifact_type")?; 646 let reason = bounded_text(row, "reason_code")?; 647 let outcome = bounded_text(row, "outcome")?; 648 let completed_at = positive_i64(row, "completed_at_unix_ms")?; 649 let exact = correlation == *request.request.correlation_id().as_bytes() 650 && method == request.request.method() 651 && connection == request.connection_id.map(|id| *id.as_bytes()) 652 && session_effect == request.session_effect 653 && provider_operation == request.provider_operation_id 654 && artifact_digest == request.artifact_sha256 655 && artifact.as_deref() == request.artifact.as_deref() 656 && artifact_kind 657 == if request.artifact.is_some() { 658 "signed_event" 659 } else { 660 "none" 661 } 662 && reason == request.reason_code 663 && outcome == "succeeded" 664 && completed_at == request.completed_at.get(); 665 if !exact { 666 return Err(CommitOperationError::Binding); 667 } 668 Ok(Some(MycNip46CommitRecord { 669 operation_id: request.request.operation_id(), 670 correlation_id: request.request.correlation_id(), 671 method, 672 connection_id: request.connection_id, 673 session_effect, 674 artifact_sha256: artifact_digest, 675 completed_at: request.completed_at, 676 })) 677 } 678 679 struct BoundProviderResult { 680 operation_id: Option<[u8; 32]>, 681 artifact: Option<Box<[u8]>>, 682 artifact_sha256: Option<[u8; 32]>, 683 } 684 685 fn bind_provider_result( 686 work: &MycNip46Work, 687 response: Option<&MycVerifiedProviderResponse>, 688 ) -> Result<BoundProviderResult, MycNip46CommitError> { 689 let Some(operation) = work.provider_operation() else { 690 if response.is_some() { 691 return Err(MycNip46CommitError::new( 692 MycNip46CommitErrorKind::InvalidBinding, 693 )); 694 } 695 return Ok(BoundProviderResult { 696 operation_id: None, 697 artifact: None, 698 artifact_sha256: None, 699 }); 700 }; 701 let response = response 702 .ok_or_else(|| MycNip46CommitError::new(MycNip46CommitErrorKind::InvalidBinding))?; 703 if response.operation_id() != operation.operation_id() 704 || response.correlation_id() != operation.correlation_id() 705 || response.instance() != operation.instance() 706 || response.role() != operation.role() 707 || response.capability() != operation.input().capability() 708 || !response.matches_operation(operation) 709 { 710 return Err(MycNip46CommitError::new( 711 MycNip46CommitErrorKind::InvalidBinding, 712 )); 713 } 714 let artifact: Option<Box<[u8]>> = response.signed_event_bytes().map(Box::<[u8]>::from); 715 if artifact 716 .as_deref() 717 .is_some_and(|bytes| bytes.is_empty() || bytes.len() > MYC_PROVIDER_OUTPUT_MAX_BYTES) 718 { 719 return Err(MycNip46CommitError::new( 720 MycNip46CommitErrorKind::InvalidBinding, 721 )); 722 } 723 let digest = artifact 724 .as_deref() 725 .map(|bytes| <[u8; 32]>::from(Sha256::digest(bytes))); 726 Ok(BoundProviderResult { 727 operation_id: Some(*operation.operation_id().as_bytes()), 728 artifact, 729 artifact_sha256: digest, 730 }) 731 } 732 733 fn exact_bytes( 734 row: &sqlx::sqlite::SqliteRow, 735 column: &str, 736 ) -> Result<[u8; 32], CommitOperationError> { 737 row.try_get::<Option<Vec<u8>>, _>(column) 738 .map_err(|_| CommitOperationError::Binding)? 739 .ok_or(CommitOperationError::Binding)? 740 .try_into() 741 .map_err(|_| CommitOperationError::Binding) 742 } 743 744 fn optional_exact_bytes( 745 row: &sqlx::sqlite::SqliteRow, 746 column: &str, 747 type_column: &str, 748 ) -> Result<Option<[u8; 32]>, CommitOperationError> { 749 let value = row 750 .try_get::<Option<Vec<u8>>, _>(column) 751 .map_err(|_| CommitOperationError::Binding)?; 752 let kind = row 753 .try_get::<String, _>(type_column) 754 .map_err(|_| CommitOperationError::Binding)?; 755 match (kind.as_str(), value) { 756 ("null", None) => Ok(None), 757 ("blob", Some(bytes)) => bytes 758 .try_into() 759 .map(Some) 760 .map_err(|_| CommitOperationError::Binding), 761 _ => Err(CommitOperationError::Binding), 762 } 763 } 764 765 fn optional_bounded_blob( 766 row: &sqlx::sqlite::SqliteRow, 767 column: &str, 768 type_column: &str, 769 ) -> Result<Option<Box<[u8]>>, CommitOperationError> { 770 let value = row 771 .try_get::<Option<Vec<u8>>, _>(column) 772 .map_err(|_| CommitOperationError::Binding)?; 773 let kind = row 774 .try_get::<String, _>(type_column) 775 .map_err(|_| CommitOperationError::Binding)?; 776 match (kind.as_str(), value) { 777 ("null", None) => Ok(None), 778 ("blob", Some(bytes)) 779 if !bytes.is_empty() && bytes.len() <= MYC_PROVIDER_OUTPUT_MAX_BYTES => 780 { 781 Ok(Some(bytes.into_boxed_slice())) 782 } 783 _ => Err(CommitOperationError::Binding), 784 } 785 } 786 787 fn bounded_text<'row>( 788 row: &'row sqlx::sqlite::SqliteRow, 789 column: &str, 790 ) -> Result<&'row str, CommitOperationError> { 791 row.try_get::<Option<&str>, _>(column) 792 .map_err(|_| CommitOperationError::Binding)? 793 .ok_or(CommitOperationError::Binding) 794 } 795 796 fn positive_i64(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<u64, CommitOperationError> { 797 row.try_get::<i64, _>(column) 798 .map_err(|_| CommitOperationError::Binding) 799 .and_then(|value| u64::try_from(value).map_err(|_| CommitOperationError::Binding)) 800 .and_then(|value| { 801 (value != 0) 802 .then_some(value) 803 .ok_or(CommitOperationError::Binding) 804 }) 805 } 806 807 fn to_i64(value: u64) -> Result<i64, CommitOperationError> { 808 i64::try_from(value).map_err(|_| CommitOperationError::Binding) 809 } 810 811 fn require_one(rows: u64) -> Result<(), CommitOperationError> { 812 (rows == 1) 813 .then_some(()) 814 .ok_or(CommitOperationError::Storage) 815 } 816 817 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 818 const fn map_repository_error(error: RepositoryOperationError) -> CommitOperationError { 819 match error { 820 RepositoryOperationError::Binding => CommitOperationError::Binding, 821 RepositoryOperationError::Storage => CommitOperationError::Storage, 822 } 823 } 824 825 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 826 fn map_transaction_error( 827 error: ServiceSqliteTransactionError<CommitOperationError>, 828 ) -> MycStateRepositoryError { 829 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 830 return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); 831 } 832 let kind = match error.operation_error() { 833 Some(CommitOperationError::Binding) => MycStateRepositoryErrorKind::Binding, 834 Some(CommitOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, 835 }; 836 MycStateRepositoryError::new(kind) 837 }