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 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#[cfg(test)]
206#[path = "../../tests/inspector.rs"]
207mod tests;