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;
Where the evaluation endpoint is (digest §1).
24pub const ENDPOINT: &str = "https://api.typesafe.ai/v1/systemone";
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.
34const IDLE_LIMIT: Duration = Duration::from_secs(300);
The response cap (pg_typesafe's 8 MB, as postjevsql's).
36const 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.
39const 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.
56impl Endpoint {
The API itself, [ENDPOINT].
Something other than the API: a relay, or a test server.
63 pub fn parse(url: &str) -> Result<Self, String> { 64 let uri: http::Uri = url.parse().map_err(|e| format!("endpoint {url}: {e}"))?; 65 if uri.scheme() != Some(&http::uri::Scheme::HTTPS) { 66 return Err(format!("the endpoint must be https: {url}")); 67 } 68 let host = uri.host().ok_or_else(|| format!("the endpoint has no host: {url}"))?; 69 let bare = host.trim_start_matches('[').trim_end_matches(']').to_owned(); 70 let server_name = ServerName::try_from(bare.clone()).map_err(|e| format!("endpoint host {host}: {e}"))?; 71 Ok(Endpoint { port: uri.port_u16().unwrap_or(443), host: bare, server_name, ca_file: None, uri }) 72 }
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.
111impl HttpsTransport { 112 pub fn new(endpoint: Endpoint, tls: Arc<ClientConfig>) -> Self { 113 HttpsTransport { 114 endpoint, 115 tls, 116 warm: Mutex::new(None), 117 requests: Bucket::new(VENDOR_REQUESTS_PER_MINUTE, VENDOR_REQUESTS_PER_MINUTE / 60.0), 118 tokens: Bucket::new(VENDOR_TOKENS_PER_SECOND, VENDOR_TOKENS_PER_SECOND), 119 } 120 } 121 122 fn take_warm(&self) -> Option<SendRequest<Full<Bytes>>> { 123 let mut warm = self.warm.lock().unwrap_or_else(PoisonError::into_inner); 124 match warm.take() { 125 Some(w) if !w.sender.is_closed() && w.used.elapsed() < IDLE_LIMIT => { 126 let sender = w.sender.clone(); 127 *warm = Some(w); 128 Some(sender) 129 } 130 _ => None, 131 } 132 } 133 134 fn keep(&self, sender: SendRequest<Full<Bytes>>) { 135 *self.warm.lock().unwrap_or_else(PoisonError::into_inner) = Some(Warm { sender, used: Instant::now() }); 136 } 137 138 fn forget(&self) { 139 *self.warm.lock().unwrap_or_else(PoisonError::into_inner) = None; 140 } 141 142 async fn connection(&self) -> Result<SendRequest<Full<Bytes>>, TransportError> { 143 if let Some(sender) = self.take_warm() { 144 return Ok(sender); 145 } 146 let sender = connect(&self.endpoint, self.tls.clone()).await?; 147 self.keep(sender.clone()); 148 Ok(sender) 149 } 150} 151 152impl 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).
162 async fn send(&self, request: HttpRequest) -> Result<HttpResponse, TransportError> { 163 let mut sender = self.connection().await?; 164 let mut builder = http::Request::post(self.endpoint.uri.clone()); 165 *builder.headers_mut().expect("a fresh builder has headers") = request.headers; 166 let request = builder 167 .body(Full::new(request.body)) 168 .map_err(|e| TransportError::new(TransportErrorKind::Config, format!("building the request: {e}")))?; 169 170 let response = match sender.try_send_request(request).await { 171 Ok(response) => response, 172 Err(mut e) => { 173 self.forget(); 174 // hyper hands the request back when it never left. 175 let kind = if e.take_message().is_some() || never_processed(e.error()) { 176 TransportErrorKind::NotSent 177 } else { 178 TransportErrorKind::Interrupted 179 }; 180 return Err(TransportError::new(kind, format!("sending to {}: {}", self.endpoint.uri, e.error()))); 181 } 182 }; 183 let (parts, body) = response.into_parts(); 184 let body = match Limited::new(body, MAX_BODY).collect().await { 185 Ok(collected) => collected.to_bytes(), 186 Err(e) if e.is::<LengthLimitError>() => { 187 return Err(TransportError::new( 188 TransportErrorKind::TooLarge, 189 format!("Jev's response exceeded {} MiB", MAX_BODY >> 20), 190 )); 191 } 192 Err(e) => { 193 self.forget(); 194 return Err(TransportError::new(TransportErrorKind::Interrupted, format!("reading Jev's response: {e}"))); 195 } 196 }; 197 self.keep(sender); 198 Ok(HttpResponse { status: parts.status.as_u16(), headers: parts.headers, body }) 199 } 200 201 fn distrust_connection(&self) { 202 self.forget(); 203 } 204}
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.
209fn never_processed(error: &hyper::Error) -> bool { 210 let mut source: Option<&(dyn std::error::Error + 'static)> = std::error::Error::source(error); 211 while let Some(e) = source { 212 if let Some(h2) = e.downcast_ref::<h2::Error>() { 213 return h2.reason() == Some(h2::Reason::REFUSED_STREAM) || (h2.is_go_away() && h2.is_remote()); 214 } 215 source = e.source(); 216 } 217 false 218}
220async fn connect(endpoint: &Endpoint, tls: Arc<ClientConfig>) -> Result<SendRequest<Full<Bytes>>, TransportError> { 221 use TransportErrorKind::{Connect, Tls}; 222 let addrs: Vec<SocketAddr> = tokio::net::lookup_host((endpoint.host.as_str(), endpoint.port)) 223 .await 224 .map_err(|e| TransportError::new(Connect, format!("resolving {}: {e}", endpoint.host)))? 225 .collect(); 226 let mut failures = Vec::new(); 227 let mut tcp = None; 228 for (i, addr) in addrs.iter().enumerate() { 229 let attempt = TcpStream::connect(*addr); 230 let result = if i + 1 < addrs.len() { 231 tokio::time::timeout(CONNECT_PER_ADDRESS, attempt).await.unwrap_or_else(|_| { 232 Err(std::io::Error::new(std::io::ErrorKind::TimedOut, "no answer within 2 s")) 233 }) 234 } else { 235 attempt.await 236 }; 237 match result { 238 Ok(stream) => { 239 tcp = Some(stream); 240 break; 241 } 242 Err(e) => failures.push(format!("{addr}: {e}")), 243 } 244 } 245 let tcp = tcp.ok_or_else(|| { 246 let why = if failures.is_empty() { "no addresses".to_owned() } else { failures.join("; ") }; 247 TransportError::new(Connect, format!("connecting to {}: {why}", endpoint.uri)) 248 })?; 249 250 let tls = TlsConnector::from(tls).connect(endpoint.server_name.clone(), tcp).await.map_err(|e| { 251 // rustls refusing (a certificate, a protocol alert) will refuse 252 // again; an I/O failure mid-handshake is the network. 253 let refused = e.get_ref().is_some_and(|inner| inner.is::<rustls::Error>()); 254 TransportError::new(if refused { Tls } else { Connect }, format!("TLS with {}: {e}", endpoint.uri)) 255 })?; 256 if tls.get_ref().1.alpn_protocol() != Some(b"h2") { 257 return Err(TransportError::new(Tls, format!("{} did not negotiate HTTP/2", endpoint.uri))); 258 } 259 // hyper's fixed receive windows are below an 8 MB answer; BDP-based 260 // adaptive windows grow to what the path carries. 261 let (sender, connection) = hyper::client::conn::http2::Builder::new(TokioExecutor::new()) 262 .adaptive_window(true) 263 .timer(TokioTimer::new()) 264 .keep_alive_interval(PING_INTERVAL) 265 .keep_alive_timeout(PING_TIMEOUT) 266 .keep_alive_while_idle(true) 267 .handshake(TokioIo::new(tls)) 268 .await 269 .map_err(|e| TransportError::new(Connect, format!("HTTP/2 with {}: {e}", endpoint.uri)))?; 270 tokio::spawn(async move { 271 let _ = connection.await; 272 }); 273 Ok(sender) 274}
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.
292 async fn take(&self, amount: f64) { 293 loop { 294 let wait = { 295 let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner); 296 let elapsed = state.1.elapsed().as_secs_f64(); 297 state.0 = (state.0 + elapsed * self.refill_per_sec).min(self.capacity); 298 state.1 = Instant::now(); 299 if state.0 >= amount { 300 state.0 -= amount; 301 None 302 } else { 303 Some(Duration::from_secs_f64((amount - state.0) / self.refill_per_sec)) 304 } 305 }; 306 match wait { 307 None => return, 308 Some(wait) => tokio::time::sleep(wait).await, 309 } 310 } 311 } 312}