Skip to main content

rivetkit_core/actor/
queue.rs

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	/// Removes all messages from the queue and resets the size counter.
521	pub async fn reset(&self) -> Result<()> {
522		self.ensure_initialized().await?;
523
524		// Serialize against receivers before touching metadata. Lock order matches
525		// try_receive_batch (receive lock then metadata) so there is no deadlock.
526		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	/// TS parity: queue-manager.ts keeps `enqueueAndWait` completion waits
742	/// alive across actor aborts; the surrounding tracked user task owns
743	/// shutdown cancellation.
744	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// Test shim keeps moved tests in crate-root tests/ with private-module access.
1002#[cfg(test)]
1003#[path = "../../tests/queue.rs"]
1004mod tests;