Skip to main content

rivetkit_core/inspector/
mod.rs

1use std::sync::Arc;
2use std::sync::Weak;
3use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering};
4
5use parking_lot::RwLock;
6
7pub mod auth;
8pub(crate) mod protocol;
9pub mod tabs;
10
11pub use auth::{InspectorAuth, init_inspector_token, set_test_inspector_token_override};
12pub use tabs::{BUILTIN_TAB_IDS, InspectorTabEntry, validate_inspector_tabs};
13
14type InspectorListener = Arc<dyn Fn(InspectorSignal) + Send + Sync>;
15
16#[derive(Clone, Debug, Default)]
17pub struct Inspector(Arc<InspectorInner>);
18
19struct InspectorInner {
20	state_revision: AtomicU64,
21	connections_revision: AtomicU64,
22	queue_revision: AtomicU64,
23	active_connections: AtomicU32,
24	queue_size: AtomicU32,
25	connected_clients: AtomicUsize,
26	next_listener_id: AtomicU64,
27	// Forced-sync: subscriptions are created/dropped from sync paths and
28	// listener callbacks are cloned before invocation.
29	listeners: RwLock<Vec<(u64, InspectorListener)>>,
30}
31
32#[allow(clippy::enum_variant_names)]
33#[derive(Clone, Copy, Debug, PartialEq, Eq)]
34pub(crate) enum InspectorSignal {
35	StateUpdated,
36	ConnectionsUpdated,
37	QueueUpdated,
38	WorkflowHistoryUpdated,
39}
40
41pub(crate) struct InspectorSubscription {
42	inspector: Weak<InspectorInner>,
43	listener_id: u64,
44}
45
46impl std::fmt::Debug for InspectorInner {
47	fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
48		f.debug_struct("InspectorInner")
49			.field(
50				"state_revision",
51				&self.state_revision.load(Ordering::SeqCst),
52			)
53			.field(
54				"connections_revision",
55				&self.connections_revision.load(Ordering::SeqCst),
56			)
57			.field(
58				"queue_revision",
59				&self.queue_revision.load(Ordering::SeqCst),
60			)
61			.field(
62				"active_connections",
63				&self.active_connections.load(Ordering::SeqCst),
64			)
65			.field("queue_size", &self.queue_size.load(Ordering::SeqCst))
66			.field(
67				"connected_clients",
68				&self.connected_clients.load(Ordering::SeqCst),
69			)
70			.finish()
71	}
72}
73
74impl Default for InspectorInner {
75	fn default() -> Self {
76		Self {
77			state_revision: AtomicU64::new(0),
78			connections_revision: AtomicU64::new(0),
79			queue_revision: AtomicU64::new(0),
80			active_connections: AtomicU32::new(0),
81			queue_size: AtomicU32::new(0),
82			connected_clients: AtomicUsize::new(0),
83			next_listener_id: AtomicU64::new(1),
84			listeners: RwLock::new(Vec::new()),
85		}
86	}
87}
88
89impl Drop for InspectorSubscription {
90	fn drop(&mut self) {
91		let Some(inspector) = self.inspector.upgrade() else {
92			return;
93		};
94		let connected_clients = {
95			let mut listeners = inspector.listeners.write();
96			listeners.retain(|(listener_id, _)| *listener_id != self.listener_id);
97			listeners.len()
98		};
99		inspector
100			.connected_clients
101			.store(connected_clients, Ordering::SeqCst);
102	}
103}
104
105#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
106pub struct InspectorSnapshot {
107	pub state_revision: u64,
108	pub connections_revision: u64,
109	pub queue_revision: u64,
110	pub active_connections: u32,
111	pub queue_size: u32,
112	pub connected_clients: usize,
113}
114
115impl Inspector {
116	pub fn new() -> Self {
117		Self::default()
118	}
119
120	pub fn snapshot(&self) -> InspectorSnapshot {
121		InspectorSnapshot {
122			state_revision: self.0.state_revision.load(Ordering::SeqCst),
123			connections_revision: self.0.connections_revision.load(Ordering::SeqCst),
124			queue_revision: self.0.queue_revision.load(Ordering::SeqCst),
125			active_connections: self.0.active_connections.load(Ordering::SeqCst),
126			queue_size: self.0.queue_size.load(Ordering::SeqCst),
127			connected_clients: self.0.connected_clients.load(Ordering::SeqCst),
128		}
129	}
130
131	pub(crate) fn subscribe(&self, listener: InspectorListener) -> InspectorSubscription {
132		let listener_id = self.0.next_listener_id.fetch_add(1, Ordering::SeqCst);
133		let connected_clients = {
134			let mut listeners = self.0.listeners.write();
135			listeners.push((listener_id, listener));
136			listeners.len()
137		};
138		self.set_connected_clients(connected_clients);
139
140		InspectorSubscription {
141			inspector: Arc::downgrade(&self.0),
142			listener_id,
143		}
144	}
145
146	pub(crate) fn record_state_updated(&self) {
147		self.0.state_revision.fetch_add(1, Ordering::SeqCst);
148		self.notify(InspectorSignal::StateUpdated);
149	}
150
151	pub(crate) fn record_connections_updated(&self, active_connections: u32) {
152		self.0
153			.active_connections
154			.store(active_connections, Ordering::SeqCst);
155		self.0.connections_revision.fetch_add(1, Ordering::SeqCst);
156		self.notify(InspectorSignal::ConnectionsUpdated);
157	}
158
159	pub(crate) fn record_queue_updated(&self, queue_size: u32) {
160		self.0.queue_size.store(queue_size, Ordering::SeqCst);
161		self.0.queue_revision.fetch_add(1, Ordering::SeqCst);
162		self.notify(InspectorSignal::QueueUpdated);
163	}
164
165	pub(crate) fn record_workflow_history_updated(&self) {
166		self.notify(InspectorSignal::WorkflowHistoryUpdated);
167	}
168
169	pub(crate) fn set_connected_clients(&self, connected_clients: usize) {
170		self.0
171			.connected_clients
172			.store(connected_clients, Ordering::SeqCst);
173	}
174
175	fn notify(&self, signal: InspectorSignal) {
176		if self.0.connected_clients.load(Ordering::SeqCst) == 0 {
177			return;
178		}
179
180		let listeners = {
181			let listeners = self.0.listeners.read();
182			listeners
183				.iter()
184				.map(|(_, listener)| listener.clone())
185				.collect::<Vec<_>>()
186		};
187
188		for listener in listeners {
189			listener(signal);
190		}
191	}
192}
193
194pub fn decode_request_payload(payload: &[u8], advertised_version: u16) -> anyhow::Result<Vec<u8>> {
195	let message = protocol::decode_client_payload(payload, advertised_version)?;
196	protocol::encode_client_payload_current(&message)
197}
198
199pub fn encode_response_payload(payload: &[u8], target_version: u16) -> anyhow::Result<Vec<u8>> {
200	let message = protocol::decode_current_server_payload(payload)?;
201	protocol::encode_server_payload(&message, target_version)
202}
203
204// Test shim keeps moved tests in crate-root tests/ with private-module access.
205#[cfg(test)]
206#[path = "../../tests/inspector.rs"]
207mod tests;