ports.rsannotatedports.rssource129 lines · 4.6 KB · raw
1use std::fmt;
2use std::future::Future;
3use std::time::{Duration, SystemTime};

The clock's type. On native targets this is std::time::Instant itself (web-time re-exports it). On wasm32-unknown-unknown, std's Instant::now panics and the type has no other constructor, so a runtime there (a Cloudflare Worker, a browser) could not implement [Runtime::now] at all without web-time's performance.now()-backed one. SystemTime stays std's: it can be built from UNIX_EPOCH plus a duration on any target.

11pub use web_time::Instant;
13use bytes::Bytes;
14use jev_protocol::Usage;

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}
23#[derive(Clone, Debug)]
24pub struct HttpResponse {
25    pub status: u16,
26    pub headers: http::HeaderMap,
27    pub body: Bytes,
28}

Sends one request and returns the whole response. It does not retry, time out or interpret statuses: that is the client's.

The futures are not required to be Send: an adapter may be tied to one thread, as the Postgres backend's is.

35pub trait Transport {

Waits until request may be sent, as a rate limit shared with other senders decides. Called before every attempt, outside its timeout and the retry budget, so a queue for the account's limit is never mistaken for a slow server. The default sends at once.

40    fn admit(&self, _request: &HttpRequest) -> impl Future<Output = ()> {
41        async {}
42    }
44    fn send(&self, request: HttpRequest) -> impl Future<Output = Result<HttpResponse, TransportError>>;

The client gave up on an exchange (it timed out). The connection it ran on is suspect: the next send should not reuse it.

48    fn distrust_connection(&self);
49}

Time and randomness, so the policy can run on virtual time in tests.

52pub trait Runtime {
53    fn now(&self) -> Instant;

Wall-clock time, for an HTTP-date Retry-After.

55    fn wall_clock(&self) -> SystemTime;
56    fn sleep(&self, duration: Duration) -> impl Future<Output = ()>;

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>>;

Uniform in [0, 1).

60    fn jitter(&self) -> f64;
61}

What the client did, as it happens: for counters and logs. It must not fail or block, because the client does not wait for it.

65pub trait Observer {
66    fn event(&self, event: Event<'_>);
67}

No observer.

70impl Observer for () {
71    fn event(&self, _: Event<'_>) {}
72}
74#[derive(Clone, Copy)]
75#[non_exhaustive]
76pub enum Event<'a> {

An attempt is going out. retry is 0 for the first attempt and for its "never sent" resends.

79    Sending { retry: u32 },

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> },

The request never left; it is resent at once on a fresh connection.

83    Redialing { cause: &'a TransportError },

A verified answer, after attempts attempts.

85    Answered { request_id: Option<&'a str>, attempts: u32, usage: Usage },
86}
88#[derive(Clone, Debug, PartialEq, Eq)]
89pub struct TransportError {
90    pub kind: TransportErrorKind,
91    pub message: String,
92}

What happened, as far as retrying is concerned.

95#[derive(Clone, Copy, Debug, PartialEq, Eq)]
96#[non_exhaustive]
97pub enum TransportErrorKind {

The request never reached the server: a connection that had already closed, or a stream refused (REFUSED_STREAM, or above a GOAWAY's last stream id). Resent at once on a new connection; it is not a retry, because nothing was billed.

102    NotSent,

Resolving or connecting failed. Retried.

104    Connect,

The exchange broke after the request may have been sent. Retried, and possibly billed twice (no idempotency key exists).

107    Interrupted,

TLS refused: bad certificate, no HTTP/2. Retrying cannot help.

109    Tls,

The response exceeded the size cap. Already produced and billed, so it is never retried.

112    TooLarge,

Our own settings are unusable. Never retried.

114    Config,
115}
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 {}