use bytes::Bytes;
use ruststream::{AckError, Headers, IncomingMessage, Positioned};
use sea_streamer_types::{Buffer as _, Message as _, SharedMessage};
use crate::wire;
pub const SEQUENCE_HEADER: &str = "stream-sequence";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FilePosition {
Beginning,
End,
Sequence(u64),
Timestamp(u64),
}
impl FilePosition {
#[must_use]
pub const fn beginning() -> Self {
Self::Beginning
}
#[must_use]
pub const fn end() -> Self {
Self::End
}
#[must_use]
pub const fn sequence(sequence: u64) -> Self {
Self::Sequence(sequence)
}
#[must_use]
pub const fn timestamp(millis: u64) -> Self {
Self::Timestamp(millis)
}
}
pub struct SeaMessage {
payload: Bytes,
headers: Headers,
stream: String,
sequence: u64,
}
impl std::fmt::Debug for SeaMessage {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SeaMessage")
.field("stream", &self.stream)
.field("sequence", &self.sequence)
.field("payload_len", &self.payload.len())
.finish_non_exhaustive()
}
}
impl SeaMessage {
pub(crate) fn new(message: &SharedMessage) -> Self {
let (mut headers, payload) = wire::decode(message.message().as_bytes());
let sequence = message.sequence();
headers.insert(SEQUENCE_HEADER, sequence.to_string());
Self {
payload,
headers,
stream: message.stream_key().name().to_owned(),
sequence,
}
}
#[must_use]
pub fn stream(&self) -> &str {
&self.stream
}
}
impl Positioned for SeaMessage {
type Position = FilePosition;
fn position(&self) -> FilePosition {
FilePosition::Sequence(self.sequence)
}
}
impl IncomingMessage for SeaMessage {
fn payload(&self) -> &[u8] {
&self.payload
}
fn headers(&self) -> &Headers {
&self.headers
}
async fn ack(self) -> Result<(), AckError> {
Err(AckError::Unsupported)
}
async fn nack(self, _requeue: bool) -> Result<(), AckError> {
Err(AckError::Unsupported)
}
}