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;

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}
15#[derive(Clone, Debug)]
16pub struct HttpResponse {
17    pub status: u16,
18    pub headers: http::HeaderMap,
19    pub body: Bytes,
20}

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.

27pub 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.

32    fn admit(&self, _request: &HttpRequest) -> impl Future<Output = ()> {
33        async {}
34    }
36    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.

40    fn distrust_connection(&self);
41}

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

44pub trait Runtime {
45    fn now(&self) -> Instant;

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

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

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

Uniform in [0, 1).

52    fn jitter(&self) -> f64;
53}

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.

57pub trait Observer {
58    fn event(&self, event: Event<'_>);
59}

No observer.

62impl Observer for () {
63    fn event(&self, _: Event<'_>) {}
64}
66#[derive(Clone, Copy)]
67#[non_exhaustive]
68pub enum Event<'a> {

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

71    Sending { retry: u32 },

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

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

75    Redialing { cause: &'a TransportError },

A verified answer, after attempts attempts.

77    Answered { request_id: Option<&'a str>, attempts: u32, usage: Usage },
78}
80#[derive(Clone, Debug, PartialEq, Eq)]
81pub struct TransportError {
82    pub kind: TransportErrorKind,
83    pub message: String,
84}

What happened, as far as retrying is concerned.

87#[derive(Clone, Copy, Debug, PartialEq, Eq)]
88#[non_exhaustive]
89pub 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.

94    NotSent,

Resolving or connecting failed. Retried.

96    Connect,

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

99    Interrupted,

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

101    Tls,

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

104    TooLarge,

Our own settings are unusable. Never retried.

106    Config,
107}
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 {}