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;