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