radrootsd

JSON-RPC bridge for Radroots event publishing
git clone https://radroots.dev/git/radrootsd.git
Log | Files | Refs | README | LICENSE

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, &registry)?;
     26     register_event(&mut module, &registry)?;
     27     register_job_get(&mut module, &registry)?;
     28     register_job_list(&mut module, &registry)?;
     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 }