use std::future::Future;
use std::time::Duration;
use efema_proto::wire::{Batch, Heads, Page, Problem, ProblemCode, Watch, Written};
use efema_proto::{Cursor, StreamName, limits};
pub trait Transport: Send + Sync + 'static {
fn describe(&self) -> String;
fn limits(&self) -> impl Future<Output = Result<Limits, TransportError>> + Send;
fn append(
&self,
stream: &StreamName,
batch: &Batch,
) -> impl Future<Output = Result<Written, TransportError>> + Send;
fn read(
&self,
stream: &StreamName,
after: Option<&Cursor>,
limit: usize,
) -> impl Future<Output = Result<Page, TransportError>> + Send;
fn wait(&self, watches: &[Watch], timeout: Duration) -> impl Future<Output = Result<Heads, TransportError>> + Send;
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct Limits {
pub max_body: usize,
pub page_entries: usize,
pub wait_max: Duration,
}
impl Limits {
pub const V1: Self = Self {
max_body: limits::MAX_BODY_BYTES,
page_entries: limits::PAGE_ENTRIES_DEFAULT,
wait_max: limits::WAIT_MAX,
};
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum TransportError {
#[error("{} refused: {} ({})", .at, .problem.message, .problem.code)]
Refused {
at: String,
problem: Problem,
},
#[error("cannot reach {at}")]
Unreachable {
at: String,
#[source]
source: Box<dyn std::error::Error + Send + Sync>,
},
#[error("{at} answered something that is not the efema protocol: {reason}")]
NotProtocol {
at: String,
reason: String,
},
}
impl TransportError {
pub fn code(&self) -> Option<&ProblemCode> {
match self {
Self::Refused { problem, .. } => Some(&problem.code),
_ => None,
}
}
}