use std::{error, fmt};
use tephra::Position;
use tephra::event::{EncodeError, Event, EventRef, EventType};
use tephra::query::{AppendCondition, Query};
use tephra::writer::{AppendError, ConflictSite};
use tephra_proto::convert as wire;
use tephra_proto::tephra as pb;
#[derive(Debug)]
pub enum ConvertError {
Wire(wire::ConvertError),
Encode(EncodeError),
}
impl fmt::Display for ConvertError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ConvertError::Wire(err) => write!(f, "{err}"),
ConvertError::Encode(err) => write!(f, "invalid event: {err}"),
}
}
}
impl error::Error for ConvertError {}
impl From<wire::ConvertError> for ConvertError {
fn from(err: wire::ConvertError) -> Self {
ConvertError::Wire(err)
}
}
pub fn event_from_proto(ev: pb::EventView<'_>) -> Result<Event, ConvertError> {
let ty = EventType::new(wire::as_str(ev.r#type())?)
.map_err(|err| ConvertError::Wire(wire::ConvertError::Name(err)))?;
let tags = wire::tags_from_pb(ev.tags().iter())?;
Event::new(&ty, &tags, ev.payload()).map_err(ConvertError::Encode)
}
pub fn events_from_proto(append: pb::AppendRequestView<'_>) -> Result<Vec<Event>, ConvertError> {
let mut events = Vec::with_capacity(append.events().len());
for ev in append.events().iter() {
events.push(event_from_proto(ev)?);
}
Ok(events)
}
pub fn query_from_proto(query: pb::QueryView<'_>) -> Result<Query, ConvertError> {
Ok(wire::query_from_pb(query)?)
}
pub fn condition_from_proto(
condition: pb::AppendConditionView<'_>,
) -> Result<AppendCondition, ConvertError> {
Ok(wire::condition_from_pb(condition)?)
}
pub fn sequenced_to_proto(position: Position, event: EventRef<'_>) -> pb::SequencedEvent {
let mut out = pb::SequencedEvent::new();
out.set_position(position.get());
let mut ev = out.event_mut();
ev.set_type(event.event_type());
for tag in event.tags() {
ev.tags_mut().push(tag);
}
ev.set_payload(event.data());
out
}
pub fn append_error_to_proto(err: &AppendError) -> pb::ErrorResponse {
let mut resp = pb::ErrorResponse::new();
resp.set_message(err.to_string());
match err {
AppendError::Conflict { at } => {
resp.set_code(pb::ErrorCode::Conflict);
match at {
ConflictSite::Durable(position) => {
resp.set_conflict_position(position.get());
resp.set_retryable(false);
}
ConflictSite::SameBatch => {
resp.set_retryable(true);
}
}
}
AppendError::AfterBeyondTip { .. } => resp.set_code(pb::ErrorCode::AfterBeyondTip),
AppendError::Empty => resp.set_code(pb::ErrorCode::Empty),
AppendError::TooLarge { .. } => resp.set_code(pb::ErrorCode::TooLarge),
AppendError::Log(_) | AppendError::Corrupt(_) => resp.set_code(pb::ErrorCode::Internal),
AppendError::Shutdown => resp.set_code(pb::ErrorCode::Shutdown),
}
resp
}
pub fn bad_request(message: impl fmt::Display) -> pb::ErrorResponse {
let mut resp = pb::ErrorResponse::new();
resp.set_code(pb::ErrorCode::BadRequest);
resp.set_message(message.to_string());
resp
}
pub fn too_large(message: impl fmt::Display) -> pb::ErrorResponse {
let mut resp = pb::ErrorResponse::new();
resp.set_code(pb::ErrorCode::TooLarge);
resp.set_message(message.to_string());
resp
}
pub fn internal_error(message: impl fmt::Display) -> pb::ErrorResponse {
let mut resp = pb::ErrorResponse::new();
resp.set_code(pb::ErrorCode::Internal);
resp.set_message(message.to_string());
resp
}