rust / expert
Snippet
Implementierung eines threadsicheren asynchronen Timer-Futures von Grund auf
Das Schreiben benutzerdefinierter Futures erfordert die manuelle Implementierung des Future-Traits und die Verwaltung von Task-Wakern. Dieser asynchrone Timer startet einen Worker-Thread und aktualisiert den gemeinsamen Zustand, wodurch die Aufgabe des Async-Runtimes nach Ablauf der Zeit geweckt wird.
snippet.rs
rust
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
use std::future::Future;use std::pin::Pin;use std::sync::{Arc, Mutex};use std::task::{Context, Poll, Waker};use std::thread;use std::time::Duration;struct AsyncTimer {state: Arc<Mutex<TimerState>>,}struct TimerState {completed: bool,waker: Option<Waker>,}impl Future for AsyncTimer {type Output = ();fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {let mut state = self.state.lock().unwrap();if state.completed {Poll::Ready(())} else {state.waker = Some(cx.waker().clone());Poll::Pending}}}impl AsyncTimer {fn new(duration: Duration) -> Self {let state = Arc::new(Mutex::new(TimerState { completed: false, waker: None }));let thread_state = state.clone();thread::spawn(move || {thread::sleep(duration);let mut state = thread_state.lock().unwrap();state.completed = true;if let Some(waker) = state.waker.take() {waker.wake();}});AsyncTimer { state }}}
Erklärung
1
state.waker = Some(cx.waker().clone());
Speichert den aktuellen Waker-Kontext des Tasks, damit der Hintergrundthread den Executor informieren kann.
2
Poll::Pending
Gibt Pending zurück, um dem Runtime zu signalisieren, dass die Operation noch nicht bereit ist.
3
waker.wake();
Löst den Waker aus, wodurch der diesem Future zugeordnete Task für einen erneuten Poll-Aufruf eingeplant wird.