rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

services_hardening_presence_publication.rs (23723B)


      1 #![forbid(unsafe_code)]
      2 #![cfg(any(target_os = "linux", target_os = "macos"))]
      3 
      4 use std::{
      5     fs,
      6     os::unix::fs::PermissionsExt,
      7     path::{Path, PathBuf},
      8 };
      9 
     10 use nostr::{Keys, SecretKey};
     11 use radroots_service_host::{
     12     EntropyError, EntropySource, SystemMonotonicClock, UnixTimeSeconds, WallClock, WallClockError,
     13 };
     14 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
     15 use radroots_storage::event::SourceGeneration;
     16 use radroots_transport::BoxFuture;
     17 use rhi::{
     18     RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiConfigDocumentV1,
     19     RhiConfigProfile, RhiDecryptedIdentity, RhiEncryptedIdentityProvisioningMaterial,
     20     RhiExactPresenceSink, RhiIdentityEnvelopeBinding, RhiPreparedPresenceAttempt,
     21     RhiPresenceAttemptOutcome, RhiPresenceDesiredAuthority, RhiPresenceLeaseOwner,
     22     RhiPresenceOutboxState, RhiPresenceRetryDelayMilliseconds, RhiPresenceTargetState,
     23     RhiPresenceUnixMilliseconds, RhiStateMetadata, RhiTimeEntropyAdapters, apply_rhi_configuration,
     24     build_rhi_signed_presence_documents, initialize_rhi_state,
     25     open_rhi_state_read_write_from_config, parse_rhi_cli_v1_from, parse_rhi_config_v1,
     26     provision_rhi_encrypted_identity, resolve_rhi_runtime_context, resolve_rhi_wrapping_credential,
     27     validate_rhi_signed_presence_documents,
     28 };
     29 use sqlx::{ConnectOptions, Connection, Row, SqliteConnection, sqlite::SqliteConnectOptions};
     30 
     31 const EXAMPLE: &str = include_str!("../contracts/services_hardening/config.v1.example.toml");
     32 const CONTRACT: &str = include_str!("../contracts/services_hardening/presence_publication.v1.json");
     33 
     34 struct FixedEntropy(u8);
     35 
     36 impl EntropySource for FixedEntropy {
     37     fn fill_bytes(&self, destination: &mut [u8]) -> Result<(), EntropyError> {
     38         destination.fill(self.0);
     39         Ok(())
     40     }
     41 }
     42 
     43 struct UnavailableEntropy;
     44 
     45 impl EntropySource for UnavailableEntropy {
     46     fn fill_bytes(&self, _destination: &mut [u8]) -> Result<(), EntropyError> {
     47         Err(EntropyError::Unavailable)
     48     }
     49 }
     50 
     51 #[derive(Clone, Copy)]
     52 struct FixedWall(u64);
     53 
     54 impl WallClock for FixedWall {
     55     fn now_utc(&self) -> Result<UnixTimeSeconds, WallClockError> {
     56         Ok(UnixTimeSeconds::new(self.0))
     57     }
     58 }
     59 
     60 struct InspectingSink {
     61     database: PathBuf,
     62     expected: Box<[u8]>,
     63 }
     64 
     65 impl RhiExactPresenceSink for InspectingSink {
     66     fn submit_exact<'a>(
     67         &'a self,
     68         attempt: &'a RhiPreparedPresenceAttempt,
     69     ) -> BoxFuture<'a, RhiPresenceAttemptOutcome> {
     70         Box::pin(async move {
     71             assert_eq!(attempt.exact_signed_event_bytes(), self.expected.as_ref());
     72             let mut connection = offline_connection(&self.database).await;
     73             let row = sqlx::query(
     74                 "SELECT targets.state, outbox.exact_signed_event_bytes \
     75                  FROM presence_targets AS targets \
     76                  JOIN presence_outbox AS outbox ON outbox.outbox_id = targets.outbox_id \
     77                  WHERE targets.last_attempt_id = ? LIMIT 2",
     78             )
     79             .bind(attempt.attempt_id().as_bytes().as_slice())
     80             .fetch_one(&mut connection)
     81             .await
     82             .expect("durable pre-I/O target");
     83             assert_eq!(row.try_get::<String, _>("state").unwrap(), "submitted");
     84             assert_eq!(
     85                 row.try_get::<Vec<u8>, _>("exact_signed_event_bytes")
     86                     .unwrap(),
     87                 self.expected.as_ref()
     88             );
     89             connection.close().await.expect("sink inspection close");
     90             RhiPresenceAttemptOutcome::Accepted
     91         })
     92     }
     93 }
     94 
     95 fn runtime(root: &Path) -> rhi::RhiRuntimeContext {
     96     let invocation = parse_rhi_cli_v1_from([
     97         "rhi",
     98         "--profile",
     99         "repo-local",
    100         "--instance",
    101         "primary",
    102         "--repo-local-root",
    103         root.to_str().expect("UTF-8 root"),
    104         "run",
    105     ])
    106     .expect("invocation");
    107     resolve_rhi_runtime_context(
    108         &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
    109         &invocation,
    110     )
    111     .expect("runtime")
    112 }
    113 
    114 fn evidence(at: u64) -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) {
    115     (
    116         MigrationAppliedAtUnixSeconds::new(at).expect("migration time"),
    117         MigrationBuildIdentity::new(
    118             env!("CARGO_PKG_VERSION"),
    119             "1111111111111111111111111111111111111111",
    120             "053d0c750bf9cd683c6ea37cefe7e79617ba629f",
    121             "rustc-test",
    122             "test-target",
    123             "service-host",
    124             1,
    125             rhi::RHI_STATE_SCHEMA_VERSION,
    126             1,
    127             1,
    128             1,
    129         )
    130         .expect("build identity"),
    131     )
    132 }
    133 
    134 async fn offline_connection(database: &Path) -> SqliteConnection {
    135     let options = SqliteConnectOptions::new()
    136         .filename(database)
    137         .create_if_missing(false)
    138         .foreign_keys(false)
    139         .disable_statement_logging();
    140     SqliteConnection::connect_with(&options)
    141         .await
    142         .expect("offline connection")
    143 }
    144 
    145 fn secret() -> [u8; 32] {
    146     [1; 32]
    147 }
    148 
    149 fn configuration(runtime: &rhi::RhiRuntimeContext, public_key: &str) -> RhiConfigDocumentV1 {
    150     configuration_from(runtime, public_key, EXAMPLE)
    151 }
    152 
    153 fn configuration_from(
    154     runtime: &rhi::RhiRuntimeContext,
    155     public_key: &str,
    156     source: &str,
    157 ) -> RhiConfigDocumentV1 {
    158     let source = source
    159         .replace(
    160             "/var/lib/radroots/services/rhi/default/secrets/service.identity.ncrypt",
    161             runtime
    162                 .identity_path()
    163                 .to_str()
    164                 .expect("UTF-8 identity path"),
    165         )
    166         .replace(&"2".repeat(64), public_key);
    167     parse_rhi_config_v1(source.as_bytes(), RhiConfigProfile::RepoLocal).expect("configuration")
    168 }
    169 
    170 #[tokio::test]
    171 async fn expired_stale_generation_is_unknown_and_superseded_before_new_bytes_commit() {
    172     let root = tempfile::tempdir().expect("root");
    173     let runtime = runtime(root.path());
    174     fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
    175     fs::set_permissions(
    176         runtime.context().paths().state(),
    177         fs::Permissions::from_mode(0o700),
    178     )
    179     .expect("state mode");
    180     let public_key = Keys::new(SecretKey::from_slice(&secret()).expect("secret"))
    181         .public_key()
    182         .to_hex();
    183     let original = configuration(&runtime, &public_key);
    184     let metadata = RhiStateMetadata::new(
    185         &runtime,
    186         &original,
    187         SourceGeneration::new([0x5b; 32]).expect("source generation"),
    188         1_725_000_000_000,
    189     )
    190     .expect("metadata");
    191     let (first_at, first_build) = evidence(1_725_000_000);
    192     initialize_rhi_state(&runtime, &metadata, first_at, &first_build)
    193         .await
    194         .expect("initialize");
    195     let identity = provision(&runtime, &original, &metadata);
    196     let original_authority =
    197         RhiPresenceDesiredAuthority::from_config(&original).expect("original authority");
    198     let host = open_rhi_state_read_write_from_config(&runtime, &original, first_at, &first_build)
    199         .await
    200         .expect("writer");
    201     let original_desired = host
    202         .repositories()
    203         .desired_presence()
    204         .commit(&original_authority)
    205         .await
    206         .expect("original desired");
    207     let original_documents = build_rhi_signed_presence_documents(
    208         original_desired,
    209         &original_authority,
    210         &identity,
    211         UnixTimeSeconds::new(1_725_000_100),
    212         &FixedEntropy(0x92),
    213     )
    214     .expect("original documents");
    215     host.repositories()
    216         .presence_outbox()
    217         .commit_signed_presence(&original_documents, millis(1_725_000_100_000))
    218         .await
    219         .expect("original outbox");
    220     let lease = host
    221         .repositories()
    222         .presence_outbox()
    223         .claim_next_presence(
    224             RhiPresenceLeaseOwner::from_bytes([0x41; 16]).expect("owner"),
    225             millis(1_725_000_101_000),
    226         )
    227         .await
    228         .expect("claim")
    229         .expect("work");
    230     let submitted = host
    231         .repositories()
    232         .presence_outbox()
    233         .prepare_next_presence_target(lease, millis(1_725_000_101_000))
    234         .await
    235         .expect("submitted");
    236     drop(submitted);
    237     host.close().await.expect("close original host");
    238 
    239     let changed_source = EXAMPLE.replace("profile = true", "profile = false");
    240     let changed = configuration_from(&runtime, &public_key, &changed_source);
    241     let (second_at, second_build) = evidence(1_725_000_001);
    242     apply_rhi_configuration(&runtime, &original, &changed, second_at, &second_build)
    243         .await
    244         .expect("apply changed configuration");
    245     let host = open_rhi_state_read_write_from_config(&runtime, &changed, second_at, &second_build)
    246         .await
    247         .expect("changed writer");
    248     let changed_authority =
    249         RhiPresenceDesiredAuthority::from_config(&changed).expect("changed authority");
    250     let changed_desired = host
    251         .repositories()
    252         .desired_presence()
    253         .commit(&changed_authority)
    254         .await
    255         .expect("changed desired");
    256     assert_eq!(changed_desired.state().generation(), 2);
    257     let changed_documents = build_rhi_signed_presence_documents(
    258         changed_desired,
    259         &changed_authority,
    260         &identity,
    261         UnixTimeSeconds::new(1_725_000_200),
    262         &FixedEntropy(0x93),
    263     )
    264     .expect("changed documents");
    265     let blocked = host
    266         .repositories()
    267         .presence_outbox()
    268         .commit_signed_presence(&changed_documents, millis(1_725_000_200_000))
    269         .await
    270         .expect_err("active stale lease blocks replacement");
    271     assert_eq!(
    272         blocked.kind(),
    273         rhi::RhiPresencePublicationErrorKind::NotReady
    274     );
    275 
    276     let recovery_adapters = RhiTimeEntropyAdapters::new(
    277         FixedWall(1_725_000_120),
    278         SystemMonotonicClock::new(),
    279         UnavailableEntropy,
    280     );
    281     assert!(
    282         host.repositories()
    283             .presence_outbox()
    284             .recover_one_expired_presence(&recovery_adapters, millis(1_725_000_120_000))
    285             .await
    286             .expect("stale recovery does not sample retry entropy")
    287     );
    288     assert!(
    289         host.repositories()
    290             .presence_outbox()
    291             .commit_signed_presence(&changed_documents, millis(1_725_000_200_000))
    292             .await
    293             .expect("new generation commit")
    294             .changed()
    295     );
    296     host.close().await.expect("close changed host");
    297 
    298     let mut disabled_table = changed_source
    299         .parse::<toml::Table>()
    300         .expect("changed configuration TOML");
    301     let disabled_presence = disabled_table
    302         .get_mut("presence")
    303         .and_then(toml::Value::as_table_mut)
    304         .expect("presence table");
    305     disabled_presence.insert("enabled".to_owned(), toml::Value::Boolean(false));
    306     disabled_presence.insert("profile".to_owned(), toml::Value::Boolean(false));
    307     disabled_presence.insert(
    308         "application_handler".to_owned(),
    309         toml::Value::Boolean(false),
    310     );
    311     disabled_presence.remove("target_relay_ids");
    312     let disabled_source = toml::to_string(&disabled_table).expect("disabled configuration TOML");
    313     let disabled = configuration_from(&runtime, &public_key, &disabled_source);
    314     let (third_at, third_build) = evidence(1_725_000_002);
    315     apply_rhi_configuration(&runtime, &changed, &disabled, third_at, &third_build)
    316         .await
    317         .expect("disable presence");
    318     let host = open_rhi_state_read_write_from_config(&runtime, &disabled, third_at, &third_build)
    319         .await
    320         .expect("disabled writer");
    321     let disabled_authority =
    322         RhiPresenceDesiredAuthority::from_config(&disabled).expect("disabled authority");
    323     let disabled_desired = host
    324         .repositories()
    325         .desired_presence()
    326         .commit(&disabled_authority)
    327         .await
    328         .expect("disabled desired");
    329     assert_eq!(disabled_desired.state().generation(), 3);
    330     let disabled_documents = build_rhi_signed_presence_documents(
    331         disabled_desired,
    332         &disabled_authority,
    333         &identity,
    334         UnixTimeSeconds::new(1_725_000_300),
    335         &UnavailableEntropy,
    336     )
    337     .expect("disabled document inventory consumes no entropy");
    338     assert!(disabled_documents.documents().is_empty());
    339     assert!(
    340         host.repositories()
    341             .presence_outbox()
    342             .commit_signed_presence(&disabled_documents, millis(1_725_000_300_000))
    343             .await
    344             .expect("disabled generation commit")
    345             .changed()
    346     );
    347     assert!(
    348         !host
    349             .repositories()
    350             .presence_outbox()
    351             .commit_signed_presence(&disabled_documents, millis(1_725_000_300_000))
    352             .await
    353             .expect("disabled exact replay")
    354             .changed()
    355     );
    356     host.close().await.expect("close disabled host");
    357 
    358     let mut connection = offline_connection(runtime.artifacts().state_database()).await;
    359     let stale: (String, String) = sqlx::query_as(
    360         "SELECT outbox.state, attempts.outcome FROM presence_outbox AS outbox \
    361          JOIN presence_attempts AS attempts ON attempts.outbox_id = outbox.outbox_id \
    362          WHERE outbox.desired_generation = 1 AND outbox.document_kind = 'service_profile'",
    363     )
    364     .fetch_one(&mut connection)
    365     .await
    366     .expect("stale durable result");
    367     assert_eq!(stale, ("superseded".to_owned(), "unknown".to_owned()));
    368     let disabled_prior_state: String = sqlx::query_scalar(
    369         "SELECT state FROM presence_outbox WHERE desired_generation = 2 LIMIT 2",
    370     )
    371     .fetch_one(&mut connection)
    372     .await
    373     .expect("disabled prior generation");
    374     assert_eq!(disabled_prior_state, "superseded");
    375     let disabled_count: i64 =
    376         sqlx::query_scalar("SELECT COUNT(*) FROM presence_outbox WHERE desired_generation = 3")
    377             .fetch_one(&mut connection)
    378             .await
    379             .expect("disabled outbox count");
    380     assert_eq!(disabled_count, 0);
    381     connection.close().await.expect("offline close");
    382 }
    383 
    384 fn provision(
    385     runtime: &rhi::RhiRuntimeContext,
    386     configuration: &RhiConfigDocumentV1,
    387     metadata: &RhiStateMetadata,
    388 ) -> RhiDecryptedIdentity {
    389     fs::create_dir_all(runtime.context().paths().secrets()).expect("secrets directory");
    390     fs::set_permissions(
    391         runtime.context().paths().secrets(),
    392         fs::Permissions::from_mode(0o700),
    393     )
    394     .expect("secrets mode");
    395     let credential_path = runtime
    396         .context()
    397         .paths()
    398         .secrets()
    399         .join("service_wrapping_key");
    400     fs::write(&credential_path, [0x81; 32]).expect("credential");
    401     fs::set_permissions(&credential_path, fs::Permissions::from_mode(0o600))
    402         .expect("credential mode");
    403     let binding = RhiIdentityEnvelopeBinding::from_configuration(configuration, metadata)
    404         .expect("identity binding");
    405     let credential =
    406         resolve_rhi_wrapping_credential(runtime, &binding).expect("wrapping credential");
    407     provision_rhi_encrypted_identity(
    408         &binding,
    409         &credential,
    410         RhiEncryptedIdentityProvisioningMaterial::new(secret(), [0x42; 32], [0x43; 24], [0x44; 24])
    411             .expect("provisioning material"),
    412     )
    413     .expect("provision identity")
    414 }
    415 
    416 fn millis(value: u64) -> RhiPresenceUnixMilliseconds {
    417     RhiPresenceUnixMilliseconds::new(value).expect("milliseconds")
    418 }
    419 
    420 #[tokio::test]
    421 async fn exact_bytes_are_durable_before_io_and_unknown_recovery_retries_unchanged() {
    422     let root = tempfile::tempdir().expect("root");
    423     let runtime = runtime(root.path());
    424     fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
    425     fs::set_permissions(
    426         runtime.context().paths().state(),
    427         fs::Permissions::from_mode(0o700),
    428     )
    429     .expect("state mode");
    430     let public_key = Keys::new(SecretKey::from_slice(&secret()).expect("secret"))
    431         .public_key()
    432         .to_hex();
    433     let configuration = configuration(&runtime, &public_key);
    434     let metadata = RhiStateMetadata::new(
    435         &runtime,
    436         &configuration,
    437         SourceGeneration::new([0x5a; 32]).expect("source generation"),
    438         1_725_000_000_000,
    439     )
    440     .expect("metadata");
    441     let (applied_at, build) = evidence(1_725_000_000);
    442     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    443         .await
    444         .expect("initialize");
    445     let identity = provision(&runtime, &configuration, &metadata);
    446     let authority = RhiPresenceDesiredAuthority::from_config(&configuration).expect("authority");
    447     let host = open_rhi_state_read_write_from_config(&runtime, &configuration, applied_at, &build)
    448         .await
    449         .expect("writer");
    450     let desired = host
    451         .repositories()
    452         .desired_presence()
    453         .commit(&authority)
    454         .await
    455         .expect("desired state");
    456     let authored = UnixTimeSeconds::new(1_725_000_100);
    457     let invalid_time = build_rhi_signed_presence_documents(
    458         desired,
    459         &authority,
    460         &identity,
    461         UnixTimeSeconds::new(i64::MAX as u64 + 1),
    462         &FixedEntropy(0x91),
    463     )
    464     .expect_err("unrepresentable authored time");
    465     assert_eq!(
    466         invalid_time.kind(),
    467         rhi::RhiPresencePublicationErrorKind::InvalidInput
    468     );
    469     let documents = build_rhi_signed_presence_documents(
    470         desired,
    471         &authority,
    472         &identity,
    473         authored,
    474         &FixedEntropy(0x91),
    475     )
    476     .expect("signed presence");
    477     validate_rhi_signed_presence_documents(&documents, &authority)
    478         .expect("independently verified signed presence");
    479     let exact: Vec<Box<[u8]>> = documents
    480         .documents()
    481         .iter()
    482         .map(|document| document.signed_event_bytes().into())
    483         .collect();
    484     let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("machine contract");
    485     assert_eq!(contract["schema"], "radroots.rhi.presence-publication");
    486     assert_eq!(contract["schema_version"], 1);
    487     assert_eq!(contract["contract_version"], 1);
    488     assert_eq!(contract["step"], 205);
    489     assert_eq!(
    490         contract["reference_vector"]["profile"]["exact_signed_event_json"],
    491         String::from_utf8_lossy(&exact[0]).as_ref()
    492     );
    493     assert_eq!(
    494         contract["reference_vector"]["application_handler"]["exact_signed_event_json"],
    495         String::from_utf8_lossy(&exact[1]).as_ref()
    496     );
    497     assert_eq!(
    498         contract["reference_vector"]["profile"]["event_id"],
    499         lower_hex(documents.documents()[0].event_id())
    500     );
    501     assert_eq!(
    502         contract["reference_vector"]["application_handler"]["event_id"],
    503         lower_hex(documents.documents()[1].event_id())
    504     );
    505     assert_eq!(
    506         contract["reference_vector"]["profile"]["exact_signed_event_sha256"],
    507         lower_hex(documents.documents()[0].signed_event_sha256())
    508     );
    509     assert_eq!(
    510         contract["reference_vector"]["application_handler"]["exact_signed_event_sha256"],
    511         lower_hex(documents.documents()[1].signed_event_sha256())
    512     );
    513     let committed_at = millis(1_725_000_100_000);
    514     let committed = host
    515         .repositories()
    516         .presence_outbox()
    517         .commit_signed_presence(&documents, committed_at)
    518         .await
    519         .expect("durable exact presence");
    520     assert!(committed.changed());
    521     assert_eq!(committed.document_count(), 2);
    522 
    523     assert!(
    524         !host
    525             .repositories()
    526             .presence_outbox()
    527             .commit_signed_presence(&documents, committed_at)
    528             .await
    529             .expect("exact replay")
    530             .changed()
    531     );
    532 
    533     let adapters = RhiTimeEntropyAdapters::new(
    534         FixedWall(1_725_000_100),
    535         SystemMonotonicClock::new(),
    536         FixedEntropy(0xff),
    537     );
    538     let sink = InspectingSink {
    539         database: runtime.artifacts().state_database().to_path_buf(),
    540         expected: exact[0].clone(),
    541     };
    542     let first = host
    543         .repositories()
    544         .presence_outbox()
    545         .execute_next_presence(
    546             RhiPresenceLeaseOwner::from_bytes([0x11; 16]).expect("owner"),
    547             &adapters,
    548             &sink,
    549         )
    550         .await
    551         .expect("execute first")
    552         .expect("first work");
    553     assert_eq!(first.outbox_state(), RhiPresenceOutboxState::Complete);
    554     assert_eq!(first.target_state(), RhiPresenceTargetState::Accepted);
    555 
    556     let repository = host.repositories().presence_outbox();
    557     let started = millis(1_725_000_101_000);
    558     let second_lease = repository
    559         .claim_next_presence(
    560             RhiPresenceLeaseOwner::from_bytes([0x22; 16]).expect("owner"),
    561             started,
    562         )
    563         .await
    564         .expect("claim second")
    565         .expect("second work");
    566     let abandoned = repository
    567         .prepare_next_presence_target(second_lease, started)
    568         .await
    569         .expect("prepare second");
    570     assert_eq!(abandoned.exact_signed_event_bytes(), exact[1].as_ref());
    571     drop(abandoned);
    572     host.close().await.expect("close after cancellation");
    573 
    574     let host = open_rhi_state_read_write_from_config(&runtime, &configuration, applied_at, &build)
    575         .await
    576         .expect("reopen writer");
    577     let recovery_now = millis(1_725_000_120_000);
    578     assert!(
    579         host.repositories()
    580             .presence_outbox()
    581             .recover_one_expired_presence(&adapters, recovery_now)
    582             .await
    583             .expect("unknown recovery")
    584     );
    585     let retry_at = millis(1_725_000_150_000);
    586     let lease = host
    587         .repositories()
    588         .presence_outbox()
    589         .claim_next_presence(
    590             RhiPresenceLeaseOwner::from_bytes([0x33; 16]).expect("owner"),
    591             retry_at,
    592         )
    593         .await
    594         .expect("retry claim")
    595         .expect("retry work");
    596     let retried = host
    597         .repositories()
    598         .presence_outbox()
    599         .prepare_next_presence_target(lease, retry_at)
    600         .await
    601         .expect("retry prepare");
    602     assert_eq!(retried.exact_signed_event_bytes(), exact[1].as_ref());
    603     let accepted = host
    604         .repositories()
    605         .presence_outbox()
    606         .record_presence_outcome(
    607             &retried,
    608             millis(1_725_000_150_001),
    609             RhiPresenceAttemptOutcome::Accepted,
    610             RhiPresenceRetryDelayMilliseconds::new(0).expect("zero delay"),
    611         )
    612         .await
    613         .expect("retry accepted");
    614     assert_eq!(accepted.attempt_number(), 2);
    615     assert_eq!(accepted.outbox_state(), RhiPresenceOutboxState::Complete);
    616     host.close().await.expect("final close");
    617 
    618     let mut connection = offline_connection(runtime.artifacts().state_database()).await;
    619     let counts = sqlx::query(
    620         "SELECT (SELECT COUNT(*) FROM presence_outbox) AS outboxes, \
    621                 (SELECT COUNT(*) FROM presence_targets) AS targets, \
    622                 (SELECT COUNT(*) FROM presence_attempts) AS attempts",
    623     )
    624     .fetch_one(&mut connection)
    625     .await
    626     .expect("workflow counts");
    627     assert_eq!(counts.try_get::<i64, _>("outboxes").unwrap(), 2);
    628     assert_eq!(counts.try_get::<i64, _>("targets").unwrap(), 4);
    629     assert_eq!(counts.try_get::<i64, _>("attempts").unwrap(), 3);
    630     let bytes: Vec<Vec<u8>> =
    631         sqlx::query("SELECT exact_signed_event_bytes FROM presence_outbox ORDER BY document_kind")
    632             .fetch_all(&mut connection)
    633             .await
    634             .expect("exact bytes")
    635             .into_iter()
    636             .map(|row| row.try_get("exact_signed_event_bytes").unwrap())
    637             .collect();
    638     assert_eq!(bytes, vec![exact[1].to_vec(), exact[0].to_vec()]);
    639     connection.close().await.expect("offline close");
    640 }
    641 
    642 fn lower_hex(bytes: &[u8]) -> String {
    643     const DIGITS: &[u8; 16] = b"0123456789abcdef";
    644     let mut rendered = String::with_capacity(bytes.len() * 2);
    645     for byte in bytes {
    646         rendered.push(char::from(DIGITS[usize::from(byte >> 4)]));
    647         rendered.push(char::from(DIGITS[usize::from(byte & 0x0f)]));
    648     }
    649     rendered
    650 }