use super::session_identity::ensure_session;
use crate::transaction::sticky_cancel::StickyCancel;
use monoloop_contracts::{
ChannelId, EventEnqueueError, ExternalSessionId, SessionId, TerminalEventDelivery,
TransactionEndEvent, TransactionEvent, TransactionEventPayload, TransactionEventSender,
TransactionId,
};
use std::sync::{Arc, Mutex};
use std::time::Instant as StdInstant;
use tokio::sync::{mpsc, oneshot};
use tokio::time::{sleep_until, Instant as TokioInstant};
pub enum EventPublisherCommand {
EstablishExternal(ExternalSessionId),
Publish(Box<TransactionEventPayload>),
}
#[derive(Clone, Debug)]
pub struct OrdinaryCmdAdmit {
tx: Arc<Mutex<Option<mpsc::Sender<EventPublisherCommand>>>>,
}
impl OrdinaryCmdAdmit {
pub fn channel(capacity: usize) -> (Self, mpsc::Receiver<EventPublisherCommand>) {
let (tx, rx) = mpsc::channel(capacity.max(1));
(
Self {
tx: Arc::new(Mutex::new(Some(tx))),
},
rx,
)
}
pub fn close(&self) {
let _ = self.tx.lock().unwrap_or_else(|e| e.into_inner()).take();
}
pub fn is_open(&self) -> bool {
self.tx.lock().unwrap_or_else(|e| e.into_inner()).is_some()
}
pub async fn send(&self, cmd: EventPublisherCommand) -> Result<(), ()> {
let tx = self.tx.lock().unwrap_or_else(|e| e.into_inner()).clone();
let Some(tx) = tx else {
return Err(());
};
tx.send(cmd).await.map_err(|_| ())
}
pub fn try_send(
&self,
cmd: EventPublisherCommand,
) -> Result<(), mpsc::error::TrySendError<EventPublisherCommand>> {
let tx = self.tx.lock().unwrap_or_else(|e| e.into_inner()).clone();
let Some(tx) = tx else {
return Err(mpsc::error::TrySendError::Closed(cmd));
};
tx.try_send(cmd)
}
#[cfg(test)]
pub async fn send_after_pre_fence_hold(
&self,
cmd: EventPublisherCommand,
holding: oneshot::Sender<()>,
) -> Result<(), ()> {
let tx = self.tx.lock().unwrap_or_else(|e| e.into_inner()).clone();
let Some(tx) = tx else {
return Err(());
};
let _ = holding.send(());
tx.send(cmd).await.map_err(|_| ())
}
}
pub struct SealCommand {
pub terminal: TransactionEndEvent,
pub reply: oneshot::Sender<TerminalPublicationResult>,
pub deadline: StdInstant,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TerminalPublicationResult {
pub delivery: TerminalEventDelivery,
pub last_sequence: u64,
}
#[allow(clippy::too_many_arguments)] pub async fn run_event_publisher(
transaction_id: TransactionId,
channel_id: ChannelId,
session_id: Option<SessionId>,
event_tx: TransactionEventSender,
mut cmd_rx: mpsc::Receiver<EventPublisherCommand>,
admit: OrdinaryCmdAdmit,
mut seal_rx: mpsc::Receiver<SealCommand>,
_cancel: Arc<StickyCancel>,
deadline: StdInstant,
) -> TerminalPublicationResult {
let mut next_seq: u64 = 1;
let mut last_committed: u64 = 0;
let mut session = session_id;
let mut session_established = false;
let mut sticky_fail: Option<TerminalEventDelivery> = None;
let mut cmd_closed = false;
let mut seal_closed = false;
loop {
if let Some(fail) = sticky_fail {
return match seal_rx.recv().await {
Some(cmd) => {
admit.close();
let _ = drain_to_disconnect(&mut cmd_rx, cmd.deadline).await;
reply_sticky(cmd, fail, last_committed)
}
None => {
admit.close();
TerminalPublicationResult {
delivery: fail,
last_sequence: last_committed,
}
}
};
}
tokio::select! {
biased;
seal = seal_rx.recv(), if !seal_closed => {
match seal {
Some(cmd) => {
admit.close();
return fence_drain_and_seal(
cmd,
transaction_id,
&channel_id,
&event_tx,
&mut cmd_rx,
&mut session,
&mut next_seq,
&mut last_committed,
&mut session_established,
None,
)
.await;
}
None => {
seal_closed = true;
admit.close();
if cmd_closed {
return TerminalPublicationResult {
delivery: TerminalEventDelivery::QueueClosed,
last_sequence: last_committed,
};
}
}
}
}
cmd = cmd_rx.recv(), if !cmd_closed => {
match cmd {
Some(EventPublisherCommand::EstablishExternal(external)) => {
match establish_external(
&event_tx,
transaction_id,
&channel_id,
external,
&mut session,
&mut next_seq,
&mut last_committed,
&mut session_established,
deadline,
&mut seal_rx,
)
.await
{
Ok(()) => {}
Err(WaitEnd::Sealed { cmd, ordinary }) => {
admit.close();
let fail =
apply_ordinary_outcome(ordinary, &mut next_seq, &mut last_committed);
return fence_drain_and_seal(
cmd,
transaction_id,
&channel_id,
&event_tx,
&mut cmd_rx,
&mut session,
&mut next_seq,
&mut last_committed,
&mut session_established,
fail,
)
.await;
}
Err(WaitEnd::Failed(fail)) => sticky_fail = Some(fail),
}
}
Some(EventPublisherCommand::Publish(payload)) => {
let payload = *payload;
let sid = ensure_session(&mut session, None, transaction_id);
let seq = next_seq;
let event = TransactionEvent {
transaction_id,
channel_id: channel_id.clone(),
session_id: sid,
sequence: seq,
payload,
};
match wait_send_or_seal(&event_tx, event, deadline, &mut seal_rx).await {
Ok(()) => {
last_committed = seq;
next_seq = seq.saturating_add(1);
}
Err(WaitEnd::Sealed { cmd, ordinary }) => {
admit.close();
let fail =
apply_ordinary_outcome(ordinary, &mut next_seq, &mut last_committed);
return fence_drain_and_seal(
cmd,
transaction_id,
&channel_id,
&event_tx,
&mut cmd_rx,
&mut session,
&mut next_seq,
&mut last_committed,
&mut session_established,
fail,
)
.await;
}
Err(WaitEnd::Failed(fail)) => sticky_fail = Some(fail),
}
}
None => {
cmd_closed = true;
if seal_closed {
return TerminalPublicationResult {
delivery: TerminalEventDelivery::QueueClosed,
last_sequence: last_committed,
};
}
}
}
}
}
}
}
enum OrdinarySealOutcome {
Committed { seq: u64 },
Failed(TerminalEventDelivery),
}
enum WaitEnd {
Sealed {
cmd: SealCommand,
ordinary: OrdinarySealOutcome,
},
Failed(TerminalEventDelivery),
}
fn apply_ordinary_outcome(
ordinary: OrdinarySealOutcome,
next_seq: &mut u64,
last_committed: &mut u64,
) -> Option<TerminalEventDelivery> {
match ordinary {
OrdinarySealOutcome::Committed { seq } => {
*last_committed = seq;
*next_seq = seq.saturating_add(1);
None
}
OrdinarySealOutcome::Failed(f) => Some(f),
}
}
fn reply_sticky(
cmd: SealCommand,
fail: TerminalEventDelivery,
last_committed: u64,
) -> TerminalPublicationResult {
let result = TerminalPublicationResult {
delivery: fail,
last_sequence: last_committed,
};
let _ = cmd.reply.send(result.clone());
result
}
async fn drain_to_disconnect(
cmd_rx: &mut mpsc::Receiver<EventPublisherCommand>,
deadline: StdInstant,
) -> Result<(), TerminalEventDelivery> {
let tokio_deadline = tokio_deadline_from(deadline);
loop {
tokio::select! {
biased;
_ = sleep_until(tokio_deadline) => {
return Err(TerminalEventDelivery::DeadlineExceeded);
}
cmd = cmd_rx.recv() => {
if cmd.is_none() {
return Ok(());
}
}
}
}
}
#[allow(clippy::too_many_arguments)]
async fn fence_drain_and_seal(
cmd: SealCommand,
transaction_id: TransactionId,
channel_id: &ChannelId,
event_tx: &TransactionEventSender,
cmd_rx: &mut mpsc::Receiver<EventPublisherCommand>,
session: &mut Option<SessionId>,
next_seq: &mut u64,
last_committed: &mut u64,
session_established: &mut bool,
mut sticky_fail: Option<TerminalEventDelivery>,
) -> TerminalPublicationResult {
let seal_deadline = cmd.deadline;
let tokio_deadline = tokio_deadline_from(seal_deadline);
while sticky_fail.is_none() {
let next = tokio::select! {
biased;
_ = sleep_until(tokio_deadline) => {
sticky_fail = Some(TerminalEventDelivery::DeadlineExceeded);
break;
}
cmd = cmd_rx.recv() => cmd,
};
match next {
None => break,
Some(EventPublisherCommand::EstablishExternal(external)) => {
if *session_established || *next_seq != 1 {
continue;
}
let sid = SessionId::from_external(&external);
let seq = *next_seq;
let event = TransactionEvent {
transaction_id,
channel_id: channel_id.clone(),
session_id: sid.clone(),
sequence: seq,
payload: TransactionEventPayload::SessionEstablished {
external_session_id: external,
},
};
match enqueue_under_deadline(event_tx, event, seal_deadline).await {
Ok(()) => {
let _ = ensure_session(session, Some(sid), transaction_id);
*last_committed = seq;
*next_seq = seq.saturating_add(1);
*session_established = true;
}
Err(fail) => sticky_fail = Some(fail),
}
}
Some(EventPublisherCommand::Publish(payload)) => {
let sid = ensure_session(session, None, transaction_id);
let seq = *next_seq;
let event = TransactionEvent {
transaction_id,
channel_id: channel_id.clone(),
session_id: sid,
sequence: seq,
payload: *payload,
};
match enqueue_under_deadline(event_tx, event, seal_deadline).await {
Ok(()) => {
*last_committed = seq;
*next_seq = seq.saturating_add(1);
}
Err(fail) => sticky_fail = Some(fail),
}
}
}
}
if sticky_fail.is_some() {
let _ = drain_to_disconnect(cmd_rx, seal_deadline).await;
}
finish_seal(
cmd,
transaction_id,
channel_id,
event_tx,
session,
next_seq,
last_committed,
sticky_fail,
)
.await
}
#[allow(clippy::too_many_arguments)] async fn finish_seal(
cmd: SealCommand,
transaction_id: TransactionId,
channel_id: &ChannelId,
event_tx: &TransactionEventSender,
session: &mut Option<SessionId>,
next_seq: &mut u64,
last_committed: &mut u64,
sticky_fail: Option<TerminalEventDelivery>,
) -> TerminalPublicationResult {
if let Some(fail) = sticky_fail {
return reply_sticky(cmd, fail, *last_committed);
}
let SealCommand {
mut terminal,
reply,
deadline,
} = cmd;
let sid = ensure_session(session, terminal.session_id.clone(), transaction_id);
let seq = *next_seq;
terminal.emitted_events = seq;
terminal.session_id = Some(sid.clone());
let event = TransactionEvent {
transaction_id,
channel_id: channel_id.clone(),
session_id: sid,
sequence: seq,
payload: TransactionEventPayload::EndedEvent(terminal),
};
let delivery = match enqueue_under_deadline(event_tx, event, deadline).await {
Ok(()) => {
*last_committed = seq;
TerminalEventDelivery::Published
}
Err(fail) => fail,
};
let result = TerminalPublicationResult {
delivery,
last_sequence: *last_committed,
};
let _ = reply.send(result.clone());
result
}
fn tokio_deadline_from(deadline: StdInstant) -> TokioInstant {
let now = StdInstant::now();
if deadline > now {
TokioInstant::now() + deadline.saturating_duration_since(now)
} else {
TokioInstant::now()
}
}
fn map_send_err(err: EventEnqueueError) -> TerminalEventDelivery {
match err {
EventEnqueueError::Closed => TerminalEventDelivery::QueueClosed,
EventEnqueueError::EventTooLarge
| EventEnqueueError::ByteCapacityExceeded
| EventEnqueueError::ItemCapacityExceeded => TerminalEventDelivery::LimitExceeded,
}
}
async fn wait_send_or_seal(
event_tx: &TransactionEventSender,
event: TransactionEvent,
deadline: StdInstant,
seal_rx: &mut mpsc::Receiver<SealCommand>,
) -> Result<(), WaitEnd> {
let seq = event.sequence;
let tokio_deadline = tokio_deadline_from(deadline);
let send_fut = event_tx.send(event.clone());
tokio::pin!(send_fut);
let mut seal_alive = true;
loop {
tokio::select! {
biased;
seal = seal_rx.recv(), if seal_alive => match seal {
Some(cmd) => {
let seal_deadline = tokio_deadline_from(cmd.deadline);
tokio::select! {
biased;
_ = sleep_until(seal_deadline) => {
return Err(WaitEnd::Sealed {
cmd,
ordinary: OrdinarySealOutcome::Failed(
TerminalEventDelivery::DeadlineExceeded,
),
});
}
res = &mut send_fut => {
let ordinary = match res {
Ok(()) => OrdinarySealOutcome::Committed { seq },
Err(e) => OrdinarySealOutcome::Failed(map_send_err(e)),
};
return Err(WaitEnd::Sealed { cmd, ordinary });
}
}
}
None => {
seal_alive = false;
}
},
_ = sleep_until(tokio_deadline) => {
return Err(WaitEnd::Failed(TerminalEventDelivery::DeadlineExceeded));
}
res = &mut send_fut => {
return match res {
Ok(()) => Ok(()),
Err(e) => Err(WaitEnd::Failed(map_send_err(e))),
};
}
}
}
}
#[allow(clippy::too_many_arguments)]
async fn establish_external(
event_tx: &TransactionEventSender,
transaction_id: TransactionId,
channel_id: &ChannelId,
external: ExternalSessionId,
session: &mut Option<SessionId>,
next_seq: &mut u64,
last_committed: &mut u64,
session_established: &mut bool,
deadline: StdInstant,
seal_rx: &mut mpsc::Receiver<SealCommand>,
) -> Result<(), WaitEnd> {
if *session_established || *next_seq != 1 {
return Ok(());
}
let sid = SessionId::from_external(&external);
let seq = *next_seq;
let event = TransactionEvent {
transaction_id,
channel_id: channel_id.clone(),
session_id: sid.clone(),
sequence: seq,
payload: TransactionEventPayload::SessionEstablished {
external_session_id: external,
},
};
match wait_send_or_seal(event_tx, event, deadline, seal_rx).await {
Ok(()) => {
let _ = ensure_session(session, Some(sid), transaction_id);
*last_committed = seq;
*next_seq = seq.saturating_add(1);
*session_established = true;
Ok(())
}
Err(WaitEnd::Sealed { cmd, ordinary }) => {
if matches!(ordinary, OrdinarySealOutcome::Committed { .. }) {
let _ = ensure_session(session, Some(sid), transaction_id);
*session_established = true;
}
Err(WaitEnd::Sealed { cmd, ordinary })
}
Err(WaitEnd::Failed(fail)) => Err(WaitEnd::Failed(fail)),
}
}
async fn enqueue_under_deadline(
event_tx: &TransactionEventSender,
event: TransactionEvent,
deadline: StdInstant,
) -> Result<(), TerminalEventDelivery> {
let tokio_deadline = tokio_deadline_from(deadline);
tokio::select! {
biased;
_ = sleep_until(tokio_deadline) => Err(TerminalEventDelivery::DeadlineExceeded),
res = event_tx.send(event) => match res {
Ok(()) => Ok(()),
Err(e) => Err(map_send_err(e)),
},
}
}