use std::{
fmt::{Debug, Display},
time::Duration,
};
use bytes::Bytes;
use watermelon_proto::{ServerMessage, Subject};
use crate::client::{Client, ClientClosedError};
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum JetstreamMessageAckError {
#[error("client closed")]
ClientClosed(#[source] ClientClosedError),
#[error("message has no reply subject")]
NoReplySubject,
#[error("message already acknowledged")]
AlreadyAcknowledged,
}
pub struct JetstreamMessage {
pub message: ServerMessage,
client: Option<Client>,
acknowledged: bool,
}
impl JetstreamMessage {
#[must_use]
pub(crate) fn new(message: ServerMessage, client: Client) -> Self {
Self {
message,
client: Some(client),
acknowledged: false,
}
}
pub async fn ack(&mut self) -> Result<(), JetstreamMessageAckError> {
let reply_subject = self.take_reply_subject()?;
self.ensure_not_acked()?;
let client = self.take_client()?;
client
.publish(reply_subject)
.payload(Bytes::from_static(b"+ACK"))
.await
.map_err(JetstreamMessageAckError::ClientClosed)?;
self.mark_acknowledged();
Ok(())
}
pub async fn nak(&mut self) -> Result<(), JetstreamMessageAckError> {
let reply_subject = self.take_reply_subject()?;
self.ensure_not_acked()?;
let client = self.take_client()?;
client
.publish(reply_subject)
.payload(Bytes::from_static(b"-NAK"))
.await
.map_err(JetstreamMessageAckError::ClientClosed)?;
self.mark_acknowledged();
Ok(())
}
pub async fn nak_with_delay(
&mut self,
delay: Duration,
) -> Result<(), JetstreamMessageAckError> {
let reply_subject = self.take_reply_subject()?;
self.ensure_not_acked()?;
let payload = format!("-NAK {{\"delay\": {}}}", delay.as_nanos()).into_bytes();
let client = self.take_client()?;
client
.publish(reply_subject)
.payload(Bytes::from(payload))
.await
.map_err(JetstreamMessageAckError::ClientClosed)?;
self.mark_acknowledged();
Ok(())
}
pub async fn progress(&self) -> Result<(), JetstreamMessageAckError> {
let reply_subject = self
.message
.base
.reply_subject
.as_ref()
.ok_or(JetstreamMessageAckError::NoReplySubject)?
.clone();
let client = self
.client
.as_ref()
.ok_or(JetstreamMessageAckError::NoReplySubject)?;
client
.publish(reply_subject)
.payload(Bytes::from_static(b"+WPI"))
.await
.map_err(JetstreamMessageAckError::ClientClosed)?;
Ok(())
}
pub async fn term(&mut self) -> Result<(), JetstreamMessageAckError> {
let reply_subject = self.take_reply_subject()?;
self.ensure_not_acked()?;
let client = self.take_client()?;
client
.publish(reply_subject)
.payload(Bytes::from_static(b"+TERM"))
.await
.map_err(JetstreamMessageAckError::ClientClosed)?;
self.mark_acknowledged();
Ok(())
}
pub async fn term_with_reason(
&mut self,
reason: impl Display,
) -> Result<(), JetstreamMessageAckError> {
let reply_subject = self.take_reply_subject()?;
self.ensure_not_acked()?;
let payload = format!("+TERM {reason}").into_bytes();
let client = self.take_client()?;
client
.publish(reply_subject)
.payload(Bytes::from(payload))
.await
.map_err(JetstreamMessageAckError::ClientClosed)?;
self.mark_acknowledged();
Ok(())
}
#[must_use]
pub fn is_acked(&self) -> bool {
self.acknowledged
}
fn take_reply_subject(&mut self) -> Result<Subject, JetstreamMessageAckError> {
self.message
.base
.reply_subject
.take()
.ok_or(JetstreamMessageAckError::NoReplySubject)
}
fn take_client(&mut self) -> Result<Client, JetstreamMessageAckError> {
self.client
.take()
.ok_or(JetstreamMessageAckError::NoReplySubject)
}
fn ensure_not_acked(&self) -> Result<(), JetstreamMessageAckError> {
if self.acknowledged {
return Err(JetstreamMessageAckError::AlreadyAcknowledged);
}
Ok(())
}
fn mark_acknowledged(&mut self) {
self.acknowledged = true;
}
}
impl Debug for JetstreamMessage {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("JetstreamMessage")
.field("message", &self.message)
.field("acknowledged", &self.acknowledged)
.finish_non_exhaustive()
}
}