rs_teststand_bridge/host.rs
1//! Owning the engine on its own thread so many callers can share it.
2
3use std::sync::mpsc;
4use std::thread;
5use std::time::Duration;
6
7use tokio::sync::{broadcast, oneshot};
8
9use rs_teststand::Engine;
10
11use crate::event::{MessageEvent, PayloadPolicy};
12
13use crate::Error;
14
15/// How many messages the broadcast channel holds before the slowest subscriber
16/// starts missing them.
17///
18/// A subscriber that falls this far behind is told how many it lost rather than
19/// being silently starved.
20const MESSAGE_BACKLOG: usize = 1024;
21
22/// How long the engine thread waits between queue polls when idle.
23///
24/// Short enough that a message is forwarded promptly, long enough that an idle
25/// host is not spinning a core.
26const IDLE_POLL_INTERVAL: Duration = Duration::from_millis(20);
27
28/// How long to let the engine finish shutting down before giving up on it.
29///
30/// Bounded on purpose: a host that cannot be stopped is worse than one that
31/// stops untidily.
32const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(10);
33
34/// A unit of work to run against the engine, on the engine's own thread.
35type Job = Box<dyn FnOnce(&Engine) + Send>;
36
37/// A shared, thread-safe front end to one engine.
38///
39/// The engine is a single-threaded-apartment COM object: its wrappers are
40/// deliberately neither [`Send`] nor [`Sync`], and driving one from several
41/// threads does not work. A server, though, is inherently multi-threaded.
42///
43/// This resolves that by giving the engine a thread of its own. Callers submit
44/// work; the engine thread runs it and sends the answer back. Because the work
45/// is a closure rather than a fixed command list, the whole API is reachable
46/// without enumerating it:
47///
48/// ```no_run
49/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
50/// use rs_teststand_bridge::EngineHost;
51///
52/// let host = EngineHost::start()?;
53/// let version = host.with_engine(|engine| engine.version_string()).await??;
54/// println!("TestStand {version}");
55/// # Ok(())
56/// # }
57/// ```
58///
59/// The same thread polls the message queue and broadcasts what it finds, so a
60/// host can serve requests and forward progress at once, which is what an IPC
61/// layer needs and what a single-threaded caller cannot do.
62#[derive(Debug)]
63pub struct EngineHost {
64 /// Taken and dropped before the worker is joined. See [`Drop`].
65 jobs: Option<mpsc::Sender<Job>>,
66 messages: broadcast::Sender<MessageEvent>,
67 /// Kept so the thread is joined when the host is dropped.
68 worker: Option<thread::JoinHandle<()>>,
69}
70
71impl EngineHost {
72 /// Starts an engine on a dedicated thread.
73 ///
74 /// Message polling is enabled, so anything a sequence posts reaches
75 /// [`subscribe`](Self::subscribe) without further setup.
76 ///
77 /// # Errors
78 /// [`Error`] if the engine cannot be created.
79 pub fn start() -> Result<Self, Error> {
80 let (job_sender, job_receiver) = mpsc::channel::<Job>();
81 let (message_sender, _) = broadcast::channel(MESSAGE_BACKLOG);
82 let (ready_sender, ready_receiver) = mpsc::channel::<Result<(), Error>>();
83
84 let messages = message_sender.clone();
85 let worker = thread::Builder::new()
86 .name("rs-teststand-engine".to_owned())
87 .spawn(move || Self::run(&job_receiver, &messages, &ready_sender))
88 .map_err(|source| Error::ThreadNotStarted {
89 reason: source.to_string(),
90 })?;
91
92 // Surface a construction failure to the caller rather than leaving a
93 // dead thread behind a healthy-looking handle.
94 match ready_receiver.recv() {
95 Ok(Ok(())) => Ok(Self {
96 jobs: Some(job_sender),
97 messages: message_sender,
98 worker: Some(worker),
99 }),
100 Ok(Err(error)) => Err(error),
101 // The thread dropped its side without reporting: it is not running.
102 Err(_) => Err(Error::HostStopped),
103 }
104 }
105
106 /// Runs a closure against the engine and waits for its result.
107 ///
108 /// The closure runs on the engine's thread, so it may use any part of the
109 /// API. Its result must be [`Send`] to come back.
110 ///
111 /// The outer error means the engine thread is gone; the inner one is
112 /// whatever the closure produced.
113 ///
114 /// # Errors
115 /// [`Error`] if the engine thread has stopped.
116 pub async fn with_engine<F, T>(&self, work: F) -> Result<T, Error>
117 where
118 F: FnOnce(&Engine) -> T + Send + 'static,
119 T: Send + 'static,
120 {
121 let (result_sender, result_receiver) = oneshot::channel();
122 let job: Job = Box::new(move |engine| {
123 // A closed receiver only means the caller stopped waiting.
124 let _ = result_sender.send(work(engine));
125 });
126 self.jobs
127 .as_ref()
128 .ok_or(Error::HostStopped)?
129 .send(job)
130 .map_err(|_| Error::HostStopped)?;
131 result_receiver.await.map_err(|_| Error::ResultLost)
132 }
133
134 /// Receives every message posted from now on.
135 ///
136 /// Each subscriber gets its own copy. One that falls further behind than
137 /// the backlog is told how many it missed rather than quietly losing them.
138 #[must_use]
139 pub fn subscribe(&self) -> broadcast::Receiver<MessageEvent> {
140 self.messages.subscribe()
141 }
142
143 /// How many subscribers are currently listening.
144 #[must_use]
145 pub fn subscriber_count(&self) -> usize {
146 self.messages.receiver_count()
147 }
148
149 /// The engine thread: create, then serve work and drain the queue.
150 fn run(
151 jobs: &mpsc::Receiver<Job>,
152 messages: &broadcast::Sender<MessageEvent>,
153 ready: &mpsc::Sender<Result<(), Error>>,
154 ) {
155 let engine = match Engine::new() {
156 Ok(engine) => engine,
157 Err(error) => {
158 let _ = ready.send(Err(error.into()));
159 return;
160 }
161 };
162 // Nothing reaches the queue until this is on.
163 if let Err(error) = engine.set_ui_message_polling_enabled(true) {
164 let _ = ready.send(Err(error.into()));
165 return;
166 }
167 if ready.send(Ok(())).is_err() {
168 return;
169 }
170
171 loop {
172 // The engine is a single-threaded-apartment object and this is a
173 // background thread, where nothing pumps by default. Without this
174 // the apartment cannot service the calls an execution makes and the
175 // process aborts rather than failing cleanly.
176 if rs_teststand_sys::pump_thread_messages() {
177 Self::close_apartment(engine);
178 return;
179 }
180
181 // Messages next: a synchronous one blocks the sequence that posted
182 // it until acknowledged, so draining beats serving a request.
183 Self::drain_messages(&engine, messages);
184
185 match jobs.recv_timeout(IDLE_POLL_INTERVAL) {
186 Ok(job) => job(&engine),
187 Err(mpsc::RecvTimeoutError::Timeout) => {}
188 // Every handle is gone: drain once more so a final message is
189 // not lost, then let the engine drop with the thread.
190 Err(mpsc::RecvTimeoutError::Disconnected) => {
191 Self::drain_messages(&engine, messages);
192 // ShutDown is asynchronous and must be waited on: dropping
193 // the engine while it is still closing files and
194 // terminating executions tears COM down underneath work in
195 // progress. Bounded, so a stuck engine cannot wedge the
196 // host thread.
197 Self::close_apartment(engine);
198 return;
199 }
200 }
201 }
202 }
203
204 /// Releases the engine and leaves the apartment, in that order.
205 ///
206 /// This thread initialized a single-threaded apartment and is about to
207 /// exit. Unlike the process's main thread it genuinely detaches, so the
208 /// apartment has to be closed or the COM runtime is left believing a live
209 /// thread still owns it.
210 ///
211 /// Order is the whole point: the engine is dropped first, because
212 /// uninitializing COM while an object is still alive aborts the process.
213 fn close_apartment(engine: Engine) {
214 let _ = engine.close(SHUTDOWN_TIMEOUT);
215 }
216
217 /// Takes everything currently queued and publishes it.
218 ///
219 /// Every message is acknowledged even when nobody is subscribed: the
220 /// acknowledgement is what releases a synchronous poster, so skipping it
221 /// would stall the sequence rather than merely drop a notification.
222 fn drain_messages(engine: &Engine, messages: &broadcast::Sender<MessageEvent>) {
223 while matches!(engine.is_ui_message_queue_empty(), Ok(false)) {
224 let Ok(message) = engine.get_ui_message() else {
225 return;
226 };
227 // A conversion failure must not skip the acknowledgement below: an
228 // unacknowledged message stalls a synchronous poster.
229 if let Ok(event) = MessageEvent::from_ui_message(&message, PayloadPolicy::default()) {
230 // A send failure only means nobody is listening yet.
231 let _ = messages.send(event);
232 }
233 let _ = message.acknowledge();
234 }
235 }
236}
237
238impl Drop for EngineHost {
239 fn drop(&mut self) {
240 // Order matters, and it is not the order the fields are written in.
241 // `Drop::drop` runs before a struct's own fields are dropped, so the
242 // job sender is still alive here. The engine thread leaves its loop
243 // only when every sender is gone, so joining first waits for a
244 // disconnect that this very value is preventing, and the thread and
245 // the joiner deadlock.
246 drop(self.jobs.take());
247
248 // Now the loop can end. Joining lets the engine release its COM
249 // references on the thread that created them.
250 if let Some(worker) = self.worker.take() {
251 let _ = worker.join();
252 }
253 }
254}
255
256#[cfg(test)]
257mod tests {
258 use rs_teststand::UIMessageCode;
259
260 use crate::event::MessageEvent;
261
262 fn event(code: i32) -> MessageEvent {
263 MessageEvent {
264 code,
265 numeric: 0.0,
266 text: String::new(),
267 payload: None,
268 synchronous: false,
269 execution_id: None,
270 }
271 }
272
273 #[test]
274 fn a_sequence_posted_event_is_distinguishable_from_an_engine_one() {
275 assert!(event(UIMessageCode::USER_MESSAGE_BASE + 1).is_from_sequence());
276 assert!(!event(UIMessageCode::EndExecution.bits()).is_from_sequence());
277 }
278
279 #[test]
280 fn an_engine_event_resolves_to_its_name() {
281 assert_eq!(
282 event(UIMessageCode::EndExecution.bits()).engine_code(),
283 Some(UIMessageCode::EndExecution)
284 );
285 // A sequence's own code has no engine name, and that is not a failure.
286 assert_eq!(event(UIMessageCode::USER_MESSAGE_BASE).engine_code(), None);
287 }
288}