1//! `jev-client`'s ports on a Cloudflare Worker: the Worker's `fetch` for the 2//! [`Transport`], and its clock, timers and `Math.random` for the [`Runtime`]. 3//! 4//! A Worker cannot open the HTTP/2 connection `jev-http` holds, and cannot read 5//! or write the on-disk ledger, so this crate is the transport half only: the 6//! retry policy, verification and price stay in `jev-client` and 7//! `jev-protocol`. A spend guard, where the Worker needs one, belongs in a 8//! Durable Object beside it. 9#![forbid(unsafe_code)] 10 11use std::future::Future; 12use std::pin::pin; 13use std::time::{Duration, SystemTime}; 14 15use bytes::Bytes; 16use futures_util::future::{Either, select}; 17use jev_client::{HttpRequest, HttpResponse, Instant, Runtime, Transport, TransportError, TransportErrorKind}; 18use worker::{Delay, Fetch, Headers, Method, Request, RequestInit}; 19 20/// The most of a response body read before giving up, as `jev-http` caps it. 21pub const MAX_BODY: usize = 8 << 20; 22 23/// Posts each request to one URL with the Worker's `fetch`. 24/// 25/// `fetch` pools its own connections, out of reach of the Worker, so there is 26/// no connection to distrust and a send never reports `NotSent`: a failure 27/// before a response is `Connect`, and is retried by the client's policy. 28#[derive(Clone, Debug)] 29pub struct FetchTransport { 30 url: String, 31} 32 33impl FetchTransport { 34 /// The API itself, [`jev_protocol::ENDPOINT`]. 35 pub fn api() -> Self { 36 Self::new(jev_protocol::ENDPOINT) 37 } 38 39 /// Something other than the API: a relay, or a test server. 40 pub fn new(url: impl Into<String>) -> Self { 41 FetchTransport { url: url.into() } 42 } 43} 44 45fn error(kind: TransportErrorKind, what: &str, cause: impl std::fmt::Display) -> TransportError { 46 TransportError::new(kind, format!("{what}: {cause}")) 47} 48 49fn request(url: &str, sent: &HttpRequest) -> Result<Request, TransportError> { 50 let headers = Headers::new(); 51 for (name, value) in &sent.headers { 52 let value = value.to_str().map_err(|e| error(TransportErrorKind::Config, "header", e))?; 53 headers.set(name.as_str(), value).map_err(|e| error(TransportErrorKind::Config, "header", e))?; 54 } 55 let body = js_sys::Uint8Array::from(sent.body.as_ref()); 56 let mut init = RequestInit::new(); 57 init.with_method(Method::Post).with_headers(headers).with_body(Some(body.into())); 58 Request::new_with_init(url, &init).map_err(|e| error(TransportErrorKind::Config, "request", e)) 59} 60 61fn header_map(headers: &Headers) -> http::HeaderMap { 62 let mut map = http::HeaderMap::new(); 63 for (name, value) in headers.entries() { 64 if let (Ok(name), Ok(value)) = 65 (http::HeaderName::try_from(name.as_str()), http::HeaderValue::try_from(value.as_str())) 66 { 67 map.append(name, value); 68 } 69 } 70 map 71} 72 73impl Transport for FetchTransport { 74 async fn send(&self, sent: HttpRequest) -> Result<HttpResponse, TransportError> { 75 let request = request(&self.url, &sent)?; 76 let mut response = 77 Fetch::Request(request).send().await.map_err(|e| error(TransportErrorKind::Connect, "fetch", e))?; 78 let status = response.status_code(); 79 let headers = header_map(response.headers()); 80 let body = response.bytes().await.map_err(|e| error(TransportErrorKind::Interrupted, "reading the body", e))?; 81 if body.len() > MAX_BODY { 82 return Err(TransportError::new( 83 TransportErrorKind::TooLarge, 84 format!("the response body is {} bytes, over the {MAX_BODY}-byte cap", body.len()), 85 )); 86 } 87 Ok(HttpResponse { status, headers, body: Bytes::from(body) }) 88 } 89 90 fn distrust_connection(&self) {} 91} 92 93/// The Worker's clock and timers. 94/// 95/// Workers advance `Date.now()` and `performance.now()` only across I/O, never 96/// during CPU work. That is still right for the client: every interval it 97/// measures spans a fetch or a timer. 98#[derive(Clone, Copy, Debug, Default)] 99pub struct WorkerRuntime; 100 101impl Runtime for WorkerRuntime { 102 fn now(&self) -> Instant { 103 Instant::now() 104 } 105 106 fn wall_clock(&self) -> SystemTime { 107 SystemTime::UNIX_EPOCH + Duration::from_secs_f64(js_sys::Date::now() / 1000.0) 108 } 109 110 fn sleep(&self, duration: Duration) -> impl Future<Output = ()> { 111 Delay::from(duration) 112 } 113 114 fn timeout<F: Future>(&self, duration: Duration, future: F) -> impl Future<Output = Option<F::Output>> { 115 async move { 116 match select(pin!(future), Delay::from(duration)).await { 117 Either::Left((output, _)) => Some(output), 118 Either::Right(((), _)) => None, 119 } 120 } 121 } 122 123 fn jitter(&self) -> f64 { 124 js_sys::Math::random() 125 } 126}