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.
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.
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).
An https URL, resolved into what a connection needs.
55impl Endpoint {
The API itself, [ENDPOINT].
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.
Sends over one kept HTTP/2 connection, reopened when it dies, idles
past IDLE_LIMIT, or the client distrusts it.
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.
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).
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.
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}