jevcrates.git / jev-http / src / transport.rs

jev-client's [Transport] for a native process: one HTTP/2 connection over rustls on tokio, kept warm, and a per-process courtesy bucket under the vendor's documented ceilings. Ported from postjevsql-pg's https.rs (which drives the same exchange on a Postgres backend's executor); it sends and classifies, and retrying, budgets and statuses are jev-client's.

8use std::net::SocketAddr;
9use std::path::{Path, PathBuf};
10use std::sync::{Arc, Mutex, PoisonError};
11use std::time::{Duration, Instant};
13use bytes::Bytes;
14use http_body_util::{BodyExt, Full, LengthLimitError, Limited};
15use hyper::client::conn::http2::SendRequest;
16use hyper_util::rt::{TokioExecutor, TokioIo, TokioTimer};
17use jev_client::{HttpRequest, HttpResponse, Transport, TransportError, TransportErrorKind};
18use rustls::ClientConfig;
19use rustls::pki_types::ServerName;
20use tokio::net::TcpStream;
21use tokio_rustls::TlsConnector;
22
23pub use jev_protocol::ENDPOINT;

The vendor's own documented ceilings (digest §4): 250,000 tokens a second, 1,200 requests a minute. "Rate limits are adjusting dynamically... can change without notice" per the same section, so these are a floor to stay well under, not a number to lean on.

29const VENDOR_TOKENS_PER_SECOND: f64 = 250_000.0;
30const VENDOR_REQUESTS_PER_MINUTE: f64 = 1_200.0;

Cloudflare closes idle connections at 400 s; reconnect well before.

33const IDLE_LIMIT: Duration = Duration::from_secs(300);

The response cap (pg_typesafe's 8 MB, as postjevsql's).

35const MAX_BODY: usize = 8 << 20;

Time for each address but the last, so one blackholed address (a broken IPv6 route) cannot spend a whole attempt.

38const CONNECT_PER_ADDRESS: Duration = Duration::from_secs(2);

HTTP/2 pings: one after this long without a frame, and the connection is closed if it goes unanswered for PING_TIMEOUT (postjevsql's contract §2 values).

42const PING_INTERVAL: Duration = Duration::from_secs(20);
43const PING_TIMEOUT: Duration = Duration::from_secs(10);

An https URL, resolved into what a connection needs.

46#[derive(Clone, Debug)]
47pub struct Endpoint {
48    uri: http::Uri,
49    host: String,
50    port: u16,
51    server_name: ServerName<'static>,
52    ca_file: Option<PathBuf>,
53}
55impl Endpoint {

The API itself, [ENDPOINT].

57    pub fn api() -> Self {
58        Self::parse(ENDPOINT).expect("ENDPOINT parses")
59    }

Something other than the API: a relay, or a test server.

62    pub fn parse(url: &str) -> Result<Self, String> {
63        let uri: http::Uri = url.parse().map_err(|e| format!("endpoint {url}: {e}"))?;
64        if uri.scheme() != Some(&http::uri::Scheme::HTTPS) {
65            return Err(format!("the endpoint must be https: {url}"));
66        }
67        let host = uri.host().ok_or_else(|| format!("the endpoint has no host: {url}"))?;
68        let bare = host.trim_start_matches('[').trim_end_matches(']').to_owned();
69        let server_name = ServerName::try_from(bare.clone()).map_err(|e| format!("endpoint host {host}: {e}"))?;
70        Ok(Endpoint { port: uri.port_u16().unwrap_or(443), host: bare, server_name, ca_file: None, uri })
71    }

Also trust the certificates in this PEM file, besides the Mozilla roots: for tests (jev-mock's throwaway CA) and private endpoints, as postjevsql's jev.ca_file.

76    pub fn trusting(mut self, ca_file: impl Into<PathBuf>) -> Self {
77        self.ca_file = Some(ca_file.into());
78        self
79    }
81    pub fn ca_file(&self) -> Option<&Path> {
82        self.ca_file.as_deref()
83    }
84
85    pub fn uri(&self) -> &http::Uri {
86        &self.uri
87    }
88}
89
90struct Warm {
91    sender: SendRequest<Full<Bytes>>,
92    used: Instant,
93}

Sends over one kept HTTP/2 connection, reopened when it dies, idles past IDLE_LIMIT, or the client distrusts it.

97pub struct HttpsTransport {
98    endpoint: Endpoint,
99    tls: Arc<ClientConfig>,
100    warm: Mutex<Option<Warm>>,

This process's own approximation of the vendor's documented limits. Approximate because it is not shared across processes the way the ledger is - and that is fine, because respecting the vendor's limit is a courtesy to stay well clear of a 429, not the invariant the ledger exists to hold exactly.

106    requests: Bucket,
107    tokens: Bucket,
108}
110impl HttpsTransport {
111    pub fn new(endpoint: Endpoint, tls: Arc<ClientConfig>) -> Self {
112        HttpsTransport {
113            endpoint,
114            tls,
115            warm: Mutex::new(None),
116            requests: Bucket::new(VENDOR_REQUESTS_PER_MINUTE, VENDOR_REQUESTS_PER_MINUTE / 60.0),
117            tokens: Bucket::new(VENDOR_TOKENS_PER_SECOND, VENDOR_TOKENS_PER_SECOND),
118        }
119    }
120
121    fn take_warm(&self) -> Option<SendRequest<Full<Bytes>>> {
122        let mut warm = self.warm.lock().unwrap_or_else(PoisonError::into_inner);
123        match warm.take() {
124            Some(w) if !w.sender.is_closed() && w.used.elapsed() < IDLE_LIMIT => {
125                let sender = w.sender.clone();
126                *warm = Some(w);
127                Some(sender)
128            }
129            _ => None,
130        }
131    }
132
133    fn keep(&self, sender: SendRequest<Full<Bytes>>) {
134        *self.warm.lock().unwrap_or_else(PoisonError::into_inner) = Some(Warm { sender, used: Instant::now() });
135    }
136
137    fn forget(&self) {
138        *self.warm.lock().unwrap_or_else(PoisonError::into_inner) = None;
139    }
140
141    async fn connection(&self) -> Result<SendRequest<Full<Bytes>>, TransportError> {
142        if let Some(sender) = self.take_warm() {
143            return Ok(sender);
144        }
145        let sender = connect(&self.endpoint, self.tls.clone()).await?;
146        self.keep(sender.clone());
147        Ok(sender)
148    }
149}
150
151impl Transport for HttpsTransport {

Under the vendor's ceilings: one request, and the body's length in characters as the token count (no token is shorter than one character, so this can only overstate).

155    async fn admit(&self, request: &HttpRequest) {
156        self.requests.take(1.0).await;
157        let chars = String::from_utf8_lossy(&request.body).chars().count() as f64;
158        self.tokens.take(chars.min(VENDOR_TOKENS_PER_SECOND)).await;
159    }
161    async fn send(&self, request: HttpRequest) -> Result<HttpResponse, TransportError> {
162        let mut sender = self.connection().await?;
163        let mut builder = http::Request::post(self.endpoint.uri.clone());
164        *builder.headers_mut().expect("a fresh builder has headers") = request.headers;
165        let request = builder
166            .body(Full::new(request.body))
167            .map_err(|e| TransportError::new(TransportErrorKind::Config, format!("building the request: {e}")))?;
168
169        let response = match sender.try_send_request(request).await {
170            Ok(response) => response,
171            Err(mut e) => {
172                self.forget();
173                // hyper hands the request back when it never left.
174                let kind = if e.take_message().is_some() || never_processed(e.error()) {
175                    TransportErrorKind::NotSent
176                } else {
177                    TransportErrorKind::Interrupted
178                };
179                return Err(TransportError::new(kind, format!("sending to {}: {}", self.endpoint.uri, e.error())));
180            }
181        };
182        let (parts, body) = response.into_parts();
183        let body = match Limited::new(body, MAX_BODY).collect().await {
184            Ok(collected) => collected.to_bytes(),
185            Err(e) if e.is::<LengthLimitError>() => {
186                return Err(TransportError::new(
187                    TransportErrorKind::TooLarge,
188                    format!("Jev's response exceeded {} MiB", MAX_BODY >> 20),
189                ));
190            }
191            Err(e) => {
192                self.forget();
193                return Err(TransportError::new(TransportErrorKind::Interrupted, format!("reading Jev's response: {e}")));
194            }
195        };
196        self.keep(sender);
197        Ok(HttpResponse { status: parts.status.as_u16(), headers: parts.headers, body })
198    }
199
200    fn distrust_connection(&self) {
201        self.forget();
202    }
203}

True when h2 says the server never processed the stream: refused (REFUSED_STREAM), or cut by a GOAWAY whose last stream id is below ours. Such a request can be resent without being billed twice.

208fn never_processed(error: &hyper::Error) -> bool {
209    let mut source: Option<&(dyn std::error::Error + 'static)> = std::error::Error::source(error);
210    while let Some(e) = source {
211        if let Some(h2) = e.downcast_ref::<h2::Error>() {
212            return h2.reason() == Some(h2::Reason::REFUSED_STREAM) || (h2.is_go_away() && h2.is_remote());
213        }
214        source = e.source();
215    }
216    false
217}
219async fn connect(endpoint: &Endpoint, tls: Arc<ClientConfig>) -> Result<SendRequest<Full<Bytes>>, TransportError> {
220    use TransportErrorKind::{Connect, Tls};
221    let addrs: Vec<SocketAddr> = tokio::net::lookup_host((endpoint.host.as_str(), endpoint.port))
222        .await
223        .map_err(|e| TransportError::new(Connect, format!("resolving {}: {e}", endpoint.host)))?
224        .collect();
225    let mut failures = Vec::new();
226    let mut tcp = None;
227    for (i, addr) in addrs.iter().enumerate() {
228        let attempt = TcpStream::connect(*addr);
229        let result = if i + 1 < addrs.len() {
230            tokio::time::timeout(CONNECT_PER_ADDRESS, attempt).await.unwrap_or_else(|_| {
231                Err(std::io::Error::new(std::io::ErrorKind::TimedOut, "no answer within 2 s"))
232            })
233        } else {
234            attempt.await
235        };
236        match result {
237            Ok(stream) => {
238                tcp = Some(stream);
239                break;
240            }
241            Err(e) => failures.push(format!("{addr}: {e}")),
242        }
243    }
244    let tcp = tcp.ok_or_else(|| {
245        let why = if failures.is_empty() { "no addresses".to_owned() } else { failures.join("; ") };
246        TransportError::new(Connect, format!("connecting to {}: {why}", endpoint.uri))
247    })?;
248
249    let tls = TlsConnector::from(tls).connect(endpoint.server_name.clone(), tcp).await.map_err(|e| {
250        // rustls refusing (a certificate, a protocol alert) will refuse
251        // again; an I/O failure mid-handshake is the network.
252        let refused = e.get_ref().is_some_and(|inner| inner.is::<rustls::Error>());
253        TransportError::new(if refused { Tls } else { Connect }, format!("TLS with {}: {e}", endpoint.uri))
254    })?;
255    if tls.get_ref().1.alpn_protocol() != Some(b"h2") {
256        return Err(TransportError::new(Tls, format!("{} did not negotiate HTTP/2", endpoint.uri)));
257    }
258    // hyper's fixed receive windows are below an 8 MB answer; BDP-based
259    // adaptive windows grow to what the path carries.
260    let (sender, connection) = hyper::client::conn::http2::Builder::new(TokioExecutor::new())
261        .adaptive_window(true)
262        .timer(TokioTimer::new())
263        .keep_alive_interval(PING_INTERVAL)
264        .keep_alive_timeout(PING_TIMEOUT)
265        .keep_alive_while_idle(true)
266        .handshake(TokioIo::new(tls))
267        .await
268        .map_err(|e| TransportError::new(Connect, format!("HTTP/2 with {}: {e}", endpoint.uri)))?;
269    tokio::spawn(async move {
270        let _ = connection.await;
271    });
272    Ok(sender)
273}

A client-side approximation of a token-bucket rate limiter: capacity tokens available at once, refilling at refill_per_sec a second. Used twice, at two different units (requests, and Jev's own input tokens), to stay under the vendor's two independently documented ceilings.

279struct Bucket {
280    capacity: f64,
281    refill_per_sec: f64,
282    state: Mutex<(f64, Instant)>,
283}
285impl Bucket {
286    fn new(capacity: f64, refill_per_sec: f64) -> Self {
287        Self { capacity, refill_per_sec, state: Mutex::new((capacity, Instant::now())) }
288    }

Wait, if necessary, until amount is available, then spend it.

291    async fn take(&self, amount: f64) {
292        loop {
293            let wait = {
294                let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
295                let elapsed = state.1.elapsed().as_secs_f64();
296                state.0 = (state.0 + elapsed * self.refill_per_sec).min(self.capacity);
297                state.1 = Instant::now();
298                if state.0 >= amount {
299                    state.0 -= amount;
300                    None
301                } else {
302                    Some(Duration::from_secs_f64((amount - state.0) / self.refill_per_sec))
303                }
304            };
305            match wait {
306                None => return,
307                Some(wait) => tokio::time::sleep(wait).await,
308            }
309        }
310    }
311}