lib.rsannotatedlib.rssource126 lines · 4.6 KB · raw

jev-client's ports on a Cloudflare Worker: the Worker's fetch for the [Transport], and its clock, timers and Math.random for the [Runtime].

A Worker cannot open the HTTP/2 connection jev-http holds, and cannot read or write the on-disk ledger, so this crate is the transport half only: the retry policy, verification and price stay in jev-client and jev-protocol. A spend guard, where the Worker needs one, belongs in a Durable Object beside it.

9#![forbid(unsafe_code)]
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};

The most of a response body read before giving up, as jev-http caps it.

21pub const MAX_BODY: usize = 8 << 20;

Posts each request to one URL with the Worker's fetch.

fetch pools its own connections, out of reach of the Worker, so there is no connection to distrust and a send never reports NotSent: a failure before a response is Connect, and is retried by the client's policy.

28#[derive(Clone, Debug)]
29pub struct FetchTransport {
30    url: String,
31}
33impl FetchTransport {

The API itself, [jev_protocol::ENDPOINT].

35    pub fn api() -> Self {
36        Self::new(jev_protocol::ENDPOINT)
37    }

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}
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}

The Worker's clock and timers.

Workers advance Date.now() and performance.now() only across I/O, never during CPU work. That is still right for the client: every interval it measures spans a fetch or a timer.

98#[derive(Clone, Copy, Debug, Default)]
99pub struct WorkerRuntime;
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}