radrootsd

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

server.rs (17861B)


      1 #![forbid(unsafe_code)]
      2 
      3 use std::net::SocketAddr;
      4 
      5 use anyhow::Result;
      6 use jsonrpsee::server::middleware::rpc::{
      7     Batch, MethodResponse, Notification, Request, RpcServiceBuilder, RpcServiceT,
      8 };
      9 use jsonrpsee::server::{
     10     BatchRequestConfig, HttpBody, HttpRequest, RpcModule, ServerBuilder, ServerConfigBuilder,
     11     ServerHandle,
     12 };
     13 use jsonrpsee::types::{ErrorObject, Id};
     14 
     15 use crate::app::config::RpcConfig;
     16 use crate::core::transport_publish::TransportPublishStore;
     17 use crate::transport::jsonrpc::RpcContext;
     18 use crate::transport::jsonrpc::auth;
     19 
     20 #[derive(Clone)]
     21 struct RejectPublishNotifications<S> {
     22     service: S,
     23 }
     24 
     25 impl<S> RpcServiceT for RejectPublishNotifications<S>
     26 where
     27     S: RpcServiceT<
     28             MethodResponse = MethodResponse,
     29             NotificationResponse = MethodResponse,
     30             BatchResponse = MethodResponse,
     31         > + Clone
     32         + Send
     33         + Sync
     34         + 'static,
     35 {
     36     type MethodResponse = MethodResponse;
     37     type NotificationResponse = MethodResponse;
     38     type BatchResponse = MethodResponse;
     39 
     40     fn call<'a>(
     41         &self,
     42         request: Request<'a>,
     43     ) -> impl Future<Output = Self::MethodResponse> + Send + 'a {
     44         self.service.call(request)
     45     }
     46 
     47     fn batch<'a>(
     48         &self,
     49         requests: Batch<'a>,
     50     ) -> impl Future<Output = Self::BatchResponse> + Send + 'a {
     51         self.service.batch(requests)
     52     }
     53 
     54     fn notification<'a>(
     55         &self,
     56         notification: Notification<'a>,
     57     ) -> impl Future<Output = Self::NotificationResponse> + Send + 'a {
     58         let service = self.service.clone();
     59         async move {
     60             if notification.method_name().starts_with("transport.publish.") {
     61                 MethodResponse::error(
     62                     Id::Null,
     63                     ErrorObject::owned(
     64                         -32600,
     65                         "transport publish notifications are not accepted",
     66                         None::<()>,
     67                     ),
     68                 )
     69             } else {
     70                 service.notification(notification).await
     71             }
     72         }
     73     }
     74 }
     75 
     76 pub async fn start_server(
     77     addr: SocketAddr,
     78     rpc_cfg: &RpcConfig,
     79     transport_publish_store: TransportPublishStore,
     80     root: RpcModule<RpcContext>,
     81 ) -> Result<ServerHandle> {
     82     let mut builder = ServerConfigBuilder::new()
     83         .max_request_body_size(rpc_cfg.max_request_body_size)
     84         .max_response_body_size(rpc_cfg.max_response_body_size)
     85         .max_connections(rpc_cfg.max_connections)
     86         .max_subscriptions_per_connection(rpc_cfg.max_subscriptions_per_connection)
     87         .set_message_buffer_capacity(rpc_cfg.message_buffer_capacity);
     88 
     89     if let Some(limit) = rpc_cfg.batch_request_limit {
     90         let cfg = if limit == 0 {
     91             BatchRequestConfig::Disabled
     92         } else {
     93             BatchRequestConfig::Limit(limit)
     94         };
     95         builder = builder.set_batch_request_config(cfg);
     96     }
     97 
     98     let server_cfg = builder.build();
     99     let rpc_middleware =
    100         RpcServiceBuilder::new().layer_fn(|service| RejectPublishNotifications { service });
    101     let server = ServerBuilder::with_config(server_cfg)
    102         .set_rpc_middleware(rpc_middleware)
    103         .set_http_middleware(tower::ServiceBuilder::new().map_request(
    104             move |mut request: HttpRequest<HttpBody>| {
    105                 let transport_publish_auth = auth::authorize_transport_publish_request(
    106                     request
    107                         .headers()
    108                         .get("authorization")
    109                         .and_then(|value| value.to_str().ok()),
    110                     &transport_publish_store,
    111                 );
    112                 request.extensions_mut().insert(transport_publish_auth);
    113                 request
    114             },
    115         ))
    116         .build(addr)
    117         .await?;
    118     Ok(server.start(root))
    119 }
    120 
    121 #[cfg(test)]
    122 mod tests {
    123     use super::start_server;
    124     use crate::app::config::{
    125         Nip46Config, NostrRelayUrlPolicy, RpcConfig, TransportPublishConfig,
    126         TransportPublishNostrConfig,
    127     };
    128     use crate::app::identity_storage::DaemonIdentity;
    129     use crate::core::Radrootsd;
    130     use crate::core::transport_publish::{
    131         PublishJobVisibility, PublishPrincipalInit, PublishRelayResolveFuture,
    132         PublishRelayResolver, generate_bearer_token, hash_bearer_token,
    133     };
    134     use crate::host_nostr::Timestamp;
    135     use crate::transport::jsonrpc::methods;
    136     use crate::transport::jsonrpc::{MethodRegistry, RpcContext};
    137     use crate::transport::relay_publish::MockRelayPublishAdapter as RadrootsMockRelayPublishAdapter;
    138     use jsonrpsee::server::RpcModule;
    139     use nostr::JsonUtil;
    140     use nostr::{EventBuilder, Kind, Tag};
    141     use radroots_protocol::radrootsd::transport_publish::v5::{
    142         NostrTargetSourcePolicy, TargetPolicyName,
    143     };
    144     use serde_json::Value;
    145     use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpListener};
    146     use std::sync::Arc;
    147     use tokio::io::{AsyncReadExt, AsyncWriteExt};
    148 
    149     const RELAY_PRIMARY: &str = "ws://localhost:7777";
    150     const RELAY_PUBLIC: &str = "wss://relay.example.com";
    151 
    152     fn unused_addr() -> SocketAddr {
    153         let listener = TcpListener::bind("127.0.0.1:0").expect("bind local addr");
    154         listener.local_addr().expect("local addr")
    155     }
    156 
    157     fn signed_event_json(identity: &DaemonIdentity) -> String {
    158         // The JSON-RPC server consumes an already-signed transport fixture.
    159         EventBuilder::new(Kind::Custom(30_402), "{}")
    160             .tag(Tag::identifier("listing-1"))
    161             .custom_created_at(Timestamp::from_secs(1_700_000_000))
    162             .sign_with_keys(identity.keys())
    163             .expect("signed event")
    164             .as_json()
    165     }
    166 
    167     async fn post_json(addr: SocketAddr, body: &str, token: Option<&str>) -> String {
    168         let mut stream = tokio::net::TcpStream::connect(addr).await.expect("connect");
    169         let auth_header = token
    170             .map(|token| format!("Authorization: Bearer {token}\r\n"))
    171             .unwrap_or_default();
    172         let request = format!(
    173             "POST / HTTP/1.1\r\nHost: {addr}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n{auth_header}\r\n{body}",
    174             body.len()
    175         );
    176         stream
    177             .write_all(request.as_bytes())
    178             .await
    179             .expect("write request");
    180         let mut bytes = Vec::new();
    181         stream.read_to_end(&mut bytes).await.expect("read response");
    182         String::from_utf8(bytes).expect("response utf8")
    183     }
    184 
    185     fn publish_server_state_with_config(
    186         transport_publish_config: TransportPublishConfig,
    187         resolver: Option<Arc<dyn PublishRelayResolver>>,
    188     ) -> (
    189         Radrootsd,
    190         String,
    191         DaemonIdentity,
    192         RadrootsMockRelayPublishAdapter,
    193     ) {
    194         let identity = DaemonIdentity::generate();
    195         let mut state = Radrootsd::new(
    196             identity.clone(),
    197             transport_publish_config,
    198             Nip46Config::default(),
    199         )
    200         .expect("state");
    201         let adapter = RadrootsMockRelayPublishAdapter::new();
    202         let mut transport_publish = state.transport_publish.clone();
    203         if let Some(resolver) = resolver {
    204             transport_publish = transport_publish.with_relay_resolver(resolver);
    205         }
    206         state.transport_publish = transport_publish.with_publisher(Arc::new(adapter.clone()));
    207         let token = generate_bearer_token();
    208         state
    209             .transport_publish
    210             .store
    211             .create_principal(PublishPrincipalInit {
    212                 label: "tester".to_owned(),
    213                 token_hash: hash_bearer_token(token.as_str()),
    214                 allowed_pubkeys: vec![identity.public_key_hex()],
    215                 allowed_kinds: vec![30_402],
    216                 allowed_target_policies: vec![TargetPolicyName::Nostr],
    217                 allowed_explicit_transport_kinds: Vec::new(),
    218                 allowed_nostr_source_policies: vec![NostrTargetSourcePolicy::DaemonDefaultOnly],
    219                 allow_request_targets: false,
    220                 job_visibility: PublishJobVisibility::Own,
    221                 expires_at_unix: None,
    222             })
    223             .expect("principal");
    224         (state, token, identity, adapter)
    225     }
    226 
    227     fn publish_server_state() -> (
    228         Radrootsd,
    229         String,
    230         DaemonIdentity,
    231         RadrootsMockRelayPublishAdapter,
    232     ) {
    233         publish_server_state_with_config(
    234             TransportPublishConfig {
    235                 nostr: TransportPublishNostrConfig {
    236                     daemon_default_relays: vec![RELAY_PRIMARY.to_owned()],
    237                     relay_url_policy: NostrRelayUrlPolicy::Localhost,
    238                     ..TransportPublishNostrConfig::default()
    239                 },
    240                 ..TransportPublishConfig::default()
    241             },
    242             None,
    243         )
    244     }
    245 
    246     struct StaticPublishRelayResolver {
    247         addresses: Vec<IpAddr>,
    248     }
    249 
    250     impl StaticPublishRelayResolver {
    251         fn forbidden_localhost() -> Self {
    252             Self {
    253                 addresses: vec![IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1))],
    254             }
    255         }
    256     }
    257 
    258     impl PublishRelayResolver for StaticPublishRelayResolver {
    259         fn resolve<'a>(
    260             &'a self,
    261             _url: &'a crate::transport::relay_publish::RelayUrl,
    262         ) -> PublishRelayResolveFuture<'a> {
    263             Box::pin(async move { Ok(self.addresses.clone()) })
    264         }
    265     }
    266 
    267     async fn start_publish_server(
    268         state: Radrootsd,
    269         rpc_cfg: RpcConfig,
    270     ) -> (SocketAddr, jsonrpsee::server::ServerHandle) {
    271         let addr = unused_addr();
    272         let store = state.transport_publish.store.clone();
    273         let registry = MethodRegistry::default();
    274         let ctx = RpcContext::new(state);
    275         let mut root = RpcModule::new(ctx.clone());
    276         methods::register_all(&mut root, ctx, registry).expect("register methods");
    277         let handle = start_server(addr, &rpc_cfg, store, root)
    278             .await
    279             .expect("start server");
    280         (addr, handle)
    281     }
    282 
    283     fn json_response_body(response: &str) -> Value {
    284         let (_headers, body) = response.split_once("\r\n\r\n").expect("http body");
    285         serde_json::from_str(body).expect("json response body")
    286     }
    287 
    288     #[tokio::test]
    289     async fn raw_http_publish_event_get_and_list_preserve_signed_event() {
    290         let (state, token, identity, adapter) = publish_server_state();
    291         let event_json = signed_event_json(&identity);
    292         let (addr, handle) = start_publish_server(state, RpcConfig::default()).await;
    293         let publish = format!(
    294             r#"{{
    295                 "jsonrpc":"2.0",
    296                 "method":"transport.publish.event",
    297                 "params":{{
    298                     "raw_event_json":{},
    299                     "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}},
    300                     "delivery_policy":{{"mode":"any"}},
    301                     "idempotency_key":"raw-http-idem"
    302                 }},
    303                 "id":1
    304             }}"#,
    305             serde_json::to_string(&event_json).expect("raw event param")
    306         );
    307         let publish_response = post_json(addr, publish.as_str(), Some(token.as_str())).await;
    308         let publish_value = json_response_body(publish_response.as_str());
    309         let job_id = publish_value["result"]["job"]["job_id"]
    310             .as_str()
    311             .expect("job id")
    312             .to_owned();
    313         assert_eq!(publish_value["result"]["deduplicated"], false);
    314         assert_eq!(
    315             publish_value["result"]["job"]["status"],
    316             "delivery_satisfied"
    317         );
    318 
    319         let get = format!(
    320             r#"{{
    321                 "jsonrpc":"2.0",
    322                 "method":"transport.publish.job.get",
    323                 "params":{{"job_id":"{job_id}"}},
    324                 "id":2
    325             }}"#
    326         );
    327         let get_response = post_json(addr, get.as_str(), Some(token.as_str())).await;
    328         let get_value = json_response_body(get_response.as_str());
    329         assert_eq!(get_value["result"]["job_id"], job_id);
    330         assert_eq!(get_value["result"]["status"], "delivery_satisfied");
    331 
    332         let list = r#"{
    333             "jsonrpc":"2.0",
    334             "method":"transport.publish.job.list",
    335             "params":{"limit":10},
    336             "id":3
    337         }"#;
    338         let list_response = post_json(addr, list, Some(token.as_str())).await;
    339         let list_value = json_response_body(list_response.as_str());
    340         let jobs = list_value["result"].as_array().expect("jobs");
    341         assert_eq!(jobs.len(), 1);
    342         assert_eq!(jobs[0]["job_id"], job_id);
    343         handle.stop().expect("stop server");
    344 
    345         assert_eq!(adapter.captured_raw_events(), vec![event_json]);
    346     }
    347 
    348     #[tokio::test]
    349     async fn raw_http_publish_event_rejects_public_relay_forbidden_dns_destination() {
    350         let (state, token, identity, adapter) = publish_server_state_with_config(
    351             TransportPublishConfig {
    352                 nostr: TransportPublishNostrConfig {
    353                     daemon_default_relays: vec![RELAY_PUBLIC.to_owned()],
    354                     relay_url_policy: NostrRelayUrlPolicy::Public,
    355                     ..TransportPublishNostrConfig::default()
    356                 },
    357                 ..TransportPublishConfig::default()
    358             },
    359             Some(Arc::new(StaticPublishRelayResolver::forbidden_localhost())),
    360         );
    361         let event_json = signed_event_json(&identity);
    362         let (addr, handle) = start_publish_server(state, RpcConfig::default()).await;
    363         let publish = format!(
    364             r#"{{
    365                 "jsonrpc":"2.0",
    366                 "method":"transport.publish.event",
    367                 "params":{{
    368                     "raw_event_json":{},
    369                     "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}},
    370                     "delivery_policy":{{"mode":"any"}},
    371                     "idempotency_key":"raw-http-public-dns-reject"
    372                 }},
    373                 "id":1
    374             }}"#,
    375             serde_json::to_string(&event_json).expect("raw event param")
    376         );
    377         let publish_response = post_json(addr, publish.as_str(), Some(token.as_str())).await;
    378         handle.stop().expect("stop server");
    379 
    380         let publish_value = json_response_body(publish_response.as_str());
    381         let job = &publish_value["result"]["job"];
    382         assert_eq!(publish_value["result"]["deduplicated"], false);
    383         assert_eq!(job["status"], "delivery_unsatisfied_terminal");
    384         assert_eq!(job["last_error"], "delivery_unsatisfied");
    385         let targets = job["targets"].as_array().expect("target outcomes");
    386         assert_eq!(targets.len(), 1);
    387         assert_eq!(targets[0]["endpoint_uri"], RELAY_PUBLIC);
    388         assert_eq!(targets[0]["source"], "daemon_default");
    389         assert_eq!(targets[0]["outcome_kind"], "target_rejected");
    390         assert_eq!(targets[0]["attempted"], false);
    391         assert!(adapter.captured_raw_events().is_empty());
    392     }
    393 
    394     #[tokio::test]
    395     async fn publish_notifications_do_not_create_jobs() {
    396         let (state, token, identity, _adapter) = publish_server_state();
    397         let store = state.transport_publish.store.clone();
    398         let (addr, handle) = start_publish_server(state, RpcConfig::default()).await;
    399         let notification = format!(
    400             r#"{{
    401                 "jsonrpc":"2.0",
    402                 "method":"transport.publish.event",
    403                 "params":{{
    404                     "raw_event_json":{},
    405                     "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}},
    406                     "delivery_policy":{{"mode":"any"}}
    407                 }}
    408             }}"#,
    409             serde_json::to_string(&signed_event_json(&identity)).expect("raw event param")
    410         );
    411         let response = post_json(addr, notification.as_str(), Some(token.as_str())).await;
    412         handle.stop().expect("stop server");
    413 
    414         assert!(
    415             response.contains("transport publish notifications are not accepted")
    416                 || response.ends_with("\r\n\r\n")
    417         );
    418         let principal = store
    419             .principal_for_token_hash(hash_bearer_token(token.as_str()).as_str())
    420             .expect("principal lookup")
    421             .expect("principal");
    422         assert!(
    423             store
    424                 .list_jobs_for_principal(&principal, 10)
    425                 .expect("jobs")
    426                 .is_empty()
    427         );
    428     }
    429 
    430     #[tokio::test]
    431     async fn batch_requests_are_disabled_by_default() {
    432         let (state, token, identity, _adapter) = publish_server_state();
    433         let store = state.transport_publish.store.clone();
    434         let (addr, handle) = start_publish_server(state, RpcConfig::default()).await;
    435         let batch = format!(
    436             r#"[{{
    437                 "jsonrpc":"2.0",
    438                 "method":"transport.publish.event",
    439                 "params":{{
    440                     "raw_event_json":{},
    441                     "target_policy":{{"kind":"nostr","source_policy":"daemon_default_only","relay_urls":[]}},
    442                     "delivery_policy":{{"mode":"any"}}
    443                 }},
    444                 "id":1
    445             }}]"#,
    446             serde_json::to_string(&signed_event_json(&identity)).expect("raw event param")
    447         );
    448         let response = post_json(addr, batch.as_str(), Some(token.as_str())).await;
    449         handle.stop().expect("stop server");
    450 
    451         assert!(
    452             response.contains("Batched requests are not supported by this server"),
    453             "{response}"
    454         );
    455         let principal = store
    456             .principal_for_token_hash(hash_bearer_token(token.as_str()).as_str())
    457             .expect("principal lookup")
    458             .expect("principal");
    459         assert!(
    460             store
    461                 .list_jobs_for_principal(&principal, 10)
    462                 .expect("jobs")
    463                 .is_empty()
    464         );
    465     }
    466 }