use monoloop_contracts::{
ChannelId, CompletionCallback, EventDeliveryOutcome, SessionId, TransactionEnd,
TransactionEndKind, TransactionId, TransactionUsage,
};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
#[derive(Debug)]
pub struct EventSequencer {
next: AtomicU64,
}
impl EventSequencer {
pub fn new() -> Self {
Self {
next: AtomicU64::new(1),
}
}
pub fn peek_next(&self) -> u64 {
self.next.load(Ordering::SeqCst)
}
pub fn allocate(&self) -> u64 {
self.next.fetch_add(1, Ordering::SeqCst)
}
pub fn last_allocated(&self) -> u64 {
self.next.load(Ordering::SeqCst).saturating_sub(1)
}
}
impl Default for EventSequencer {
fn default() -> Self {
Self::new()
}
}
pub struct FinalizationPayload {
pub callback: Box<dyn CompletionCallback>,
pub channel_id: ChannelId,
pub session_id: Option<SessionId>,
pub transaction_id: TransactionId,
}
pub struct FinalizationGuard {
claimed: AtomicBool,
payload: Mutex<Option<FinalizationPayload>>,
sequencer: Arc<EventSequencer>,
callback_scheduled: AtomicBool,
}
impl FinalizationGuard {
pub fn new(
transaction_id: TransactionId,
channel_id: ChannelId,
session_id: Option<SessionId>,
callback: Box<dyn CompletionCallback>,
sequencer: Arc<EventSequencer>,
) -> Arc<Self> {
Arc::new(Self {
claimed: AtomicBool::new(false),
payload: Mutex::new(Some(FinalizationPayload {
callback,
channel_id,
session_id,
transaction_id,
})),
sequencer,
callback_scheduled: AtomicBool::new(false),
})
}
pub fn sequencer(&self) -> &Arc<EventSequencer> {
&self.sequencer
}
pub fn set_session_id(&self, session_id: SessionId) {
if let Ok(mut g) = self.payload.lock() {
if let Some(p) = g.as_mut() {
p.session_id = Some(session_id);
}
}
}
pub fn try_claim(&self) -> Option<FinalizationPayload> {
if self
.claimed
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.is_err()
{
return None;
}
self.payload.lock().ok().and_then(|mut g| g.take())
}
pub fn is_claimed(&self) -> bool {
self.claimed.load(Ordering::SeqCst)
}
pub fn mark_callback_scheduled(&self) {
self.callback_scheduled.store(true, Ordering::SeqCst);
}
pub fn callback_was_scheduled(&self) -> bool {
self.callback_scheduled.load(Ordering::SeqCst)
}
}
pub fn build_transaction_end(
payload: &FinalizationPayload,
kind: TransactionEndKind,
prior: Option<TransactionEndKind>,
event_delivery: EventDeliveryOutcome,
emitted_events: u64,
) -> TransactionEnd {
TransactionEnd {
transaction_id: payload.transaction_id,
session_id: payload.session_id.clone(),
channel_id: payload.channel_id.clone(),
kind,
prior_terminal_cause: prior,
event_delivery,
emitted_events,
usage: TransactionUsage::default(),
diagnostics: vec![],
}
}
pub fn bound_diagnostics(
mut diagnostics: Vec<monoloop_contracts::TransactionDiagnostic>,
max_count: usize,
max_message_bytes: usize,
) -> Vec<monoloop_contracts::TransactionDiagnostic> {
if diagnostics.len() > max_count {
diagnostics.truncate(max_count.max(1));
}
for d in &mut diagnostics {
if let Some(ref mut msg) = d.diagnostic.message {
if msg.len() > max_message_bytes {
msg.truncate(max_message_bytes);
while !msg.is_char_boundary(msg.len()) {
msg.pop();
}
}
}
}
diagnostics
}