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