mod client_request;
mod precondition;
pub mod request_options;
use crate::{
error::ClientError,
event::{Event, EventCandidate, ManagementEvent},
};
use client_request::{
ClientRequest, ListEventTypesRequest, ListSubjectsRequest, ObserveEventsRequest,
OneShotRequest, PingRequest, ReadEventsRequest, RegisterEventSchemaRequest,
RunEventqlQueryRequest, StreamingRequest, VerifyApiTokenRequest, WriteEventsRequest,
list_event_types::EventType,
};
use futures::Stream;
pub use precondition::Precondition;
use reqwest;
use url::Url;
#[derive(Debug)]
pub struct Client {
base_url: Url,
api_token: String,
reqwest: reqwest::Client,
}
impl Client {
pub fn new(base_url: Url, api_token: impl Into<String>) -> Self {
Client {
base_url,
api_token: api_token.into(),
reqwest: reqwest::Client::new(),
}
}
#[must_use]
pub fn get_base_url(&self) -> &Url {
&self.base_url
}
#[must_use]
pub fn get_api_token(&self) -> &str {
&self.api_token
}
fn build_request<R: ClientRequest>(
&self,
endpoint: &R,
) -> Result<reqwest::RequestBuilder, ClientError> {
let url = self
.base_url
.join(endpoint.url_path())
.map_err(ClientError::URLParseError)?;
let request = match endpoint.method() {
reqwest::Method::GET => self.reqwest.get(url),
reqwest::Method::POST => self.reqwest.post(url),
_ => return Err(ClientError::InvalidRequestMethod),
}
.bearer_auth(&self.api_token);
let request = if let Some(body) = endpoint.body() {
request
.header("Content-Type", "application/json")
.json(&body?)
} else {
request
};
Ok(request)
}
async fn request_oneshot<R: OneShotRequest>(
&self,
endpoint: R,
) -> Result<R::Response, ClientError> {
let response = self.build_request(&endpoint)?.send().await?;
if response.status().is_success() {
let result = response.json().await?;
endpoint.validate_response(&result)?;
Ok(result)
} else {
Err(ClientError::DBApiError(
response.status(),
response.text().await.unwrap_or_default(),
))
}
}
async fn request_streaming<R: StreamingRequest>(
&self,
endpoint: R,
) -> Result<impl Stream<Item = Result<R::ItemType, ClientError>>, ClientError> {
let response = self.build_request(&endpoint)?.send().await?;
if response.status().is_success() {
Ok(R::build_stream(response))
} else {
Err(ClientError::DBApiError(
response.status(),
response.text().await.unwrap_or_default(),
))
}
}
pub async fn ping(&self) -> Result<(), ClientError> {
let _ = self.request_oneshot(PingRequest).await?;
Ok(())
}
pub async fn read_events<'a>(
&self,
subject: &'a str,
options: Option<request_options::ReadEventsOptions<'a>>,
) -> Result<impl Stream<Item = Result<Event, ClientError>>, ClientError> {
let response = self
.request_streaming(ReadEventsRequest { subject, options })
.await?;
Ok(response)
}
pub async fn observe_events<'a>(
&self,
subject: &'a str,
options: Option<request_options::ObserveEventsOptions<'a>>,
) -> Result<impl Stream<Item = Result<Event, ClientError>>, ClientError> {
let response = self
.request_streaming(ObserveEventsRequest { subject, options })
.await?;
Ok(response)
}
pub async fn verify_api_token(&self) -> Result<(), ClientError> {
let _ = self.request_oneshot(VerifyApiTokenRequest).await?;
Ok(())
}
pub async fn register_event_schema(
&self,
event_type: &str,
schema: &serde_json::Value,
) -> Result<ManagementEvent, ClientError> {
self.request_oneshot(RegisterEventSchemaRequest::try_new(event_type, schema)?)
.await
}
pub async fn list_subjects(
&self,
base_subject: Option<&str>,
) -> Result<impl Stream<Item = Result<String, ClientError>>, ClientError> {
let response = self
.request_streaming(ListSubjectsRequest {
base_subject: base_subject.unwrap_or("/"),
})
.await?;
Ok(response)
}
pub async fn list_event_types(
&self,
) -> Result<impl Stream<Item = Result<EventType, ClientError>>, ClientError> {
let response = self.request_streaming(ListEventTypesRequest).await?;
Ok(response)
}
pub async fn write_events(
&self,
events: Vec<EventCandidate>,
preconditions: Vec<Precondition>,
) -> Result<Vec<Event>, ClientError> {
self.request_oneshot(WriteEventsRequest {
events,
preconditions,
})
.await
}
pub async fn run_eventql_query(
&self,
query: &str,
) -> Result<impl Stream<Item = Result<serde_json::Value, ClientError>>, ClientError> {
let response = self
.request_streaming(RunEventqlQueryRequest { query })
.await?;
Ok(response)
}
}