ports.rsannotatedports.rssource121 lines · 4.1 KB · raw
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 {}