1use std::fmt; 2use std::future::Future; 3use std::time::{Duration, SystemTime}; 4 5/// The clock's type. On native targets this is `std::time::Instant` itself 6/// (web-time re-exports it). On wasm32-unknown-unknown, std's `Instant::now` 7/// panics and the type has no other constructor, so a runtime there (a 8/// Cloudflare Worker, a browser) could not implement [`Runtime::now`] at all 9/// without web-time's `performance.now()`-backed one. `SystemTime` stays std's: 10/// it can be built from `UNIX_EPOCH` plus a duration on any target. 11pub use web_time::Instant; 12 13use bytes::Bytes; 14use jev_protocol::Usage; 15 16/// One POST to the System One endpoint the transport was built for. 17#[derive(Clone, Debug)] 18pub struct HttpRequest { 19 pub headers: http::HeaderMap, 20 pub body: Bytes, 21} 22 23#[derive(Clone, Debug)] 24pub struct HttpResponse { 25 pub status: u16, 26 pub headers: http::HeaderMap, 27 pub body: Bytes, 28} 29 30/// Sends one request and returns the whole response. It does not retry, 31/// time out or interpret statuses: that is the client's. 32/// 33/// The futures are not required to be `Send`: an adapter may be tied to 34/// one thread, as the Postgres backend's is. 35pub trait Transport { 36 /// Waits until `request` may be sent, as a rate limit shared with 37 /// other senders decides. Called before every attempt, outside its 38 /// timeout and the retry budget, so a queue for the account's limit 39 /// is never mistaken for a slow server. The default sends at once. 40 fn admit(&self, _request: &HttpRequest) -> impl Future<Output = ()> { 41 async {} 42 } 43 44 fn send(&self, request: HttpRequest) -> impl Future<Output = Result<HttpResponse, TransportError>>; 45 46 /// The client gave up on an exchange (it timed out). The connection 47 /// it ran on is suspect: the next send should not reuse it. 48 fn distrust_connection(&self); 49} 50 51/// Time and randomness, so the policy can run on virtual time in tests. 52pub trait Runtime { 53 fn now(&self) -> Instant; 54 /// Wall-clock time, for an HTTP-date `Retry-After`. 55 fn wall_clock(&self) -> SystemTime; 56 fn sleep(&self, duration: Duration) -> impl Future<Output = ()>; 57 /// `future`, or `None` if `duration` passes first (dropping it). 58 fn timeout<F: Future>(&self, duration: Duration, future: F) -> impl Future<Output = Option<F::Output>>; 59 /// Uniform in [0, 1). 60 fn jitter(&self) -> f64; 61} 62 63/// What the client did, as it happens: for counters and logs. It must not 64/// fail or block, because the client does not wait for it. 65pub trait Observer { 66 fn event(&self, event: Event<'_>); 67} 68 69/// No observer. 70impl Observer for () { 71 fn event(&self, _: Event<'_>) {} 72} 73 74#[derive(Clone, Copy)] 75#[non_exhaustive] 76pub enum Event<'a> { 77 /// An attempt is going out. `retry` is 0 for the first attempt and 78 /// for its "never sent" resends. 79 Sending { retry: u32 }, 80 /// Attempt `attempt` failed with `cause` and is retried after `delay`. 81 Retrying { attempt: u32, delay: Duration, cause: &'a dyn fmt::Display, request_id: Option<&'a str> }, 82 /// The request never left; it is resent at once on a fresh connection. 83 Redialing { cause: &'a TransportError }, 84 /// A verified answer, after `attempts` attempts. 85 Answered { request_id: Option<&'a str>, attempts: u32, usage: Usage }, 86} 87 88#[derive(Clone, Debug, PartialEq, Eq)] 89pub struct TransportError { 90 pub kind: TransportErrorKind, 91 pub message: String, 92} 93 94/// What happened, as far as retrying is concerned. 95#[derive(Clone, Copy, Debug, PartialEq, Eq)] 96#[non_exhaustive] 97pub enum TransportErrorKind { 98 /// The request never reached the server: a connection that had 99 /// already closed, or a stream refused (REFUSED_STREAM, or above a 100 /// GOAWAY's last stream id). Resent at once on a new connection; it 101 /// is not a retry, because nothing was billed. 102 NotSent, 103 /// Resolving or connecting failed. Retried. 104 Connect, 105 /// The exchange broke after the request may have been sent. Retried, 106 /// and possibly billed twice (no idempotency key exists). 107 Interrupted, 108 /// TLS refused: bad certificate, no HTTP/2. Retrying cannot help. 109 Tls, 110 /// The response exceeded the size cap. Already produced and billed, 111 /// so it is never retried. 112 TooLarge, 113 /// Our own settings are unusable. Never retried. 114 Config, 115} 116 117impl TransportError { 118 pub fn new(kind: TransportErrorKind, message: impl Into<String>) -> Self { 119 TransportError { kind, message: message.into() } 120 } 121} 122 123impl fmt::Display for TransportError { 124 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { 125 f.write_str(&self.message) 126 } 127} 128 129impl std::error::Error for TransportError {}