transport_publish.rs (16126B)
1 use anyhow::Result; 2 use jsonrpsee::server::RpcModule; 3 use radroots_protocol::radrootsd::transport_publish::v5::{ 4 Capabilities, EventRequest, METHOD_CAPABILITIES, METHOD_EVENT, METHOD_JOB_GET, METHOD_JOB_LIST, 5 }; 6 use serde::Deserialize; 7 8 use crate::core::transport_publish::TransportPublishError; 9 use crate::transport::jsonrpc::auth::require_publish_principal; 10 use crate::transport::jsonrpc::{MethodRegistry, RpcContext, RpcError}; 11 12 #[derive(Debug, Deserialize)] 13 struct JobGetParams { 14 job_id: String, 15 } 16 17 #[derive(Debug, Deserialize)] 18 #[serde(deny_unknown_fields)] 19 struct JobListParams { 20 limit: Option<usize>, 21 } 22 23 pub fn module(ctx: RpcContext, registry: MethodRegistry) -> Result<RpcModule<RpcContext>> { 24 let mut module = RpcModule::new(ctx); 25 register_capabilities(&mut module, ®istry)?; 26 register_event(&mut module, ®istry)?; 27 register_job_get(&mut module, ®istry)?; 28 register_job_list(&mut module, ®istry)?; 29 Ok(module) 30 } 31 32 fn register_capabilities( 33 module: &mut RpcModule<RpcContext>, 34 registry: &MethodRegistry, 35 ) -> Result<()> { 36 registry.track(METHOD_CAPABILITIES); 37 module.register_async_method(METHOD_CAPABILITIES, |_params, ctx, extensions| async move { 38 require_publish_principal(&extensions)?; 39 Ok::<Capabilities, RpcError>(Capabilities::v5( 40 ctx.state.transport_publish.config.max_event_bytes, 41 ctx.state.transport_publish.config.max_targets_per_request, 42 )) 43 })?; 44 Ok(()) 45 } 46 47 fn register_event(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { 48 registry.track(METHOD_EVENT); 49 module.register_async_method(METHOD_EVENT, |params, ctx, extensions| async move { 50 let principal = require_publish_principal(&extensions)?; 51 let request: EventRequest = params 52 .parse() 53 .map_err(|error| RpcError::InvalidParams(error.to_string()))?; 54 ctx.state 55 .transport_publish 56 .publish_event(&principal, request) 57 .await 58 .map_err(rpc_error_from_transport_publish) 59 })?; 60 Ok(()) 61 } 62 63 fn register_job_get(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { 64 registry.track(METHOD_JOB_GET); 65 module.register_async_method(METHOD_JOB_GET, |params, ctx, extensions| async move { 66 let principal = require_publish_principal(&extensions)?; 67 let params: JobGetParams = params 68 .parse() 69 .map_err(|error| RpcError::InvalidParams(error.to_string()))?; 70 let job_id = params.job_id.trim(); 71 if job_id.is_empty() { 72 return Err(RpcError::InvalidParams("missing job_id".to_owned())); 73 } 74 ctx.state 75 .transport_publish 76 .store 77 .job_by_id_for_principal(job_id, &principal) 78 .map_err(|error| RpcError::Other(error.to_string()))? 79 .ok_or_else(|| RpcError::Other(format!("unknown publish job: {job_id}"))) 80 })?; 81 Ok(()) 82 } 83 84 fn register_job_list(module: &mut RpcModule<RpcContext>, registry: &MethodRegistry) -> Result<()> { 85 registry.track(METHOD_JOB_LIST); 86 module.register_async_method(METHOD_JOB_LIST, |params, ctx, extensions| async move { 87 let principal = require_publish_principal(&extensions)?; 88 let params = if params.len_bytes() == 0 || params.as_str() == Some("[]") { 89 JobListParams { limit: None } 90 } else { 91 params 92 .parse::<JobListParams>() 93 .map_err(|error| RpcError::InvalidParams(error.to_string()))? 94 }; 95 if params.limit == Some(0) { 96 return Err(RpcError::InvalidParams( 97 "limit must be greater than zero".to_owned(), 98 )); 99 } 100 let configured_limit = ctx.state.transport_publish.config.job_list_limit; 101 let limit = params 102 .limit 103 .unwrap_or(configured_limit) 104 .min(configured_limit); 105 ctx.state 106 .transport_publish 107 .store 108 .list_jobs_for_principal(&principal, limit) 109 .map_err(|error| RpcError::Other(error.to_string())) 110 })?; 111 Ok(()) 112 } 113 114 fn rpc_error_from_transport_publish(error: TransportPublishError) -> RpcError { 115 match error { 116 TransportPublishError::InvalidScope(message) => RpcError::Unauthorized(message), 117 TransportPublishError::InvalidSignedEvent(message) => RpcError::InvalidParams(message), 118 TransportPublishError::EventWire(_) 119 | TransportPublishError::SignedEvent(_) 120 | TransportPublishError::Relay(_) => RpcError::InvalidParams(error.to_string()), 121 TransportPublishError::IdempotencyConflict(_) => RpcError::Other(error.to_string()), 122 other => RpcError::Other(other.to_string()), 123 } 124 } 125 126 #[cfg(test)] 127 mod tests { 128 use super::module; 129 use std::sync::Arc; 130 131 use crate::app::config::{Nip46Config, TransportPublishConfig, TransportPublishNostrConfig}; 132 use crate::app::identity_storage::DaemonIdentity; 133 use crate::core::Radrootsd; 134 use crate::core::transport_publish::{ 135 PublishJobVisibility, PublishPrincipalInit, generate_bearer_token, hash_bearer_token, 136 }; 137 use crate::host_nostr::Timestamp; 138 use crate::transport::jsonrpc::auth::{ 139 TransportPublishAuthorization, authorize_transport_publish_request, 140 }; 141 use crate::transport::jsonrpc::{MethodRegistry, RpcContext}; 142 use crate::transport::relay_publish::MockRelayPublishAdapter as RadrootsMockRelayPublishAdapter; 143 use jsonrpsee::server::RpcModule; 144 use nostr::JsonUtil; 145 use nostr::{EventBuilder, Kind, Tag}; 146 use radroots_protocol::radrootsd::transport_publish::v5::{ 147 NostrTargetSourcePolicy, TargetPolicyName, 148 }; 149 150 fn signed_event(identity: &DaemonIdentity) -> String { 151 // This method accepts an already-signed wire event; construct the test 152 // fixture at that explicit low-level interoperability boundary. 153 let event = EventBuilder::new(Kind::Custom(30_402), "{}") 154 .tag(Tag::identifier("listing-1")) 155 .custom_created_at(Timestamp::from_secs(1_700_000_000)) 156 .sign_with_keys(identity.keys()) 157 .expect("signed event"); 158 event.as_json() 159 } 160 161 fn module_with_principal_and_config( 162 admin: bool, 163 transport_publish_config: TransportPublishConfig, 164 ) -> (RpcModule<RpcContext>, RpcContext, String, String) { 165 let identity = DaemonIdentity::generate(); 166 let signed_event = signed_event(&identity); 167 let state = Radrootsd::new( 168 identity.clone(), 169 transport_publish_config, 170 Nip46Config::default(), 171 ) 172 .expect("state"); 173 let mut state = state; 174 state.transport_publish = state 175 .transport_publish 176 .clone() 177 .with_publisher(Arc::new(RadrootsMockRelayPublishAdapter::new())); 178 let token = generate_bearer_token(); 179 let principal = state 180 .transport_publish 181 .store 182 .create_principal(PublishPrincipalInit { 183 label: "tester".to_owned(), 184 token_hash: hash_bearer_token(token.as_str()), 185 allowed_pubkeys: vec![identity.public_key_hex()], 186 allowed_kinds: vec![30_402], 187 allowed_target_policies: vec![TargetPolicyName::Nostr], 188 allowed_explicit_transport_kinds: Vec::new(), 189 allowed_nostr_source_policies: vec![NostrTargetSourcePolicy::DaemonDefaultOnly], 190 allow_request_targets: false, 191 job_visibility: if admin { 192 PublishJobVisibility::Admin 193 } else { 194 PublishJobVisibility::Own 195 }, 196 expires_at_unix: None, 197 }) 198 .expect("principal"); 199 let registry = MethodRegistry::default(); 200 let ctx = RpcContext::new(state); 201 let mut module = module(ctx.clone(), registry).expect("module"); 202 module 203 .extensions_mut() 204 .insert(TransportPublishAuthorization::Authorized(principal)); 205 (module, ctx, token, signed_event) 206 } 207 208 fn module_with_principal(admin: bool) -> (RpcModule<RpcContext>, RpcContext, String, String) { 209 module_with_principal_and_config( 210 admin, 211 TransportPublishConfig { 212 nostr: TransportPublishNostrConfig { 213 daemon_default_relays: vec!["wss://relay.example.com".to_owned()], 214 ..TransportPublishNostrConfig::default() 215 }, 216 ..TransportPublishConfig::default() 217 }, 218 ) 219 } 220 221 #[tokio::test] 222 async fn publish_event_records_job_and_deduplicates_idempotency() { 223 let (module, _ctx, _token, event) = module_with_principal(false); 224 let request = format!( 225 r#"{{ 226 "jsonrpc":"2.0", 227 "method":"transport.publish.event", 228 "params":{{ 229 "raw_event_json":{}, 230 "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, 231 "delivery_policy":{{"mode":"any"}}, 232 "idempotency_key":"idem-1" 233 }}, 234 "id":1 235 }}"#, 236 serde_json::to_string(&event).expect("event json") 237 ); 238 let (response, _stream) = module 239 .raw_json_request(request.as_str(), 1) 240 .await 241 .expect("request"); 242 assert!(response.get().contains("\"deduplicated\":false")); 243 let (response, _stream) = module 244 .raw_json_request(request.as_str(), 1) 245 .await 246 .expect("request"); 247 assert!(response.get().contains("\"deduplicated\":true")); 248 } 249 250 #[tokio::test] 251 async fn publish_event_rejects_principal_scope_gap() { 252 let (module, _ctx, _token, _pubkey) = module_with_principal(false); 253 let other_identity = DaemonIdentity::generate(); 254 let event = signed_event(&other_identity); 255 let request = format!( 256 r#"{{ 257 "jsonrpc":"2.0", 258 "method":"transport.publish.event", 259 "params":{{ 260 "raw_event_json":{}, 261 "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, 262 "delivery_policy":{{"mode":"any"}} 263 }}, 264 "id":1 265 }}"#, 266 serde_json::to_string(&event).expect("event json") 267 ); 268 let (response, _stream) = module 269 .raw_json_request(request.as_str(), 1) 270 .await 271 .expect("request"); 272 assert!(response.get().contains("unauthorized")); 273 } 274 275 #[tokio::test] 276 async fn publish_job_list_rejects_malformed_and_zero_limits() { 277 let (module, _ctx, _token, _event) = module_with_principal(false); 278 let malformed = r#"{ 279 "jsonrpc":"2.0", 280 "method":"transport.publish.job.list", 281 "params":"bad", 282 "id":1 283 }"#; 284 let (response, _stream) = module 285 .raw_json_request(malformed, 1) 286 .await 287 .expect("malformed request"); 288 assert!(response.get().contains("\"code\":-32602")); 289 290 let zero = r#"{ 291 "jsonrpc":"2.0", 292 "method":"transport.publish.job.list", 293 "params":{"limit":0}, 294 "id":1 295 }"#; 296 let (response, _stream) = module 297 .raw_json_request(zero, 1) 298 .await 299 .expect("zero request"); 300 assert!(response.get().contains("\"code\":-32602")); 301 assert!(response.get().contains("limit must be greater than zero")); 302 } 303 304 #[tokio::test] 305 async fn publish_job_list_rejects_unknown_fields() { 306 let (module, _ctx, _token, _event) = module_with_principal(false); 307 for params in [ 308 r#"{"cursor":"next"}"#, 309 r#"{"status":"publishing"}"#, 310 r#"{"limit":1,"extra":true}"#, 311 ] { 312 let request = format!( 313 r#"{{ 314 "jsonrpc":"2.0", 315 "method":"transport.publish.job.list", 316 "params":{params}, 317 "id":1 318 }}"# 319 ); 320 let (response, _stream) = module 321 .raw_json_request(request.as_str(), 1) 322 .await 323 .expect("unknown field request"); 324 assert!( 325 response.get().contains("\"code\":-32602"), 326 "{}", 327 response.get() 328 ); 329 } 330 } 331 332 #[tokio::test] 333 async fn publish_job_list_uses_configured_limit_when_omitted_and_caps_positive_limits() { 334 let mut config = TransportPublishConfig { 335 nostr: TransportPublishNostrConfig { 336 daemon_default_relays: vec!["wss://relay.example.com".to_owned()], 337 ..TransportPublishNostrConfig::default() 338 }, 339 ..TransportPublishConfig::default() 340 }; 341 config.job_list_limit = 1; 342 let (module, _ctx, _token, event) = module_with_principal_and_config(false, config); 343 for idempotency_key in ["idem-list-1", "idem-list-2"] { 344 let request = format!( 345 r#"{{ 346 "jsonrpc":"2.0", 347 "method":"transport.publish.event", 348 "params":{{ 349 "raw_event_json":{}, 350 "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}}, 351 "delivery_policy":{{"mode":"any"}}, 352 "idempotency_key":"{idempotency_key}" 353 }}, 354 "id":1 355 }}"#, 356 serde_json::to_string(&event).expect("event json") 357 ); 358 let (response, _stream) = module 359 .raw_json_request(request.as_str(), 1) 360 .await 361 .expect("publish request"); 362 assert!(response.get().contains("\"deduplicated\":false")); 363 } 364 365 let omitted = r#"{ 366 "jsonrpc":"2.0", 367 "method":"transport.publish.job.list", 368 "id":1 369 }"#; 370 let (response, _stream) = module 371 .raw_json_request(omitted, 1) 372 .await 373 .expect("omitted request"); 374 let value: serde_json::Value = 375 serde_json::from_str(response.get()).expect("omitted response json"); 376 assert_eq!(value["result"].as_array().expect("jobs").len(), 1); 377 378 let empty_array = r#"{ 379 "jsonrpc":"2.0", 380 "method":"transport.publish.job.list", 381 "params":[], 382 "id":1 383 }"#; 384 let (response, _stream) = module 385 .raw_json_request(empty_array, 1) 386 .await 387 .expect("empty array request"); 388 let value: serde_json::Value = 389 serde_json::from_str(response.get()).expect("empty array response json"); 390 assert_eq!(value["result"].as_array().expect("jobs").len(), 1); 391 392 let over_limit = r#"{ 393 "jsonrpc":"2.0", 394 "method":"transport.publish.job.list", 395 "params":{"limit":50}, 396 "id":1 397 }"#; 398 let (response, _stream) = module 399 .raw_json_request(over_limit, 1) 400 .await 401 .expect("over limit request"); 402 let value: serde_json::Value = 403 serde_json::from_str(response.get()).expect("over limit response json"); 404 assert_eq!(value["result"].as_array().expect("jobs").len(), 1); 405 } 406 407 #[test] 408 fn http_auth_finds_principal_from_hashed_token() { 409 let (_module, ctx, token, _pubkey) = module_with_principal(false); 410 let header = format!("Bearer {token}"); 411 let auth = authorize_transport_publish_request( 412 Some(header.as_str()), 413 &ctx.state.transport_publish.store, 414 ); 415 assert!(matches!(auth, TransportPublishAuthorization::Authorized(_))); 416 } 417 }