ridl-rt 0.6.0

The vocabulary that code generated from ridl and a runtime agree on: identity, time, the envelope, samples, payload traits, interaction descriptors, ports, and contract and transport errors.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
//! `block_on`, `noop_waker` and `flag_waker` under the `std` feature (driftsys/ridl#568).
//!
//! The whole file is gated on the feature, because `just test` builds the
//! workspace with default features and the module does not exist there. The
//! tests run through the `cargo test -p ridl-rt --all-features` line of
//! `just test`, and again under `just compat-check`.
//!
//! Every timing assertion here is a lower bound — "the wait lasted at least
//! the delay" — or a bound so wide that a loaded machine cannot cross it, so
//! the tests do not depend on scheduling latency. A poll-count ceiling is wide
//! in the same way: a parked wait polls once before the wait and once per wake
//! or timeout, so a ceiling of eight catches a wait that spins, and tolerates
//! the few extra wakes a platform may deliver. That each park is given only
//! the time left until the deadline is not pinned by any test here. The
//! spurious-wake test asserts that at least one such wake polled the future,
//! not how many did, because the thread that sends them is not scheduled at a
//! fixed cadence.
#![cfg(feature = "std")]

use std::future::{poll_fn, ready, Future};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::task::{Poll, Waker};
use std::thread;
use std::time::{Duration, Instant};

use ridl_rt::task::{block_on, flag_waker, noop_waker};

/// A future that stays pending until `done` is set, and that stores the waker
/// of its last poll so another thread can wake it. `polls` counts the polls.
struct Gate {
    done: AtomicBool,
    waker: Mutex<Option<Waker>>,
    polls: AtomicUsize,
}

impl Gate {
    fn new() -> Arc<Gate> {
        Arc::new(Gate {
            done: AtomicBool::new(false),
            waker: Mutex::new(None),
            polls: AtomicUsize::new(0),
        })
    }

    fn future(self: &Arc<Self>) -> impl Future<Output = u32> {
        let gate = Arc::clone(self);
        poll_fn(move |cx| {
            gate.polls.fetch_add(1, Ordering::SeqCst);
            if gate.done.load(Ordering::SeqCst) {
                Poll::Ready(42)
            } else {
                *gate.waker.lock().unwrap() = Some(cx.waker().clone());
                Poll::Pending
            }
        })
    }

    /// Complete the future and wake whoever polled it last.
    fn open(&self) {
        self.done.store(true, Ordering::SeqCst);
        if let Some(waker) = self.waker.lock().unwrap().take() {
            waker.wake();
        }
    }
}

#[test]
fn block_on_returns_the_output_when_another_thread_wakes_the_future() {
    let gate = Gate::new();
    let delay = Duration::from_millis(50);
    let start = Instant::now();
    let opener = {
        let gate = Arc::clone(&gate);
        thread::spawn(move || {
            thread::sleep(delay);
            gate.open();
        })
    };

    // With no deadline, a lost wake would park this thread forever. The
    // rescue thread turns that into a failure: after `limit` it completes the
    // future and unparks this thread directly, without the waker, so the
    // assertion on the elapsed time fails instead of the test never ending.
    let limit = Duration::from_secs(5);
    let stop = Arc::new(AtomicBool::new(false));
    let rescue = {
        let gate = Arc::clone(&gate);
        let stop = Arc::clone(&stop);
        let waiter = thread::current();
        thread::spawn(move || {
            while !stop.load(Ordering::SeqCst) {
                if start.elapsed() > limit {
                    gate.done.store(true, Ordering::SeqCst);
                    waiter.unpark();
                    return;
                }
                thread::sleep(Duration::from_millis(10));
            }
        })
    };

    let out = block_on(gate.future(), None);
    let elapsed = start.elapsed();
    stop.store(true, Ordering::SeqCst);
    opener.join().unwrap();
    rescue.join().unwrap();

    assert_eq!(out, Some(42));
    assert!(
        elapsed >= delay,
        "returned after {elapsed:?}, before the opener ran"
    );
    assert!(
        elapsed < limit,
        "the wake was lost: the rescue thread ended the wait at {elapsed:?}"
    );
    let polls = gate.polls.load(Ordering::SeqCst);
    assert!(
        polls >= 2,
        "the future is polled once before the wait and once after the wake"
    );
    assert!(
        polls <= 8,
        "a parked wait polls on a wake only, not in a loop: {polls} polls"
    );
}

#[test]
fn block_on_with_a_far_deadline_returns_the_output_when_woken_before_it() {
    let gate = Gate::new();
    let delay = Duration::from_millis(20);
    let far = Duration::from_secs(30);
    let start = Instant::now();
    let opener = {
        let gate = Arc::clone(&gate);
        thread::spawn(move || {
            thread::sleep(delay);
            gate.open();
        })
    };

    let out = block_on(gate.future(), Some(start + far));
    let elapsed = start.elapsed();
    opener.join().unwrap();

    assert_eq!(out, Some(42));
    assert!(elapsed >= delay);
    assert!(elapsed < far, "the wake, not the deadline, ended the wait");
}

#[test]
fn a_wake_during_a_poll_is_not_lost() {
    // The future wakes its own waker by reference inside its first poll, before
    // the thread parks, and is ready on its second poll. The unpark token has
    // to survive from the poll to the park that follows it, or the wait runs
    // to the deadline.
    let polls = Arc::new(AtomicUsize::new(0));
    let fut = {
        let polls = Arc::clone(&polls);
        poll_fn(move |cx| {
            let n = polls.fetch_add(1, Ordering::SeqCst) + 1;
            if n == 1 {
                cx.waker().wake_by_ref();
                Poll::Pending
            } else {
                Poll::Ready(n)
            }
        })
    };
    let limit = Duration::from_secs(5);

    let start = Instant::now();
    let out = block_on(fut, Some(start + limit));
    let elapsed = start.elapsed();

    assert_eq!(out, Some(2));
    assert!(
        elapsed < Duration::from_secs(1),
        "the wake from inside the poll was lost: the wait took {elapsed:?}"
    );
}

#[test]
fn a_wake_between_a_poll_and_the_park_is_not_lost_without_a_deadline() {
    // With no deadline, `block_on` waits with `thread::park`, not
    // `park_timeout`, so a lost wake parks this thread forever. The future,
    // on its first poll, hands a clone of its waker to another thread, which
    // wakes it by value; the poll joins that thread before it returns
    // `Pending`, so the wake lands after the poll and before the park that
    // follows it. The unpark token has to survive that gap. The future is
    // ready on its second poll.
    let polls = Arc::new(AtomicUsize::new(0));
    let fut = {
        let polls = Arc::clone(&polls);
        poll_fn(move |cx| {
            let n = polls.fetch_add(1, Ordering::SeqCst) + 1;
            if n == 1 {
                let waker = cx.waker().clone();
                thread::spawn(move || waker.wake()).join().unwrap();
                Poll::Pending
            } else {
                Poll::Ready(n)
            }
        })
    };

    // The rescue thread turns a lost wake into a failure: after `limit` it
    // unparks this thread directly, without the waker, so the assertion on
    // the elapsed time fails instead of the test never ending.
    let limit = Duration::from_secs(5);
    let stop = Arc::new(AtomicBool::new(false));
    let start = Instant::now();
    let rescue = {
        let stop = Arc::clone(&stop);
        let waiter = thread::current();
        thread::spawn(move || {
            while !stop.load(Ordering::SeqCst) {
                if start.elapsed() > limit {
                    waiter.unpark();
                    return;
                }
                thread::sleep(Duration::from_millis(10));
            }
        })
    };

    let out = block_on(fut, None);
    let elapsed = start.elapsed();
    stop.store(true, Ordering::SeqCst);
    rescue.join().unwrap();

    assert_eq!(out, Some(2));
    assert!(
        elapsed < limit,
        "the wake was lost: the rescue thread ended the wait at {elapsed:?}"
    );
}

#[test]
fn block_on_returns_none_when_the_deadline_passes() {
    let gate = Gate::new();
    let bound = Duration::from_millis(50);

    let start = Instant::now();
    let out = block_on(gate.future(), Some(start + bound));
    let elapsed = start.elapsed();

    assert_eq!(out, None);
    assert!(
        elapsed >= bound,
        "returned after {elapsed:?}, before the deadline"
    );
    assert!(
        elapsed < bound + Duration::from_millis(500),
        "the wait ended long after the deadline: {elapsed:?}"
    );
    let polls = gate.polls.load(Ordering::SeqCst);
    assert!(
        polls >= 2,
        "the future is polled once before the wait and once at the deadline"
    );
    assert!(
        polls <= 8,
        "a parked wait polls on a wake or the deadline only: {polls} polls"
    );
}

#[test]
fn block_on_polls_once_even_when_the_deadline_has_already_passed() {
    let past = Instant::now() - Duration::from_secs(1);
    assert_eq!(block_on(ready(7u8), Some(past)), Some(7));

    let gate = Gate::new();
    assert_eq!(block_on(gate.future(), Some(past)), None);
    assert_eq!(gate.polls.load(Ordering::SeqCst), 1);
}

#[test]
fn block_on_returns_the_output_of_a_ready_poll_after_the_deadline() {
    // The future is pending while the deadline is ahead and ready once it has
    // passed, and it never wakes its waker. So the first poll is `Pending`,
    // the park runs to the deadline, and the poll after that park, which runs
    // when `Instant::now() >= deadline`, is `Ready`. Its output is returned
    // as `Some`: the deadline turns a `Pending` poll into `None`, not a
    // `Ready` one.
    let bound = Duration::from_millis(100);
    let start = Instant::now();
    let deadline = start + bound;
    let polls = Arc::new(AtomicUsize::new(0));
    let fut = {
        let polls = Arc::clone(&polls);
        poll_fn(move |_cx| {
            polls.fetch_add(1, Ordering::SeqCst);
            if Instant::now() >= deadline {
                Poll::Ready(42)
            } else {
                Poll::Pending
            }
        })
    };

    let out = block_on(fut, Some(deadline));
    let elapsed = start.elapsed();

    assert_eq!(
        out,
        Some(42),
        "a poll that is ready after the deadline returns its output"
    );
    assert!(
        elapsed >= bound,
        "returned after {elapsed:?}, before the deadline"
    );
    assert!(
        polls.load(Ordering::SeqCst) >= 2,
        "the ready poll is the one after the park, not the first"
    );
}

#[test]
fn a_spurious_unpark_does_not_end_the_wait_early() {
    let gate = Gate::new();
    let bound = Duration::from_millis(100);
    let stop = Arc::new(AtomicBool::new(false));
    let waiter = thread::current();
    let nagger = {
        let stop = Arc::clone(&stop);
        thread::spawn(move || {
            while !stop.load(Ordering::SeqCst) {
                waiter.unpark();
                thread::sleep(Duration::from_millis(5));
            }
        })
    };

    let start = Instant::now();
    let out = block_on(gate.future(), Some(start + bound));
    let elapsed = start.elapsed();
    stop.store(true, Ordering::SeqCst);
    nagger.join().unwrap();

    assert_eq!(out, None);
    assert!(
        elapsed >= bound,
        "a spurious unpark ended the wait at {elapsed:?}"
    );
    assert!(
        gate.polls.load(Ordering::SeqCst) >= 3,
        "a spurious unpark polls the future once, and the wait continues"
    );
}

#[test]
fn noop_waker_can_be_cloned_and_woken_without_effect() {
    let waker = noop_waker();
    let clone = waker.clone();
    assert!(waker.will_wake(&clone));
    assert!(
        !waker.will_wake(&noop_waker()),
        "two calls give two wakers, which is why one is created per loop"
    );

    waker.wake_by_ref();
    clone.wake_by_ref();
    clone.wake();

    // A future polled with it makes progress on the caller's polls only: the
    // waker does not poll, and waking it is not observable from the future.
    let gate = Gate::new();
    let mut fut = std::pin::pin!(gate.future());
    let mut cx = std::task::Context::from_waker(&waker);
    assert_eq!(fut.as_mut().poll(&mut cx), Poll::Pending);
    gate.open();
    assert_eq!(fut.as_mut().poll(&mut cx), Poll::Ready(42));
}

#[test]
fn waking_a_noop_waker_from_another_thread_does_not_poll_the_future() {
    // A thread wakes a no-op waker every 5 ms while this thread waits on a
    // pending future under a 100 ms deadline. A waker that unparked this
    // thread would add a poll per wake; a no-op one adds none, so the wait
    // polls only before the park and at the deadline.
    let gate = Gate::new();
    let bound = Duration::from_millis(100);
    let stop = Arc::new(AtomicBool::new(false));
    let waker = noop_waker();
    let nagger = {
        let stop = Arc::clone(&stop);
        let waker = waker.clone();
        thread::spawn(move || {
            while !stop.load(Ordering::SeqCst) {
                waker.wake_by_ref();
                thread::sleep(Duration::from_millis(5));
            }
        })
    };

    let start = Instant::now();
    let out = block_on(gate.future(), Some(start + bound));
    stop.store(true, Ordering::SeqCst);
    nagger.join().unwrap();

    assert_eq!(out, None);
    let polls = gate.polls.load(Ordering::SeqCst);
    assert!(
        (2..=8).contains(&polls),
        "a no-op wake unparked the waiting thread: {polls} polls in {bound:?}"
    );
}

#[test]
fn a_flag_waker_sets_its_flag_on_each_kind_of_wake_and_take_clears_it() {
    let (waker, woken) = flag_waker();
    assert!(!woken.take(), "the flag starts clear");

    waker.wake_by_ref();
    assert!(woken.take(), "a wake by reference sets the flag");
    assert!(!woken.take(), "take clears the flag");

    let clone = waker.clone();
    clone.wake();
    assert!(woken.take(), "a wake of a clone, by value, sets the flag");

    let remote = waker.clone();
    thread::spawn(move || remote.wake_by_ref()).join().unwrap();
    assert!(woken.take(), "a wake from another thread sets the flag");
    assert!(!woken.take());

    let (other, other_woken) = flag_waker();
    other.wake_by_ref();
    assert!(!woken.take(), "each call gives its own flag");
    assert!(other_woken.take());
}

#[test]
fn a_frame_loop_over_a_flag_waker_polls_again_while_the_future_wakes_itself() {
    // A future that wakes itself on each of its first three polls and is ready
    // on the fourth. Under a flag waker one frame polls it four times; under a
    // no-op waker one poll is all a frame gets.
    let polls = AtomicUsize::new(0);
    let fut = poll_fn(|cx| {
        if polls.fetch_add(1, Ordering::SeqCst) < 3 {
            cx.waker().wake_by_ref();
            Poll::Pending
        } else {
            Poll::Ready(())
        }
    });
    let mut fut = std::pin::pin!(fut);
    let (waker, woken) = flag_waker();
    let mut cx = std::task::Context::from_waker(&waker);
    let mut ready = false;
    for _ in 0..8 {
        if fut.as_mut().poll(&mut cx).is_ready() {
            ready = true;
            break;
        }
        if !woken.take() {
            break;
        }
    }
    assert!(
        ready,
        "the frame polled until the future stopped waking itself"
    );
    assert_eq!(polls.load(Ordering::SeqCst), 4);
}