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