1use std::collections::BTreeSet;
2use std::fmt;
3use std::future::pending;
4use std::sync::Arc;
5use std::sync::atomic::Ordering;
6use std::time::Duration;
7
8use crate::time::{Instant, SystemTime, UNIX_EPOCH, sleep};
9
10use anyhow::{Context, Result};
11use rivet_error::RivetError;
12use rivetkit_actor_persist::{generated::v4 as persist_v4, versioned as persist_versioned};
13use serde::{Deserialize, Serialize};
14#[cfg(not(target_arch = "wasm32"))]
15use tokio::runtime::{Builder, Handle};
16use tokio::sync::oneshot;
17use tokio_util::sync::CancellationToken;
18
19use crate::actor::config::ActorConfig;
20use crate::actor::context::ActorContext;
21use crate::actor::internal_storage;
22use crate::actor::persist::{
23 decode_latest_with_embedded_version, encode_latest_with_embedded_version,
24};
25use crate::actor::task_types::UserTaskKind;
26#[cfg(target_arch = "wasm32")]
27use crate::error::ActorRuntime;
28
29#[derive(Clone, Debug, Default)]
30pub struct QueueNextOpts {
31 pub names: Option<Vec<String>>,
32 pub timeout: Option<Duration>,
33 pub signal: Option<CancellationToken>,
34 pub completable: bool,
35}
36
37#[derive(Clone, Debug, Default)]
38pub struct QueueWaitOpts {
39 pub timeout: Option<Duration>,
40 pub signal: Option<CancellationToken>,
41 pub completable: bool,
42}
43
44#[derive(Clone, Debug, Default)]
45pub struct EnqueueAndWaitOpts {
46 pub timeout: Option<Duration>,
47 pub signal: Option<CancellationToken>,
48}
49
50#[derive(Clone, Debug)]
51pub struct QueueNextBatchOpts {
52 pub names: Option<Vec<String>>,
53 pub count: u32,
54 pub timeout: Option<Duration>,
55 pub signal: Option<CancellationToken>,
56 pub completable: bool,
57}
58
59impl Default for QueueNextBatchOpts {
60 fn default() -> Self {
61 Self {
62 names: None,
63 count: 1,
64 timeout: None,
65 signal: None,
66 completable: false,
67 }
68 }
69}
70
71#[derive(Clone, Debug, Default)]
72pub struct QueueTryNextOpts {
73 pub names: Option<Vec<String>>,
74 pub completable: bool,
75}
76
77#[derive(Clone, Debug)]
78pub struct QueueTryNextBatchOpts {
79 pub names: Option<Vec<String>>,
80 pub count: u32,
81 pub completable: bool,
82}
83
84impl Default for QueueTryNextBatchOpts {
85 fn default() -> Self {
86 Self {
87 names: None,
88 count: 1,
89 completable: false,
90 }
91 }
92}
93
94pub(super) type QueueWaitActivityCallback = Arc<dyn Fn() + Send + Sync>;
95pub(super) type QueueInspectorUpdateCallback = Arc<dyn Fn(u32) + Send + Sync>;
96
97#[derive(Clone, Debug)]
98pub struct QueueMessage {
99 pub id: u64,
100 pub name: String,
101 pub body: Vec<u8>,
102 pub created_at: i64,
103 completion: Option<CompletionHandle>,
104}
105
106#[derive(Clone, Debug)]
107pub struct CompletableQueueMessage {
108 pub id: u64,
109 pub name: String,
110 pub body: Vec<u8>,
111 pub created_at: i64,
112 completion: CompletionHandle,
113}
114
115#[derive(Clone)]
116struct CompletionHandle(Arc<CompletionHandleInner>);
117
118struct CompletionHandleInner {
119 ctx: ActorContext,
120 message_id: u64,
121 completed: std::sync::atomic::AtomicBool,
122}
123
124pub(crate) type QueueMetadata = persist_v4::QueueMetadata;
125pub(crate) type PersistedQueueMessage = persist_v4::QueueMessage;
126
127#[cfg(test)]
128pub(crate) fn encode_queue_metadata(metadata: &QueueMetadata) -> Result<Vec<u8>> {
129 encode_latest_with_embedded_version::<persist_versioned::QueueMetadata>(
130 metadata.clone(),
131 rivetkit_actor_persist::CURRENT_VERSION,
132 "queue metadata",
133 )
134}
135
136pub(crate) fn decode_queue_metadata(payload: &[u8]) -> Result<QueueMetadata> {
137 let metadata = decode_latest_with_embedded_version::<persist_versioned::QueueMetadata>(
138 payload,
139 "queue metadata",
140 )?;
141 Ok(metadata)
142}
143
144pub(crate) fn encode_queue_message(message: &PersistedQueueMessage) -> Result<Vec<u8>> {
145 encode_latest_with_embedded_version::<persist_versioned::QueueMessage>(
146 message.clone(),
147 rivetkit_actor_persist::CURRENT_VERSION,
148 "queue message",
149 )
150}
151
152pub(crate) fn decode_queue_message(payload: &[u8]) -> Result<PersistedQueueMessage> {
153 let message = decode_latest_with_embedded_version::<persist_versioned::QueueMessage>(
154 payload,
155 "queue message",
156 )?;
157 Ok(message)
158}
159
160#[derive(RivetError, Serialize, Deserialize)]
161#[error(
162 "queue",
163 "full",
164 "Queue is full",
165 "Queue is full. Limit is {limit} messages."
166)]
167struct QueueFull {
168 limit: u32,
169}
170
171#[derive(RivetError, Serialize, Deserialize)]
172#[error(
173 "queue",
174 "message_too_large",
175 "Queue message is too large",
176 "Queue message too large ({size} bytes). Limit is {limit} bytes."
177)]
178struct QueueMessageTooLarge {
179 size: usize,
180 limit: u32,
181}
182
183#[derive(RivetError)]
184#[error("queue", "already_completed", "Queue message was already completed")]
185struct QueueAlreadyCompleted;
186
187#[derive(RivetError, Serialize, Deserialize)]
188#[error(
189 "queue",
190 "complete_not_configured",
191 "Queue message does not support completion",
192 "Queue '{name}' does not support completion responses."
193)]
194struct QueueCompleteNotConfigured {
195 name: String,
196}
197
198#[derive(RivetError)]
199#[error("actor", "aborted", "Actor aborted")]
200struct QueueActorAborted;
201
202#[derive(RivetError, Serialize, Deserialize)]
203#[error(
204 "queue",
205 "timed_out",
206 "Queue wait timed out",
207 "Queue wait timed out after {timeout_ms} ms."
208)]
209struct QueueWaitTimedOut {
210 timeout_ms: u64,
211}
212
213#[derive(RivetError, Serialize, Deserialize)]
214#[error(
215 "queue",
216 "completion_waiter_conflict",
217 "Queue completion waiter conflict",
218 "Queue completion waiter is already registered for message {message_id}."
219)]
220struct QueueCompletionWaiterConflict {
221 message_id: u64,
222}
223
224#[derive(RivetError)]
225#[error(
226 "queue",
227 "completion_waiter_dropped",
228 "Queue completion waiter dropped before response"
229)]
230struct QueueCompletionWaiterDropped;
231
232impl ActorContext {
233 pub async fn send(&self, name: &str, body: &[u8]) -> Result<QueueMessage> {
234 self.enqueue_message(name, body, None).await
235 }
236
237 pub async fn enqueue_and_wait(
238 &self,
239 name: &str,
240 body: &[u8],
241 opts: EnqueueAndWaitOpts,
242 ) -> Result<Option<Vec<u8>>> {
243 let (sender, receiver) = oneshot::channel();
244 let message = self.enqueue_message(name, body, Some(sender)).await?;
245 let result = self
246 .wait_for_completion_response(message.id, receiver, opts.timeout, opts.signal.as_ref())
247 .await;
248 self.remove_completion_waiter(message.id).await;
249 result
250 }
251
252 async fn enqueue_message(
253 &self,
254 name: &str,
255 body: &[u8],
256 completion_waiter: Option<oneshot::Sender<Option<Vec<u8>>>>,
257 ) -> Result<QueueMessage> {
258 self.ensure_initialized().await?;
259
260 let created_at = current_timestamp_ms()?;
261 let persisted = PersistedQueueMessage {
262 name: name.to_owned(),
263 body: body.to_vec(),
264 created_at,
265 failure_count: None,
266 available_at: None,
267 in_flight: None,
268 in_flight_at: None,
269 };
270 let encoded_message = encode_queue_message(&persisted).context("encode queue message")?;
271
272 let config = self.config();
273 if encoded_message.len() > config.max_queue_message_size as usize {
274 return Err(QueueMessageTooLarge {
275 size: encoded_message.len(),
276 limit: config.max_queue_message_size,
277 }
278 .build());
279 }
280
281 let mut metadata = self.0.queue_metadata.lock().await;
282 if metadata.size >= config.max_queue_size {
283 return Err(QueueFull {
284 limit: config.max_queue_size,
285 }
286 .build());
287 }
288
289 let id = if metadata.next_id == 0 {
290 1
291 } else {
292 metadata.next_id
293 };
294 metadata.next_id = id.saturating_add(1);
295 metadata.size = metadata.size.saturating_add(1);
296 let registered_completion_waiter = if let Some(waiter) = completion_waiter {
297 if self
298 .0
299 .queue_completion_waiters
300 .insert_async(id, waiter)
301 .await
302 .is_err()
303 {
304 metadata.next_id = id;
305 metadata.size = metadata.size.saturating_sub(1);
306 return Err(QueueCompletionWaiterConflict { message_id: id }.build());
307 }
308 true
309 } else {
310 false
311 };
312
313 let persist_result =
314 internal_storage::persist_queue_message(self.sql(), id, metadata.next_id, &persisted)
315 .await;
316
317 if let Err(error) = persist_result {
318 metadata.next_id = id;
319 metadata.size = metadata.size.saturating_sub(1);
320 if registered_completion_waiter {
321 self.remove_completion_waiter(id).await;
322 }
323 return Err(error).context("persist queue message");
324 }
325
326 let queue_size = metadata.size;
327 drop(metadata);
328 self.0.metrics.add_queue_messages_sent(1);
329 self.0
330 .metrics
331 .set_queue_depth(self.0.queue_metadata.lock().await.size);
332 self.notify_inspector_update(queue_size);
333 self.0.queue_notify.notify_waiters();
334
335 Ok(QueueMessage {
336 id,
337 name: name.to_owned(),
338 body: body.to_vec(),
339 created_at,
340 completion: None,
341 })
342 }
343
344 pub async fn next(&self, opts: QueueNextOpts) -> Result<Option<QueueMessage>> {
345 let mut messages = self
346 .next_batch(QueueNextBatchOpts {
347 names: opts.names,
348 count: 1,
349 timeout: opts.timeout,
350 signal: opts.signal,
351 completable: opts.completable,
352 })
353 .await?;
354 Ok(messages.pop())
355 }
356
357 pub async fn next_batch(&self, opts: QueueNextBatchOpts) -> Result<Vec<QueueMessage>> {
358 self.ensure_initialized().await?;
359
360 let count = opts.count.max(1);
361 let deadline = opts.timeout.map(|timeout| Instant::now() + timeout);
362 let names = normalize_names(opts.names);
363
364 loop {
365 let messages = self
366 .try_receive_batch(names.as_ref(), count, opts.completable)
367 .await?;
368 if !messages.is_empty() {
369 return Ok(messages);
370 }
371
372 let remaining_timeout =
373 deadline.map(|deadline| deadline.saturating_duration_since(Instant::now()));
374 if matches!(remaining_timeout, Some(timeout) if timeout.is_zero()) {
375 return Ok(Vec::new());
376 }
377
378 let wait_guard = ActiveQueueWaitGuard::new(self);
379 let result = self
380 .wait_for_message(remaining_timeout, opts.signal.as_ref())
381 .await;
382 drop(wait_guard);
383
384 match result {
385 WaitOutcome::Notified => continue,
386 WaitOutcome::TimedOut => return Ok(Vec::new()),
387 WaitOutcome::Aborted => return Err(QueueActorAborted.build()),
388 }
389 }
390 }
391
392 pub async fn wait_for_names(
393 &self,
394 names: Vec<String>,
395 opts: QueueWaitOpts,
396 ) -> Result<QueueMessage> {
397 self.ensure_initialized().await?;
398
399 let deadline = opts.timeout.map(|timeout| Instant::now() + timeout);
400 let names = normalize_names(Some(names));
401
402 loop {
403 if let Some(message) = self
404 .try_receive_batch(names.as_ref(), 1, opts.completable)
405 .await?
406 .into_iter()
407 .next()
408 {
409 return Ok(message);
410 }
411
412 let remaining_timeout =
413 deadline.map(|deadline| deadline.saturating_duration_since(Instant::now()));
414 if let Some(timeout) = remaining_timeout
415 && timeout.is_zero()
416 {
417 return Err(QueueWaitTimedOut {
418 timeout_ms: opts.timeout.map(duration_ms).unwrap_or(0),
419 }
420 .build());
421 }
422
423 let wait_guard = ActiveQueueWaitGuard::new(self);
424 let result = self
425 .wait_for_message(remaining_timeout, opts.signal.as_ref())
426 .await;
427 drop(wait_guard);
428
429 match result {
430 WaitOutcome::Notified => continue,
431 WaitOutcome::TimedOut => {
432 return Err(QueueWaitTimedOut {
433 timeout_ms: opts.timeout.map(duration_ms).unwrap_or(0),
434 }
435 .build());
436 }
437 WaitOutcome::Aborted => return Err(QueueActorAborted.build()),
438 }
439 }
440 }
441
442 pub async fn wait_for_names_available(
443 &self,
444 names: Vec<String>,
445 opts: QueueWaitOpts,
446 ) -> Result<()> {
447 self.ensure_initialized().await?;
448
449 let deadline = opts.timeout.map(|timeout| Instant::now() + timeout);
450 let names = normalize_names(Some(names));
451
452 loop {
453 if internal_storage::has_queue_messages(self.sql(), names.as_ref())
454 .await
455 .context("check for matching sqlite queue messages")?
456 {
457 return Ok(());
458 }
459
460 let remaining_timeout =
461 deadline.map(|deadline| deadline.saturating_duration_since(Instant::now()));
462 if let Some(timeout) = remaining_timeout
463 && timeout.is_zero()
464 {
465 return Err(QueueWaitTimedOut {
466 timeout_ms: opts.timeout.map(duration_ms).unwrap_or(0),
467 }
468 .build());
469 }
470
471 let wait_guard = ActiveQueueWaitGuard::new(self);
472 let result = self
473 .wait_for_message(remaining_timeout, opts.signal.as_ref())
474 .await;
475 drop(wait_guard);
476
477 match result {
478 WaitOutcome::Notified => continue,
479 WaitOutcome::TimedOut => {
480 return Err(QueueWaitTimedOut {
481 timeout_ms: opts.timeout.map(duration_ms).unwrap_or(0),
482 }
483 .build());
484 }
485 WaitOutcome::Aborted => return Err(QueueActorAborted.build()),
486 }
487 }
488 }
489
490 pub fn try_next(&self, opts: QueueTryNextOpts) -> Result<Option<QueueMessage>> {
491 let mut messages = self.try_next_batch(QueueTryNextBatchOpts {
492 names: opts.names,
493 count: 1,
494 completable: opts.completable,
495 })?;
496 Ok(messages.pop())
497 }
498
499 pub fn try_next_batch(&self, opts: QueueTryNextBatchOpts) -> Result<Vec<QueueMessage>> {
500 self.block_on(async {
501 self.ensure_initialized().await?;
502 self.try_receive_batch(
503 normalize_names(opts.names).as_ref(),
504 opts.count.max(1),
505 opts.completable,
506 )
507 .await
508 })
509 }
510
511 pub async fn inspect_messages(&self) -> Result<Vec<QueueMessage>> {
512 self.ensure_initialized().await?;
513 self.list_messages().await
514 }
515
516 pub fn max_size(&self) -> u32 {
517 self.config().max_queue_size
518 }
519
520 pub async fn reset(&self) -> Result<()> {
522 self.ensure_initialized().await?;
523
524 let _receive_guard = self.0.queue_receive_lock.lock().await;
527
528 let mut metadata = self.0.queue_metadata.lock().await;
529
530 internal_storage::reset_queue(self.sql())
531 .await
532 .context("delete all sqlite queue messages")?;
533
534 metadata.size = 0;
535
536 self.0.queue_completion_waiters.clear_async().await;
537
538 drop(metadata);
539
540 self.0.metrics.set_queue_depth(0);
541 self.notify_inspector_update(0);
542 self.0.queue_notify.notify_waiters();
543
544 Ok(())
545 }
546
547 pub(crate) fn configure_queue(&self, config: ActorConfig) {
548 *self.0.queue_config.lock() = config;
549 }
550
551 pub(crate) fn set_wait_activity_callback(&self, callback: Option<Arc<dyn Fn() + Send + Sync>>) {
552 *self.0.queue_wait_activity_callback.lock() = callback;
553 }
554
555 pub(crate) fn set_inspector_update_callback(
556 &self,
557 callback: Option<Arc<dyn Fn(u32) + Send + Sync>>,
558 ) {
559 *self.0.queue_inspector_update_callback.lock() = callback;
560 }
561
562 async fn ensure_initialized(&self) -> Result<()> {
563 self.0
564 .queue_initialize
565 .get_or_try_init(|| async {
566 let metadata = internal_storage::load_queue_metadata(self.sql())
567 .await
568 .context("load queue metadata from sqlite")?;
569 let mut state = self.0.queue_metadata.lock().await;
570 *state = metadata;
571 self.0.metrics.set_queue_depth(state.size);
572 Ok(())
573 })
574 .await
575 .map(|_| ())
576 }
577
578 async fn try_receive_batch(
579 &self,
580 names: Option<&BTreeSet<String>>,
581 count: u32,
582 completable: bool,
583 ) -> Result<Vec<QueueMessage>> {
584 let _receive_guard = self.0.queue_receive_lock.lock().await;
585
586 let selected = self.list_messages_matching(names, count).await?;
587
588 if selected.is_empty() {
589 return Ok(Vec::new());
590 }
591
592 if completable {
593 let queue_size = self.0.queue_metadata.lock().await.size;
594 self.0
595 .metrics
596 .add_queue_messages_received(selected.len().try_into().unwrap_or(u64::MAX));
597 self.notify_inspector_update(queue_size);
598 return Ok(selected
599 .into_iter()
600 .map(|message| self.attach_completion(message))
601 .collect());
602 }
603
604 self.remove_messages(selected.iter().map(|message| message.id).collect())
605 .await?;
606 self.0
607 .metrics
608 .add_queue_messages_received(selected.len().try_into().unwrap_or(u64::MAX));
609
610 Ok(selected)
611 }
612
613 async fn list_messages(&self) -> Result<Vec<QueueMessage>> {
614 let messages: Vec<QueueMessage> = internal_storage::load_queue_messages(self.sql())
615 .await
616 .context("list sqlite queue messages")?
617 .into_iter()
618 .map(queue_message_from_row)
619 .collect();
620
621 let actual_size = messages.len().try_into().unwrap_or(u32::MAX);
622 let mut metadata = self.0.queue_metadata.lock().await;
623 if metadata.size != actual_size {
624 metadata.size = actual_size;
625 }
626 if metadata.next_id == 0 {
627 metadata.next_id = messages
628 .last()
629 .map(|message| message.id.saturating_add(1))
630 .unwrap_or(1);
631 }
632
633 Ok(messages)
634 }
635
636 async fn list_messages_matching(
637 &self,
638 names: Option<&BTreeSet<String>>,
639 limit: u32,
640 ) -> Result<Vec<QueueMessage>> {
641 internal_storage::load_queue_messages_matching(self.sql(), names, limit)
642 .await
643 .context("list matching sqlite queue messages")
644 .map(|rows| rows.into_iter().map(queue_message_from_row).collect())
645 }
646
647 fn attach_completion(&self, mut message: QueueMessage) -> QueueMessage {
648 message.completion = Some(CompletionHandle::new(self.clone(), message.id));
649 message
650 }
651
652 async fn remove_messages(&self, message_ids: Vec<u64>) -> Result<()> {
653 if message_ids.is_empty() {
654 return Ok(());
655 }
656
657 let deleted_count = message_ids.len();
658 internal_storage::delete_queue_messages(self.sql(), &message_ids)
659 .await
660 .context("delete sqlite queue messages")?;
661
662 let queue_size = {
663 let mut metadata = self.0.queue_metadata.lock().await;
664 metadata.size = metadata.size.saturating_sub(deleted_count as u32);
665 metadata.size
666 };
667 self.0
668 .metrics
669 .set_queue_depth(self.0.queue_metadata.lock().await.size);
670 self.notify_inspector_update(queue_size);
671 Ok(())
672 }
673
674 async fn complete_message_by_id(
675 &self,
676 message_id: u64,
677 response: Option<Vec<u8>>,
678 ) -> Result<()> {
679 self.remove_messages(vec![message_id]).await?;
680 if let Some(waiter) = self.remove_completion_waiter(message_id).await {
681 let _ = waiter.send(response);
682 }
683 Ok(())
684 }
685
686 async fn remove_completion_waiter(
687 &self,
688 message_id: u64,
689 ) -> Option<oneshot::Sender<Option<Vec<u8>>>> {
690 self.0
691 .queue_completion_waiters
692 .remove_async(&message_id)
693 .await
694 .map(|(_, waiter)| waiter)
695 }
696
697 async fn wait_for_message(
698 &self,
699 timeout: Option<Duration>,
700 signal: Option<&CancellationToken>,
701 ) -> WaitOutcome {
702 let actor_abort_signal = self.0.queue_abort_signal.lock().clone();
703 if signal.is_some_and(CancellationToken::is_cancelled) {
704 return WaitOutcome::Aborted;
705 }
706 if actor_abort_signal.is_cancelled() {
707 return WaitOutcome::Aborted;
708 }
709
710 let notified = self.0.queue_notify.notified();
711 let actor_aborted = async {
712 actor_abort_signal.cancelled().await;
713 };
714 let external_aborted = async {
715 if let Some(signal) = signal {
716 signal.cancelled().await;
717 } else {
718 pending::<()>().await;
719 }
720 };
721
722 match timeout {
723 Some(timeout) => {
724 tokio::select! {
725 _ = notified => WaitOutcome::Notified,
726 _ = actor_aborted => WaitOutcome::Aborted,
727 _ = external_aborted => WaitOutcome::Aborted,
728 _ = sleep(timeout) => WaitOutcome::TimedOut,
729 }
730 }
731 None => {
732 tokio::select! {
733 _ = notified => WaitOutcome::Notified,
734 _ = actor_aborted => WaitOutcome::Aborted,
735 _ = external_aborted => WaitOutcome::Aborted,
736 }
737 }
738 }
739 }
740
741 async fn wait_for_completion_response(
745 &self,
746 message_id: u64,
747 mut receiver: oneshot::Receiver<Option<Vec<u8>>>,
748 timeout: Option<Duration>,
749 signal: Option<&CancellationToken>,
750 ) -> Result<Option<Vec<u8>>> {
751 if signal.is_some_and(CancellationToken::is_cancelled) {
752 return Err(QueueActorAborted.build());
753 }
754
755 let external_aborted = async {
756 if let Some(signal) = signal {
757 signal.cancelled().await;
758 } else {
759 pending::<()>().await;
760 }
761 };
762
763 let wait_result = match timeout {
764 Some(timeout) => {
765 tokio::select! {
766 response = &mut receiver => CompletionWaitOutcome::Response(response),
767 _ = external_aborted => CompletionWaitOutcome::Aborted,
768 _ = sleep(timeout) => CompletionWaitOutcome::TimedOut,
769 }
770 }
771 None => {
772 tokio::select! {
773 response = &mut receiver => CompletionWaitOutcome::Response(response),
774 _ = external_aborted => CompletionWaitOutcome::Aborted,
775 }
776 }
777 };
778
779 match wait_result {
780 CompletionWaitOutcome::Response(Ok(response)) => Ok(response),
781 CompletionWaitOutcome::Response(Err(_)) => Err(QueueCompletionWaiterDropped.build())
782 .context(format!("wait for queue completion on message {message_id}")),
783 CompletionWaitOutcome::TimedOut => Err(QueueWaitTimedOut {
784 timeout_ms: timeout.map(duration_ms).unwrap_or(0),
785 }
786 .build()),
787 CompletionWaitOutcome::Aborted => Err(QueueActorAborted.build()),
788 }
789 }
790
791 fn block_on<T>(&self, future: impl std::future::Future<Output = Result<T>>) -> Result<T> {
792 #[cfg(not(target_arch = "wasm32"))]
793 {
794 if let Ok(handle) = Handle::try_current() {
795 tokio::task::block_in_place(|| handle.block_on(future))
796 } else {
797 Builder::new_current_thread()
798 .enable_all()
799 .build()
800 .context("build temporary runtime for queue operation")?
801 .block_on(future)
802 }
803 }
804
805 #[cfg(target_arch = "wasm32")]
806 {
807 drop(future);
808 Err(ActorRuntime::InvalidOperation {
809 operation: "queue.try_next_batch".to_owned(),
810 reason: "synchronous queue receive requires native runtime support".to_owned(),
811 }
812 .build())
813 }
814 }
815
816 fn config(&self) -> ActorConfig {
817 self.0.queue_config.lock().clone()
818 }
819
820 #[cfg(test)]
821 pub(crate) fn queue_config_for_tests(&self) -> ActorConfig {
822 self.config()
823 }
824
825 fn notify_wait_activity(&self) {
826 if let Some(callback) = self.0.queue_wait_activity_callback.lock().clone() {
827 callback();
828 }
829 }
830
831 fn notify_inspector_update(&self, queue_size: u32) {
832 if let Some(callback) = self.0.queue_inspector_update_callback.lock().clone() {
833 callback(queue_size);
834 }
835 }
836}
837
838impl QueueMessage {
839 pub async fn complete(self, response: Option<Vec<u8>>) -> Result<()> {
840 let completable = self.into_completable()?;
841 completable.complete(response).await
842 }
843
844 pub fn into_completable(self) -> Result<CompletableQueueMessage> {
845 let completion = self.completion.clone().ok_or_else(|| {
846 QueueCompleteNotConfigured {
847 name: self.name.clone(),
848 }
849 .build()
850 })?;
851
852 Ok(CompletableQueueMessage {
853 id: self.id,
854 name: self.name,
855 body: self.body,
856 created_at: self.created_at,
857 completion,
858 })
859 }
860
861 pub fn is_completable(&self) -> bool {
862 self.completion.is_some()
863 }
864}
865
866impl CompletableQueueMessage {
867 pub async fn complete(self, response: Option<Vec<u8>>) -> Result<()> {
868 self.completion.complete(response).await
869 }
870
871 pub fn into_message(self) -> QueueMessage {
872 QueueMessage {
873 id: self.id,
874 name: self.name,
875 body: self.body,
876 created_at: self.created_at,
877 completion: Some(self.completion),
878 }
879 }
880}
881
882impl CompletionHandle {
883 fn new(ctx: ActorContext, message_id: u64) -> Self {
884 Self(Arc::new(CompletionHandleInner {
885 ctx,
886 message_id,
887 completed: std::sync::atomic::AtomicBool::new(false),
888 }))
889 }
890
891 async fn complete(&self, response: Option<Vec<u8>>) -> Result<()> {
892 if self.0.completed.swap(true, Ordering::SeqCst) {
893 return Err(QueueAlreadyCompleted.build());
894 }
895
896 if let Err(error) = self
897 .0
898 .ctx
899 .complete_message_by_id(self.0.message_id, response)
900 .await
901 {
902 self.0.completed.store(false, Ordering::SeqCst);
903 return Err(error);
904 }
905
906 Ok(())
907 }
908}
909
910impl fmt::Debug for CompletionHandle {
911 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
912 f.debug_struct("CompletionHandle")
913 .field("message_id", &self.0.message_id)
914 .field("completed", &self.0.completed.load(Ordering::SeqCst))
915 .finish()
916 }
917}
918
919struct ActiveQueueWaitGuard<'a> {
920 ctx: &'a ActorContext,
921 started_at: Instant,
922}
923
924impl<'a> ActiveQueueWaitGuard<'a> {
925 fn new(ctx: &'a ActorContext) -> Self {
926 ctx.0.active_queue_wait_count.fetch_add(1, Ordering::SeqCst);
927 ctx.0.metrics.begin_user_task(UserTaskKind::QueueWait);
928 ctx.notify_wait_activity();
929 Self {
930 ctx,
931 started_at: Instant::now(),
932 }
933 }
934}
935
936impl Drop for ActiveQueueWaitGuard<'_> {
937 fn drop(&mut self) {
938 self.ctx
939 .0
940 .metrics
941 .end_user_task(UserTaskKind::QueueWait, self.started_at.elapsed());
942 let previous = self
943 .ctx
944 .0
945 .active_queue_wait_count
946 .fetch_sub(1, Ordering::SeqCst);
947 if previous == 0 {
948 self.ctx
949 .0
950 .active_queue_wait_count
951 .store(0, Ordering::SeqCst);
952 }
953 self.ctx.notify_wait_activity();
954 }
955}
956
957enum WaitOutcome {
958 Notified,
959 TimedOut,
960 Aborted,
961}
962
963enum CompletionWaitOutcome {
964 Response(Result<Option<Vec<u8>>, oneshot::error::RecvError>),
965 TimedOut,
966 Aborted,
967}
968
969fn normalize_names(names: Option<Vec<String>>) -> Option<BTreeSet<String>> {
970 names.and_then(|names| {
971 let normalized = names.into_iter().collect::<BTreeSet<_>>();
972 if normalized.is_empty() {
973 None
974 } else {
975 Some(normalized)
976 }
977 })
978}
979
980fn queue_message_from_row(row: internal_storage::QueueMessageRow) -> QueueMessage {
981 QueueMessage {
982 id: row.id,
983 name: row.message.name,
984 body: row.message.body,
985 created_at: row.message.created_at,
986 completion: None,
987 }
988}
989
990fn current_timestamp_ms() -> Result<i64> {
991 let now = SystemTime::now()
992 .duration_since(UNIX_EPOCH)
993 .context("current time is before unix epoch")?;
994 i64::try_from(now.as_millis()).context("queue timestamp exceeds i64")
995}
996
997fn duration_ms(duration: Duration) -> u64 {
998 duration.as_millis().try_into().unwrap_or(u64::MAX)
999}
1000
1001#[cfg(test)]
1003#[path = "../../tests/queue.rs"]
1004mod tests;