nip46_wave_080_b.rs (25948B)
1 //! Native test-only transaction qualification for RCLD-RSHR-080 wave 080-b. 2 3 use std::{error::Error as _, fs, os::unix::fs::PermissionsExt}; 4 5 use nostr::{JsonUtil as _, Kind, Tag, Timestamp, UnsignedEvent as NostrUnsignedEvent}; 6 use radroots_nostr_connect::message::Request; 7 use sha2::{Digest, Sha256}; 8 use sqlx::{ConnectOptions as _, Connection as _, Row as _, sqlite::SqliteConnectOptions}; 9 10 use crate::{ 11 MycAuditCorrelationId, MycConnectionAdmissionPolicy, MycConnectionNonce, 12 MycConnectionOperatorDecision, MycConnectionPolicyGeneration, MycLocalSignerUntrustedResponse, 13 MycNip46CommitAdmission, MycNip46CommitRequest, MycNip46ResponseCommitAdmission, 14 MycNip46ResponseCommitRequest, MycNip46SessionEffect, MycProviderCorrelationId, 15 MycProviderDeadlineUnixMs, MycProviderOperation, MycProviderOperationId, 16 MycProviderOperationInput, MycProviderResponseObservedAtUnixMs, MycProviderRole, 17 MycRateRelayId, MycSignerRequestAdmission, MycSignerRequestMethod, MycStateRepositoryErrorKind, 18 initialize_myc_state, open_myc_state_read_write, prepare_myc_nip46_work, 19 provider_local_signer::{ProtectedWireHex, WireProviderResult}, 20 }; 21 22 use super::nip46_wave_080_a::{ 23 Encryption, OBSERVED_AT_SECONDS, PROVIDER_DEADLINE_MS, RECEIVED_AT_MS, configuration, 24 connect_request, connection_time, keys, metadata, migration_evidence, permissions, 25 prepared_request, runtime, unsigned_sign_event, untrusted_response, 26 }; 27 use crate::provider_verification::verify_encrypted_provider_response; 28 use crate::state_response::MycNip46PendingResponseCommitRequest; 29 30 pub(crate) async fn active_connection( 31 repository: &crate::MycStateRepository<'_>, 32 config: &crate::MycConfigDocumentV1, 33 ) -> ( 34 crate::MycNip46Work, 35 crate::MycConnectionRecord, 36 crate::MycConnectionDecisionRecord, 37 ) { 38 let prepared = prepared_request( 39 config, 40 &keys(10), 41 "step147-connect", 42 connect_request(), 43 Encryption::Nip44V2, 44 OBSERVED_AT_SECONDS + 20, 45 20, 46 RECEIVED_AT_MS + 20, 47 ); 48 let admitted = repository 49 .admit_signer_request(prepared.signer_request()) 50 .await 51 .expect("connect request admission"); 52 let work = prepare_myc_nip46_work( 53 prepared, 54 admitted.record().clone(), 55 None, 56 config.provider_contract(), 57 connection_time(RECEIVED_AT_MS + 1_000), 58 None, 59 ) 60 .expect("connect work"); 61 let connection_request = work 62 .connection_admission_request( 63 MycConnectionPolicyGeneration::new(1).expect("policy generation"), 64 MycConnectionNonce::from_injected_entropy([0x51; 32]), 65 connection_time(RECEIVED_AT_MS + 1_000), 66 None, 67 MycConnectionAdmissionPolicy::ExplicitApproval, 68 MycRateRelayId::new("primary").expect("relay ID"), 69 ) 70 .expect("connection request"); 71 let pending = repository 72 .admit_connection(&connection_request) 73 .await 74 .expect("pending connection"); 75 let record = pending.record().expect("pending record"); 76 let active = repository 77 .decide_pending_connection( 78 record.operation_id(), 79 record.connection().expect("connection").id(), 80 MycConnectionPolicyGeneration::new(1).expect("policy generation"), 81 connection_time(RECEIVED_AT_MS + 2_000), 82 MycAuditCorrelationId::new([0x52; 32]), 83 MycConnectionOperatorDecision::Approve { 84 granted_permissions: permissions(), 85 authorized_until: Some(connection_time(RECEIVED_AT_MS + 100_000)), 86 }, 87 ) 88 .await 89 .expect("approved connection"); 90 let decision = repository 91 .read_connection_decision(record.operation_id()) 92 .await 93 .expect("terminal connect decision"); 94 (work, active, decision) 95 } 96 97 async fn pending_connection( 98 repository: &crate::MycStateRepository<'_>, 99 config: &crate::MycConfigDocumentV1, 100 ) -> (crate::MycNip46Work, crate::MycConnectionDecisionRecord) { 101 let prepared = prepared_request( 102 config, 103 &keys(10), 104 "pending-response-connect", 105 connect_request(), 106 Encryption::Nip44V2, 107 OBSERVED_AT_SECONDS + 40, 108 40, 109 RECEIVED_AT_MS + 40, 110 ); 111 let admitted = repository 112 .admit_signer_request(prepared.signer_request()) 113 .await 114 .expect("pending request admission"); 115 let work = prepare_myc_nip46_work( 116 prepared, 117 admitted.record().clone(), 118 None, 119 config.provider_contract(), 120 connection_time(RECEIVED_AT_MS + 4_000), 121 None, 122 ) 123 .expect("pending connect work"); 124 let connection_request = work 125 .connection_admission_request( 126 MycConnectionPolicyGeneration::new(1).expect("policy generation"), 127 crate::MycConnectionNonce::from_injected_entropy([0x81; 32]), 128 connection_time(RECEIVED_AT_MS + 4_000), 129 None, 130 MycConnectionAdmissionPolicy::ExplicitApproval, 131 MycRateRelayId::new("primary").expect("relay"), 132 ) 133 .expect("pending connection request"); 134 let admission = repository 135 .admit_connection(&connection_request) 136 .await 137 .expect("pending connection admission"); 138 let decision = admission.record().expect("pending decision").clone(); 139 assert_eq!( 140 decision.decision(), 141 crate::MycConnectionDecision::PendingApproval 142 ); 143 (work, decision) 144 } 145 146 fn pending_response_request( 147 config: &crate::MycConfigDocumentV1, 148 work: &crate::MycNip46Work, 149 decision: &crate::MycConnectionDecisionRecord, 150 ) -> (MycNip46PendingResponseCommitRequest, Vec<u8>) { 151 let unsigned = NostrUnsignedEvent::new( 152 keys(2).public_key(), 153 Timestamp::from_secs(OBSERVED_AT_SECONDS + 41), 154 Kind::Custom(24_133), 155 vec![Tag::public_key(keys(10).public_key())], 156 "encrypted-pending-response", 157 ); 158 let operation = MycProviderOperation::new( 159 config 160 .provider_contract() 161 .binding(MycProviderRole::Transport) 162 .expect("transport binding"), 163 MycProviderOperationId::from_bytes([0xa1; 32]), 164 MycProviderCorrelationId::from_bytes([0xa2; 32]), 165 MycProviderDeadlineUnixMs::new(PROVIDER_DEADLINE_MS).expect("provider deadline"), 166 MycProviderOperationInput::sign_event(unsigned.as_json().as_bytes()) 167 .expect("response signing input"), 168 ) 169 .expect("pending response operation"); 170 let signed = unsigned 171 .sign_with_keys(&keys(2)) 172 .expect("signed pending response"); 173 let bytes = serde_json::to_vec(&signed).expect("canonical pending response"); 174 let verified = verify_encrypted_provider_response( 175 config 176 .provider_contract() 177 .binding(MycProviderRole::Transport) 178 .expect("transport binding"), 179 &operation, 180 MycProviderResponseObservedAtUnixMs::new(RECEIVED_AT_MS + 41_001).expect("response time"), 181 WireProviderResult::SignEvent { 182 payload_hex: ProtectedWireHex::from_bytes(&bytes), 183 }, 184 ) 185 .expect("verified pending response"); 186 let request = MycNip46PendingResponseCommitRequest::new( 187 work, 188 decision, 189 &operation, 190 &verified, 191 crate::MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 41_003).expect("commit time"), 192 ) 193 .expect("pending response request"); 194 (request, bytes) 195 } 196 197 #[tokio::test] 198 async fn pending_approval_response_and_delivery_job_commit_atomically_and_replay_exactly() { 199 let directory = tempfile::tempdir().expect("temporary root"); 200 let runtime = runtime(directory.path()); 201 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 202 fs::set_permissions( 203 runtime.context().paths().state(), 204 fs::Permissions::from_mode(0o700), 205 ) 206 .expect("state mode"); 207 let metadata = metadata(&runtime); 208 let config = configuration(); 209 let (applied_at, build) = migration_evidence(); 210 initialize_myc_state(&runtime, &metadata, applied_at, &build) 211 .await 212 .expect("state initialization"); 213 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 214 .await 215 .expect("writable state"); 216 let repository = host.repository(); 217 let (work, decision) = pending_connection(&repository, &config).await; 218 let (request, response_bytes) = pending_response_request(&config, &work, &decision); 219 let debug = format!("{request:?}"); 220 assert!(!debug.contains("encrypted-pending-response")); 221 222 let rollback = repository 223 .commit_nip46_pending_response(&request.fail_after_response_for_test()) 224 .await 225 .expect_err("pending response without delivery must roll back"); 226 assert_eq!( 227 rollback.kind(), 228 crate::MycStateRepositoryErrorKind::Transaction 229 ); 230 let committed = repository 231 .commit_nip46_pending_response(&request) 232 .await 233 .expect("pending response commit"); 234 assert_eq!(committed.signed_response_bytes(), response_bytes); 235 assert_eq!( 236 committed.operation_id(), 237 work.request_record().operation_id() 238 ); 239 let replay = repository 240 .commit_nip46_pending_response(&request) 241 .await 242 .expect("exact pending response replay"); 243 assert_eq!(replay, committed); 244 let by_job = repository 245 .read_nip46_response(committed.delivery_job().id()) 246 .await 247 .expect("pending response read") 248 .expect("retained pending response"); 249 assert_eq!(by_job, committed); 250 repository 251 .verify_delivery_invariants() 252 .await 253 .expect("pending response satisfies delivery invariants"); 254 repository 255 .recover_delivery_state( 256 crate::MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 41_004).expect("recovery time"), 257 crate::MycDeliveryRecoveryEntropy::from_injected_entropy([0xa3; 32]), 258 ) 259 .await 260 .expect("pending response survives restart recovery"); 261 repository 262 .verify_delivery_invariants() 263 .await 264 .expect("recovered pending response satisfies delivery invariants"); 265 host.close().await.expect("host close"); 266 267 let options = SqliteConnectOptions::new() 268 .filename(runtime.artifacts().state_database()) 269 .create_if_missing(false) 270 .disable_statement_logging(); 271 let mut connection = sqlx::SqliteConnection::connect_with(&options) 272 .await 273 .expect("inspection connection"); 274 assert_eq!( 275 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_pending_responses") 276 .fetch_one(&mut connection) 277 .await 278 .expect("pending response count"), 279 1 280 ); 281 assert_eq!( 282 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_signed_responses") 283 .fetch_one(&mut connection) 284 .await 285 .expect("terminal response count"), 286 0 287 ); 288 assert_eq!( 289 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_operation_commits") 290 .fetch_one(&mut connection) 291 .await 292 .expect("terminal operation count"), 293 0 294 ); 295 assert_eq!( 296 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_jobs") 297 .fetch_one(&mut connection) 298 .await 299 .expect("delivery job count"), 300 1 301 ); 302 connection.close().await.expect("inspection close"); 303 } 304 305 pub(crate) fn atomic_response_request( 306 config: &crate::MycConfigDocumentV1, 307 completion: &MycNip46CommitRequest, 308 ) -> (MycNip46ResponseCommitRequest, Vec<u8>) { 309 let unsigned = NostrUnsignedEvent::new( 310 keys(2).public_key(), 311 Timestamp::from_secs(OBSERVED_AT_SECONDS + 3), 312 Kind::Custom(24_133), 313 vec![Tag::public_key(keys(10).public_key())], 314 "encrypted-response", 315 ); 316 let operation = MycProviderOperation::new( 317 config 318 .provider_contract() 319 .binding(MycProviderRole::Transport) 320 .expect("transport binding"), 321 MycProviderOperationId::from_bytes([0x91; 32]), 322 MycProviderCorrelationId::from_bytes([0x92; 32]), 323 MycProviderDeadlineUnixMs::new(PROVIDER_DEADLINE_MS).expect("provider deadline"), 324 MycProviderOperationInput::sign_event(unsigned.as_json().as_bytes()) 325 .expect("response signing input"), 326 ) 327 .expect("response operation"); 328 let signed = unsigned.sign_with_keys(&keys(2)).expect("signed response"); 329 let bytes = serde_json::to_vec(&signed).expect("canonical response"); 330 let verified = verify_encrypted_provider_response( 331 config 332 .provider_contract() 333 .binding(MycProviderRole::Transport) 334 .expect("transport binding"), 335 &operation, 336 MycProviderResponseObservedAtUnixMs::new(RECEIVED_AT_MS + 3_001).expect("response time"), 337 WireProviderResult::SignEvent { 338 payload_hex: ProtectedWireHex::from_bytes(&bytes), 339 }, 340 ) 341 .expect("verified response"); 342 let request = MycNip46ResponseCommitRequest::new( 343 completion, 344 &operation, 345 &verified, 346 crate::MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 3_003).expect("commit time"), 347 ) 348 .expect("atomic response request"); 349 (request, bytes) 350 } 351 352 #[tokio::test] 353 async fn verified_signed_response_commits_completion_and_outbox_exactly_once() { 354 let directory = tempfile::tempdir().expect("temporary root"); 355 let runtime = runtime(directory.path()); 356 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 357 fs::set_permissions( 358 runtime.context().paths().state(), 359 fs::Permissions::from_mode(0o700), 360 ) 361 .expect("state mode"); 362 let metadata = metadata(&runtime); 363 let config = configuration(); 364 let (applied_at, build) = migration_evidence(); 365 initialize_myc_state(&runtime, &metadata, applied_at, &build) 366 .await 367 .expect("state initialization"); 368 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 369 .await 370 .expect("writable state"); 371 let repository = host.repository(); 372 let (connect_work, active, connect_decision) = active_connection(&repository, &config).await; 373 let connect_commit = MycNip46CommitRequest::new( 374 &connect_work, 375 Some(&connect_decision), 376 None, 377 connection_time(RECEIVED_AT_MS + 2_001), 378 ) 379 .expect("connect completion"); 380 let invalid_binding = MycNip46CommitRequest::new( 381 &connect_work, 382 None, 383 None, 384 connection_time(RECEIVED_AT_MS + 2_001), 385 ) 386 .expect_err("connect completion requires its terminal decision"); 387 assert_eq!( 388 invalid_binding.kind(), 389 crate::MycNip46CommitErrorKind::InvalidBinding 390 ); 391 assert!(invalid_binding.source().is_none()); 392 assert!(!format!("{invalid_binding:?}").contains("step147-connect")); 393 let invalid_time = MycNip46CommitRequest::new( 394 &connect_work, 395 Some(&connect_decision), 396 None, 397 connection_time(RECEIVED_AT_MS), 398 ) 399 .expect_err("completion time precedes request admission"); 400 assert_eq!( 401 invalid_time.kind(), 402 crate::MycNip46CommitErrorKind::InvalidTime 403 ); 404 assert!(invalid_time.source().is_none()); 405 let before_terminal_decision = MycNip46CommitRequest::new( 406 &connect_work, 407 Some(&connect_decision), 408 None, 409 connection_time(RECEIVED_AT_MS + 1_500), 410 ) 411 .expect_err("completion time precedes terminal decision"); 412 assert_eq!( 413 before_terminal_decision.kind(), 414 crate::MycNip46CommitErrorKind::InvalidTime 415 ); 416 assert_eq!( 417 repository 418 .commit_nip46_operation(&connect_commit) 419 .await 420 .expect("connect completion commit") 421 .record() 422 .session_effect(), 423 MycNip46SessionEffect::ConnectionAdmitted 424 ); 425 let (partial_response, _) = atomic_response_request(&config, &connect_commit); 426 assert_eq!( 427 repository 428 .commit_nip46_response(&partial_response) 429 .await 430 .expect_err("completion-only legacy state is not repaired") 431 .kind(), 432 MycStateRepositoryErrorKind::Binding 433 ); 434 435 let prepared = prepared_request( 436 &config, 437 &keys(10), 438 "step147-sign", 439 Request::SignEvent(unsigned_sign_event()), 440 Encryption::Nip44V2, 441 OBSERVED_AT_SECONDS + 21, 442 21, 443 RECEIVED_AT_MS + 21, 444 ); 445 let admitted = repository 446 .admit_signer_request(prepared.signer_request()) 447 .await 448 .expect("sign request admission"); 449 assert!(matches!(admitted, MycSignerRequestAdmission::Admitted(_))); 450 let work = prepare_myc_nip46_work( 451 prepared, 452 admitted.record().clone(), 453 Some(active), 454 config.provider_contract(), 455 connection_time(RECEIVED_AT_MS + 3_000), 456 Some(MycProviderDeadlineUnixMs::new(PROVIDER_DEADLINE_MS).expect("provider deadline")), 457 ) 458 .expect("sign work"); 459 let operation = work.provider_operation().expect("provider operation"); 460 let unsigned: NostrUnsignedEvent = 461 serde_json::from_slice(operation.input().bytes().expect("canonical unsigned event")) 462 .expect("unsigned event"); 463 let signed = unsigned.sign_with_keys(&keys(3)).expect("signed event"); 464 let signed_bytes = serde_json::to_vec(&signed).expect("canonical signed event"); 465 let response: MycLocalSignerUntrustedResponse = untrusted_response( 466 operation, 467 hex::encode(operation.correlation_id().as_bytes()), 468 WireProviderResult::SignEvent { 469 payload_hex: ProtectedWireHex::from_bytes(&signed_bytes), 470 }, 471 ); 472 let verified = response 473 .verify( 474 config 475 .provider_contract() 476 .binding(MycProviderRole::User) 477 .expect("user binding"), 478 operation, 479 MycProviderResponseObservedAtUnixMs::new(RECEIVED_AT_MS + 3_001) 480 .expect("response time"), 481 ) 482 .expect("verified signed event"); 483 let request = MycNip46CommitRequest::new( 484 &work, 485 None, 486 Some(&verified), 487 connection_time(RECEIVED_AT_MS + 3_002), 488 ) 489 .expect("completion request"); 490 let request_debug = format!("{request:?}"); 491 assert!(!request_debug.contains(&hex::encode(operation.operation_id().as_bytes()))); 492 assert!(!request_debug.contains(&hex::encode(operation.correlation_id().as_bytes()))); 493 assert!(!request_debug.contains(String::from_utf8_lossy(&signed_bytes).as_ref())); 494 let (atomic_request, response_bytes) = atomic_response_request(&config, &request); 495 assert_eq!( 496 repository 497 .commit_nip46_response(&atomic_request.fail_after_completion_for_test()) 498 .await 499 .expect_err("completion-edge failure rolls back") 500 .kind(), 501 MycStateRepositoryErrorKind::Transaction 502 ); 503 assert_eq!( 504 repository 505 .commit_nip46_response(&atomic_request.fail_after_response_for_test()) 506 .await 507 .expect_err("response-edge failure rolls back") 508 .kind(), 509 MycStateRepositoryErrorKind::Transaction 510 ); 511 let committed = repository 512 .commit_nip46_response(&atomic_request) 513 .await 514 .expect("atomic response commit"); 515 assert!(matches!( 516 committed, 517 MycNip46ResponseCommitAdmission::Committed(_) 518 )); 519 assert_eq!( 520 committed.record().completion().method(), 521 MycSignerRequestMethod::SignEvent 522 ); 523 assert_eq!( 524 committed.record().completion().artifact_sha256(), 525 Some(&<[u8; 32]>::from(Sha256::digest(&signed_bytes))) 526 ); 527 assert_eq!( 528 committed.record().response().signed_response_bytes(), 529 response_bytes 530 ); 531 assert!( 532 committed 533 .record() 534 .response() 535 .delivery_job() 536 .targets() 537 .iter() 538 .all(|target| target.attempt_count() == 0 539 && target.status() == crate::MycDeliveryTargetStatus::Pending) 540 ); 541 let record_debug = format!("{:?}", committed.record()); 542 assert!(!record_debug.contains(&hex::encode( 543 committed.record().completion().operation_id().as_bytes() 544 ))); 545 assert!(!record_debug.contains(&hex::encode( 546 committed.record().completion().correlation_id().as_bytes() 547 ))); 548 assert!(!record_debug.contains(&hex::encode(Sha256::digest(&signed_bytes)))); 549 let replay = repository 550 .commit_nip46_response(&atomic_request) 551 .await 552 .expect("exact atomic replay"); 553 assert!(matches!( 554 replay, 555 MycNip46ResponseCommitAdmission::ExactReplay(_) 556 )); 557 let retry = repository 558 .read_nip46_response(committed.record().response().delivery_job().id()) 559 .await 560 .expect("response retry read") 561 .expect("retained response"); 562 assert_eq!(retry.signed_response_bytes(), response_bytes); 563 host.close().await.expect("host close"); 564 565 let options = SqliteConnectOptions::new() 566 .filename(runtime.artifacts().state_database()) 567 .create_if_missing(false) 568 .disable_statement_logging(); 569 let mut connection = sqlx::SqliteConnection::connect_with(&options) 570 .await 571 .expect("inspection connection"); 572 let row = sqlx::query( 573 "SELECT provider_artifact, provider_artifact_sha256 FROM nip46_operation_commits \ 574 WHERE method = 'sign_event'", 575 ) 576 .fetch_one(&mut connection) 577 .await 578 .expect("committed artifact"); 579 assert_eq!(row.get::<Vec<u8>, _>("provider_artifact"), signed_bytes); 580 assert_eq!( 581 row.get::<Vec<u8>, _>("provider_artifact_sha256"), 582 Sha256::digest(&signed_bytes).as_slice() 583 ); 584 assert_eq!( 585 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_jobs") 586 .fetch_one(&mut connection) 587 .await 588 .expect("delivery job count"), 589 1 590 ); 591 assert_eq!( 592 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_signed_responses") 593 .fetch_one(&mut connection) 594 .await 595 .expect("response count"), 596 1 597 ); 598 connection.close().await.expect("inspection close"); 599 } 600 601 #[tokio::test] 602 async fn session_revocation_and_completion_roll_back_together_then_replay_exactly() { 603 let directory = tempfile::tempdir().expect("temporary root"); 604 let runtime = runtime(directory.path()); 605 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 606 fs::set_permissions( 607 runtime.context().paths().state(), 608 fs::Permissions::from_mode(0o700), 609 ) 610 .expect("state mode"); 611 let metadata = metadata(&runtime); 612 let config = configuration(); 613 let (applied_at, build) = migration_evidence(); 614 initialize_myc_state(&runtime, &metadata, applied_at, &build) 615 .await 616 .expect("state initialization"); 617 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 618 .await 619 .expect("writable state"); 620 let repository = host.repository(); 621 let (_connect_work, active, _connect_decision) = active_connection(&repository, &config).await; 622 623 let prepared = prepared_request( 624 &config, 625 &keys(10), 626 "step147-logout", 627 Request::Logout, 628 Encryption::Nip44V2, 629 OBSERVED_AT_SECONDS + 22, 630 22, 631 RECEIVED_AT_MS + 22, 632 ); 633 let admitted = repository 634 .admit_signer_request(prepared.signer_request()) 635 .await 636 .expect("logout request admission"); 637 let work = crate::MycNip46Work::local_for_test( 638 admitted.record().clone(), 639 active.clone(), 640 MycSignerRequestMethod::Logout, 641 ); 642 let failed = 643 MycNip46CommitRequest::new(&work, None, None, connection_time(RECEIVED_AT_MS + 3_100)) 644 .expect("logout completion") 645 .fail_after_session_effect_for_test(); 646 assert_eq!( 647 repository 648 .commit_nip46_operation(&failed) 649 .await 650 .expect_err("injected transaction failure") 651 .kind(), 652 MycStateRepositoryErrorKind::Transaction 653 ); 654 655 let request = 656 MycNip46CommitRequest::new(&work, None, None, connection_time(RECEIVED_AT_MS + 3_100)) 657 .expect("logout completion retry"); 658 let committed = repository 659 .commit_nip46_operation(&request) 660 .await 661 .expect("rollback preserved retry authority"); 662 assert_eq!( 663 committed.record().session_effect(), 664 MycNip46SessionEffect::ConnectionRevoked 665 ); 666 assert!(matches!( 667 repository 668 .commit_nip46_operation(&request) 669 .await 670 .expect("logout replay"), 671 MycNip46CommitAdmission::ExactReplay(_) 672 )); 673 host.close().await.expect("host close"); 674 675 let options = SqliteConnectOptions::new() 676 .filename(runtime.artifacts().state_database()) 677 .create_if_missing(false) 678 .disable_statement_logging(); 679 let mut connection = sqlx::SqliteConnection::connect_with(&options) 680 .await 681 .expect("inspection connection"); 682 assert_eq!( 683 sqlx::query_scalar::<_, String>("SELECT status FROM connections WHERE connection_id = ?") 684 .bind(active.id().as_bytes().as_slice()) 685 .fetch_one(&mut connection) 686 .await 687 .expect("connection status"), 688 "expired" 689 ); 690 assert_eq!( 691 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nip46_operation_commits") 692 .fetch_one(&mut connection) 693 .await 694 .expect("completion count"), 695 1 696 ); 697 connection.close().await.expect("inspection close"); 698 }