pub mod list_event_types;
mod list_subjects;
mod observe_events;
mod ping;
mod read_event_type;
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_event_type::ReadEventTypeRequest;
pub use read_events::ReadEventsRequest;
pub use register_event_schema::RegisterEventSchemaRequest;
pub use run_eventql_query::RunEventqlQueryRequest;
use serde_json::value::RawValue;
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)]
struct StreamLineItem {
#[serde(rename = "type")]
ty: String,
payload: Box<RawValue>,
}
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(|line| Ok(serde_json::from_str::<StreamLineItem>(line?.as_str())?))
.filter_map(|o| async {
match o {
Ok(StreamLineItem { payload, ty }) => match ty.as_str() {
ty if ty == Self::ITEM_TYPE_NAME => {
Some(serde_json::from_str(payload.get()).map_err(ClientError::from))
}
"error" => Some(Err(ClientError::DBError(payload.get().to_string()))),
"heartbeat" => None,
other => Some(Err(ClientError::InvalidResponseType(format!(
"Expected type {}, but got {}",
Self::ITEM_TYPE_NAME,
other
)))),
},
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)
}
}