rhi

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

services_hardening_trade_persistence.rs (17434B)


      1 #![forbid(unsafe_code)]
      2 #![cfg(any(target_os = "linux", target_os = "macos"))]
      3 
      4 use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path};
      5 
      6 use nostr::secp256k1::{Keypair, Message};
      7 use nostr::{EventId, Keys, SECP256K1};
      8 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
      9 use radroots_storage::event::SourceGeneration;
     10 use rhi::{
     11     RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION, RadrootsHostEnvironment, RadrootsPathResolver,
     12     RadrootsPlatform, RhiConfigProfile, RhiStateMetadata, RhiTradeEvidencePersistenceErrorKind,
     13     RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy,
     14     RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceObservation,
     15     admit_rhi_trade_mutation_event, initialize_rhi_state, open_rhi_state_inspection,
     16     open_rhi_state_read_write, parse_rhi_cli_v1_from, parse_rhi_config_v1,
     17     resolve_rhi_runtime_context,
     18 };
     19 use serde_json::Value;
     20 use sqlx::{ConnectOptions, Connection, SqliteConnection, sqlite::SqliteConnectOptions};
     21 
     22 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml");
     23 const CONTRACT: &str =
     24     include_str!("../contracts/services_hardening/trade_evidence_persistence.v1.json");
     25 const VECTOR: &str = include_str!("../contracts/conformance/vectors/trade_ingest_proposal.v1.json");
     26 
     27 fn runtime(root: &Path) -> rhi::RhiRuntimeContext {
     28     let invocation = parse_rhi_cli_v1_from([
     29         "rhi",
     30         "--profile",
     31         "repo-local",
     32         "--instance",
     33         "primary",
     34         "--repo-local-root",
     35         root.to_str().expect("UTF-8 root"),
     36         "run",
     37     ])
     38     .expect("invocation");
     39     resolve_rhi_runtime_context(
     40         &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
     41         &invocation,
     42     )
     43     .expect("runtime")
     44 }
     45 
     46 fn configuration() -> rhi::RhiConfigDocumentV1 {
     47     parse_rhi_config_v1(CONFIG.as_bytes(), RhiConfigProfile::RepoLocal).expect("configuration")
     48 }
     49 
     50 fn metadata(
     51     runtime: &rhi::RhiRuntimeContext,
     52     configuration: &rhi::RhiConfigDocumentV1,
     53 ) -> RhiStateMetadata {
     54     RhiStateMetadata::new(
     55         runtime,
     56         configuration,
     57         SourceGeneration::new([0x5a; 32]).expect("generation"),
     58         1_725_000_000_000,
     59     )
     60     .expect("metadata")
     61 }
     62 
     63 fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) {
     64     (
     65         MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time"),
     66         MigrationBuildIdentity::new(
     67             env!("CARGO_PKG_VERSION"),
     68             "1111111111111111111111111111111111111111",
     69             "053d0c750bf9cd683c6ea37cefe7e79617ba629f",
     70             "rustc-test",
     71             "test-target",
     72             "service-host",
     73             1,
     74             rhi::RHI_STATE_SCHEMA_VERSION,
     75             1,
     76             1,
     77             1,
     78         )
     79         .expect("build identity"),
     80     )
     81 }
     82 
     83 fn prepare_state_directory(runtime: &rhi::RhiRuntimeContext) {
     84     fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
     85     fs::set_permissions(
     86         runtime.context().paths().state(),
     87         fs::Permissions::from_mode(0o700),
     88     )
     89     .expect("state mode");
     90 }
     91 
     92 fn signed_wire(auxiliary: u8) -> Vec<u8> {
     93     let vector: Value = serde_json::from_str(VECTOR).expect("vector");
     94     let mut event: Value =
     95         serde_json::from_str(vector["raw_json"].as_str().expect("raw event")).expect("event JSON");
     96     let keys = Keys::parse("10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5")
     97         .expect("approved fixture keys");
     98     let event_id = EventId::from_hex(event["id"].as_str().expect("event id")).expect("event id");
     99     let message = Message::from_digest(event_id.to_bytes());
    100     let keypair = Keypair::from_secret_key(SECP256K1, keys.secret_key());
    101     let signature = SECP256K1.sign_schnorr_with_aux_rand(&message, &keypair, &[auxiliary; 32]);
    102     event["sig"] = signature.to_string().into();
    103     serde_json::to_vec(&event).expect("event JSON")
    104 }
    105 
    106 fn admitted(
    107     configuration: &rhi::RhiConfigDocumentV1,
    108     wire: &[u8],
    109     observed_at: u64,
    110 ) -> rhi::RhiAdmittedTradeMutationEvent {
    111     admit_rhi_trade_mutation_event(
    112         RhiTradeMutationAdmissionLimits::from_config(configuration).expect("limits"),
    113         wire,
    114         RhiTradeMutationObservedAtUnixSeconds::new(observed_at).expect("observation time"),
    115         RhiTradeMutationAuthoredTimePolicy::new(0).expect("time policy"),
    116     )
    117     .expect("admitted event")
    118 }
    119 
    120 async fn offline_connection(runtime: &rhi::RhiRuntimeContext) -> SqliteConnection {
    121     SqliteConnection::connect_with(
    122         &SqliteConnectOptions::new()
    123             .filename(runtime.artifacts().state_database())
    124             .create_if_missing(false)
    125             .disable_statement_logging(),
    126     )
    127     .await
    128     .expect("offline connection")
    129 }
    130 
    131 async fn insert_mutation(
    132     connection: &mut SqliteConnection,
    133     event: &rhi::RhiAdmittedTradeMutationEvent,
    134     canonical_content: &[u8],
    135 ) {
    136     let mutation = event.mutation();
    137     sqlx::query(
    138         r#"INSERT INTO trade_mutations (
    139             mutation_id, trade_id, contract_id, schema_version, event_kind,
    140             author_pubkey, canonical_content
    141         ) VALUES (?, ?, ?, ?, ?, ?, ?)"#,
    142     )
    143     .bind(event.mutation_id().as_bytes().as_slice())
    144     .bind(mutation.trade_id.as_bytes().as_slice())
    145     .bind(mutation.mutation_kind().contract_id())
    146     .bind(i64::from(mutation.schema_version))
    147     .bind(i64::from(event.event_kind()))
    148     .bind(mutation.author_pubkey.as_bytes().as_slice())
    149     .bind(canonical_content)
    150     .execute(connection)
    151     .await
    152     .expect("seed mutation");
    153 }
    154 
    155 fn decode_hex<const N: usize>(value: &str) -> [u8; N] {
    156     assert_eq!(value.len(), N * 2);
    157     let mut bytes = [0_u8; N];
    158     for (index, byte) in bytes.iter_mut().enumerate() {
    159         *byte = u8::from_str_radix(&value[index * 2..index * 2 + 2], 16).expect("hex byte");
    160     }
    161     bytes
    162 }
    163 
    164 #[test]
    165 fn machine_contract_freezes_three_separate_immutable_facts() {
    166     let contract: Value = serde_json::from_str(CONTRACT).expect("contract");
    167     assert_eq!(
    168         contract["schema"],
    169         "radroots.rhi.trade-evidence-persistence.v1"
    170     );
    171     assert_eq!(contract["contract_version"], 1);
    172     assert_eq!(RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION, 1);
    173     assert_eq!(
    174         contract["facts"]["canonical_mutation"]["table"],
    175         "trade_mutations"
    176     );
    177     assert_eq!(
    178         contract["facts"]["canonical_mutation"]["event_authored_time_column"],
    179         "absent_event_fact_only"
    180     );
    181     assert_eq!(contract["facts"]["signed_event"]["table"], "nostr_events");
    182     assert_eq!(
    183         contract["facts"]["signed_event"]["identity"],
    184         serde_json::json!(["verified_event_id", "verified_event_signature"])
    185     );
    186     assert_eq!(
    187         contract["facts"]["source_observation"]["table"],
    188         "relay_observations"
    189     );
    190     assert_eq!(contract["effects"]["checkpoint"], false);
    191     assert_eq!(contract["effects"]["dirty_generation"], false);
    192 }
    193 
    194 #[tokio::test]
    195 async fn signed_events_mutations_and_observations_are_atomic_distinct_and_idempotent() {
    196     let root = tempfile::tempdir().expect("root");
    197     let runtime = runtime(root.path());
    198     prepare_state_directory(&runtime);
    199     let configuration = configuration();
    200     let metadata = metadata(&runtime, &configuration);
    201     let (applied_at, build) = migration_evidence();
    202     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    203         .await
    204         .expect("initialize");
    205     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    206         .await
    207         .expect("writer");
    208 
    209     let first_wire = signed_wire(1);
    210     let second_wire = signed_wire(2);
    211     let first_json: Value = serde_json::from_slice(&first_wire).expect("first JSON");
    212     let second_json: Value = serde_json::from_slice(&second_wire).expect("second JSON");
    213     assert_eq!(first_json["id"], second_json["id"]);
    214     assert_ne!(first_json["sig"], second_json["sig"]);
    215 
    216     let first = admitted(&configuration, &first_wire, 1_784_347_200);
    217     let first_observation =
    218         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &first)
    219             .expect("observation");
    220     let outcome = host
    221         .repositories()
    222         .persist_trade_evidence(first, first_observation)
    223         .await
    224         .expect("first persistence");
    225     assert!(outcome.mutation_inserted());
    226     assert!(outcome.signed_event_inserted());
    227     assert!(outcome.observation_inserted());
    228 
    229     let second = admitted(&configuration, &second_wire, 1_784_347_201);
    230     let second_observation =
    231         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &second)
    232             .expect("second observation");
    233     let outcome = host
    234         .repositories()
    235         .persist_trade_evidence(second, second_observation)
    236         .await
    237         .expect("second signature");
    238     assert!(!outcome.mutation_inserted());
    239     assert!(outcome.signed_event_inserted());
    240     assert!(outcome.observation_inserted());
    241 
    242     let replay = admitted(&configuration, &first_wire, 1_784_347_200);
    243     let replay_observation =
    244         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &replay)
    245             .expect("replay observation");
    246     let outcome = host
    247         .repositories()
    248         .persist_trade_evidence(replay, replay_observation)
    249         .await
    250         .expect("exact replay");
    251     assert!(!outcome.mutation_inserted());
    252     assert!(!outcome.signed_event_inserted());
    253     assert!(!outcome.observation_inserted());
    254 
    255     let later = admitted(&configuration, &first_wire, 1_784_347_202);
    256     let later_observation =
    257         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &later)
    258             .expect("later observation");
    259     let outcome = host
    260         .repositories()
    261         .persist_trade_evidence(later, later_observation)
    262         .await
    263         .expect("later observation");
    264     assert!(!outcome.mutation_inserted());
    265     assert!(!outcome.signed_event_inserted());
    266     assert!(outcome.observation_inserted());
    267 
    268     host.close().await.expect("close");
    269     let mut connection = offline_connection(&runtime).await;
    270     for (query, table, expected) in [
    271         (
    272             "SELECT COUNT(*) FROM trade_mutations",
    273             "trade_mutations",
    274             1_i64,
    275         ),
    276         ("SELECT COUNT(*) FROM nostr_events", "nostr_events", 2),
    277         (
    278             "SELECT COUNT(*) FROM relay_observations",
    279             "relay_observations",
    280             3,
    281         ),
    282     ] {
    283         let count = sqlx::query_scalar::<_, i64>(query)
    284             .fetch_one(&mut connection)
    285             .await
    286             .expect("count");
    287         assert_eq!(count, expected, "{table}");
    288     }
    289     assert!(
    290         sqlx::query("UPDATE trade_mutations SET mutation_id = mutation_id")
    291             .execute(&mut connection)
    292             .await
    293             .is_err()
    294     );
    295     assert!(
    296         sqlx::query("DELETE FROM nostr_events")
    297             .execute(&mut connection)
    298             .await
    299             .is_err()
    300     );
    301     assert!(
    302         sqlx::query("DELETE FROM relay_observations")
    303             .execute(&mut connection)
    304             .await
    305             .is_err()
    306     );
    307     connection.close().await.expect("offline close");
    308 
    309     let inspection = open_rhi_state_inspection(&runtime, &metadata)
    310         .await
    311         .expect("inspection");
    312     let rejected = admitted(&configuration, &first_wire, 1_784_347_203);
    313     let observation =
    314         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &rejected)
    315             .expect("observation");
    316     let error = inspection
    317         .repositories()
    318         .persist_trade_evidence(rejected, observation)
    319         .await
    320         .expect_err("inspection cannot persist");
    321     assert_eq!(
    322         error.kind(),
    323         RhiTradeEvidencePersistenceErrorKind::InvalidMode
    324     );
    325     inspection.close().await.expect("inspection close");
    326 }
    327 
    328 #[tokio::test]
    329 async fn durable_mutation_conflict_rolls_back_the_event_and_observation() {
    330     let root = tempfile::tempdir().expect("root");
    331     let runtime = runtime(root.path());
    332     prepare_state_directory(&runtime);
    333     let configuration = configuration();
    334     let metadata = metadata(&runtime, &configuration);
    335     let (applied_at, build) = migration_evidence();
    336     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    337         .await
    338         .expect("initialize");
    339 
    340     let wire = signed_wire(1);
    341     let event = admitted(&configuration, &wire, 1_784_347_200);
    342     let mut connection = offline_connection(&runtime).await;
    343     insert_mutation(&mut connection, &event, br#"{}"#).await;
    344     connection.close().await.expect("offline close");
    345 
    346     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    347         .await
    348         .expect("writer");
    349     let observation =
    350         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &event)
    351             .expect("observation");
    352     let error = host
    353         .repositories()
    354         .persist_trade_evidence(event, observation)
    355         .await
    356         .expect_err("conflicting mutation");
    357     assert_eq!(
    358         error.kind(),
    359         RhiTradeEvidencePersistenceErrorKind::MutationConflict
    360     );
    361     assert!(Error::source(&error).is_none());
    362     host.close().await.expect("close");
    363 
    364     let mut connection = offline_connection(&runtime).await;
    365     assert_eq!(
    366         sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nostr_events")
    367             .fetch_one(&mut connection)
    368             .await
    369             .expect("event count"),
    370         0
    371     );
    372     assert_eq!(
    373         sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM relay_observations")
    374             .fetch_one(&mut connection)
    375             .await
    376             .expect("observation count"),
    377         0
    378     );
    379     connection.close().await.expect("offline close");
    380 }
    381 
    382 #[tokio::test]
    383 async fn durable_signed_event_conflict_rolls_back_the_observation() {
    384     let root = tempfile::tempdir().expect("root");
    385     let runtime = runtime(root.path());
    386     prepare_state_directory(&runtime);
    387     let configuration = configuration();
    388     let metadata = metadata(&runtime, &configuration);
    389     let (applied_at, build) = migration_evidence();
    390     initialize_rhi_state(&runtime, &metadata, applied_at, &build)
    391         .await
    392         .expect("initialize");
    393 
    394     let wire = signed_wire(1);
    395     let wire_json: Value = serde_json::from_slice(&wire).expect("wire JSON");
    396     let event = admitted(&configuration, &wire, 1_784_347_200);
    397     let mut connection = offline_connection(&runtime).await;
    398     insert_mutation(
    399         &mut connection,
    400         &event,
    401         wire_json["content"].as_str().expect("content").as_bytes(),
    402     )
    403     .await;
    404     sqlx::query(
    405         r#"INSERT INTO nostr_events (
    406             event_id, event_signature, mutation_id, author_pubkey, event_kind,
    407             authored_at_unix_s, canonical_event_json
    408         ) VALUES (?, ?, ?, ?, ?, ?, ?)"#,
    409     )
    410     .bind(event.event_id().as_bytes().as_slice())
    411     .bind(decode_hex::<64>(wire_json["sig"].as_str().expect("signature")).as_slice())
    412     .bind(event.mutation_id().as_bytes().as_slice())
    413     .bind(event.mutation().author_pubkey.as_bytes().as_slice())
    414     .bind(i64::from(event.event_kind()))
    415     .bind(i64::try_from(event.authored_at_unix_seconds()).expect("authored time"))
    416     .bind(br#"{}"#.as_slice())
    417     .execute(&mut connection)
    418     .await
    419     .expect("seed conflicting event");
    420     connection.close().await.expect("offline close");
    421 
    422     let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
    423         .await
    424         .expect("writer");
    425     let observation =
    426         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &event)
    427             .expect("observation");
    428     let error = host
    429         .repositories()
    430         .persist_trade_evidence(event, observation)
    431         .await
    432         .expect_err("conflicting signed event");
    433     assert_eq!(
    434         error.kind(),
    435         RhiTradeEvidencePersistenceErrorKind::SignedEventConflict
    436     );
    437     assert!(Error::source(&error).is_none());
    438     host.close().await.expect("close");
    439 
    440     let mut connection = offline_connection(&runtime).await;
    441     assert_eq!(
    442         sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM relay_observations")
    443             .fetch_one(&mut connection)
    444             .await
    445             .expect("observation count"),
    446         0
    447     );
    448     connection.close().await.expect("offline close");
    449 }
    450 
    451 #[test]
    452 fn observation_construction_and_public_diagnostics_are_closed_and_redacted() {
    453     let configuration = configuration();
    454     let wire = signed_wire(1);
    455     let event = admitted(&configuration, &wire, 1_784_347_200);
    456     let missing = RhiTradeSourceObservation::from_config(&configuration, "missing", &event)
    457         .expect_err("unknown source");
    458     assert_eq!(
    459         missing.kind(),
    460         RhiTradeEvidencePersistenceErrorKind::InvalidObservation
    461     );
    462     assert!(Error::source(&missing).is_none());
    463 
    464     let observation =
    465         RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &event)
    466             .expect("observation");
    467     let rendered = format!("{observation:?} {missing} {missing:?}");
    468     let wire: Value = serde_json::from_slice(&wire).expect("wire JSON");
    469     assert!(!rendered.contains("trade-primary"));
    470     assert!(!rendered.contains(wire["id"].as_str().expect("id")));
    471     assert!(!rendered.contains(wire["sig"].as_str().expect("signature")));
    472 }