Skip to main content

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}