Line data Source code
1 : use std::panic::AssertUnwindSafe;
2 : use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU32, Ordering};
3 : use std::sync::{Arc, OnceLock, mpsc};
4 : use std::thread::JoinHandle;
5 :
6 : use tracing::{debug, error, trace};
7 :
8 : use crate::time::Timing;
9 :
10 : /// Stores status of the running threads. Threads are registered in the runtime upon creation
11 : /// and deregistered upon termination.
12 : pub struct Runtime {
13 : // stores handles to all threads that are currently running
14 : threads: Vec<ThreadHandle>,
15 : // stores current time and pending wakeups
16 : clock: Arc<Timing>,
17 : // thread counter
18 : thread_counter: AtomicU32,
19 : // Thread step counter -- how many times all threads has been actually
20 : // stepped (note that all world/time/executor/thread have slightly different
21 : // meaning of steps). For observability.
22 : pub step_counter: u64,
23 : }
24 :
25 : impl Runtime {
26 : /// Init new runtime, no running threads.
27 528 : pub fn new(clock: Arc<Timing>) -> Self {
28 528 : Self {
29 528 : threads: Vec::new(),
30 528 : clock,
31 528 : thread_counter: AtomicU32::new(0),
32 528 : step_counter: 0,
33 528 : }
34 528 : }
35 :
36 : /// Spawn a new thread and register it in the runtime.
37 19806 : pub fn spawn<F>(&mut self, f: F) -> ExternalHandle
38 19806 : where
39 19806 : F: FnOnce() + Send + 'static,
40 19806 : {
41 19806 : let (tx, rx) = mpsc::channel();
42 19806 :
43 19806 : let clock = self.clock.clone();
44 19806 : let tid = self.thread_counter.fetch_add(1, Ordering::SeqCst);
45 19806 : debug!("spawning thread-{}", tid);
46 :
47 19806 : let join = std::thread::spawn(move || {
48 19806 : let _guard = tracing::info_span!("", tid).entered();
49 19806 :
50 19806 : let res = std::panic::catch_unwind(AssertUnwindSafe(|| {
51 19806 : with_thread_context(|ctx| {
52 19806 : assert!(ctx.clock.set(clock).is_ok());
53 19806 : ctx.id.store(tid, Ordering::SeqCst);
54 19806 : tx.send(ctx.clone()).expect("failed to send thread context");
55 19806 : // suspend thread to put it to `threads` in sleeping state
56 19806 : ctx.yield_me(0);
57 19806 : });
58 19806 :
59 19806 : // start user-provided function
60 19806 : f();
61 19806 : }));
62 19806 : debug!("thread finished");
63 :
64 19748 : if let Err(e) = res {
65 19728 : with_thread_context(|ctx| {
66 19728 : if !ctx.allow_panic.load(std::sync::atomic::Ordering::SeqCst) {
67 0 : error!("thread panicked, terminating the process: {:?}", e);
68 0 : std::process::exit(1);
69 19728 : }
70 19728 :
71 19728 : debug!("thread panicked: {:?}", e);
72 19728 : let mut result = ctx.result.lock();
73 19728 : if result.0 == -1 {
74 19296 : *result = (256, format!("thread panicked: {:?}", e));
75 19296 : }
76 19728 : });
77 19728 : }
78 :
79 19748 : with_thread_context(|ctx| {
80 19748 : ctx.finish_me();
81 19748 : });
82 19806 : });
83 19806 :
84 19806 : let ctx = rx.recv().expect("failed to receive thread context");
85 19806 : let handle = ThreadHandle::new(ctx.clone(), join);
86 19806 :
87 19806 : self.threads.push(handle);
88 19806 :
89 19806 : ExternalHandle { ctx }
90 19806 : }
91 :
92 : /// Returns true if there are any unfinished activity, such as running thread or pending events.
93 : /// Otherwise returns false, which means all threads are blocked forever.
94 418254 : pub fn step(&mut self) -> bool {
95 418254 : trace!("runtime step");
96 :
97 : // have we run any thread?
98 418254 : let mut ran = false;
99 418254 :
100 2087733 : self.threads.retain(|thread: &ThreadHandle| {
101 2087733 : let res = thread.ctx.wakeup.compare_exchange(
102 2087733 : PENDING_WAKEUP,
103 2087733 : NO_WAKEUP,
104 2087733 : Ordering::SeqCst,
105 2087733 : Ordering::SeqCst,
106 2087733 : );
107 2087733 : if res.is_err() {
108 : // thread has no pending wakeups, leaving as is
109 1798773 : return true;
110 288960 : }
111 288960 : ran = true;
112 288960 :
113 288960 : trace!("entering thread-{}", thread.ctx.tid());
114 288960 : let status = thread.step();
115 288960 : self.step_counter += 1;
116 288960 : trace!(
117 0 : "out of thread-{} with status {:?}",
118 0 : thread.ctx.tid(),
119 : status
120 : );
121 :
122 288960 : if status == Status::Sleep {
123 269212 : true
124 : } else {
125 19748 : trace!("thread has finished");
126 : // removing the thread from the list
127 19748 : false
128 : }
129 2087733 : });
130 418254 :
131 418254 : if !ran {
132 221271 : trace!("no threads were run, stepping clock");
133 221271 : if let Some(ctx_to_wake) = self.clock.step() {
134 220723 : trace!("waking up thread-{}", ctx_to_wake.tid());
135 220723 : ctx_to_wake.inc_wake();
136 : } else {
137 548 : return false;
138 : }
139 196983 : }
140 :
141 417706 : true
142 418254 : }
143 :
144 : /// Kill all threads. This is done by setting a flag in each thread context and waking it up.
145 1008 : pub fn crash_all_threads(&mut self) {
146 2863 : for thread in self.threads.iter() {
147 2863 : thread.ctx.crash_stop();
148 2863 : }
149 :
150 : // all threads should be finished after a few steps
151 1512 : while !self.threads.is_empty() {
152 504 : self.step();
153 504 : }
154 1008 : }
155 : }
156 :
157 : impl Drop for Runtime {
158 503 : fn drop(&mut self) {
159 503 : debug!("dropping the runtime");
160 503 : self.crash_all_threads();
161 503 : }
162 : }
163 :
164 : #[derive(Clone)]
165 : pub struct ExternalHandle {
166 : ctx: Arc<ThreadContext>,
167 : }
168 :
169 : impl ExternalHandle {
170 : /// Returns true if thread has finished execution.
171 433007 : pub fn is_finished(&self) -> bool {
172 433007 : let status = self.ctx.mutex.lock();
173 433007 : *status == Status::Finished
174 433007 : }
175 :
176 : /// Returns exitcode and message, which is available after thread has finished execution.
177 427 : pub fn result(&self) -> (i32, String) {
178 427 : let result = self.ctx.result.lock();
179 427 : result.clone()
180 427 : }
181 :
182 : /// Returns thread id.
183 16 : pub fn id(&self) -> u32 {
184 16 : self.ctx.id.load(Ordering::SeqCst)
185 16 : }
186 :
187 : /// Sets a flag to crash thread on the next wakeup.
188 16781 : pub fn crash_stop(&self) {
189 16781 : self.ctx.crash_stop();
190 16781 : }
191 : }
192 :
193 : struct ThreadHandle {
194 : ctx: Arc<ThreadContext>,
195 : _join: JoinHandle<()>,
196 : }
197 :
198 : impl ThreadHandle {
199 : /// Create a new [`ThreadHandle`] and wait until thread will enter [`Status::Sleep`] state.
200 19806 : fn new(ctx: Arc<ThreadContext>, join: JoinHandle<()>) -> Self {
201 19806 : let mut status = ctx.mutex.lock();
202 : // wait until thread will go into the first yield
203 19838 : while *status != Status::Sleep {
204 32 : ctx.condvar.wait(&mut status);
205 32 : }
206 19806 : drop(status);
207 19806 :
208 19806 : Self { ctx, _join: join }
209 19806 : }
210 :
211 : /// Allows thread to execute one step of its execution.
212 : /// Returns [`Status`] of the thread after the step.
213 288960 : fn step(&self) -> Status {
214 288960 : let mut status = self.ctx.mutex.lock();
215 288960 : assert!(matches!(*status, Status::Sleep));
216 :
217 288960 : *status = Status::Running;
218 288960 : self.ctx.condvar.notify_all();
219 :
220 577920 : while *status == Status::Running {
221 288960 : self.ctx.condvar.wait(&mut status);
222 288960 : }
223 :
224 288960 : *status
225 288960 : }
226 : }
227 :
228 : #[derive(Clone, Copy, Debug, PartialEq, Eq)]
229 : enum Status {
230 : /// Thread is running.
231 : Running,
232 : /// Waiting for event to complete, will be resumed by the executor step, once wakeup flag is set.
233 : Sleep,
234 : /// Thread finished execution.
235 : Finished,
236 : }
237 :
238 : const NO_WAKEUP: u8 = 0;
239 : const PENDING_WAKEUP: u8 = 1;
240 :
241 : pub struct ThreadContext {
242 : id: AtomicU32,
243 : // used to block thread until it is woken up
244 : mutex: parking_lot::Mutex<Status>,
245 : condvar: parking_lot::Condvar,
246 : // used as a flag to indicate runtime that thread is ready to be woken up
247 : wakeup: AtomicU8,
248 : clock: OnceLock<Arc<Timing>>,
249 : // execution result, set by exit() call
250 : result: parking_lot::Mutex<(i32, String)>,
251 : // determines if process should be killed on receiving panic
252 : allow_panic: AtomicBool,
253 : // acts as a signal that thread should crash itself on the next wakeup
254 : crash_request: AtomicBool,
255 : }
256 :
257 : impl ThreadContext {
258 20334 : pub(crate) fn new() -> Self {
259 20334 : Self {
260 20334 : id: AtomicU32::new(0),
261 20334 : mutex: parking_lot::Mutex::new(Status::Running),
262 20334 : condvar: parking_lot::Condvar::new(),
263 20334 : wakeup: AtomicU8::new(NO_WAKEUP),
264 20334 : clock: OnceLock::new(),
265 20334 : result: parking_lot::Mutex::new((-1, String::new())),
266 20334 : allow_panic: AtomicBool::new(false),
267 20334 : crash_request: AtomicBool::new(false),
268 20334 : }
269 20334 : }
270 : }
271 :
272 : // Functions for executor to control thread execution.
273 : impl ThreadContext {
274 : /// Set atomic flag to indicate that thread is ready to be woken up.
275 684780 : fn inc_wake(&self) {
276 684780 : self.wakeup.store(PENDING_WAKEUP, Ordering::SeqCst);
277 684780 : }
278 :
279 : /// Internal function used for event queues.
280 179596 : pub(crate) fn schedule_wakeup(self: &Arc<Self>, after_ms: u64) {
281 179596 : self.clock
282 179596 : .get()
283 179596 : .unwrap()
284 179596 : .schedule_wakeup(after_ms, self.clone());
285 179596 : }
286 :
287 1 : fn tid(&self) -> u32 {
288 1 : self.id.load(Ordering::SeqCst)
289 1 : }
290 :
291 19644 : fn crash_stop(&self) {
292 19644 : let status = self.mutex.lock();
293 19644 : if *status == Status::Finished {
294 5 : debug!(
295 0 : "trying to crash thread-{}, which is already finished",
296 0 : self.tid()
297 : );
298 5 : return;
299 19639 : }
300 19639 : assert!(matches!(*status, Status::Sleep));
301 19639 : drop(status);
302 19639 :
303 19639 : self.allow_panic.store(true, Ordering::SeqCst);
304 19639 : self.crash_request.store(true, Ordering::SeqCst);
305 19639 : // set a wakeup
306 19639 : self.inc_wake();
307 : // it will panic on the next wakeup
308 19644 : }
309 : }
310 :
311 : // Internal functions.
312 : impl ThreadContext {
313 : /// Blocks thread until it's woken up by the executor. If `after_ms` is 0, is will be
314 : /// woken on the next step. If `after_ms` > 0, wakeup is scheduled after that time.
315 : /// Otherwise wakeup is not scheduled inside `yield_me`, and should be arranged before
316 : /// calling this function.
317 289018 : fn yield_me(self: &Arc<Self>, after_ms: i64) {
318 289018 : let mut status = self.mutex.lock();
319 289018 : assert!(matches!(*status, Status::Running));
320 :
321 289018 : match after_ms.cmp(&0) {
322 238012 : std::cmp::Ordering::Less => {
323 238012 : // block until something wakes us up
324 238012 : }
325 21310 : std::cmp::Ordering::Equal => {
326 21310 : // tell executor that we are ready to be woken up
327 21310 : self.inc_wake();
328 21310 : }
329 29696 : std::cmp::Ordering::Greater => {
330 29696 : // schedule wakeup
331 29696 : self.clock
332 29696 : .get()
333 29696 : .unwrap()
334 29696 : .schedule_wakeup(after_ms as u64, self.clone());
335 29696 : }
336 : }
337 :
338 289018 : *status = Status::Sleep;
339 289018 : self.condvar.notify_all();
340 :
341 : // wait until executor wakes us up
342 578036 : while *status != Status::Running {
343 289018 : self.condvar.wait(&mut status);
344 289018 : }
345 :
346 289018 : if self.crash_request.load(Ordering::SeqCst) {
347 19296 : panic!("crashed by request");
348 269722 : }
349 269722 : }
350 :
351 : /// Called only once, exactly before thread finishes execution.
352 19748 : fn finish_me(&self) {
353 19748 : let mut status = self.mutex.lock();
354 19748 : assert!(matches!(*status, Status::Running));
355 :
356 19748 : *status = Status::Finished;
357 19748 : {
358 19748 : let mut result = self.result.lock();
359 19748 : if result.0 == -1 {
360 20 : *result = (0, "finished normally".to_owned());
361 19728 : }
362 : }
363 19748 : self.condvar.notify_all();
364 19748 : }
365 : }
366 :
367 : /// Invokes the given closure with a reference to the current thread [`ThreadContext`].
368 : #[inline(always)]
369 1876737 : fn with_thread_context<T>(f: impl FnOnce(&Arc<ThreadContext>) -> T) -> T {
370 1876737 : thread_local!(static THREAD_DATA: Arc<ThreadContext> = Arc::new(ThreadContext::new()));
371 1876737 : THREAD_DATA.with(f)
372 1876737 : }
373 :
374 : /// Waker is used to wake up threads that are blocked on condition.
375 : /// It keeps track of contexts [`Arc<ThreadContext>`] and can increment the counter
376 : /// of several contexts to send a notification.
377 : pub struct Waker {
378 : // contexts that are waiting for a notification
379 : contexts: parking_lot::Mutex<smallvec::SmallVec<[Arc<ThreadContext>; 8]>>,
380 : }
381 :
382 : impl Default for Waker {
383 0 : fn default() -> Self {
384 0 : Self::new()
385 0 : }
386 : }
387 :
388 : impl Waker {
389 81572 : pub fn new() -> Self {
390 81572 : Self {
391 81572 : contexts: parking_lot::Mutex::new(smallvec::SmallVec::new()),
392 81572 : }
393 81572 : }
394 :
395 : /// Subscribe current thread to receive a wake notification later.
396 829354 : pub fn wake_me_later(&self) {
397 829354 : with_thread_context(|ctx| {
398 829354 : self.contexts.lock().push(ctx.clone());
399 829354 : });
400 829354 : }
401 :
402 : /// Wake up all threads that are waiting for a notification and clear the list.
403 125505 : pub fn wake_all(&self) {
404 125505 : let mut v = self.contexts.lock();
405 423108 : for ctx in v.iter() {
406 423108 : ctx.inc_wake();
407 423108 : }
408 125505 : v.clear();
409 125505 : }
410 : }
411 :
412 : /// See [`ThreadContext::yield_me`].
413 269212 : pub fn yield_me(after_ms: i64) {
414 269212 : with_thread_context(|ctx| ctx.yield_me(after_ms))
415 269212 : }
416 :
417 : /// Get current time.
418 717929 : pub fn now() -> u64 {
419 717929 : with_thread_context(|ctx| ctx.clock.get().unwrap().now())
420 717929 : }
421 :
422 432 : pub fn exit(code: i32, msg: String) {
423 432 : with_thread_context(|ctx| {
424 432 : ctx.allow_panic.store(true, Ordering::SeqCst);
425 432 : let mut result = ctx.result.lock();
426 432 : *result = (code, msg);
427 432 : panic!("exit");
428 432 : });
429 : }
430 :
431 528 : pub(crate) fn get_thread_ctx() -> Arc<ThreadContext> {
432 528 : with_thread_context(|ctx| ctx.clone())
433 528 : }
434 :
435 : /// Trait for polling channels until they have something.
436 : pub trait PollSome {
437 : /// Schedule wakeup for message arrival.
438 : fn wake_me(&self);
439 :
440 : /// Check if channel has a ready message.
441 : fn has_some(&self) -> bool;
442 : }
443 :
444 : /// Blocks current thread until one of the channels has a ready message. Returns
445 : /// index of the channel that has a message. If timeout is reached, returns None.
446 : ///
447 : /// Negative timeout means block forever. Zero timeout means check channels and return
448 : /// immediately. Positive timeout means block until timeout is reached.
449 107951 : pub fn epoll_chans(chans: &[Box<dyn PollSome>], timeout: i64) -> Option<usize> {
450 107951 : let deadline = if timeout < 0 {
451 78087 : 0
452 : } else {
453 29864 : now() + timeout as u64
454 : };
455 :
456 : loop {
457 1031664 : for chan in chans {
458 826548 : chan.wake_me()
459 : }
460 :
461 670628 : for (i, chan) in chans.iter().enumerate() {
462 670628 : if chan.has_some() {
463 86781 : return Some(i);
464 583847 : }
465 : }
466 :
467 99942 : if timeout < 0 {
468 67469 : // block until wakeup
469 67469 : yield_me(-1);
470 67469 : } else {
471 32473 : let current_time = now();
472 32473 : if current_time >= deadline {
473 2777 : return None;
474 29696 : }
475 29696 :
476 29696 : yield_me((deadline - current_time) as i64);
477 : }
478 : }
479 89558 : }
|