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