lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

radrootsd.rs (10700B)


      1 use core::fmt;
      2 use core::time::Duration;
      3 
      4 use radroots_event::draft::SignedEvent;
      5 use radroots_protocol::radrootsd::transport_publish::v5::{
      6     DeliveryPolicy as TransportPublishDeliveryPolicy, Error as TransportPublishProtocolError,
      7     EventRequest as TransportPublishEventRequest, EventResponse as TransportPublishEventResponse,
      8     METHOD_EVENT, TargetPolicy as TransportPublishTargetPolicy,
      9 };
     10 use reqwest::header::{AUTHORIZATION, CONTENT_TYPE, HeaderMap, HeaderValue};
     11 use serde::{Deserialize, Serialize, de::DeserializeOwned};
     12 use serde_json::{Value, json};
     13 
     14 pub const SDK_RADROOTSD_PUBLISH_REQUEST_ID: &str = "radroots-sdk-transport-publish-event";
     15 pub const SDK_RADROOTSD_PUBLISH_MAX_TARGETS: usize = 20;
     16 
     17 #[derive(Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
     18 pub enum RadrootsdAuth {
     19     #[default]
     20     None,
     21     BearerToken(String),
     22 }
     23 
     24 impl fmt::Debug for RadrootsdAuth {
     25     fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
     26         match self {
     27             Self::None => f.write_str("None"),
     28             Self::BearerToken(_) => f.write_str("BearerToken(<redacted>)"),
     29         }
     30     }
     31 }
     32 
     33 #[derive(Clone, Debug, PartialEq, Eq)]
     34 pub struct RadrootsdPublishConfig {
     35     pub endpoint: String,
     36     pub auth: RadrootsdAuth,
     37     pub timeout: Duration,
     38 }
     39 
     40 impl RadrootsdPublishConfig {
     41     pub fn new(endpoint: impl Into<String>) -> Self {
     42         Self {
     43             endpoint: endpoint.into(),
     44             auth: RadrootsdAuth::None,
     45             timeout: Duration::from_secs(10),
     46         }
     47     }
     48 
     49     pub fn with_auth(mut self, auth: RadrootsdAuth) -> Self {
     50         self.auth = auth;
     51         self
     52     }
     53 
     54     pub fn with_timeout(mut self, timeout: Duration) -> Self {
     55         self.timeout = timeout;
     56         self
     57     }
     58 }
     59 
     60 #[derive(Clone, Debug, PartialEq, Eq)]
     61 pub struct RadrootsdPublishAdapter {
     62     config: RadrootsdPublishConfig,
     63 }
     64 
     65 impl RadrootsdPublishAdapter {
     66     pub fn new(config: RadrootsdPublishConfig) -> Self {
     67         Self { config }
     68     }
     69 
     70     #[cfg(test)]
     71     pub fn config(&self) -> &RadrootsdPublishConfig {
     72         &self.config
     73     }
     74 
     75     pub async fn publish_signed_event(
     76         &self,
     77         request: RadrootsdPublishRequest,
     78     ) -> Result<TransportPublishEventResponse, RadrootsdError> {
     79         let event_identity =
     80             RadrootsdPublishEventIdentity::from_signed_event(&request.signed_event);
     81         let request = request.into_protocol_request();
     82         request
     83             .validate(SDK_RADROOTSD_PUBLISH_MAX_TARGETS)
     84             .map_err(RadrootsdError::from_protocol)?;
     85         let response = publish_event(
     86             self.config.endpoint.as_str(),
     87             &self.config.auth,
     88             &request,
     89             self.config.timeout,
     90         )
     91         .await?;
     92         validate_transport_publish_response_for_request(&request, &event_identity, &response)?;
     93         Ok(response)
     94     }
     95 }
     96 
     97 #[derive(Clone, Debug, PartialEq, Eq)]
     98 pub struct RadrootsdPublishRequest {
     99     pub signed_event: SignedEvent,
    100     pub target_policy: TransportPublishTargetPolicy,
    101     pub delivery_policy: TransportPublishDeliveryPolicy,
    102     pub idempotency_key: Option<String>,
    103     pub timeout_ms: Option<u64>,
    104 }
    105 
    106 impl RadrootsdPublishRequest {
    107     fn into_protocol_request(self) -> TransportPublishEventRequest {
    108         TransportPublishEventRequest {
    109             raw_event_json: self.signed_event.raw_json().to_owned(),
    110             target_policy: self.target_policy,
    111             delivery_policy: self.delivery_policy,
    112             idempotency_key: self.idempotency_key,
    113             timeout_ms: self.timeout_ms,
    114         }
    115     }
    116 }
    117 
    118 #[derive(Clone, Debug, PartialEq, Eq)]
    119 struct RadrootsdPublishEventIdentity {
    120     event_id: String,
    121     pubkey: String,
    122     kind: u32,
    123 }
    124 
    125 impl RadrootsdPublishEventIdentity {
    126     fn from_signed_event(event: &SignedEvent) -> Self {
    127         Self {
    128             event_id: event.id_str().to_owned(),
    129             pubkey: event.pubkey().to_hex().to_owned(),
    130             kind: event.kind(),
    131         }
    132     }
    133 }
    134 
    135 #[derive(Debug, Clone, PartialEq, Eq)]
    136 pub enum RadrootsdError {
    137     InvalidAuthHeader(String),
    138     InvalidRequest(String),
    139     Http(String),
    140     JsonRpc { code: i64, message: String },
    141     MalformedResponse(String),
    142 }
    143 
    144 impl RadrootsdError {
    145     fn from_protocol(error: TransportPublishProtocolError) -> Self {
    146         Self::InvalidRequest(error.to_string())
    147     }
    148 }
    149 
    150 impl fmt::Display for RadrootsdError {
    151     fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
    152         match self {
    153             Self::InvalidAuthHeader(value) => {
    154                 write!(f, "invalid radrootsd bearer token header: {value}")
    155             }
    156             Self::InvalidRequest(value) => f.write_str(value),
    157             Self::Http(value) => f.write_str(value),
    158             Self::MalformedResponse(value) => f.write_str(value),
    159             Self::JsonRpc { code, message } => {
    160                 write!(f, "radrootsd jsonrpc failed {code}: {message}")
    161             }
    162         }
    163     }
    164 }
    165 
    166 impl std::error::Error for RadrootsdError {}
    167 
    168 #[derive(Debug, Deserialize)]
    169 struct JsonRpcEnvelope<T> {
    170     jsonrpc: Option<String>,
    171     id: Option<Value>,
    172     result: Option<T>,
    173     error: Option<JsonRpcError>,
    174 }
    175 
    176 #[derive(Debug, Deserialize)]
    177 struct JsonRpcError {
    178     code: i64,
    179     message: String,
    180 }
    181 
    182 pub async fn publish_event(
    183     endpoint: &str,
    184     auth: &RadrootsdAuth,
    185     request: &TransportPublishEventRequest,
    186     timeout: Duration,
    187 ) -> Result<TransportPublishEventResponse, RadrootsdError> {
    188     jsonrpc_call(
    189         endpoint,
    190         auth,
    191         SDK_RADROOTSD_PUBLISH_REQUEST_ID,
    192         METHOD_EVENT,
    193         request,
    194         timeout,
    195     )
    196     .await
    197 }
    198 
    199 fn auth_headers(auth: &RadrootsdAuth) -> Result<HeaderMap, RadrootsdError> {
    200     let mut headers = HeaderMap::new();
    201     match auth {
    202         RadrootsdAuth::None => Ok(headers),
    203         RadrootsdAuth::BearerToken(token) => {
    204             let header = format!("Bearer {token}");
    205             let value = HeaderValue::from_str(header.as_str())
    206                 .map_err(|err| RadrootsdError::InvalidAuthHeader(err.to_string()))?;
    207             headers.insert(AUTHORIZATION, value);
    208             Ok(headers)
    209         }
    210     }
    211 }
    212 
    213 #[cfg(test)]
    214 pub fn publish_event_request_json(
    215     request: &TransportPublishEventRequest,
    216 ) -> Result<Value, RadrootsdError> {
    217     Ok(serde_json::to_value(request).expect("radrootsd transport publish request serializes"))
    218 }
    219 
    220 fn http_status_error(status: reqwest::StatusCode, body: &str) -> RadrootsdError {
    221     let body_summary = if body.is_empty() {
    222         "response body empty".to_owned()
    223     } else {
    224         format!("response body omitted ({} bytes)", body.len())
    225     };
    226     RadrootsdError::Http(format!(
    227         "radrootsd returned http {}: {}",
    228         status.as_u16(),
    229         body_summary
    230     ))
    231 }
    232 
    233 fn decode_jsonrpc_response<R>(
    234     method: &str,
    235     expected_id: &str,
    236     body: &str,
    237 ) -> Result<R, RadrootsdError>
    238 where
    239     R: DeserializeOwned,
    240 {
    241     let envelope: JsonRpcEnvelope<R> = serde_json::from_str(body).map_err(|err| {
    242         RadrootsdError::MalformedResponse(format!("decode radrootsd {method} response: {err}"))
    243     })?;
    244     if envelope.jsonrpc.as_deref() != Some("2.0") {
    245         return Err(RadrootsdError::MalformedResponse(format!(
    246             "radrootsd {method} returned invalid jsonrpc version"
    247         )));
    248     }
    249     let expected_id_value = Value::String(expected_id.to_owned());
    250     if envelope.id.as_ref() != Some(&expected_id_value) {
    251         return Err(RadrootsdError::MalformedResponse(format!(
    252             "radrootsd {method} returned mismatched jsonrpc id"
    253         )));
    254     }
    255     match (envelope.result, envelope.error) {
    256         (Some(result), None) => Ok(result),
    257         (None, Some(error)) => Err(RadrootsdError::JsonRpc {
    258             code: error.code,
    259             message: error.message,
    260         }),
    261         (Some(_), Some(error)) => Err(RadrootsdError::MalformedResponse(format!(
    262             "radrootsd {method} returned result and error: {} {}",
    263             error.code, error.message
    264         ))),
    265         (None, None) => Err(RadrootsdError::MalformedResponse(format!(
    266             "radrootsd {method} returned neither result nor error"
    267         ))),
    268     }
    269 }
    270 
    271 async fn jsonrpc_call<P, R>(
    272     endpoint: &str,
    273     auth: &RadrootsdAuth,
    274     request_id: &str,
    275     method: &str,
    276     params: &P,
    277     timeout: Duration,
    278 ) -> Result<R, RadrootsdError>
    279 where
    280     P: Serialize + ?Sized,
    281     R: DeserializeOwned,
    282 {
    283     let client = reqwest::Client::builder()
    284         .timeout(timeout)
    285         .build()
    286         .map_err(|err| RadrootsdError::Http(format!("build radrootsd client: {err}")))?;
    287     let mut request_builder = client
    288         .post(endpoint)
    289         .headers(auth_headers(auth)?)
    290         .json(&json!({
    291             "jsonrpc": "2.0",
    292             "id": request_id,
    293             "method": method,
    294             "params": params,
    295         }));
    296 
    297     request_builder = request_builder.header(CONTENT_TYPE, "application/json");
    298 
    299     let response = request_builder
    300         .send()
    301         .await
    302         .map_err(|err| RadrootsdError::Http(format!("send radrootsd {method} request: {err}")))?;
    303     let status = response.status();
    304     let body = response
    305         .text()
    306         .await
    307         .map_err(|err| RadrootsdError::Http(format!("read radrootsd response body: {err}")))?;
    308 
    309     if !status.is_success() {
    310         return Err(http_status_error(status, body.as_str()));
    311     }
    312 
    313     decode_jsonrpc_response(method, request_id, body.as_str())
    314 }
    315 
    316 fn validate_transport_publish_response_for_request(
    317     request: &TransportPublishEventRequest,
    318     event_identity: &RadrootsdPublishEventIdentity,
    319     response: &TransportPublishEventResponse,
    320 ) -> Result<(), RadrootsdError> {
    321     response.job.validate().map_err(|error| {
    322         RadrootsdError::MalformedResponse(format!(
    323             "radrootsd transport publish response invalid: {error}"
    324         ))
    325     })?;
    326     if response.job.event_id != event_identity.event_id {
    327         return Err(response_mismatch("event_id"));
    328     }
    329     if response.job.pubkey != event_identity.pubkey {
    330         return Err(response_mismatch("pubkey"));
    331     }
    332     if response.job.event_kind != event_identity.kind {
    333         return Err(response_mismatch("event_kind"));
    334     }
    335     if response.job.delivery_policy != request.delivery_policy {
    336         return Err(response_mismatch("delivery_policy"));
    337     }
    338     if response.job.target_policy != request.target_policy {
    339         return Err(response_mismatch("target_policy"));
    340     }
    341     Ok(())
    342 }
    343 
    344 fn response_mismatch(field: &str) -> RadrootsdError {
    345     RadrootsdError::MalformedResponse(format!(
    346         "radrootsd transport publish response {field} does not match request"
    347     ))
    348 }
    349 
    350 #[cfg(test)]
    351 #[path = "../../tests/unit/adapters_radrootsd_tests.rs"]
    352 mod tests;