use anyhow::Result;
use async_trait::async_trait;
use futures::stream::BoxStream;
use serde::{Deserialize, Serialize};
use std::{collections::BTreeMap, fmt::Display, num::NonZeroU64};
use crate::{
event::{Event, EventKey, Metadata},
language::Query,
scalars::StreamId,
tags::TagSet,
AppId, LamportTimestamp, Offset, OffsetMap, Payload, Timestamp,
};
#[derive(Debug, Serialize, Deserialize, Clone, Eq, PartialEq)]
#[serde(rename_all = "kebab-case")]
pub enum Order {
Asc,
Desc,
StreamAsc,
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct QueryRequest {
pub lower_bound: Option<OffsetMap>,
pub upper_bound: Option<OffsetMap>,
pub query: Query,
pub order: Order,
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct SubscribeRequest {
pub lower_bound: Option<OffsetMap>,
pub query: Query,
}
#[derive(Debug, Serialize, Deserialize, Clone, Ord, PartialOrd, Eq, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct EventResponse<T> {
pub lamport: LamportTimestamp,
pub stream: StreamId,
pub offset: Offset,
pub timestamp: Timestamp,
pub tags: TagSet,
pub app_id: AppId,
pub payload: T,
}
impl<T> From<Event<T>> for EventResponse<T> {
fn from(env: Event<T>) -> Self {
let EventKey {
lamport,
stream,
offset,
} = env.key;
let Metadata {
timestamp,
tags,
app_id,
} = env.meta;
let payload = env.payload;
EventResponse {
lamport,
stream,
offset,
timestamp,
tags,
app_id,
payload,
}
}
}
impl EventResponse<Payload> {
pub fn extract<'a, T>(&'a self) -> Result<EventResponse<T>, serde_cbor::Error>
where
T: Deserialize<'a> + Clone,
{
Ok(EventResponse {
stream: self.stream,
lamport: self.lamport,
offset: self.offset,
timestamp: self.timestamp,
tags: self.tags.clone(),
app_id: self.app_id.clone(),
payload: self.payload.extract::<T>()?,
})
}
}
impl<T> std::fmt::Display for EventResponse<T> {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
use chrono::TimeZone;
let time = chrono::Local.timestamp_millis(self.timestamp.as_i64() / 1000);
write!(
f,
"Event at {} ({}, stream ID {})",
time.to_rfc3339_opts(chrono::SecondsFormat::Millis, false),
self.lamport,
self.stream,
)
}
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct OffsetMapResponse {
pub offsets: OffsetMap,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct PublishEvent {
pub tags: TagSet,
pub payload: Payload,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct PublishRequest {
pub data: Vec<PublishEvent>,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct PublishResponseKey {
pub lamport: LamportTimestamp,
pub stream: StreamId,
pub offset: Offset,
pub timestamp: Timestamp,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct PublishResponse {
pub data: Vec<PublishResponseKey>,
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialOrd, Eq, PartialEq)]
#[serde(rename_all = "camelCase")]
pub enum StartFrom {
LowerBound(OffsetMap),
}
#[derive(Debug, Clone, Serialize, Deserialize, Ord, PartialOrd, Eq, PartialEq, Hash)]
pub struct SessionId(Box<str>);
impl Display for SessionId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&*self.0)
}
}
impl From<&str> for SessionId {
fn from(s: &str) -> Self {
Self(s.into())
}
}
impl From<String> for SessionId {
fn from(s: String) -> Self {
Self(s.into())
}
}
impl SessionId {
pub fn as_str(&self) -> &str {
&*self.0
}
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct SubscribeMonotonicRequest {
pub session: SessionId,
pub query: Query,
#[serde(flatten)]
pub from: StartFrom,
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)]
#[serde(rename_all = "camelCase", tag = "type")]
pub enum SubscribeMonotonicResponse {
#[serde(rename_all = "camelCase")]
Event {
#[serde(flatten)]
event: EventResponse<Payload>,
caught_up: bool,
},
#[serde(rename_all = "camelCase")]
Offsets(OffsetMapResponse),
#[serde(rename_all = "camelCase")]
TimeTravel { new_start: EventKey },
#[serde(rename_all = "camelCase")]
Diagnostic(Diagnostic),
#[serde(other)]
FutureCompat,
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)]
#[serde(rename_all = "camelCase", tag = "type")]
pub enum QueryResponse {
#[serde(rename_all = "camelCase")]
Event(EventResponse<Payload>),
#[serde(rename_all = "camelCase")]
Offsets(OffsetMapResponse),
#[serde(rename_all = "camelCase")]
Diagnostic(Diagnostic),
#[serde(other)]
FutureCompat,
}
#[derive(Debug, PartialEq, Serialize, Deserialize, Clone)]
#[serde(rename_all = "camelCase", tag = "type")]
pub enum SubscribeResponse {
#[serde(rename_all = "camelCase")]
Event(EventResponse<Payload>),
#[serde(rename_all = "camelCase")]
Offsets(OffsetMapResponse),
#[serde(rename_all = "camelCase")]
Diagnostic(Diagnostic),
#[serde(other)]
FutureCompat,
}
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct Diagnostic {
pub severity: Severity,
pub message: String,
}
impl Diagnostic {
pub fn warn(message: String) -> Self {
Self {
severity: Severity::Warning,
message,
}
}
pub fn error(message: String) -> Self {
Self {
severity: Severity::Error,
message,
}
}
}
#[derive(Debug, Serialize, Deserialize, Clone, Copy, PartialEq)]
#[serde(rename_all = "camelCase")]
pub enum Severity {
Warning,
Error,
#[serde(other)]
FutureCompat,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct OffsetsResponse {
pub present: OffsetMap,
pub to_replicate: BTreeMap<StreamId, NonZeroU64>,
}
#[async_trait]
pub trait EventService: Clone + Send {
async fn offsets(&self) -> Result<OffsetsResponse>;
async fn publish(&self, request: PublishRequest) -> Result<PublishResponse>;
async fn query(&self, request: QueryRequest) -> Result<BoxStream<'static, QueryResponse>>;
async fn subscribe(&self, request: SubscribeRequest) -> Result<BoxStream<'static, SubscribeResponse>>;
async fn subscribe_monotonic(
&self,
request: SubscribeMonotonicRequest,
) -> Result<BoxStream<'static, SubscribeMonotonicResponse>>;
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn future_compat() {
assert_eq!(
serde_json::from_str::<QueryResponse>(r#"{"type":"fromTheFuture","x":42}"#).unwrap(),
QueryResponse::FutureCompat
);
assert_eq!(
serde_json::from_str::<SubscribeResponse>(r#"{"type":"fromTheFuture","x":42}"#).unwrap(),
SubscribeResponse::FutureCompat
);
assert_eq!(
serde_json::from_str::<SubscribeMonotonicResponse>(r#"{"type":"fromTheFuture","x":42}"#).unwrap(),
SubscribeMonotonicResponse::FutureCompat
);
}
}