monoloop_loop/transaction/
events.rs1use super::executor_spawn::try_spawn;
4use super::finalization::EventSequencer;
5use monoloop_contracts::{
6 ChannelId, EventDeliveryError, SessionId, TransactionEvent, TransactionEventPayload,
7 TransactionEventSink, TransactionId,
8};
9use std::panic::{catch_unwind, AssertUnwindSafe};
10use std::sync::atomic::{AtomicUsize, Ordering};
11use std::sync::Arc;
12use std::time::Duration;
13use tokio::runtime::Handle;
14use tokio::sync::{mpsc, Mutex};
15
16pub struct QueuedEvent {
18 pub event: TransactionEvent,
20 pub ack: Option<tokio::sync::oneshot::Sender<Result<(), EventDeliveryError>>>,
22 pub approx_bytes: usize,
24}
25
26impl QueuedEvent {
27 pub fn new(
29 event: TransactionEvent,
30 ack: Option<tokio::sync::oneshot::Sender<Result<(), EventDeliveryError>>>,
31 ) -> Self {
32 let approx_bytes = estimate_event_bytes(&event);
33 Self {
34 event,
35 ack,
36 approx_bytes,
37 }
38 }
39}
40
41fn estimate_event_bytes(event: &TransactionEvent) -> usize {
42 serde_json::to_vec(event)
44 .map(|b| b.len().max(64))
45 .unwrap_or(256)
46}
47
48#[derive(Clone)]
50pub struct BoundedEventSender {
51 tx: mpsc::Sender<QueuedEvent>,
52 queued_bytes: Arc<AtomicUsize>,
53 max_bytes: usize,
54}
55
56impl BoundedEventSender {
57 pub fn new(tx: mpsc::Sender<QueuedEvent>, max_bytes: usize) -> Self {
59 Self {
60 tx,
61 queued_bytes: Arc::new(AtomicUsize::new(0)),
62 max_bytes: max_bytes.max(1),
63 }
64 }
65
66 pub async fn send(&self, item: QueuedEvent) -> Result<(), EventQueueFull> {
72 let bytes = item.approx_bytes;
73 loop {
74 let cur = self.queued_bytes.load(Ordering::SeqCst);
75 if cur.saturating_add(bytes) > self.max_bytes {
76 return Err(EventQueueFull::Bytes);
77 }
78 if self
79 .queued_bytes
80 .compare_exchange(cur, cur + bytes, Ordering::SeqCst, Ordering::SeqCst)
81 .is_ok()
82 {
83 break;
84 }
85 }
86 let mut reservation = ByteReservation {
88 counter: &self.queued_bytes,
89 bytes,
90 released: false,
91 };
92 match self.tx.send(item).await {
93 Ok(()) => {
94 reservation.released = true;
95 Ok(())
96 }
97 Err(_) => {
98 reservation.release();
99 Err(EventQueueFull::Closed)
100 }
101 }
102 }
103
104 pub fn byte_counter(&self) -> Arc<AtomicUsize> {
106 Arc::clone(&self.queued_bytes)
107 }
108}
109
110#[derive(Clone)]
115pub struct OrderedEventPublisher {
116 order: Arc<Mutex<()>>,
117 event_tx: BoundedEventSender,
118 sequencer: Arc<EventSequencer>,
119}
120
121impl OrderedEventPublisher {
122 pub fn new(event_tx: BoundedEventSender, sequencer: Arc<EventSequencer>) -> Self {
124 Self {
125 order: Arc::new(Mutex::new(())),
126 event_tx,
127 sequencer,
128 }
129 }
130
131 pub fn sequencer(&self) -> &Arc<EventSequencer> {
133 &self.sequencer
134 }
135
136 pub async fn publish(
138 &self,
139 transaction_id: TransactionId,
140 channel_id: ChannelId,
141 session_id: SessionId,
142 payload: TransactionEventPayload,
143 ) -> Result<u64, EventQueueFull> {
144 self.publish_inner(transaction_id, channel_id, session_id, payload, None)
145 .await
146 }
147
148 pub async fn publish_terminal(
150 &self,
151 transaction_id: TransactionId,
152 channel_id: ChannelId,
153 session_id: SessionId,
154 payload: TransactionEventPayload,
155 ack: tokio::sync::oneshot::Sender<Result<(), EventDeliveryError>>,
156 ) -> Result<u64, EventQueueFull> {
157 self.publish_inner(transaction_id, channel_id, session_id, payload, Some(ack))
158 .await
159 }
160
161 async fn publish_inner(
162 &self,
163 transaction_id: TransactionId,
164 channel_id: ChannelId,
165 session_id: SessionId,
166 payload: TransactionEventPayload,
167 ack: Option<tokio::sync::oneshot::Sender<Result<(), EventDeliveryError>>>,
168 ) -> Result<u64, EventQueueFull> {
169 let _guard = self.order.lock().await;
171 let seq = self.sequencer.peek_next();
172 let event = TransactionEvent {
173 transaction_id,
174 channel_id,
175 session_id,
176 sequence: seq,
177 payload,
178 };
179 self.event_tx.send(QueuedEvent::new(event, ack)).await?;
180 let got = self.sequencer.allocate();
181 debug_assert_eq!(got, seq);
182 Ok(seq)
183 }
184}
185
186#[derive(Clone, Copy, Debug, PartialEq, Eq)]
188pub enum EventQueueFull {
189 Bytes,
191 Closed,
193}
194
195struct ByteReservation<'a> {
197 counter: &'a AtomicUsize,
198 bytes: usize,
199 released: bool,
200}
201
202impl ByteReservation<'_> {
203 fn release(&mut self) {
204 if !self.released {
205 self.counter.fetch_sub(self.bytes, Ordering::SeqCst);
206 self.released = true;
207 }
208 }
209}
210
211impl Drop for ByteReservation<'_> {
212 fn drop(&mut self) {
213 self.release();
214 }
215}
216
217pub fn spawn_delivery_task(
219 executor: &Handle,
220 mut rx: mpsc::Receiver<QueuedEvent>,
221 sink: Arc<dyn TransactionEventSink>,
222 on_fail: mpsc::Sender<()>,
223 byte_counter: Arc<AtomicUsize>,
224 deliver_deadline: Duration,
225) -> Result<tokio::task::JoinHandle<()>, ()> {
226 let executor_child = executor.clone();
227 try_spawn(executor, async move {
228 while let Some(item) = rx.recv().await {
229 let bytes = item.approx_bytes;
230 let result =
232 deliver_isolated(&executor_child, &sink, item.event, deliver_deadline).await;
233 byte_counter.fetch_sub(
234 bytes.min(byte_counter.load(Ordering::SeqCst)),
235 Ordering::SeqCst,
236 );
237 let ok = result.is_ok();
238 if let Some(ack) = item.ack {
239 let _ = ack.send(if ok {
240 Ok(())
241 } else {
242 Err(EventDeliveryError::Failed)
243 });
244 }
245 if !ok {
246 let _ = on_fail.try_send(());
247 while let Some(rest) = rx.recv().await {
248 byte_counter.fetch_sub(
249 rest.approx_bytes.min(byte_counter.load(Ordering::SeqCst)),
250 Ordering::SeqCst,
251 );
252 if let Some(ack) = rest.ack {
253 let _ = ack.send(Err(EventDeliveryError::Failed));
254 }
255 }
256 break;
257 }
258 }
259 })
260}
261
262async fn deliver_isolated(
264 executor: &Handle,
265 sink: &Arc<dyn TransactionEventSink>,
266 event: TransactionEvent,
267 deadline: Duration,
268) -> Result<(), EventDeliveryError> {
269 let deliver_fut = catch_unwind(AssertUnwindSafe(|| sink.deliver(event)));
270 let fut = match deliver_fut {
271 Ok(f) => f,
272 Err(_) => return Err(EventDeliveryError::Failed),
273 };
274 let handle = match try_spawn(executor, fut) {
276 Ok(h) => h,
277 Err(()) => return Err(EventDeliveryError::Failed),
278 };
279 let abort = handle.abort_handle();
280 match tokio::time::timeout(deadline, handle).await {
281 Ok(Ok(r)) => r,
282 Ok(Err(_)) => Err(EventDeliveryError::Failed),
283 Err(_) => {
284 abort.abort();
285 Err(EventDeliveryError::Failed)
286 }
287 }
288}