rs_teststand_websocket/server/mod.rs
1//! A bidirectional `WebSocket` bridge between a host and its user interfaces.
2//!
3//! The shape a real station takes: an orchestrator owns the engine, one or more
4//! panels connect, and the **same connection** carries progress out and requests
5//! back. A panel does not poll, and does not open a second channel to ask for
6//! something.
7//!
8//! Everything is JSON in `WebSocket` **text** frames, opcode `0x1` in RFC 6455.
9//! The protocol frames each message itself, so there is no terminator to agree
10//! on, unlike the line transport in `rs-teststand-bridge`, where CRLF exists precisely because a
11//! raw socket has no record boundary.
12//!
13//! # Threads
14//!
15//! Two, and the split is the point:
16//!
17//! - **The engine thread**, whichever thread called [`WebSocketBridge::bind`]
18//! and owns the engine. It publishes events and drains commands. Engine
19//! wrappers are neither [`Send`] nor [`Sync`], so nothing here can take one
20//! even by accident.
21//! - **The server thread**, a runtime of its own, accepting panels and moving
22//! bytes. It never sees the engine; what crosses between the two is
23//! `MessageEvent`, `Command` and `Response` from `rs-teststand-bridge`, all
24//! plain data.
25//!
26//! Replies go out as an `Ack`, a fixed five-field record, rather
27//! than as the `Response` enum whose fields vary by variant. A client sorts the
28//! two kinds of traffic on `command`: an acknowledgement always carries one and
29//! an event never does.
30
31use std::net::{SocketAddr, TcpListener as StdListener};
32use std::sync::mpsc;
33use std::thread;
34
35// Four things are easy to get wrong in a tokio websocket server, and each is
36// answered deliberately here rather than by accident. Changing this file means
37// keeping them true.
38//
39// A lagging receiver. `broadcast` drops messages for a receiver that falls
40// behind and reports `Lagged`. That is treated as fatal for the panel rather
41// than ignored: it is disconnected, because silently missing messages is worse
42// than a close it can react to. `EVENT_BACKLOG` bounds what one slow panel can
43// hold open.
44//
45// Cancellation in `select!`. The macro drops the futures it was polling when a
46// branch wins, so a branch future that had consumed something would lose it.
47// Both branches here poll cancel-safe futures. The bodies are safe for a
48// different reason: once a branch is chosen its body runs to completion, so the
49// `send` calls inside are never cut short.
50//
51// Split halves. `split` produces a read and a write half that cannot be
52// recombined, so both stay in this one task rather than being handed out.
53//
54// Locks across await points. There are none. State moves through channels.
55
56use tokio::sync::broadcast;
57
58use rs_teststand_bridge::{Command, Error, MessageEvent, Response};
59
60/// How many events the fan-out holds before the slowest panel misses some.
61///
62/// A panel that falls this far behind is disconnected rather than allowed to
63/// hold the buffer open: one that stopped reading is not a reason for the
64/// station to grow memory without limit.
65const EVENT_BACKLOG: usize = 256;
66
67/// Largest message a panel may send, in bytes.
68///
69/// Commands are small. A sequence path and a lookup string are the biggest
70/// parts of one, so a megabyte is generous by a wide margin. Without a limit a
71/// single frame can make the host allocate until it dies, and that needs no
72/// malice: a client with a loop bug reaches the same place.
73///
74/// This bounds what one panel can make the host hold. `EVENT_BACKLOG` bounds
75/// what a slow panel can make it keep.
76const MAX_MESSAGE_BYTES: usize = 1024 * 1024;
77
78/// Largest single frame accepted, in bytes.
79///
80/// Kept at the message limit. A message can arrive as several frames, so
81/// capping only the message would still let one frame be assembled unbounded
82/// before the total is known.
83const MAX_FRAME_BYTES: usize = MAX_MESSAGE_BYTES;
84
85/// Most panels served at once.
86///
87/// A host serves an orchestrator and the few panels a person has open, so this
88/// is far above normal use. It exists because nothing else stops a client that
89/// reconnects in a loop from opening sockets until the host runs out of them,
90/// and a station that has stopped answering is worse than one that refused a
91/// connection.
92///
93/// Refusing is deliberate rather than queueing: a panel told no can back off
94/// and return, while one left waiting cannot tell a busy host from a dead one.
95const MAX_CLIENTS: usize = 64;
96
97/// What travels out to the panels: an event, or an answer to one of them.
98mod accept;
99mod options;
100mod origin;
101mod page;
102mod session;
103
104pub use options::Options;
105
106use accept::serve;
107
108#[derive(Debug, Clone)]
109enum Outbound {
110 /// Broadcast to everyone.
111 Event(Box<MessageEvent>),
112 /// Addressed to the panel that asked.
113 Reply {
114 /// Which connection the answer belongs to.
115 client: u64,
116 /// The answer.
117 response: Box<Response>,
118 },
119 /// The host is going away, so every session should close.
120 ///
121 /// Sent when the bridge is dropped. Without it the server thread outlives
122 /// the bridge and the sockets it accepted stay open, so a panel keeps
123 /// waiting on a read that will never complete and looks connected to a host
124 /// that no longer exists.
125 Shutdown,
126}
127
128/// A command, with the panel that sent it.
129///
130/// The identity matters: a reply goes to the panel that asked, not to every
131/// panel watching.
132#[derive(Debug, Clone)]
133pub struct Request {
134 /// Which connection this arrived on.
135 pub client: u64,
136 /// What was asked.
137 pub command: Command,
138}
139
140/// Accepts panels, broadcasts events to them, and collects their commands.
141///
142/// Built on the engine's thread and used from there; the server runs elsewhere
143/// and shares nothing but data.
144#[derive(Debug)]
145pub struct WebSocketBridge {
146 outbound: broadcast::Sender<Outbound>,
147 commands: mpsc::Receiver<Request>,
148 address: SocketAddr,
149}
150
151impl WebSocketBridge {
152 /// Binds a listener and starts serving in the background.
153 ///
154 /// The socket is bound synchronously, so a port already in use is reported
155 /// to the caller that can do something about it rather than disappearing
156 /// into a thread.
157 ///
158 /// # Errors
159 /// [`Error::Transport`] if the address cannot be bound, or
160 /// [`Error::ThreadNotStarted`] if the server thread cannot be created.
161 pub fn bind(address: &str) -> Result<Self, Error> {
162 Self::bind_with(address, Options::default())
163 }
164
165 /// Binds, and also serves `page` to a browser asking for the root.
166 ///
167 /// One address for both, so opening the host's address is the whole setup
168 /// and the panel need not be found on disk. It also gives the page and the
169 /// socket the same origin, which is what makes an origin check meaningful;
170 /// a page loaded from a file has an origin of `null`.
171 ///
172 /// A browser request never consumes a connection slot or a subscription: it
173 /// is answered and closed before either is taken.
174 ///
175 /// # Errors
176 /// [`Error::Transport`] if the address cannot be bound, or
177 /// [`Error::ThreadNotStarted`] if the server thread will not start.
178 pub fn bind_with_page(address: &str, page: Option<String>) -> Result<Self, Error> {
179 let options = page.map_or_else(Options::default, |html| Options::default().page(html));
180 Self::bind_with(address, options)
181 }
182
183 /// Binds with everything the host wants to decide up front.
184 ///
185 /// The way to reach the origin allowlist. See [`Options`].
186 ///
187 /// # Errors
188 /// [`Error::Transport`] if the address cannot be bound, or
189 /// [`Error::ThreadNotStarted`] if the server thread will not start.
190 pub fn bind_with(address: &str, options: Options) -> Result<Self, Error> {
191 let listener = StdListener::bind(address)?;
192 listener.set_nonblocking(true)?;
193 let address = listener.local_addr()?;
194
195 let (outbound, _) = broadcast::channel(EVENT_BACKLOG);
196 let (command_sender, commands) = mpsc::channel();
197
198 let publisher = outbound.clone();
199 // A runtime of its own, on a thread of its own. The engine's thread must
200 // not host one: it is a single-threaded apartment that has to stay free
201 // to pump its own message queue.
202 thread::Builder::new()
203 .name("rs-teststand-websocket".to_owned())
204 .spawn(move || {
205 let Ok(runtime) = tokio::runtime::Builder::new_current_thread()
206 .enable_all()
207 .build()
208 else {
209 return;
210 };
211 runtime.block_on(serve(listener, publisher, command_sender, options));
212 })
213 .map_err(|error| Error::ThreadNotStarted {
214 reason: error.to_string(),
215 })?;
216
217 Ok(Self {
218 outbound,
219 commands,
220 address,
221 })
222 }
223
224 /// Tells every session to close when the bridge goes away.
225 ///
226 /// The sessions own the sockets, so this is the only way to reach them.
227 /// A send failure means nobody is listening, which is the same outcome.
228 fn shutdown(&self) {
229 let _ = self.outbound.send(Outbound::Shutdown);
230 }
231
232 /// The address actually bound, which resolves a port of `0`.
233 #[must_use]
234 pub const fn address(&self) -> SocketAddr {
235 self.address
236 }
237
238 /// How many panels are connected.
239 #[must_use]
240 pub fn client_count(&self) -> usize {
241 self.outbound.receiver_count()
242 }
243
244 /// Sends an event to every connected panel.
245 ///
246 /// Never blocks and never fails for want of an audience: with nobody
247 /// connected the event is dropped, because a station must not stop testing
248 /// because no one is watching.
249 pub fn publish(&self, event: &MessageEvent) {
250 let _ = self.outbound.send(Outbound::Event(Box::new(event.clone())));
251 }
252
253 /// Answers one request, addressed to the panel that made it.
254 pub fn reply(&self, request: &Request, response: &Response) {
255 let _ = self.outbound.send(Outbound::Reply {
256 client: request.client,
257 response: Box::new(response.clone()),
258 });
259 }
260
261 /// Takes the next command, if one is waiting.
262 ///
263 /// Non-blocking on purpose. The engine thread has its own queue to pump and
264 /// cannot afford to wait here; it drains what has arrived and gets on with
265 /// the run.
266 #[must_use]
267 pub fn next_command(&self) -> Option<Request> {
268 self.commands.try_recv().ok()
269 }
270}
271
272impl Drop for WebSocketBridge {
273 /// Closes the sessions rather than abandoning them.
274 ///
275 /// The server runs on its own thread and does not stop when this type is
276 /// dropped. Without telling the sessions to close, a panel is left holding
277 /// a socket to a host that has gone, blocked on a read that never returns.
278 fn drop(&mut self) {
279 self.shutdown();
280 }
281}