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