Skip to main content

rivetkit_core/
lib.rs

1#[cfg(all(feature = "native-runtime", feature = "wasm-runtime"))]
2compile_error!(
3	"`native-runtime` and `wasm-runtime` are mutually exclusive. Enable exactly one rivetkit-core runtime."
4);
5
6pub mod actor;
7#[cfg(feature = "native-runtime")]
8pub mod engine_process;
9pub mod error;
10pub mod inspector;
11pub mod inspector_bundle;
12pub mod metrics_endpoint;
13pub mod registry;
14pub mod runtime;
15pub mod serverless;
16#[cfg(feature = "native-runtime")]
17pub mod serverless_http;
18#[cfg(any(test, feature = "test-support"))]
19pub mod testing;
20pub(crate) mod time {
21	use std::fmt;
22	use std::future::Future;
23	use std::time::Duration;
24
25	#[cfg(target_arch = "wasm32")]
26	use futures::FutureExt;
27	#[cfg(target_arch = "wasm32")]
28	use wasm_bindgen::{JsCast, JsValue};
29	#[cfg(target_arch = "wasm32")]
30	use wasm_bindgen_futures::JsFuture;
31
32	#[cfg(not(target_arch = "wasm32"))]
33	pub use std::time::{Instant, SystemTime, UNIX_EPOCH};
34	#[cfg(target_arch = "wasm32")]
35	pub use web_time::{Instant, SystemTime, UNIX_EPOCH};
36
37	#[derive(Debug, Clone, Copy)]
38	pub struct TimeoutError;
39
40	impl fmt::Display for TimeoutError {
41		fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
42			f.write_str("operation timed out")
43		}
44	}
45
46	impl std::error::Error for TimeoutError {}
47
48	#[cfg(not(target_arch = "wasm32"))]
49	pub fn tokio_deadline(deadline: Instant) -> tokio::time::Instant {
50		deadline.into()
51	}
52
53	#[cfg(target_arch = "wasm32")]
54	pub async fn sleep(duration: Duration) {
55		let delay_ms = duration.as_millis().min(u32::MAX as u128) as f64;
56		let promise = js_sys::Promise::new(&mut |resolve, _reject| {
57			let global = js_sys::global();
58			let set_timeout = js_sys::Reflect::get(&global, &JsValue::from_str("setTimeout"))
59				.ok()
60				.and_then(|value| value.dyn_into::<js_sys::Function>().ok());
61
62			if let Some(set_timeout) = set_timeout {
63				let _ = set_timeout.call2(&global, &resolve, &JsValue::from_f64(delay_ms));
64			} else {
65				let _ = resolve.call0(&JsValue::UNDEFINED);
66			}
67		});
68
69		let _ = JsFuture::from(promise).await;
70	}
71
72	#[cfg(not(target_arch = "wasm32"))]
73	pub async fn sleep(duration: Duration) {
74		tokio::time::sleep(duration).await;
75	}
76
77	#[cfg(not(target_arch = "wasm32"))]
78	pub async fn sleep_until(deadline: Instant) {
79		tokio::time::sleep_until(tokio_deadline(deadline)).await;
80	}
81
82	#[cfg(target_arch = "wasm32")]
83	pub async fn sleep_until(deadline: Instant) {
84		let remaining = deadline
85			.checked_duration_since(Instant::now())
86			.unwrap_or(Duration::ZERO);
87		sleep(remaining).await;
88	}
89
90	#[cfg(not(target_arch = "wasm32"))]
91	pub async fn timeout<F>(duration: Duration, future: F) -> Result<F::Output, TimeoutError>
92	where
93		F: Future,
94	{
95		tokio::time::timeout(duration, future)
96			.await
97			.map_err(|_| TimeoutError)
98	}
99
100	#[cfg(target_arch = "wasm32")]
101	pub async fn timeout<F>(duration: Duration, future: F) -> Result<F::Output, TimeoutError>
102	where
103		F: Future,
104	{
105		futures::pin_mut!(future);
106		let timer = sleep(duration);
107		futures::pin_mut!(timer);
108
109		futures::select! {
110			result = future.fuse() => Ok(result),
111			_ = timer.fuse() => Err(TimeoutError),
112		}
113	}
114}
115pub mod types;
116pub mod websocket;
117pub use actor::{kv, sqlite};
118
119pub use actor::action::ActionDispatchError;
120pub use actor::config::{
121	ActionDefinition, ActorConfig, ActorConfigInput, ActorConfigOverrides, CanHibernateWebSocket,
122};
123pub use actor::connection::ConnHandle;
124pub use actor::context::{
125	ActorContext, ActorKv, ActorWorkRegion, KeepAwakeRegion, WebSocketCallbackRegion,
126};
127pub use actor::factory::{ActorEntryFn, ActorFactory};
128pub use actor::lifecycle_hooks::{ActorEvents, ActorStart, Reply};
129pub use actor::messages::{
130	ActorEvent, ActorHttpResponse, QueueSendResult, QueueSendStatus, Request, Response,
131	SerializeStateReason, StateDelta, StreamingResponse, WorkflowKvWrite,
132};
133pub use actor::queue::{
134	CompletableQueueMessage, EnqueueAndWaitOpts, QueueMessage, QueueNextBatchOpts, QueueNextOpts,
135	QueueTryNextBatchOpts, QueueTryNextOpts, QueueWaitOpts,
136};
137pub use actor::sqlite::{
138	BindParam, ColumnValue, ExecResult, ExecuteResult, QueryResult, SqliteBackend,
139	SqliteBatchStatement, SqliteDb, SqliteTransaction,
140};
141pub use actor::state::RequestSaveOpts;
142pub use actor::task::{
143	ActionDispatchResult, ActorTask, DispatchCommand, HttpDispatchResult, LifecycleCommand,
144	LifecycleEvent, LifecycleState,
145};
146pub use actor::task_types::ShutdownKind;
147pub use actor::work_registry::{ActorWorkKind, ActorWorkPolicy};
148pub use error::ActorLifecycle;
149pub use inspector::{Inspector, InspectorSnapshot};
150pub use registry::{CoreRegistry, EngineSpawnMode, ServeConfig};
151pub use rivet_envoy_client::config::{
152	HTTP_BODY_STREAM_CHANNEL_CAPACITY, HttpRequestBodyStream, ResponseChunk,
153};
154pub use runtime::{RuntimeBoxFuture, RuntimeSpawner, boxed_runtime_future};
155pub use serverless::{CoreServerlessRuntime, ServerlessRequest, ServerlessResponse};
156pub use types::{
157	ActorKey, ActorKeySegment, ConnId, ListOpts, SaveStateOpts, WsMessage, format_actor_key,
158};
159pub use websocket::WebSocket;