pub mod list_event_types;
mod list_subjects;
mod observe_events;
mod ping;
mod read_events;
mod register_event_schema;
mod run_eventql_query;
mod verify_api_token;
mod write_events;
pub use list_event_types::ListEventTypesRequest;
pub use list_subjects::ListSubjectsRequest;
pub use observe_events::ObserveEventsRequest;
pub use ping::PingRequest;
pub use read_events::ReadEventsRequest;
pub use register_event_schema::RegisterEventSchemaRequest;
pub use run_eventql_query::RunEventqlQueryRequest;
use serde_json::Value;
pub use verify_api_token::VerifyApiTokenRequest;
pub use write_events::WriteEventsRequest;
use crate::error::ClientError;
use futures::{
Stream,
stream::{StreamExt, TryStreamExt},
};
use futures_util::io;
use reqwest::Method;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use tokio::io::{AsyncBufReadExt, BufReader};
use tokio_stream::wrappers::LinesStream;
use tokio_util::io::StreamReader;
pub trait ClientRequest {
const URL_PATH: &'static str;
const METHOD: Method;
fn url_path(&self) -> &'static str {
Self::URL_PATH
}
fn method(&self) -> Method {
Self::METHOD
}
fn body(&self) -> Option<Result<impl Serialize, ClientError>> {
None::<Result<(), _>>
}
}
pub trait OneShotRequest: ClientRequest {
type Response: DeserializeOwned;
fn validate_response(&self, _response: &Self::Response) -> Result<(), ClientError> {
Ok(())
}
}
#[derive(Deserialize, Debug)]
#[serde(tag = "type", content = "payload", rename_all = "camelCase")]
enum StreamLineItem<T> {
Error { error: String },
Heartbeat(Value),
#[serde(untagged)]
Ok {
#[serde(rename = "type")]
ty: String,
payload: T,
},
}
pub trait StreamingRequest: ClientRequest {
type ItemType: DeserializeOwned;
const ITEM_TYPE_NAME: &'static str;
fn build_stream(
response: reqwest::Response,
) -> impl Stream<Item = Result<Self::ItemType, ClientError>> {
Box::pin(
Self::lines_stream(response)
.map(|maybe_line| {
let line = maybe_line?;
Ok(serde_json::from_str::<StreamLineItem<Self::ItemType>>(
line.as_str(),
)?)
})
.filter_map(|o| async {
match o {
Ok(StreamLineItem::Error { error }) => {
Some(Err(ClientError::DBError(error)))
}
Ok(StreamLineItem::Heartbeat(_value)) => None,
Ok(StreamLineItem::Ok { payload, ty }) if ty == Self::ITEM_TYPE_NAME => {
Some(Ok(payload))
}
Ok(StreamLineItem::Ok { ty, .. }) => {
Some(Err(ClientError::InvalidResponseType(ty)))
}
Err(e) => Some(Err(e)),
}
}),
)
}
fn lines_stream(
response: reqwest::Response,
) -> impl Stream<Item = Result<String, ClientError>> {
let bytes = response
.bytes_stream()
.map_err(|err| io::Error::other(format!("Failed to read response stream: {err}")));
let stream_reader = StreamReader::new(bytes);
LinesStream::new(BufReader::new(stream_reader).lines()).map_err(ClientError::from)
}
}