adapters_radrootsd_tests.rs (36626B)
1 use super::*; 2 use radroots_event::wire::Nip01EventWire; 3 use radroots_protocol::radrootsd::transport_publish::v5::{ 4 DeliveryPolicy as TransportPublishDeliveryPolicy, EventRequest as TransportPublishEventRequest, 5 EventResponse as TransportPublishEventResponse, Job as TransportPublishJobView, 6 JobStatus as TransportPublishJobStatus, 7 NostrTargetSourcePolicy as NostrPublishTargetSourcePolicy, 8 OutcomeKind as TransportPublishOutcomeKind, 9 ReticulumBehavior as TransportPublishReticulumBehavior, Target as TransportPublishTarget, 10 TargetOutcome as TransportPublishTargetOutcome, TargetPolicy as TransportPublishTargetPolicy, 11 TargetSource as TransportPublishTargetSource, 12 }; 13 const RADROOTS_RETICULUM_ENDPOINT_URI: &str = "reticulum:local"; 14 use std::io::{Read, Write}; 15 use std::net::TcpListener; 16 use std::thread::JoinHandle; 17 18 const SIGNED_EVENT_PUBLIC_KEY: &str = 19 "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 20 21 struct RecordedHttpRequest { 22 request_line: String, 23 headers: Vec<(String, String)>, 24 body: String, 25 } 26 27 fn spawn_http_server( 28 status: &str, 29 response_body: &str, 30 ) -> (String, JoinHandle<RecordedHttpRequest>) { 31 let listener = TcpListener::bind("127.0.0.1:0").expect("bind test server"); 32 let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr")); 33 let status = status.to_owned(); 34 let response_body = response_body.to_owned(); 35 let content_length = response_body.len(); 36 let handle = std::thread::spawn(move || { 37 let (mut stream, _) = listener.accept().expect("accept"); 38 let mut request = Vec::new(); 39 let mut buffer = [0u8; 1024]; 40 loop { 41 let read = stream.read(&mut buffer).expect("read request"); 42 if read == 0 { 43 break; 44 } 45 request.extend_from_slice(&buffer[..read]); 46 if request.windows(4).any(|window| window == b"\r\n\r\n") { 47 let headers_end = request 48 .windows(4) 49 .position(|window| window == b"\r\n\r\n") 50 .expect("headers end") 51 + 4; 52 let header_text = String::from_utf8_lossy(&request[..headers_end]); 53 let request_content_length = header_text 54 .lines() 55 .find_map(|line| { 56 let (name, value) = line.split_once(':')?; 57 name.eq_ignore_ascii_case("content-length") 58 .then(|| value.trim().parse::<usize>().expect("content length")) 59 }) 60 .unwrap_or(0); 61 while request.len() < headers_end + request_content_length { 62 let read = stream.read(&mut buffer).expect("read body"); 63 if read == 0 { 64 break; 65 } 66 request.extend_from_slice(&buffer[..read]); 67 } 68 break; 69 } 70 } 71 let request_text = String::from_utf8_lossy(&request); 72 let (headers_text, body) = request_text.split_once("\r\n\r\n").expect("request body"); 73 let mut header_lines = headers_text.lines(); 74 let request_line = header_lines.next().expect("request line").to_owned(); 75 let headers = header_lines 76 .filter_map(|line| { 77 let (name, value) = line.split_once(':')?; 78 Some((name.to_ascii_lowercase(), value.trim().to_owned())) 79 }) 80 .collect::<Vec<_>>(); 81 let response = format!( 82 "HTTP/1.1 {status}\r\ncontent-type: application/json\r\ncontent-length: {content_length}\r\nconnection: close\r\n\r\n{response_body}", 83 ); 84 stream 85 .write_all(response.as_bytes()) 86 .expect("write response"); 87 RecordedHttpRequest { 88 request_line, 89 headers, 90 body: body.to_owned(), 91 } 92 }); 93 (endpoint, handle) 94 } 95 96 fn signed_event() -> SignedEvent { 97 let mut wire = Nip01EventWire { 98 id: String::new(), 99 pubkey: SIGNED_EVENT_PUBLIC_KEY.to_owned(), 100 created_at: 1_700_000_000, 101 kind: 30_402, 102 tags: vec![vec!["d".to_owned(), "listing-1".to_owned()]], 103 content: "{\"name\":\"carrots\"}".to_owned(), 104 sig: "c".repeat(128), 105 extra: Default::default(), 106 }; 107 wire.id = wire.computed_event_id().expect("event id").into_string(); 108 let raw_json = serde_json::to_string(&wire).expect("raw event json"); 109 SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") 110 } 111 112 fn publish_request() -> TransportPublishEventRequest { 113 TransportPublishEventRequest { 114 raw_event_json: signed_event().raw_json().to_owned(), 115 target_policy: TransportPublishTargetPolicy::nostr( 116 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 117 vec!["wss://relay.example.com".to_owned()], 118 ), 119 delivery_policy: TransportPublishDeliveryPolicy::Any, 120 idempotency_key: Some("idem-1".to_owned()), 121 timeout_ms: Some(10_000), 122 } 123 } 124 125 fn signed_event_id() -> String { 126 signed_event().id_str().to_owned() 127 } 128 129 fn signed_event_pubkey() -> String { 130 signed_event().pubkey().to_hex().to_owned() 131 } 132 133 fn job_status_for_outcome(outcome_kind: TransportPublishOutcomeKind) -> TransportPublishJobStatus { 134 if outcome_kind.counts_toward_accepted_delivery() { 135 TransportPublishJobStatus::DeliverySatisfied 136 } else if outcome_kind.is_retryable() { 137 TransportPublishJobStatus::DeliveryUnsatisfiedRetryable 138 } else if outcome_kind.is_terminal_failure() { 139 TransportPublishJobStatus::DeliveryUnsatisfiedTerminal 140 } else if outcome_kind == TransportPublishOutcomeKind::DeferredUntilImplemented { 141 TransportPublishJobStatus::DeliveryDeferred 142 } else { 143 TransportPublishJobStatus::DeliveryUnsatisfiedRetryable 144 } 145 } 146 147 fn job(outcome_kind: TransportPublishOutcomeKind) -> TransportPublishJobView { 148 let status = job_status_for_outcome(outcome_kind); 149 TransportPublishJobView { 150 job_id: "job-1".to_owned(), 151 status, 152 terminal: matches!( 153 status, 154 TransportPublishJobStatus::DeliverySatisfied 155 | TransportPublishJobStatus::DeliveryUnsatisfiedTerminal 156 | TransportPublishJobStatus::DeliveryDeferred 157 | TransportPublishJobStatus::DeliveryDeferredUntilImplemented 158 | TransportPublishJobStatus::Rejected 159 ), 160 delivery_satisfied: status == TransportPublishJobStatus::DeliverySatisfied, 161 event_id: signed_event_id(), 162 pubkey: signed_event_pubkey(), 163 event_kind: 30_402, 164 target_policy: TransportPublishTargetPolicy::nostr( 165 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 166 vec!["wss://relay.example.com".to_owned()], 167 ), 168 delivery_policy: TransportPublishDeliveryPolicy::Any, 169 target_count: 1, 170 acknowledged_count: usize::from(outcome_kind.counts_toward_accepted_delivery()), 171 retryable_count: usize::from(outcome_kind.is_retryable()), 172 terminal_count: usize::from(outcome_kind.is_terminal_failure()), 173 requested_at_ms: 1_700_000_000_000, 174 completed_at_ms: Some(1_700_000_000_100), 175 last_error: None, 176 targets: vec![TransportPublishTargetOutcome { 177 transport_kind: "nostr".to_owned(), 178 endpoint_uri: "wss://relay.example.com".to_owned(), 179 target_scope: None, 180 target_label: None, 181 source: TransportPublishTargetSource::Request, 182 attempted: true, 183 outcome_kind, 184 message: Some("relay outcome".to_owned()), 185 latency_ms: Some(7), 186 }], 187 } 188 } 189 190 fn explicit_nostr_job( 191 endpoints: Vec<String>, 192 delivery_policy: TransportPublishDeliveryPolicy, 193 ) -> TransportPublishJobView { 194 let targets = endpoints 195 .iter() 196 .map(|endpoint| TransportPublishTargetOutcome { 197 transport_kind: "nostr".to_owned(), 198 endpoint_uri: endpoint.clone(), 199 target_scope: None, 200 target_label: None, 201 source: TransportPublishTargetSource::Request, 202 attempted: true, 203 outcome_kind: TransportPublishOutcomeKind::Accepted, 204 message: Some("relay outcome".to_owned()), 205 latency_ms: Some(7), 206 }) 207 .collect::<Vec<_>>(); 208 TransportPublishJobView { 209 job_id: "job-explicit-nostr".to_owned(), 210 status: TransportPublishJobStatus::DeliverySatisfied, 211 terminal: true, 212 delivery_satisfied: true, 213 event_id: signed_event_id(), 214 pubkey: signed_event_pubkey(), 215 event_kind: 30_402, 216 target_policy: TransportPublishTargetPolicy::explicit_targets( 217 endpoints 218 .iter() 219 .map(|endpoint| TransportPublishTarget::nostr(endpoint.as_str())) 220 .collect::<Vec<_>>(), 221 ), 222 delivery_policy, 223 target_count: targets.len(), 224 acknowledged_count: targets.len(), 225 retryable_count: 0, 226 terminal_count: 0, 227 requested_at_ms: 1_700_000_000_000, 228 completed_at_ms: Some(1_700_000_000_100), 229 last_error: None, 230 targets, 231 } 232 } 233 234 fn reticulum_deferred_job() -> TransportPublishJobView { 235 TransportPublishJobView { 236 job_id: "job-reticulum".to_owned(), 237 status: TransportPublishJobStatus::DeliveryDeferred, 238 terminal: true, 239 delivery_satisfied: false, 240 event_id: signed_event_id(), 241 pubkey: signed_event_pubkey(), 242 event_kind: 30_402, 243 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 244 TransportPublishTarget::reticulum( 245 TransportPublishReticulumBehavior::DeferDeliveryPlans, 246 ), 247 ]), 248 delivery_policy: TransportPublishDeliveryPolicy::Any, 249 target_count: 1, 250 acknowledged_count: 0, 251 retryable_count: 0, 252 terminal_count: 0, 253 requested_at_ms: 1_700_000_000_000, 254 completed_at_ms: Some(1_700_000_000_100), 255 last_error: Some("delivery_deferred_until_implemented".to_owned()), 256 targets: vec![TransportPublishTargetOutcome { 257 transport_kind: "reticulum".to_owned(), 258 endpoint_uri: RADROOTS_RETICULUM_ENDPOINT_URI.to_owned(), 259 target_scope: None, 260 target_label: None, 261 source: TransportPublishTargetSource::Reticulum, 262 attempted: false, 263 outcome_kind: TransportPublishOutcomeKind::DeferredUntilImplemented, 264 message: Some("reticulum deferred until implemented".to_owned()), 265 latency_ms: None, 266 }], 267 } 268 } 269 270 fn publish_response_json_for_job(job: TransportPublishJobView) -> String { 271 serde_json::json!({ 272 "jsonrpc": "2.0", 273 "id": SDK_RADROOTSD_PUBLISH_REQUEST_ID, 274 "result": { 275 "deduplicated": false, 276 "job": job 277 } 278 }) 279 .to_string() 280 } 281 282 fn publish_response_json() -> String { 283 publish_response_json_for_job(job(TransportPublishOutcomeKind::Accepted)) 284 } 285 286 fn reticulum_deferred_response_json() -> String { 287 serde_json::json!({ 288 "jsonrpc": "2.0", 289 "id": SDK_RADROOTSD_PUBLISH_REQUEST_ID, 290 "result": { 291 "deduplicated": false, 292 "job": reticulum_deferred_job() 293 } 294 }) 295 .to_string() 296 } 297 298 fn assert_message(error: RadrootsdError, fragment: &str) { 299 let message = error.to_string(); 300 assert!( 301 message.contains(fragment), 302 "expected {message:?} to contain {fragment:?}" 303 ); 304 } 305 306 #[test] 307 fn auth_headers_omit_or_redact_bearer_authorization() { 308 let none = auth_headers(&RadrootsdAuth::None).expect("none auth"); 309 assert!(!none.contains_key(AUTHORIZATION)); 310 assert_eq!(format!("{:?}", RadrootsdAuth::None), "None"); 311 312 let bearer = auth_headers(&RadrootsdAuth::BearerToken("sdk-token".into())).expect("bearer"); 313 assert_eq!( 314 bearer 315 .get(AUTHORIZATION) 316 .expect("authorization") 317 .to_str() 318 .expect("authorization str"), 319 "Bearer sdk-token" 320 ); 321 322 let error = auth_headers(&RadrootsdAuth::BearerToken("bad\ntoken".into())).expect_err("error"); 323 assert!(matches!(error, RadrootsdError::InvalidAuthHeader(_))); 324 assert_eq!( 325 format!("{:?}", RadrootsdAuth::BearerToken("token-secret".into())), 326 "BearerToken(<redacted>)" 327 ); 328 assert!( 329 RadrootsdError::InvalidAuthHeader("bad header".to_owned()) 330 .to_string() 331 .contains("invalid radrootsd bearer token header") 332 ); 333 assert_eq!( 334 RadrootsdError::InvalidRequest("invalid request".to_owned()).to_string(), 335 "invalid request" 336 ); 337 assert_eq!( 338 RadrootsdError::Http("http failed".to_owned()).to_string(), 339 "http failed" 340 ); 341 assert_eq!( 342 RadrootsdError::MalformedResponse("bad envelope".to_owned()).to_string(), 343 "bad envelope" 344 ); 345 } 346 347 #[test] 348 fn radrootsd_publish_config_builders_preserve_typed_runtime_options() { 349 let config = RadrootsdPublishConfig::new("http://127.0.0.1:8080/rpc") 350 .with_auth(RadrootsdAuth::BearerToken("sdk-token".to_owned())) 351 .with_timeout(Duration::from_millis(250)); 352 let adapter = RadrootsdPublishAdapter::new(config.clone()); 353 354 assert_eq!(adapter.config(), &config); 355 assert_eq!(adapter.config().endpoint, "http://127.0.0.1:8080/rpc"); 356 assert_eq!( 357 adapter.config().auth, 358 RadrootsdAuth::BearerToken("sdk-token".to_owned()) 359 ); 360 assert_eq!(adapter.config().timeout, Duration::from_millis(250)); 361 } 362 363 #[test] 364 fn publish_event_request_json_uses_signed_event_contract() { 365 let value = publish_event_request_json(&publish_request()).expect("request json"); 366 367 let raw_event_json = value["raw_event_json"].as_str().expect("raw event json"); 368 let raw_event: serde_json::Value = serde_json::from_str(raw_event_json).expect("raw event"); 369 assert_eq!(raw_event["id"], signed_event_id()); 370 assert_eq!(raw_event["pubkey"], SIGNED_EVENT_PUBLIC_KEY); 371 assert_eq!(raw_event["kind"], 30_402); 372 assert_eq!(value["target_policy"]["kind"], "nostr"); 373 assert_eq!( 374 value["target_policy"]["source_policy"], 375 "request_then_author_write_then_daemon_default" 376 ); 377 assert_eq!( 378 value["target_policy"]["relay_urls"][0], 379 "wss://relay.example.com" 380 ); 381 assert_eq!(value["delivery_policy"]["mode"], "any"); 382 assert_eq!(value["idempotency_key"], "idem-1"); 383 let rendered = value.to_string(); 384 assert!(!rendered.contains("signer_session_id")); 385 assert!(!rendered.contains("bridge.")); 386 } 387 388 #[test] 389 fn decode_jsonrpc_response_validates_envelope_and_errors() { 390 let response: TransportPublishEventResponse = decode_jsonrpc_response( 391 METHOD_EVENT, 392 SDK_RADROOTSD_PUBLISH_REQUEST_ID, 393 publish_response_json().as_str(), 394 ) 395 .expect("response"); 396 assert_eq!(response.job.event_id, signed_event_id()); 397 398 let error = decode_jsonrpc_response::<TransportPublishEventResponse>( 399 METHOD_EVENT, 400 SDK_RADROOTSD_PUBLISH_REQUEST_ID, 401 r#"{"jsonrpc":"2.0","id":"radroots-sdk-transport-publish-event","error":{"code":-32001,"message":"principal unauthorized"}}"#, 402 ) 403 .expect_err("jsonrpc error"); 404 assert!(matches!( 405 error, 406 RadrootsdError::JsonRpc { code: -32001, .. } 407 )); 408 assert_message(error, "principal unauthorized"); 409 410 assert!(matches!( 411 decode_jsonrpc_response::<serde_json::Value>(METHOD_EVENT, "expected", "not json"), 412 Err(RadrootsdError::MalformedResponse(_)) 413 )); 414 assert!(matches!( 415 decode_jsonrpc_response::<serde_json::Value>( 416 METHOD_EVENT, 417 "expected", 418 r#"{"jsonrpc":"2.0","id":"other","result":{}}"# 419 ), 420 Err(RadrootsdError::MalformedResponse(_)) 421 )); 422 assert!(matches!( 423 decode_jsonrpc_response::<serde_json::Value>( 424 METHOD_EVENT, 425 "expected", 426 r#"{"jsonrpc":"2.0","id":"expected"}"# 427 ), 428 Err(RadrootsdError::MalformedResponse(_)) 429 )); 430 assert!(matches!( 431 decode_jsonrpc_response::<serde_json::Value>( 432 METHOD_EVENT, 433 "expected", 434 r#"{"jsonrpc":"1.0","id":"expected","result":{}}"# 435 ), 436 Err(RadrootsdError::MalformedResponse(_)) 437 )); 438 assert!(matches!( 439 decode_jsonrpc_response::<serde_json::Value>( 440 METHOD_EVENT, 441 "expected", 442 r#"{"jsonrpc":"2.0","id":"expected","result":{},"error":{"code":-32002,"message":"both"}}"# 443 ), 444 Err(RadrootsdError::MalformedResponse(_)) 445 )); 446 } 447 448 #[tokio::test] 449 async fn publish_event_posts_transport_publish_jsonrpc() { 450 let (endpoint, handle) = spawn_http_server("200 OK", publish_response_json().as_str()); 451 452 let receipt = publish_event( 453 endpoint.as_str(), 454 &RadrootsdAuth::BearerToken("sdk-token".into()), 455 &publish_request(), 456 Duration::from_secs(2), 457 ) 458 .await 459 .expect("publish"); 460 461 assert_eq!(receipt.job.event_id, signed_event_id()); 462 let recorded = handle.join().expect("server thread"); 463 assert_eq!(recorded.request_line, "POST /rpc HTTP/1.1"); 464 assert!( 465 recorded 466 .headers 467 .iter() 468 .any(|(name, value)| name == "authorization" && value == "Bearer sdk-token") 469 ); 470 let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); 471 assert_eq!(body["method"], METHOD_EVENT); 472 assert_eq!(body["id"], SDK_RADROOTSD_PUBLISH_REQUEST_ID); 473 let raw_event_json = body["params"]["raw_event_json"] 474 .as_str() 475 .expect("raw event json"); 476 let raw_event: serde_json::Value = serde_json::from_str(raw_event_json).expect("raw event"); 477 assert_eq!(raw_event["content"], "{\"name\":\"carrots\"}"); 478 assert_eq!(body["params"]["target_policy"]["kind"], "nostr"); 479 assert_eq!( 480 body["params"]["target_policy"]["relay_urls"][0], 481 "wss://relay.example.com" 482 ); 483 } 484 485 #[tokio::test] 486 async fn publish_signed_event_posts_typed_radrootsd_request() { 487 let mut response_job = explicit_nostr_job( 488 vec!["wss://relay.example.com".to_owned()], 489 TransportPublishDeliveryPolicy::All, 490 ); 491 response_job.target_policy = TransportPublishTargetPolicy::explicit_targets(vec![ 492 TransportPublishTarget::nostr("wss://relay.example.com") 493 .with_scope("farm.local") 494 .with_label("Farm relay"), 495 ]); 496 response_job.targets[0].target_scope = Some("farm.local".to_owned()); 497 response_job.targets[0].target_label = Some("Farm relay".to_owned()); 498 let response_json = publish_response_json_for_job(response_job); 499 let (endpoint, handle) = spawn_http_server("200 OK", response_json.as_str()); 500 let adapter = RadrootsdPublishAdapter::new( 501 RadrootsdPublishConfig::new(endpoint) 502 .with_auth(RadrootsdAuth::BearerToken("sdk-token".into())), 503 ); 504 505 let receipt = adapter 506 .publish_signed_event(RadrootsdPublishRequest { 507 signed_event: signed_event(), 508 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 509 TransportPublishTarget::nostr("wss://relay.example.com") 510 .with_scope("farm.local") 511 .with_label("Farm relay"), 512 ]), 513 delivery_policy: TransportPublishDeliveryPolicy::All, 514 idempotency_key: Some("idem-typed".to_owned()), 515 timeout_ms: Some(7_000), 516 }) 517 .await 518 .expect("typed publish"); 519 520 assert!(receipt.job.delivery_satisfied); 521 assert_eq!( 522 receipt.job.targets[0].target_scope.as_deref(), 523 Some("farm.local") 524 ); 525 assert_eq!( 526 receipt.job.targets[0].target_label.as_deref(), 527 Some("Farm relay") 528 ); 529 let recorded = handle.join().expect("server thread"); 530 assert!( 531 recorded 532 .headers 533 .iter() 534 .any(|(name, value)| name == "authorization" && value == "Bearer sdk-token") 535 ); 536 let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); 537 assert_eq!(body["params"]["delivery_policy"]["mode"], "all"); 538 assert_eq!(body["params"]["target_policy"]["kind"], "explicit_targets"); 539 assert_eq!( 540 body["params"]["target_policy"]["targets"][0]["endpoint_uri"], 541 "wss://relay.example.com" 542 ); 543 assert_eq!( 544 body["params"]["target_policy"]["targets"][0]["target_scope"], 545 "farm.local" 546 ); 547 assert_eq!( 548 body["params"]["target_policy"]["targets"][0]["target_label"], 549 "Farm relay" 550 ); 551 assert_eq!(body["params"]["idempotency_key"], "idem-typed"); 552 assert_eq!(body["params"]["timeout_ms"], 7_000); 553 } 554 555 #[tokio::test] 556 async fn publish_signed_event_preserves_typed_reticulum_behavior() { 557 let response_json = reticulum_deferred_response_json(); 558 let (endpoint, handle) = spawn_http_server("200 OK", response_json.as_str()); 559 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 560 561 let response = adapter 562 .publish_signed_event(RadrootsdPublishRequest { 563 signed_event: signed_event(), 564 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 565 TransportPublishTarget::reticulum( 566 TransportPublishReticulumBehavior::DeferDeliveryPlans, 567 ), 568 ]), 569 delivery_policy: TransportPublishDeliveryPolicy::Any, 570 idempotency_key: Some("idem-reticulum".to_owned()), 571 timeout_ms: None, 572 }) 573 .await 574 .expect("typed Reticulum publish request"); 575 assert_eq!( 576 response.job.status, 577 TransportPublishJobStatus::DeliveryDeferred 578 ); 579 assert!(!response.job.delivery_satisfied); 580 581 let recorded = handle.join().expect("server thread"); 582 let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); 583 assert_eq!(body["params"]["target_policy"]["kind"], "explicit_targets"); 584 assert_eq!( 585 body["params"]["target_policy"]["targets"][0]["transport_kind"], 586 "reticulum" 587 ); 588 assert_eq!( 589 body["params"]["target_policy"]["targets"][0]["endpoint_uri"], 590 RADROOTS_RETICULUM_ENDPOINT_URI 591 ); 592 assert_eq!( 593 body["params"]["target_policy"]["targets"][0]["reticulum_behavior"], 594 "defer_delivery_plans" 595 ); 596 } 597 598 #[tokio::test] 599 async fn publish_signed_event_rejects_mismatched_daemon_event_identity() { 600 let mut response_job = job(TransportPublishOutcomeKind::Accepted); 601 response_job.event_id = "0".repeat(64); 602 let response_json = publish_response_json_for_job(response_job); 603 let (endpoint, _handle) = spawn_http_server("200 OK", response_json.as_str()); 604 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 605 606 let error = adapter 607 .publish_signed_event(RadrootsdPublishRequest { 608 signed_event: signed_event(), 609 target_policy: TransportPublishTargetPolicy::nostr( 610 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 611 vec!["wss://relay.example.com".to_owned()], 612 ), 613 delivery_policy: TransportPublishDeliveryPolicy::Any, 614 idempotency_key: Some("idem-mismatch".to_owned()), 615 timeout_ms: None, 616 }) 617 .await 618 .expect_err("mismatched response"); 619 620 assert!(matches!(error, RadrootsdError::MalformedResponse(_))); 621 assert_message(error, "event_id"); 622 } 623 624 #[tokio::test] 625 async fn publish_signed_event_rejects_mismatched_daemon_pubkey_and_kind() { 626 for (field, response_job) in [ 627 { 628 let mut response_job = job(TransportPublishOutcomeKind::Accepted); 629 response_job.pubkey = "0".repeat(64); 630 ("pubkey", response_job) 631 }, 632 { 633 let mut response_job = job(TransportPublishOutcomeKind::Accepted); 634 response_job.event_kind = 30_403; 635 ("event_kind", response_job) 636 }, 637 ] { 638 let response_json = publish_response_json_for_job(response_job); 639 let (endpoint, _handle) = spawn_http_server("200 OK", response_json.as_str()); 640 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 641 642 let error = adapter 643 .publish_signed_event(RadrootsdPublishRequest { 644 signed_event: signed_event(), 645 target_policy: TransportPublishTargetPolicy::nostr( 646 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 647 vec!["wss://relay.example.com".to_owned()], 648 ), 649 delivery_policy: TransportPublishDeliveryPolicy::Any, 650 idempotency_key: Some(format!("idem-mismatch-{field}")), 651 timeout_ms: None, 652 }) 653 .await 654 .expect_err("mismatched response"); 655 656 assert!(matches!(error, RadrootsdError::MalformedResponse(_))); 657 assert_message(error, field); 658 } 659 } 660 661 #[tokio::test] 662 async fn publish_signed_event_rejects_mismatched_daemon_delivery_policy() { 663 let mut response_job = job(TransportPublishOutcomeKind::Accepted); 664 response_job.delivery_policy = TransportPublishDeliveryPolicy::All; 665 let response_json = publish_response_json_for_job(response_job); 666 let (endpoint, _handle) = spawn_http_server("200 OK", response_json.as_str()); 667 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 668 669 let error = adapter 670 .publish_signed_event(RadrootsdPublishRequest { 671 signed_event: signed_event(), 672 target_policy: TransportPublishTargetPolicy::nostr( 673 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 674 vec!["wss://relay.example.com".to_owned()], 675 ), 676 delivery_policy: TransportPublishDeliveryPolicy::Any, 677 idempotency_key: Some("idem-delivery-policy-mismatch".to_owned()), 678 timeout_ms: None, 679 }) 680 .await 681 .expect_err("delivery policy mismatch"); 682 683 assert!(matches!(error, RadrootsdError::MalformedResponse(_))); 684 assert_message(error, "delivery_policy"); 685 } 686 687 #[tokio::test] 688 async fn publish_signed_event_rejects_mismatched_explicit_target_response() { 689 let (endpoint, _handle) = spawn_http_server("200 OK", publish_response_json().as_str()); 690 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 691 692 let error = adapter 693 .publish_signed_event(RadrootsdPublishRequest { 694 signed_event: signed_event(), 695 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 696 TransportPublishTarget::reticulum( 697 TransportPublishReticulumBehavior::DeferDeliveryPlans, 698 ), 699 ]), 700 delivery_policy: TransportPublishDeliveryPolicy::Any, 701 idempotency_key: Some("idem-target-mismatch".to_owned()), 702 timeout_ms: None, 703 }) 704 .await 705 .expect_err("target mismatch"); 706 707 assert!(matches!(error, RadrootsdError::MalformedResponse(_))); 708 assert_message(error, "target_policy"); 709 } 710 711 #[tokio::test] 712 async fn publish_signed_event_accepts_reordered_explicit_target_outcomes() { 713 let mut response_job = explicit_nostr_job( 714 vec![ 715 "wss://relay-a.example.com".to_owned(), 716 "wss://relay-b.example.com".to_owned(), 717 ], 718 TransportPublishDeliveryPolicy::Any, 719 ); 720 response_job.targets.reverse(); 721 let response_json = publish_response_json_for_job(response_job); 722 let (endpoint, _handle) = spawn_http_server("200 OK", response_json.as_str()); 723 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 724 725 let response = adapter 726 .publish_signed_event(RadrootsdPublishRequest { 727 signed_event: signed_event(), 728 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 729 TransportPublishTarget::nostr("wss://relay-a.example.com"), 730 TransportPublishTarget::nostr("wss://relay-b.example.com"), 731 ]), 732 delivery_policy: TransportPublishDeliveryPolicy::Any, 733 idempotency_key: Some("idem-explicit-reordered-outcomes".to_owned()), 734 timeout_ms: None, 735 }) 736 .await 737 .expect("reordered explicit target outcomes"); 738 739 assert!(response.job.delivery_satisfied); 740 assert_eq!( 741 response.job.targets[0].endpoint_uri, 742 "wss://relay-b.example.com" 743 ); 744 assert_eq!( 745 response.job.targets[1].endpoint_uri, 746 "wss://relay-a.example.com" 747 ); 748 } 749 750 #[tokio::test] 751 async fn publish_signed_event_rejects_mismatched_explicit_target_outcomes() { 752 let mut response_job = explicit_nostr_job( 753 vec!["wss://relay.example.com".to_owned()], 754 TransportPublishDeliveryPolicy::Any, 755 ); 756 response_job.targets[0].endpoint_uri = "wss://relay-other.example.com".to_owned(); 757 let response_json = publish_response_json_for_job(response_job); 758 let (endpoint, _handle) = spawn_http_server("200 OK", response_json.as_str()); 759 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 760 761 let error = adapter 762 .publish_signed_event(RadrootsdPublishRequest { 763 signed_event: signed_event(), 764 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 765 TransportPublishTarget::nostr("wss://relay.example.com"), 766 ]), 767 delivery_policy: TransportPublishDeliveryPolicy::Any, 768 idempotency_key: Some("idem-explicit-outcome-mismatch".to_owned()), 769 timeout_ms: None, 770 }) 771 .await 772 .expect_err("explicit target outcome mismatch"); 773 774 assert!(matches!(error, RadrootsdError::MalformedResponse(_))); 775 assert_message(error, "explicit target policy"); 776 } 777 778 #[tokio::test] 779 async fn publish_signed_event_rejects_mismatched_scoped_explicit_target_outcomes() { 780 let mut response_job = explicit_nostr_job( 781 vec!["wss://relay.example.com".to_owned()], 782 TransportPublishDeliveryPolicy::Any, 783 ); 784 response_job.target_policy = TransportPublishTargetPolicy::explicit_targets(vec![ 785 TransportPublishTarget::nostr("wss://relay.example.com").with_scope("farm.local"), 786 ]); 787 response_job.targets[0].target_scope = Some("farm.remote".to_owned()); 788 let response_json = publish_response_json_for_job(response_job); 789 let (endpoint, _handle) = spawn_http_server("200 OK", response_json.as_str()); 790 let adapter = RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new(endpoint)); 791 792 let error = adapter 793 .publish_signed_event(RadrootsdPublishRequest { 794 signed_event: signed_event(), 795 target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ 796 TransportPublishTarget::nostr("wss://relay.example.com").with_scope("farm.local"), 797 ]), 798 delivery_policy: TransportPublishDeliveryPolicy::Any, 799 idempotency_key: Some("idem-scoped-explicit-outcome-mismatch".to_owned()), 800 timeout_ms: None, 801 }) 802 .await 803 .expect_err("scoped explicit target outcome mismatch"); 804 805 assert!(matches!(error, RadrootsdError::MalformedResponse(_))); 806 assert_message(error, "explicit target policy"); 807 } 808 809 #[tokio::test] 810 async fn publish_event_http_errors_omit_body_and_token_material() { 811 let body = "{\"error\":\"token-secret content carrots\"}"; 812 let (endpoint, _handle) = spawn_http_server("503 Service Unavailable", body); 813 814 let error = publish_event( 815 endpoint.as_str(), 816 &RadrootsdAuth::BearerToken("token-secret".into()), 817 &publish_request(), 818 Duration::from_secs(2), 819 ) 820 .await 821 .expect_err("http error"); 822 let message = error.to_string(); 823 824 assert!(message.contains("503")); 825 assert!(message.contains("response body omitted")); 826 assert!(!message.contains("token-secret")); 827 assert!(!message.contains("carrots")); 828 } 829 830 #[tokio::test] 831 async fn publish_event_empty_http_error_reports_empty_body() { 832 let (endpoint, _handle) = spawn_http_server("500 Internal Server Error", ""); 833 834 let error = publish_event( 835 endpoint.as_str(), 836 &RadrootsdAuth::None, 837 &publish_request(), 838 Duration::from_secs(2), 839 ) 840 .await 841 .expect_err("http error"); 842 843 assert!(error.to_string().contains("response body empty")); 844 } 845 846 #[tokio::test] 847 async fn publish_signed_event_rejects_invalid_target_requests_before_http() { 848 let adapter = 849 RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new("http://127.0.0.1:9/rpc")); 850 let base = RadrootsdPublishRequest { 851 signed_event: signed_event(), 852 target_policy: TransportPublishTargetPolicy::nostr( 853 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 854 vec!["wss://relay.example.com".to_owned()], 855 ), 856 delivery_policy: TransportPublishDeliveryPolicy::Any, 857 idempotency_key: Some("idem-1".to_owned()), 858 timeout_ms: Some(1_000), 859 }; 860 861 let mut invalid_quorum = base.clone(); 862 invalid_quorum.delivery_policy = TransportPublishDeliveryPolicy::Quorum { quorum: 0 }; 863 let mut too_many_targets = base.clone(); 864 too_many_targets.target_policy = TransportPublishTargetPolicy::nostr( 865 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 866 (0..=SDK_RADROOTSD_PUBLISH_MAX_TARGETS) 867 .map(|index| format!("wss://relay-{index}.example.com")) 868 .collect(), 869 ); 870 let mut empty_endpoint_uri = base.clone(); 871 empty_endpoint_uri.target_policy = TransportPublishTargetPolicy::nostr( 872 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 873 vec![" ".to_owned()], 874 ); 875 let mut nostr_reticulum_behavior = base.clone(); 876 nostr_reticulum_behavior.target_policy = 877 TransportPublishTargetPolicy::explicit_targets(vec![TransportPublishTarget { 878 transport_kind: "nostr".to_owned(), 879 endpoint_uri: "wss://relay.example.com".to_owned(), 880 target_scope: None, 881 target_label: None, 882 reticulum_behavior: Some(TransportPublishReticulumBehavior::RejectDeliveryAttempts), 883 }]); 884 let mut explicit_radrootsd_target = base.clone(); 885 explicit_radrootsd_target.target_policy = 886 TransportPublishTargetPolicy::explicit_targets(vec![TransportPublishTarget { 887 transport_kind: "radrootsd".to_owned(), 888 endpoint_uri: "radrootsd-execution:publish".to_owned(), 889 target_scope: None, 890 target_label: None, 891 reticulum_behavior: None, 892 }]); 893 let radrootsd_error = adapter 894 .publish_signed_event(explicit_radrootsd_target) 895 .await 896 .expect_err("explicit radrootsd target"); 897 assert!(matches!( 898 radrootsd_error, 899 RadrootsdError::InvalidRequest(message) 900 if message.contains("transport target 0 kind must be canonical lowercase") 901 )); 902 let mut empty_idempotency = base; 903 empty_idempotency.idempotency_key = Some(" ".to_owned()); 904 905 for request in [ 906 invalid_quorum, 907 too_many_targets, 908 empty_endpoint_uri, 909 nostr_reticulum_behavior, 910 empty_idempotency, 911 ] { 912 assert!(matches!( 913 adapter.publish_signed_event(request).await, 914 Err(RadrootsdError::InvalidRequest(_)) 915 )); 916 } 917 } 918 919 #[tokio::test] 920 async fn adapter_rejects_invalid_request_before_transport() { 921 let adapter = 922 RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new("http://127.0.0.1:9/rpc")); 923 let request = RadrootsdPublishRequest { 924 signed_event: signed_event(), 925 target_policy: TransportPublishTargetPolicy::nostr( 926 NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, 927 Vec::new(), 928 ), 929 delivery_policy: TransportPublishDeliveryPolicy::Quorum { quorum: 0 }, 930 idempotency_key: None, 931 timeout_ms: None, 932 }; 933 934 let error = adapter 935 .publish_signed_event(request) 936 .await 937 .expect_err("invalid request"); 938 939 assert!(matches!(error, RadrootsdError::InvalidRequest(_))); 940 }