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 }