Kapitel 1: Das mentale Modell der Asynchronität: Future, Waker und Executor als Trio
Asynchrone Programmierung ist in Rust keine Bibliothek, sondern ein sprachweites Protokoll. Dass Tokio zu einer produktionsreifen Laufzeitumgebung werden konnte, liegt nicht daran, dass es Future erfunden hat, sondern daran, dass es die Randbedingungen jedes einzelnen Vertrags dieses Protokolls präzise implementiert. Dieses Kapitel springt nicht vorschnell in den Scheduler-Code von Tokio, sondern erklärt zunächst gründlich die Verantwortungsgrenzen und den umgekehrten Kontrollfluss des „Trio" — Future, Waker, Executor. Erst wenn man versteht, wie diese drei ineinandergreifen, finden die Zusammensetzung der Runtime, das Work-Stealing-Scheduling und der I/O-Treiber in den folgenden Kapiteln einen Anknüpfungspunkt.
1.1 Vom Blockieren zum Polling: Warum Rust poll statt Callbacks wählt
Intuitives Modell
Stellen Sie sich vor, Sie bestellen in einem Restaurant ein Gericht, das frisch zubereitet werden muss. Callback-basierte Asynchronität (wie der frühe Node.js-Stil) entspricht dem Hinterlassen Ihrer Telefonnummer, und der Koch ruft Sieaktiv an— die Kontrolle liegt beim Koch, Ihr Code reagiert nur passiv. Polling-basierte Asynchronität (Rusts Wahl) entspricht dem Erhalt eines Abholscheins, Sieentscheiden selbstwann Sie zum Fenster gehen und fragen „Ist es fertig?": Wenn nicht, machen Sie etwas anderes, wenn ja, holen Sie es ab.
Dieser Unterschied erscheint geringfügig, bestimmt aber die Form des gesamten Systems. Im Callback-Modell muss jede asynchrone Operation eine Closure mitführen, die angibt, „was nach Abschluss zu tun ist". Closures verschachteln sich Schicht für Schicht und bilden die Callback-Hölle, und das Abbrechen von Operationen ist extrem schwierig — Sie können einen bereits registrierten Callback nicht „zurückziehen". Im Polling-Modell ist ein Future nur ein Zustandsautomat,pollist eine reine Abfrageaktion, ohne Voranschreiten werden keine Ressourcen verbraucht, Abbrechen ist einfach drop, sauber und ordentlich.
Der Kernvertrag des Polling-Modells
Das von der Rust-Standardbibliothek definierteFutureTrait hat nur zwei Elemente: einepollMethode, einenOutputassoziierten Typ. Tokio definiert dieses Trait nicht neu, sondern verwendet direkt die Implementierung der Standardbibliothek wieder. Dies zeigt sich deutlich im Quellcode:
// tokio/src/future/mod.rs
cfg_not_trace! {
cfg_rt! {
pub(crate) use std::future::Future;
}
}📎 tokio/src/future/mod.rs:24-28
Dieser Code offenbart eine wichtige Tatsache: Wenn dastracingFeature nicht aktiviert ist, ist das interneFuturevon Tokio ein Alias fürstd::future::Futureohne jegliche Umhüllung. Nur wenntracingaktiviert ist, wird es durchInstrumentedFutureersetzt:
cfg_trace! {
mod trace;
#[allow(unused_imports)]
pub(crate) use trace::InstrumentedFuture as Future;
}📎 tokio/src/future/mod.rs:18-22
Dieses Design „standardmäßig zero-cost, Instrumentierung nach Bedarf" ist Tokios beständige Philosophie: Der Kernpfad führt keine zusätzliche Abstraktionsebene ein, Beobachtbarkeit wird als optionales Feature hinzugefügt.InstrumentedFutureDie Existenz von zeigt, dass das Tokio-Team der Ansicht ist, dass die Instrumentierungskosten von tracing nicht von allen Nutzern getragen werden sollten.
Drei implizite Einschränkungen des poll-Vertrags
pollDie Signatur derfn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>Methode ist
. In dieser Signatur verbergen sich drei Verträge; die Verletzung eines jeden führt zu undefiniertem Verhalten oder logischen Fehlern: Pin<&mut Self>Vertrag eins: Pin garantiert Selbstreferenz-Sicherheit.
bedeutet, dass ein Future, sobald es gepollt wurde, seine Speicheradresse nicht mehr ändern darf. Der Grund ist, dass ein async-Block nach der Kompilierung einen Zustandsautomaten mit Selbstreferenzen erzeugt — lokale Variablen können Referenzen auf andere Felder innerhalb desselben Zustandsautomaten halten. Wenn eine Verschiebung erlaubt wäre, würden diese Referenzen baumeln.Vertrag zwei: Pending muss bereits ein Aufwecken registriert haben.pollWennPoll::Pendingzurückgibtcx.waker()Holt und speichert den Waker oder hat den Waker bereits bei einer Ereignisquelle registriert. Andernfalls wird der Executor niemals erfahren, wann dieses Future erneut gepollt werden kann, was dazu führt, dass die Aufgabe dauerhaft hängen bleibt.
Vertrag drei: Nach Ready sollte nicht erneut gepollt werden.SobaldpollzurückgibtPoll::Readyist es ein logischer Fehler, dasselbe Future erneut zu pollen (obwohl dies kein UB verursacht, ist das Verhalten undefiniert). Der Executor ist dafür verantwortlich, die Aufgabe nach Erhalt von Ready nicht erneut zu planen.
Von diesen drei Verträgen ist Vertrag zwei die fehleranfälligste Stelle und auch der grundlegende Grund für die Existenz des Wakers.
1.2 Waker: Der Träger des umgekehrten Kontrollflusses
Intuitives Modell
Der Waker ist der „Vibrations-Piepser", den dir das Restaurant gibt. Du musst nicht ständig am Fenster stehen und fragen „Ist es fertig?" – das würde nur deine Zeit verschwenden. Du musst nur beim ersten Gang zum Fenster den Piepser dem Koch geben (Waker registrieren) und dann beruhigt andere Dinge tun. Wenn das Essen fertig ist, drückt der Koch den Knopf, der Piepser vibriert (ruftwakeauf), du erhältst das Signal und gehst dann zum Fenster, um das Essen abzuholen (erneut pollen).
Ohne den Waker hätte der Executor nur zwei Möglichkeiten: entweder alle Aufgaben beschäftigt zu pollen (verschwendet CPU) oder Aufgaben, die bereits Pending zurückgegeben haben, niemals zu pollen (Aufgaben verhungern). Der Waker ist der einzige Mechanismus, um dieses Patt zu durchbrechen.
Speicherlayout und Vtable-Design des Wakers
Der Waker ist ein Standardbibliothekstyp, aber sein Design hat direkt die Aufgabenstruktur von Tokio beeinflusst.Wakerist im Wesentlichen ein fetter Zeiger: eineRawWakerStruktur, die einen Datenzeiger und einen Vtable-Zeiger enthält.
// 标准库中的定义(非 Tokio 源码,此处为背景说明)
pub struct RawWaker {
data: *const (),
vtable: &'static RawWakerVTable,
}
pub struct RawWakerVTable {
clone: unsafe fn(*const ()) -> RawWaker,
wake: unsafe fn(*const ()),
wake_by_ref: unsafe fn(*const ()),
drop: unsafe fn(*const ()),
}Das Raffinierte an diesem Design ist:Wakerselbst kümmert sich nicht darum, was „Aufwecken" konkret bedeutet. Es ist nur ein Träger für vier Funktionszeiger. Tokio kann einen Waker bereitstellen, dessenwakeFunktion die Aufgabe erneut in die Scheduling-Warteschlange einreiht; während eine andere Laufzeitumgebung (zum BeispielfuturesCrateblock_on) eine völlig andere Waker-Implementierung bereitstellen kann. Dieses „Daten + Vtable"-Muster ermöglicht es, den Waker zwischen verschiedenen Laufzeitumgebungen zu übergeben, ohne die Semantik zu verlieren.
wakeundwake_by_refDer Unterschied ist entscheidend:wakekonsumiert die Eigentümerschaft des Wakers (nach dem Aufruf wird der Waker gedroppt), währendwake_by_refnur ausleiht. Der Executor implementiert üblicherweisewake_by_refals „Aufgabe als bereit markieren und in die Warteschlange einreihen", währendwakezusätzlich die Dekrementierung des Referenzzählers behandelt. In der Aufgabenstruktur von Tokio zeigt der Datenzeiger des Wakers auf den Referenzzählerkopf der Aufgabe; jedes Klonen erhöht den Zähler, jedes Droppen verringert ihn, und wenn der Zähler null erreicht, wird der Aufgabenspeicher freigegeben.
Der vollständige Ablauf des Aufweckens
Das folgende Sequenzdiagramm zeigt die vollständige Kette einer TCP-Leseoperation von der Initiierung bis zum Aufwecken. Beachten Sie, wie der Waker vom Aufgabenkontext bis zum I/O-Treiber weitergegeben wird:
sequenceDiagram
participant App as 应用任务
participant Exec as 调度器 Worker
participant Future as TcpStream::read Future
participant Reactor as I/O 驱动 (epoll)
participant Kernel as 操作系统内核
App->>Future: poll(cx) 携带 Waker
Future->>Reactor: 注册可读兴趣 + 保存 Waker
Reactor->>Kernel: epoll_ctl(ADD, fd, EPOLLIN)
Future-->>Exec: 返回 Poll::Pending
Note over Exec: 任务挂起,Worker 去执行其他任务
Kernel-->>Reactor: epoll_wait 返回 fd 就绪
Reactor->>Reactor: 查找 fd 对应的 Waker
Reactor->>Exec: waker.wake_by_ref()
Note over Exec: 任务重新入队
Exec->>Future: 再次 poll(cx)
Future->>Kernel: read(fd, buf) 非阻塞读取
Kernel-->>Future: 返回数据
Future-->>App: 返回 Poll::Ready(n)Der Schlüssel in diesem Diagramm ist:Der Waker ist der einzige Kanal, der vom Reactor zurück zum Executor gelangen kann. Der Reactor besitzt keine anderen Informationen über die Aufgabe; er weiß nur „wenn dieser fd bereit ist, rufe diesen Waker auf". Diese Entkopplung ermöglicht es, den I/O-Treiber unabhängig vom Scheduler zu implementieren; beide kommunizieren nur über die schmale Schnittstelle des Wakers.
Falsches Aufwecken: Die Grauzone des Vertrags
Die Dokumentation von Tokio erkennt ausdrücklich die Existenz von falschem Aufwecken an:
Normally, tasks are scheduled only if they have been woken by calling wake on their waker. However, this is not guaranteed, and Tokio may schedule tasks that have not been woken under some circumstances.
📎 tokio/src/runtime/mod.rs:306-309
Das bedeutet, dasspollDie Implementierung muss tolerieren können, dass „ohne aufgeweckt zu werden erneut gepollt wird". Ein korrektes Future sollte nach der Rückgabe von Pending, selbst wenn kein Ereignis eingetreten ist, bei erneutem Polling wieder Pending zurückgeben und nicht panicen oder fehlerhafte Ergebnisse liefern. Diese Einschränkung scheint locker zu sein, stellt aber tatsächlich Anforderungen an das Design des Zustandsautomaten: Man darf nicht annehmen, dass „zwischen zwei Polls unbedingt ein Ereignis stattfindet".
1.3 Executor: Von Future zur Aufgabenkapselung
Intuitives Modell
Der Executor ist der Disponent des Restaurants. Er hat einen Stapel Bestellungen (Aufgabenwarteschlange) und entscheidet, welche Bestellung zuerst bearbeitet wird und wer sie bearbeitet. Wenn der Piepser vibriert, reiht er die entsprechende Bestellung wieder in die Warteschlange ein. Ohne Disponent wüssten die Köche nicht, welches Gericht sie zubereiten sollen, noch wann sie die Arbeit wechseln sollen.
Doch die Aufgaben des Executors gehen weit über „Future pollen" hinaus. Er muss drei Kernprobleme lösen:Lebenszyklusverwaltung von Aufgaben(Erstellen, Planen, Abschließen, Abbrechen),Fairness-Garantie(verhindern, dass eine Aufgabe andere aushungert),Ressourcentreiber-Integration(wie I/O- und Timer-Ereignisse in Aufwecken umgewandelt werden).
Speicherlayout der Aufgabe: Vom Future zur Task
Beim Aufruf vontokio::spawnwird das übergebene Future nicht direkt in die Warteschlange gestellt. Es wird in eineTaskStruktur verpackt, die einen Referenzzählerkopf, Scheduling-Metadaten und das Future selbst enthält. Dieser Verpackungsprozess hat eine entscheidende Optimierungsentscheidung:
/// Boundary value to prevent stack overflow caused by a large-sized
/// Future being placed in the stack.
pub(crate) const BOX_FUTURE_THRESHOLD: usize = if cfg!(debug_assertions) {
2048
} else {
16384
};
pub(crate) struct AutoBox<T>(std::marker::PhantomData<T>);
impl<T> AutoBox<T> {
/// `true` if a value of type `T` is larger than [`BOX_FUTURE_THRESHOLD`].
pub(crate) const SHOULD_BOX: bool = std::mem::size_of::<T>() > BOX_FUTURE_THRESHOLD;
}📎 tokio/src/runtime/mod.rs:649-673
Dieser Code löst ein sehr konkretes Problem: Wenn das Future zu groß ist (über 16KB, im Debug-Modus 2KB), würde das direkte Inlining in die Task-Struktur zu Stack-Überlauf oder Speicherverschwendung führen.AutoBoxentscheidet durch die Kompilierzeit-KonstanteSHOULD_BOX, ob das Future geboxt wird.
In den Kommentaren wird besonders betont, „assozierte Konstanten statt Laufzeit-ifDer Grund: Wenn zur Laufzeit entschieden wird, instanziiert der Compiler für jedenTgleichzeitig den Code beider Zweige (einen fürT, einen fürPin<Box<T>>), was zu Code-Aufblähung führt. Bei konstanten Zweigen schneidet der Monomorphisierungs-Sammler die unerreichbaren Zweige weg und generiert Code nur für die tatsächlich verwendeten Typen. Dies ist eine typische Optimierung, bei der das Typsystem Laufzeitentscheidungen ersetzt.
Fairness der Planung: Die magischen Zahlen 31 und 61
In der Scheduler-Dokumentation von Tokio ist eine formale Fairness-Garantie definiert:
If the total number of tasks does not grow without bound, and no task is blocking the thread, then it is guaranteed that tasks are scheduled fairly.
📎 tokio/src/runtime/mod.rs:279-281
Die Implementierung dieser Garantie hängt von zwei Schlüsselparametern ab. Für die current-thread-Laufzeit:
The runtime will prefer to choose the next task to schedule from the local queue, and will only pick a task from the global queue if the local queue is empty, or if it has picked a task from the local queue 31 times in a row.
📎 tokio/src/runtime/mod.rs:328-333
The runtime will check for new IO or timer events whenever there are no tasks ready to be scheduled, or when it has scheduled 61 tasks in a row.
📎 tokio/src/runtime/mod.rs:335-337
Diese beiden Zahlen (31 und 61) sind nicht willkürlich gewählt. 31 ist 2 hoch 5 minus 1 und kann schnell mit Bitoperationen geprüft werden; 61 dient dazu, sicherzustellen, dass I/O-Ereignisse nicht unbegrenzt verzögert werden – selbst wenn die Aufgabenwarteschlange niemals leer wird, muss nach jeweils 61 Planungen einmal I/O geprüft werden.
Warum 31 und nicht 32? Weil der Zähler bei 0 beginnt, bei jeder Planung um 1 erhöht wird und bei Erreichen von 31 die Prüfung der globalen Warteschlange auslöst. Die Prüfung mitcounter & 31 == 31ist effizienter als mitcounter % 32 == 0(obwohl moderne Compiler dies automatisch optimieren). Die Wahl von 61 ist subtiler: Sie muss groß genug sein, um den Overhead häufiger epoll_wait-Systemaufrufe zu vermeiden, und klein genug, um die I/O-Latenz in einem akzeptablen Bereich zu halten.
LIFO-Slot-Optimierung der Multithread-Laufzeit
Die Multithread-Laufzeit fügt der Fairness noch eine Leistungsoptimierung hinzu – den LIFO-Slot:
The multi thread runtime uses the lifo slot optimization: Whenever a task wakes up another task, the other task is added to the worker thread's lifo slot instead of being added to a queue.
📎 tokio/src/runtime/mod.rs:373-377
Die Intuition hinter dieser Optimierung ist: Wenn eine Aufgabe eine andere aufweckt, hat die aufgeweckte Aufgabe wahrscheinlich eine Datenabhängigkeit mit der aktuellen Aufgabe (z. B. im Producer-Consumer-Muster). Indem man sie in den LIFO-Slot legt, kann die aktuelle Aufgabe sie sofort nach Abschluss ausführen und die heißen Daten im CPU-Cache nutzen.
Aber der LIFO-Slot hat einen Missbrauchsschutz:
if a worker thread uses the lifo slot three times in a row, it is temporarily disabled until the worker thread has scheduled a task that didn't come from the lifo slot.
📎 tokio/src/runtime/mod.rs:380-382
Die Regel „nach drei aufeinanderfolgenden Verwendungen deaktivieren“ dient dazu, zu verhindern, dass zwei Aufgaben sich gegenseitig aufwecken und eine Livelock bilden. Wenn Aufgabe A Aufgabe B aufweckt und B wiederum A, würde der LIFO-Slot ohne diese Einschränkung dauerhaft von diesen beiden Aufgaben belegt, und andere Aufgaben kämen nie zum Zug. Die Dreifach-Begrenzung gibt anderen Aufgaben eine Chance, sich einzufügen.
Aufgabenabbruch: Die wahre Semantik von abort
JoinHandle::abortDas Verhalten von
Be aware that calls to JoinHandle::abort just schedule the task for cancellation, and will return before the cancellation has completed.
📎 tokio/src/task/mod.rs:146-148
wird oft missverstanden. Die Dokumentation stellt ausdrücklich klar:abortDas bedeutet,.awaitist nicht synchron. Es setzt lediglich ein Flag, und die Aufgabe prüft dieses Flag am nächsten.await-Punkt und beendet sich selbst. Wenn die Aufgabe gerade einen CPU-intensiven Code ohneabortausführt, wird
nicht sofort wirksam.
Note that aborting a task does not guarantee that it fails with a cancelled error, since it may complete normally first.
📎 tokio/src/task/mod.rs:134-138
〔Design-Schlussfolgerung und Architektur-Abwägung〕spawn_blockingDie Design-Motivation dieser Semantik ist: Abbruch ist eine „Best-Effort“-Operation. Tokio erzwingt nicht das Töten von Aufgaben (Rust hat keinen sicheren Mechanismus zur erzwungenen Beendigung), sondern fordert Aufgaben kooperativ auf, sich selbst zu beenden. Dies steht im Einklang mit dem Design, dass.await-Aufgaben nicht abbrechbar sind – blockierende Aufgaben haben keine
-Punkte und können das Abbruch-Flag nicht prüfen.
1.4 Design-Überlegungen: Grenzen und Kosten des Trios
Warum Future keinen Executor enthältFutureDas
-Trait von Rust enthält bewusst keine Information darüber, „wie man sich selbst plant“. Dies ist eine wohlüberlegte Entkopplungsentscheidung. Wenn ein Future seinen Executor kennen würde, dann:
1. Könnte dasselbe Future nicht auf verschiedenen Laufzeiten ausgeführt werden (z. B. Migration von Tokio zu async-std)block_on2. Könnte man es beim Testen nicht mit einem einfachen
antreibenselect!、join!3. Könnten Kombinatoren (wie
) nicht laufzeitübergreifend funktionieren
Der Waker existiert genau deshalb, um diese Entkopplung beizubehalten und dem Future dennoch zu erlauben, den Executor zu benachrichtigen. Der Waker ist ein „Fähigkeits-Token“ – das Future weiß nur „Ich kann dies aufrufen, um eine Neuplanung anzufordern“, aber nicht, wie die Planung konkret abläuft.
Die Kosten der kooperativen Planung.awaitTokios Aufgaben sind kooperativ: Eine Aufgabe gibt die Ausführung nur an
code that spends a long time without reaching an .await will prevent other tasks from running.
📎 tokio/src/lib.rs:178-179
〔Design-Schlussfolgerung und Architektur-Abwägung〕.awaitDies sind die grundlegenden Kosten der kooperativen Planung. Das Betriebssystem kann einen Thread an jeder Instruktionsgrenze unterbrechen, aber Tokio kann Aufgaben nur an.await-Punkten wechseln. Wenn eine Aufgabe eine 10-sekündige CPU-intensive Schleife ohnespawn_blockingdazwischen ausführt, werden alle anderen Aufgaben auf demselben Worker-Thread 10 Sekunden lang blockiert. Tokios Gegenstrategie ist die Bereitstellung vonblock_in_placeund
, um solche Arbeiten in einen dedizierten Thread-Pool zu verlagern. Aber das liegt in der Verantwortung des Nutzers; die Laufzeit kann dies nicht automatisch erkennen.
Randbedingungen der Fairness-Garantie
- Tokios Fairness-Garantie hat zwei Voraussetzungen: Die Gesamtzahl der Aufgaben ist nach oben beschränkt, und keine Aufgabe blockiert den Thread. Diese beiden Bedingungen werden in der Praxis häufig verletzt:
- Wenn Aufgaben ständig neue Aufgaben spawnen und nicht zurückgewonnen werden, ist die Gesamtzahl unbegrenzt und die Fairness-Garantie ungültig
〔Design-Schlussfolgerung und Architektur-Abwägung〕
Deshalb betont die Tokio-Dokumentation wiederholt: „Führen Sie keine blockierenden Operationen in asynchronen Aufgaben aus.“ Die Fairness-Garantie ist keine harte Garantie der Laufzeit, sondern eine Garantie „unter der Voraussetzung korrekter Nutzung“. Die Laufzeit erkennt Verstöße nicht, da die Erkennung selbst Overhead verursachen würde.
1.5 Zusammenfassung dieses Kapitels
Future ist eine zustandsmaschine im Pull-Stil. pollist eine reine Abfrageaktion und gibt zurückPendingmuss bereits ein Waker registriert sein, wenn zurückgegeben wirdReadydarf danach nicht erneut gepollt werden. Tokio verwendet direkt wiederstd::future::Future, ohne zusätzliches Wrapping (außer wenn tracing aktiviert ist).
Waker ist der einzige Kanal für umgekehrte Kontrollflüsse.Es erreicht Laufzeitunabhängigkeit durch das Design aus „Datenzeiger + Vtable“.wakekonsumiert Ownership,wake_by_refleiht nur aus. Falsche Weckrufe sind erlaubt, Future muss sie tolerieren.
Executor ist verantwortlich für Lebenszyklus, Fairness und Ressourcenintegration.Es verpackt Future als Task, entscheidet überAutoBoxzur Kompilierzeit, ob geboxt wird, balanciert über die beiden magischen Zahlen 31/61 die Planung zwischen lokaler und globaler Queue und optimiert über LIFO-Slots die Leistung in Szenarien mit Datenabhängigkeiten.
Diese drei Komponenten sind über schmale Schnittstellen entkoppelt: Future kennt nurpoll, Waker kennt nurwake, Executor kennt nur „poll bis Pending oder Ready“. Genau diese Entkopplung ermöglicht es Tokio, erweiterte Funktionen wie Work-Stealing-Scheduling, I/O-Treiber-Integration und kooperative Budgets zu implementieren, ohne die Definition von Future zu ändern.
Gedanken und Selbsttest dieses Kapitels
Q1: Wenn manAutoBox::SHOULD_BOXvon einer Kompilierzeit-Konstante in eine Laufzeit-if size_of::<T>() > THRESHOLDändert, welche Auswirkungen hätte das auf das kompilierte Artefakt? Warum betont Tokios Kommentar diesen Punkt besonders?
Referenzanalyse: Laut den Kommentaren zu📎 tokio/src/runtime/mod.rs:657-667, wenn Laufzeit-ifverwendet wird, instanziiert der Compiler für jedesTgleichzeitig den Code beider Zweige – einen für den Fall, dassTdirekt inlined, und einen für den FallPin<Box<T>>. Das bedeutet, dass für jeden gespawnten Future-Typ zwei Kopien des Task-Treibers (task harness) erzeugt werden, was die Binärgröße verdoppelt. Mit der assoziierten KonstanteSHOULD_BOXhingegen, da sie nach Bestimmung vonTeine Kompilierzeit-Konstante ist, entfernt der Monomorphisierungs-Sammler unerreichbare Zweige und erzeugt Code nur für den tatsächlich verwendeten Pfad. Dies ist eine typische Optimierung, bei der „das Typsystem Laufzeitentscheidungen ersetzt“, zum Preis, dassAutoBoxeine generische Struktur statt einer gewöhnlichen Funktion sein muss.
Q2: Angenommen, eine Task gibt inpollzurückPending, vergisst aber, einen Waker zu registrieren. Was passiert mit dieser Task jeweils in der current-thread-Laufzeit und der multi-thread-Laufzeit? Hat Tokio einen Mechanismus, um diesen Fall zu erkennen?
Referenzanalyse: Laut📎 tokio/src/runtime/mod.rs:306-309erlaubt Tokio falsche Weckrufe, was bedeutet, dass eine Task ohne Weckruf erneut geplant werden kann. Das heißt aber nicht, dass es sicher ist, die Waker-Registrierung zu vergessen. In der current-thread-Laufzeit geht die Laufzeit, wenn sowohl lokale als auch globale Queue leer sind, in denparkZustand über und wartet auf I/O- oder Timer-Ereignisse. Eine Task, die vergessen hat, einen Waker zu registrieren, wird niemals erneut in die Queue eingereiht und hängt dauerhaft. In der multi-thread-Laufzeit ist die Situation ähnlich, aber wenn andere Tasks kontinuierlich Weckrufe auslösen, kann diese Task durch falsche Weckrufe zufällig erneut geplant werden – worauf man sich jedoch nicht verlassen kann. Tokio hat keinen Laufzeit-Erkennungsmechanismus, um den Fall „gibt Pending zurück, aber kein Waker registriert“ zu entdecken, da dies nach jedem Poll prüfen müsste, ob der Waker verwendet wurde, was zu teuer wäre. Das liegt in der Verantwortung des Future-Implementierers.
Q3: Welches konkrete Szenario soll die Regel „LIFO-Slot nach drei aufeinanderfolgenden Verwendungen deaktivieren“ verhindern? Wenn man diese Einschränkung entfernt, bei welchem Task-Abhängigkeitsmuster würden andere Tasks verhungern?
Referenzanalyse: Laut📎 tokio/src/runtime/mod.rs:380-382wird der LIFO-Slot nach drei aufeinanderfolgenden Verwendungen vorübergehend deaktiviert, bis eine Task aus einer Nicht-LIFO-Quelle geplant wurde. Das Szenario, das diese Regel verhindert, ist: Zwei Tasks wecken sich gegenseitig und bilden eine enge Schleife. Zum Beispiel weckt Task A nach der Verarbeitung eines Datenstapels Task B, und Task B weckt sofort nach der Verarbeitung Task A. Ohne die Drei-Beschränkung würden A und B dauerhaft den LIFO-Slot besetzen, der Worker-Thread würde endlos zwischen diesen beiden Tasks wechseln, und andere Tasks in der lokalen und globalen Queue bekämen nie eine Ausführungschance. Die Drei-Beschränkung stellt sicher, dass nach jeweils drei Runden „gegenseitigen Weckens“ mindestens eine andere Task geplant wird, wodurch ein Livelock durchbrochen wird. Die Wahl dieser Zahl ist empirisch: Zu klein verringert den Nutzen der LIFO-Optimierung, zu groß erhöht die Latenz anderer Tasks.
Damit sind die Verantwortungsgrenzen und das Zusammenspiel von Future, Waker und Executor klar: Future definiert die Berechnung, Waker ist für das Wecken verantwortlich, Executor treibt die Ausführung an. Aber eine einzelne Komponente kann nicht unabhängig arbeiten; sie müssen in eine einheitliche Laufzeitumgebung zusammengesetzt werden. Im nächsten Kapitel verfolgen wir die vollständige Montagekette von Runtime::new und Builder::build, sehen, wie Scheduler, I/O-Treiber, Zeit-Treiber und Blocking-Thread-Pool in dieselbe Runtime-Instanz injiziert werden, und decken die grundlegenden Unterschiede zwischen current_thread und multi_thread in der Montagephase auf.
Kapitel 2: Die Montage der Runtime: Wie Builder Treiber, Scheduler und Thread-Pool zusammensetzt
VonBuilderbisRuntime: Eine vollständige Reise der Assemblierung
Im vorherigen Kapitel haben wir die Verantwortungsgrenzen von Future, Waker und Executor geklärt. Aber eine real nutzbare Laufzeit ist weit mehr als „ein Executor" – sie benötigt auch eine I/O-Ereignisschleife, Timer, einen Blocking-Thread-Pool, und diese Komponenten müssen dieselbe Menge von Handles und denselben Lebenszyklus teilen. Dieses Kapitel verfolgtBuilder::builddie vollständige Assemblierungskette und beantwortet eine Kernfrage:Welche Komponenten befinden sich tatsächlich im Inneren einesRuntime, wie werden sie zusammengesetzt und teilen sich Handles。
Tokios Assemblierungseinstiegspunkt istBuilder. Es selbst ist ein reiner Konfigurationscontainer, alle Felder sind „Absichtserklärungen" und halten keine Laufzeitressourcen. Die tatsächliche Ressourcenerstellung erfolgt beim Aufruf vonbuild().
Intuitives Modell: Builder ist der „Renovierungsplan", Runtime ist das „Haus nach der Übergabe"
Builderist wie ein Renovierungsplan: Man markiert darauf „wie viele Zimmer (worker_threads)", „ob Wasseranschluss (enable_io)", „ob Stromanschluss (enable_time)", „Obergrenze für ausgelagerte Hilfskräfte (max_blocking_threads)". Der Plan selbst erzeugt keine physischen Entitäten. Erst wennbuild()aufgerufen wird, baut das Bautrupp nach Plan und errichtet die „Zimmer" wie Scheduler, Treiber, Thread-Pool tatsächlich und übergibt eineRuntimeInstanz.
Ohne die Ebene vonBuildermüsste der Benutzer jede Komponente manuell new-en, manuell verdrahten, manuell Fehler-Rollback handhaben – jede Reihenfolgefehler würde zu hängenden Handles oder Ressourcenlecks führen.BuilderDer Wert von liegt darin:„Konfiguration" und „Konstruktion" vollständig zu trennen, sodass der Konstruktionsprozess zentral Validierung, Fehlerbereinigung und Handle-Sharing durchführen kann。
Speicherlayout:Builderdie Feldpartitionierung von
BuilderDie Felder von können nach Verantwortlichkeit in vier Gruppen unterteilt werden. Die erste Gruppe istForm und Schalter:kindbestimmt die Scheduler-Form,enable_io / enable_timebestimmt, ob der entsprechende Treiber erstellt wird.
📎 tokio/src/runtime/builder.rs:55-68
pub struct Builder {
kind: Kind,
name: Option<String>,
enable_io: bool,
nevents: usize,
nevents_busy: Option<usize>,
enable_time: bool,
start_paused: bool,
// ...
}Die zweite Gruppe istThread-Pool-Parameter:worker_threadsistOption<usize>,Nonebedeutet „beim Build auf Basis der CPU-Kernzahl automatisch erkennen";max_blocking_threadsStandard ist 512.
📎 tokio/src/runtime/builder.rs:73-79
worker_threads: Option<usize>,
max_blocking_threads: usize,Die dritte Gruppe istCallback-Hooks, alle sindOption<Arc<dyn Fn ...>>. Beachten Sie, dass sieArcstattBoxverwenden, weil diese Callbacks in denConfig。
📎 tokio/src/runtime/builder.rs:87-97
pub(super) after_start: Option<Callback>,
pub(super) before_stop: Option<Callback>,
pub(super) before_park: Option<Callback>,
pub(super) after_unpark: Option<Callback>,KopierenDie vierte Gruppe ist:global_queue_interval、event_interval、disable_lifo_slot、seed_generator。
📎 tokio/src/runtime/builder.rs:116-134
pub(super) global_queue_interval: Option<u32>,
pub(super) event_interval: u32,
pub(super) disable_lifo_slot: bool,
pub(super) seed_generator: RngSeedGenerator,KopierenKindHier gibt es ein bemerkenswertes Design:Copyist ein
📎 tokio/src/runtime/builder.rs:261-265
#[derive(Clone, Copy)]
pub(crate) enum Kind {
CurrentThread,
#[cfg(feature = "rt-multi-thread")]
MultiThread,
}MultiThreadKopierenrt-multi-threadDie Variante wird durch dasrtFeature gesteuert. Das bedeutet, in einem Build, bei dem nur dasKindFeature aktiviert ist, hatbuild()nur eine Variante, und dasmatchvonwird vom Compiler zu einem einzigen Zweig optimiert –。
Verwendung des Typsystems statt Laufzeitprüfung, um die Codegröße des Multithread-Schedulers zu eliminieren
Builder::newDie Philosophie der Standardwerte: Warum I/O und time standardmäßig deaktiviert sindenable_ioist der gemeinsame Einstiegspunkt für alle Konstruktionen. Es setztenable_timeundfalse。
📎 tokio/src/runtime/builder.rs:309-318
// I/O defaults to "off"
enable_io: false,
nevents: 1024,
nevents_busy: None,
// Time defaults to "off"
enable_time: false,
// The clock starts not-paused
start_paused: false,〔Design-Inferenz und Architektur-Abwägung〕#[tokio::main]Diese Standardwertwahl ist absichtlich: Das Erstellen des I/O-Treibers erfordert die Anforderung von epoll/kqueue-Handles vom Betriebssystem, das Erstellen des time-Treibers erfordert den Start der Timer-Infrastruktur. Wenn der Benutzer nur einen reinen Berechnungs-Task-Scheduler möchte (z. B. CPU-intensive async-Logik ausführen), ist das erzwungene Erstellen dieser Treiber reine Verschwendung.enable_all()。
enable_all()Das
📎 tokio/src/runtime/builder.rs:398-419
pub fn enable_all(&mut self) -> &mut Self {
#[cfg(any(
feature = "net",
all(unix, feature = "process"),
all(unix, feature = "signal")
))]
self.enable_io();
#[cfg(all(
tokio_unstable,
feature = "io-uring",
// ...
))]
self.enable_io_uring();
#[cfg(feature = "time")]
self.enable_time();
self
}aufruft. Die Implementierung vonenable_io()offenbart, wie Feature-Gating die Semantik von „alles an" beeinflusst.net、processKopierensignalBeachten Sie, dasstime feature,enable_all()nur aufgerufen wird, wenn
oder dasbuild()Feature aktiviert ist. Wenn der Benutzer nur
build()aktiviert hat, wirdkindden I/O-Treiber nicht öffnen – weil im Kompilat überhaupt kein I/O-Treiber-Code vorhanden ist.
📎 tokio/src/runtime/builder.rs:1146-1152
pub fn build(&mut self) -> io::Result<Runtime> {
match &self.kind {
Kind::CurrentThread => self.build_current_thread_runtime(),
#[cfg(feature = "rt-multi-thread")]
Kind::MultiThread => self.build_threaded_runtime(),
}
}die Verzweigung von
ist der Startpunkt der Assemblierung, es verzweigt nach
build_current_thread_runtimein zwei völlig unterschiedliche Pfade.build_current_thread_runtime_componentsKopierenRuntime。
📎 tokio/src/runtime/builder.rs:1725-1736
fn build_current_thread_runtime(&mut self) -> io::Result<Runtime> {
use crate::runtime::runtime::Scheduler;
let (scheduler, handle, blocking_pool) =
self.build_current_thread_runtime_components(None)?;
Ok(Runtime::from_parts(
Scheduler::CurrentThread(scheduler),
handle,
blocking_pool,
))
}Pfad eins: Assemblierung von current_threadbuild_current_thread_runtime_componentsselbst ist sehr dünn, es delegiert an
📎 tokio/src/runtime/builder.rs:1760-1766
let mut cfg = self.get_cfg();
cfg.timer_flavor = TimerFlavor::Traditional;
let (driver, driver_handle) = driver::Driver::new(cfg)?;
// Blocking pool
let blocking_pool = blocking::create_blocking_pool(self, self.max_blocking_threads, 0);
let blocking_spawner = blocking_pool.spawner().clone();KopierendriverDie eigentliche Assemblierungslogik befindet sich in(driver, driver_handle). Ihre Ausführungsreihenfolge ist entscheidend:?KopierenbuildDer erste Schritt erstelltErrund gibt ein Paar
zurück. Beachten Sie, dass hierspawnerFehler direkt nach oben propagiert – wenn die I/O-Treiber-Initialisierung fehlschlägt (z. B. epoll-Erstellung fehlschlägt), gibt das gesamtespawner
zurück, zu diesem Zeitpunkt wurde der Blocking-Pool noch nicht erstellt, keine Bereinigung erforderlich.
📎 tokio/src/runtime/builder.rs:1768-1770
let seed_generator_1 = self.seed_generator.next_generator();
let seed_generator_2 = self.seed_generator.next_generator();wird in den Scheduler injiziert, damit der Scheduler die Fähigkeit hat, Blocking-Tasks an den Thread-Pool zu übermitteln.seed_generator_1Der dritte Schritt generiert zwei unabhängige RNG-Seed-Generatoren.ConfigKopierenselect!〔Design-Inferenz und Architektur-Abwägung〕seed_generator_2Warum werden zwei benötigt?CurrentThread::newwird inrng_seedplatziert, für die interne Verwendung des Schedulers (z. B.
zufällige Verzweigungsreihenfolge);Configwird anCurrentThread::new。
📎 tokio/src/runtime/builder.rs:1776-1807
let (scheduler, handle) = CurrentThread::new(
driver,
driver_handle,
blocking_spawner,
seed_generator_2,
Config {
before_park: self.before_park.clone(),
after_unpark: self.after_unpark.clone(),
// ...
global_queue_interval: self.global_queue_interval,
event_interval: self.event_interval,
// ...
enable_eager_driver_handoff: false,
seed_generator: seed_generator_1,
// ...
},
local_tid,
self.name.clone(),
);gewährleistet wird.enable_eager_driver_handoffDer vierte Schritt ist der Kern: driver, driver_handle, blocking_spawner, Seeds undfalse。
📎 tokio/src/runtime/builder.rs:1795-1798
// This setting never makes sense for a current thread runtime,
// as it only configures how the I/O driver is stolen across
// workers.
enable_eager_driver_handoff: false,Dieser Kommentar verdeutlicht das Wesen dieser Option: Sie beschreibt, „wie mehrere Worker um den I/O-Treiber konkurrieren", und da current_thread nur einen einzigen Thread hat, gibt es keine Konkurrenz, weshalb sie zwangsweise deaktiviert wird. Dies ist ein typisches Beispiel dafür, dass „die Semantik eines Konfigurationselements stark von der Form abhängt" – dasselbeBuilderFeld hat in verschiedenen Formen unterschiedliche Bedeutungen.
Schließlich wirdCurrentThread::newdas vonhandlezurückgegebenescheduler::Handle::CurrentThreadinHandle。
📎 tokio/src/runtime/builder.rs:1816-1822
let handle = Handle {
inner: scheduler::Handle::CurrentThread(handle),
};
Ok((scheduler, handle, blocking_pool))Kopie
build_threaded_runtimePfad zwei: Die Assemblierung von multi_thread
📎 tokio/src/runtime/builder.rs:2185
let worker_threads = self.worker_threads.unwrap_or_else(num_cpus);Noneähnelt dem von current_thread, weist jedoch drei wesentliche Unterschiede auf. Der erste Unterschied ist die Bestimmung der Worker-Thread-Anzahl:num_cpus()KopieBuilder::newwird hier zu
aufgelöst. Dies ist der Ort, an dem die „verzögerte automatische Erkennung" greift – die Erkennung erfolgt zur build-Zeit und nicht zur
📎 tokio/src/runtime/builder.rs:2189-2192
let blocking_pool =
blocking::create_blocking_pool(self, self.max_blocking_threads + worker_threads, worker_threads);
let blocking_spawner = blocking_pool.spawner().clone();Der zweite Unterschied liegt in der Kapazitätsberechnung des blocking pool:max_blocking_threads + worker_threadsKopieself.max_blocking_threadsBeachten Sie0。
📎 tokio/src/runtime/builder.rs:1765
let blocking_pool = blocking::create_blocking_pool(self, self.max_blocking_threads, 0);Kopiemax_blocking_threads〔Design-Inferenz und Architektur-Abwägung〕worker_threadsDieser Unterschied offenbart die Kapazitätssemantik des blocking pool: Unter multi_thread istmax_blocking_threadsdie Obergrenze für „zusätzliche" Blocking-Threads; die tatsächliche Gesamt-Thread-Obergrenze ergibt sich aus der Addition der Worker-Thread-Anzahl. Der dritte Parameter (bei current_thread 0, bei multi_thread
) ist höchstwahrscheinlich ein Hinweis auf die „Anzahl reservierter Threads" oder „Anzahl initialer Threads". Dieses Design sorgt dafür, dass die Semantik vonMultiThread::newin beiden Formen konsistent bleibt: Sie beschreibt, „wie viele zusätzliche Blocking-Threads über die Kern-Worker hinaus geöffnet werden können".
📎 tokio/src/runtime/builder.rs:2198-2226
let (scheduler, handle, launch) = MultiThread::new(
worker_threads,
driver,
driver_handle,
blocking_spawner,
seed_generator_2,
Config {
// ...
enable_eager_driver_handoff: self.enable_eager_driver_handoff,
// ...
},
self.timer_flavor,
self.name.clone(),
);ein Tripel statt eines Tupels zurückgibt:launchKopieMultiThread::newDas zusätzlicheist ein „Start-Handle".ist nur für die Konstruktion der Scheduler-Struktur verantwortlich,
📎 tokio/src/runtime/builder.rs:2228-2234
let handle = Handle { inner: scheduler::Handle::MultiThread(handle) };
// Spawn the thread pool workers
let _enter = handle.enter();
launch.launch();
Ok(Runtime::from_parts(Scheduler::MultiThread(scheduler), handle, blocking_pool))handle.enter(). Der eigentliche Start erfolgt später:launch.launch()Kopie
werden tatsächlich alle Worker-Threads gespawnt. Dieses zweiphasige Design „erst konstruieren, dann starten" ist äußerst entscheidend.handle〔Design-Inferenz und Architektur-Abwägung〕handleWarum kann man nicht während der Konstruktion starten? Weil Worker-Threads, sobald sie gestartet sind, sofort mit dem Pollen von Aufgaben beginnen, und Aufgaben möglicherweise aufverweisen. WennHandlenoch nicht fertig konstruiert ist, entsteht eine Race Condition, bei der „Worker ein halbfertiges Handle halten". Das zweiphasige Design stellt sicher:。_enterWenn alle Worker-Threads starten, ist das vollständige
bereits bereit
Der Guard stellt sicher, dass sich Worker-Threads im Moment des Starts bereits im korrekten Runtime-Kontext befinden.driver::Driver::newAssemblierungs-FlussdiagrammErrDie folgende Abbildung zeigt die Assemblierungsreihenfolge, die wichtigsten Verzweigungen und die Fehlerpfade beider Pfade zusammen. Beachten Sie, dass bei einem Fehlschlag von
flowchart TD
start["Builder::build()"] --> match_kind{"self.kind?"}
match_kind -->|CurrentThread| ct_cfg["get_cfg() + timer_flavor=Traditional"]
match_kind -->|MultiThread| mt_workers["worker_threads = self.worker_threads.unwrap_or_else(num_cpus)"]
ct_cfg --> ct_driver["driver::Driver::new(cfg)?"]
mt_workers --> mt_driver["driver::Driver::new(self.get_cfg())?"]
ct_driver -->|Err| ret_err["return Err(io::Error)"]
mt_driver -->|Err| ret_err
ct_driver -->|Ok driver, driver_handle| ct_pool["create_blocking_pool(self, max_blocking_threads, 0)"]
mt_driver -->|Ok driver, driver_handle| mt_pool["create_blocking_pool(self, max_blocking_threads + worker_threads, worker_threads)"]
ct_pool --> ct_seed["next_generator() x2"]
mt_pool --> mt_seed["next_generator() x2"]
ct_seed --> ct_new["CurrentThread::new(driver, driver_handle, blocking_spawner, ...)"]
mt_seed --> mt_new["MultiThread::new(worker_threads, driver, ...) -> (scheduler, handle, launch)"]
ct_new --> ct_wrap["Handle { inner: CurrentThread(handle) }"]
mt_new --> mt_wrap["Handle { inner: MultiThread(handle) }"]
ct_wrap --> ct_rt["Runtime::from_parts(Scheduler::CurrentThread, handle, blocking_pool)"]
mt_wrap --> mt_enter["handle.enter()"]
mt_enter --> mt_launch["launch.launch() 启动 worker 线程"]
mt_launch --> mt_rt["Runtime::from_parts(Scheduler::MultiThread, handle, blocking_pool)"]zurückgegeben wird, wobei der blocking pool zu diesem Zeitpunkt noch nicht erstellt wurde.HandleKopie
Handle-Sharing:RuntimeWiescheduler、handle、blocking_poolzum „Passierschein" über Komponentengrenzen hinweg wirdhandleNach Abschluss der Assemblierung hält
📎 tokio/src/runtime/scheduler/mod.rs:29-41
#[derive(Debug, Clone)]
pub(crate) enum Handle {
#[cfg(feature = "rt")]
CurrentThread(Arc<current_thread::Handle>),
#[cfg(feature = "rt-multi-thread")]
MultiThread(Arc<multi_thread::Handle>),
#[cfg(not(feature = "rt"))]
#[allow(dead_code)]
Disabled,
}. Davon istArcder gemeinsame Kern. Sein Inneres ist eine Enum:HandleKopieHandleBeachten Sie, dass beide Variantenmatchumschließen. Das bedeutet, dass das Klonen vondriver():
📎 tokio/src/runtime/scheduler/mod.rs:53-64
pub(crate) fn driver(&self) -> &driver::Handle {
match *self {
#[cfg(feature = "rt")]
Handle::CurrentThread(ref h) => &h.driver,
#[cfg(feature = "rt-multi-thread")]
Handle::MultiThread(ref h) => &h.driver,
#[cfg(not(feature = "rt"))]
Handle::Disabled => unreachable!(),
}
}blocking_spawner()bietet eine einheitliche Zugriffsschnittstelle, die die Formunterschiede im Inneren vonmatch_flavor!kapselt. Zum Beispiel
📎 tokio/src/runtime/scheduler/mod.rs:96-98
pub(crate) fn blocking_spawner(&self) -> &blocking::Spawner {
match_flavor!(self, Handle(h) => &h.blocking_spawner)
}verwendet dasdriver()-Makro zur Eliminierung von Wiederholungen:matchKopiematch_flavor!Dieses Makro expandiert zu einemmatchwie dem obigen
. Sein Wert liegt darin: Wenn ein neuer Accessor hinzugefügt wird, der nach Form verteilen muss, genügt eine ZeileHandle, anstatt zweimal diescheduler::Handle-Verzweigung von Hand zu schreiben.
📎 tokio/src/runtime/handle.rs:13-15
pub struct Handle {
pub(crate) inner: scheduler::Handle,
}ist ein dünner Wrapper um das interneHandle:spawnKopieblock_on。spawnDas vom Benutzer erhalteneAutoBoxkann threadübergreifend geklont werden, kann
📎 tokio/src/runtime/handle.rs:197-208
pub fn spawn<F>(&self, future: F) -> JoinHandle<F::Output>
where
F: Future + Send + 'static,
F::Output: Send + 'static,
{
let fut_size = mem::size_of::<F>();
if AutoBox::<F>::SHOULD_BOX {
self.spawn_named(Box::pin(future), SpawnMeta::new_unnamed(fut_size))
} else {
self.spawn_named(future, SpawnMeta::new_unnamed(fut_size))
}
}AutoBox::<F>::SHOULD_BOXDie Implementierung vonsize_of::<F>()zeigt die Compile-Zeit-Verzweigung von
📎 tokio/src/runtime/mod.rs:668-673
pub(crate) struct AutoBox<T>(std::marker::PhantomData<T>);
impl<T> AutoBox<T> {
pub(crate) const SHOULD_BOX: bool = std::mem::size_of::<T>() > BOX_FUTURE_THRESHOLD;
}ist eine assoziierte Konstante, die durch den Vergleich vonifmit einem Schwellenwert ermittelt wird.spawn_namedKopieF〔Design-Inferenz und Architektur-Abwägung〕Pin<Box<F>>Der Kommentar erklärt, warum eine assoziierte Konstante statt einer Laufzeit-
verwendet wird: Bei einer Laufzeitprüfung würde
zweimal monomorphisiert (einmal für, einmal fürdriver -> blocking_pool -> scheduler), was dazu führt, dass für jede gespawnte Future zwei Task-Harnesses generiert werden und sich die Codegröße verdoppelt. Mit einer konstanten Verzweigung behält der Monomorphisierungs-Sammler nur den tatsächlich durchlaufenen Zweig.
Design-Überlegungen: Assemblierungsreihenfolge, Fehlerwiederherstellung und Produktions-Fallstrickelocal_tidReihenfolge ist Vertrag。build_local. Die Assemblierungsreihenfolgebuild_current_thread_local_runtimeist nicht willkürlich. Der Treiber wird zuerst erstellt, da er der einzige Schritt ist, der aufgrund unzureichender OS-Ressourcen fehlschlagen kann und bei einem Fehlschlag keine Bereinigung anderer Komponenten erfordert. blocking_pool kommt nach dem Treiber und vor dem Scheduler, da der Scheduler blocking_spawner benötigt. Wenn die Erstellung von blocking_pool fehlschlägt (was praktisch kaum vorkommt), wird der Treiber durch Drop automatisch bereinigt.
📎 tokio/src/runtime/builder.rs:1738-1751
fn build_current_thread_local_runtime(&mut self) -> io::Result<LocalRuntime> {
use crate::runtime::local_runtime::LocalRuntimeScheduler;
let tid = std::thread::current().id();
let (scheduler, handle, blocking_pool) =
self.build_current_thread_runtime_components(Some(tid))?;
Ok(LocalRuntime::from_parts(
LocalRuntimeScheduler::CurrentThread(scheduler),
handle,
blocking_pool,
))
}-Zweig von current_threadtidverwendetHandle, wobei die aktuelle Thread-ID übergeben wird:can_spawn_local_on_local_runtimeKopie
📎 tokio/src/runtime/scheduler/mod.rs:140-147
pub(crate) fn can_spawn_local_on_local_runtime(&self) -> bool {
match self {
Handle::CurrentThread(h) => h.local_tid.is_some_and(|x| std::thread::current().id() == x),
#[cfg(feature = "rt-multi-thread")]
Handle::MultiThread(_) => false,
}
}gespeichert, und später verwendetLocalRuntimesie zur Überprüfung, „ob spawn_local auf dem Owner-Thread aufgerufen wird":!SendKopielocal_tid〔Design-Inferenz und Architektur-Abwägung〕!SendDies ist der Grundpfeiler der Sicherheit von
:worker_threads(0)Die Future von。worker_threadsdarf nur auf ihrem Owner-Thread gepollt werden, und
📎 tokio/src/runtime/builder.rs:582-586
pub fn worker_threads(&mut self, val: usize) -> &mut Self {
assert!(val > 0, "Worker threads cannot be set to 0");
self.worker_threads = Some(val);
self
}Diese Assertion schlägt bereits in der Konfigurationsphase fehl, nicht erst beim Build. Der Vorteil ist eine frühere Fehlerlokalisierung, der Nachteil ist, dass der Benutzer, wenn die Thread-Anzahl aus einem dynamischen Wert der Konfigurationsdatei stammt, sie vor dem Aufruf selbst validieren muss.
Produktions-Fallstrick Zwei:max_blocking_threadsZu klein eingestellt führt zu Hängen. Die Dokumentation warnt ausdrücklich:
📎 tokio/src/runtime/builder.rs:600-601
/// It's recommended to not set this limit too low in order to avoid hanging on operations
/// requiring [`spawn_blocking`].Weil die Warteschlange des blocking pool kein Backpressure hat – Aufgaben sammeln sich an, bis ein Thread verfügbar ist. Wenn alle blockierenden Threads auf eine Operation warten, die „einen neuen blockierenden Thread benötigt, um abgeschlossen zu werden", kommt es zum Deadlock. Die Aussage in der Dokumentation „the queue does not apply any backpressure, it could potentially grow unbounded" ist genau die Fußnote zu diesem Risiko.
Produktions-Fallstrick Drei:UnhandledPanic::ShutdownRuntimeNur current_thread wird unterstützt。
📎 tokio/src/runtime/builder.rs:1374-1381
pub fn unhandled_panic(&mut self, behavior: UnhandledPanic) -> &mut Self {
if !matches!(self.kind, Kind::CurrentThread) && matches!(behavior, UnhandledPanic::ShutdownRuntime) {
panic!("UnhandledPanic::ShutdownRuntime is only supported in current thread runtime");
}
self.unhandled_panic = behavior;
self
}Der Grund für diese Einschränkung ist: Unter multi_thread erfordert „sofortiges Herunterfahren der Runtime" die Koordination des Stopps aller Worker-Threads, was implementierungstechnisch komplex und semantisch unklar ist (was passiert mit anderen Aufgaben, die gerade gepollt werden?). current_thread hat nur einen Thread, die Shutdown-Semantik ist klar.
Zusammenfassung dieses Kapitels
Dieses Kapitel hatBuilder::builddie vollständige Assemblierungskette nachverfolgt. Kernschlussfolgerungen:
1. Builderist ein reiner Konfigurationscontainer,build()erst erstellt Ressourcen. Die Assemblierungsreihenfolgedriver -> blocking_pool -> schedulerwird durch die Fehlerwiederherstellungsanforderung bestimmt.
2. Der Unterschied zwischen current_thread und multi_thread beschränkt sich nicht auf die Thread-Anzahl: Die Kapazitätsberechnung des blocking pool ist unterschiedlich (max_blocking_threads vs max_blocking_threads + worker_threads), multi_thread hat zusätzlich einenlaunchzweiphasigen Start,enable_eager_driver_handoffwird unter current_thread zwangsweise deaktiviert.
3. Handleist der Kern, der komponentenübergreifend geteilt wird, intern wirdArcverwendet, um form-spezifische Handles zu umschließen, und übermatchodermatch_flavor!Makros einheitlich zugegriffen.
4. AutoBoxverwendet assoziierte Konstanten, um zur Kompilierzeit zu entscheiden, ob das Future geboxt wird, wodurch eine Verdopplung der Codegröße vermieden wird.
5. local_tidistLocalRuntimeder Laufzeitprüfpunkt für Sicherheit.
Im nächsten Kapitel betreten wir den Lebenszyklus von Aufgaben:spawnwie ein Future in eine schedulierbare Entität umgewandelt wird,JoinHandlewie mit der Aufgaben-Zustandsmaschine interagiert wird, und die Zustandsübergänge von Aufgaben zwischenPENDING / RUNNING / COMPLETE.
Gedanken und Selbsttest dieses Kapitels
Q1: Wenn man inbuild_threaded_runtimeden Kapazitätsparameter voncreate_blocking_poolvonself.max_blocking_threads + worker_threadsaufself.max_blocking_threadsändert, in welchen Szenarien würde dies dazu führen, dass blockierende Aufgaben verhungern? Warum kann der current_thread-Pfadself.max_blocking_threads?
Referenzanalyse: Gemäß📎 tokio/src/runtime/builder.rs:2189-2192übergibt der multi_thread-Pfadself.max_blocking_threads + worker_threads, während der current_thread-Pfad📎 tokio/src/runtime/builder.rs:1765übergibtself.max_blocking_threads. Die Ursache des Unterschieds liegt darin: Unter multi_thread führen die Worker-Threads selbst auch blockierende Aufgaben aus (zum Beispielblock_in_placewandelt den Worker-Thread vorübergehend in einen blockierenden Thread um), daher muss das Gesamtbudget der blockierenden Threads die Anzahl der Worker-Threads einschließen. Wenn man stattdessen nurself.max_blocking_threadsübergibt, wennmax_blocking_threadsklein eingestellt ist (zum Beispiel 1) und bereits Worker-Threads inblock_in_placedas Budget belegen, werden neuespawn_blockingAufgaben keinen Thread zur Verfügung haben und sich in der Warteschlange ohne Backpressure ansammeln, was dazu führt, dass async-Aufgaben, die von diesen blockierenden Aufgaben abhängen, dauerhaft hängen. current_thread hat nur einen Thread und unterstützt nicht die Worker-Konvertierungssemantik vonblock_in_place, daher muss die Worker-Anzahl nicht addiert werden.
Q2: MultiThread::newgibt daslaunchHandle zurück, was die Worker-Threads tatsächlich startet, istlaunch.launch(). Was passiert, wenn man die Zeilehandle.enter()entfernt und direktlaunch.launch()aufruft?
Referenzanalyse: Gemäß📎 tokio/src/runtime/builder.rs:2230-2232gibt es vor dem Startlet _enter = handle.enter();und erst dannlaunch.launch()。handle.enter()Der Zweck ist, den thread-lokalen Kontext (thread-local) zu setzen, damit der aktuelle Thread „so aussieht", als befände er sich innerhalb der Runtime. Nach dem Start beginnen die Worker-Threads sofort mit dem Pollen von Aufgaben, und der Aufgabencode könnte APIs wieHandle::current()、tokio::spawnaufrufen, die vom Kontext abhängen. Wenn man_enterentfernt, könnte die Kontextsetzung im Moment des Starts des Worker-Threads unvollständig sein (je nachdem, oblaunchintern selbst setzt), im schlimmsten Fall würde der Initialisierungscode, der auf dem Worker-Thread ausgeführt wird, bei einem Aufruf vonHandle::current()panic verursachen (CONTEXT_MISSING_ERROR). Selbst wennlaunchintern für jeden Worker den Kontext setzt,_enterstellt ebenfalls sicher, dass „die Startaktion selbst" im richtigen Kontext stattfindet, wodurch Races während des Startvorgangs vermieden werden.
Q3: AutoBox::<F>::SHOULD_BOXverwendet assoziierte Konstanten statt Laufzeit-if size_of::<F>() > THRESHOLD. Angenommen, man ändert es zu einer Laufzeitprüfung, in welchen Fällen würde dies neben der Verdopplung der Codegröße zu Leistungseinbußen führen?
Referenzanalyse: Gemäß📎 tokio/src/runtime/mod.rs:657-673den Kommentaren vonifwürde Laufzeit-spawn_nameddazu führen, dassTfür jedesTzweimal monomorphisiert wird (Pin<Box<T>>undPin<Box<T>>jeweils einmal). Neben der Verdopplung der Codegröße zeigt sich die Leistungseinbuße in: 1) erhöhter Druck auf den Instruktionscache (i-cache), da beide Harness-Code-Sätze resident sein müssen; 2) der Compiler kann nicht optimieren, dass „tatsächlich nur ein Zweig durchlaufen wird", die Laufzeit-Branch-Vorhersage ist zwar normalerweise genau, aber der Branch selbst und die Unterschiede in der Registerzuweisung der beiden Code-Sätze summieren sich; 3) subtiler ist, dass dersize_ofPfad eine Heap-Allokation erzwingt, wenn die Laufzeitprüfung aus irgendeinem Grund (zum Beispiel
Kapitel 3: Das Leben einer Aufgabe (Teil 1): Wie spawn einen Future in eine planbare Entität verwandelt
Im vorherigen Kapitel haben wir die Montage der Runtime abgeschlossen: I/O driver, time driver, blocking pool und Scheduler werden in dieselbeRuntimeInstanz injiziert,Handleund werden zu einem gemeinsamen Handle für den threadübergreifenden Zugriff auf diese Komponenten. Doch die montierte Runtime ist zu diesem Zeitpunkt noch eine leere Hülle – sie besitzt die Engine zum Antreiben von Aufgaben, aber es gibt keine Aufgaben, die angetrieben werden könnten. Die Frage, die dieses Kapitel beantworten will, ist genau: Wenn dutokio::spawn(async { ... })eingibst, was durchläuft dann dieserasyncBlock, um von einem gewöhnlichen Rust-Code-Stück zu einer Entität zu werden, die „vom Scheduler übernommen, geweckt und gejoint werden kann"? Dies ist die erste Halbzeit von „Das Leben einer Aufgabe", wir konzentrieren uns auf die Geburt: Ausgehend vonHandle::spawndurchlaufen wirnew_taskdie Referenzzählungs-Allokation und landen beiCell<T, S>dem Speicherlayout, um schließlich klar zu sehen, wie die Aufgabe in die lokale Warteschlange eines Workers oder in die globale Injektions-Warteschlange eingereiht wird. Die zweite Halbzeit (Kapitel 4) wird dann in die Scheduling-Schleife und den poll/wake-Kreislauf eintreten.
3.1 Future ist keine Aufgabe: Was genau erzeugt ein spawn
Intuitives Modell
Stelle dirFutureals ein „Rezept" vor und eine Aufgabe als „ein Gericht, das gerade in der Küche gekocht wird". Das Rezept selbst ist statisch, kopierbar und hat keinerlei Ausführungszustand; erst wenn die Küche (der Scheduler) entscheidet „jetzt dieses Gericht zubereiten", ihm einen Herd (Worker), eine Bestellnummer (TaskId) und eine Ausgabestation (JoinHandle) zuweist, wird es zu einem „Gericht in Zubereitung". Ohne diese Verpackungsschicht kann der Scheduler nicht wissen, „bis zu welchem Schritt dieses Gericht zubereitet ist", „wer darauf wartet", „wen es nach Fertigstellung benachrichtigen soll" – er sieht nur ein Rezept und kann nichts verwalten.
Datenstruktur und Speicherlayout
Tokio verwendetTask<S>um „eine von der Runtime besessene Aufgabenreferenz" darzustellen, es ist ein transparenter Wrapper umRawTask
#[repr(transparent)]
pub(crate) struct Task<S: 'static> {
raw: RawTask,
_p: PhantomData<S>,
}📎 tokio/src/runtime/task/mod.rs:233-238
#[repr(transparent)]bedeutet, dassTask<S>undRawTaskim Speicher vollständig identisch sind, ohne zusätzlichen Overhead.PhantomData<S>ist nur eine Typmarkierung zur Kompilierungszeit, die markiert, zu welchem Scheduler-Typ diese Aufgabe gehörtS。
Was tatsächlich den gesamten Zustand der Aufgabe trägt, istCell<T, S>dessen Layout der Grundstein des gesamten Aufgabenmoduls ist:
#[repr(C)]
pub(super) struct Cell<T: Future, S> {
pub(super) header: Header,
pub(super) core: Core<T, S>,
pub(super) trailer: Trailer,
}📎 tokio/src/runtime/task/core.rs:126-136
Die drei Felder sind nach „heiß-warm-kalt" angeordnet.Headersind heiße Daten (bei jedem Scheduling, bei jedem Zustandsübergang zugegriffen),Coresind warme Daten (beim Poll zugegriffen),Trailersind kalte Daten (nur bei Erstellung und Zerstörung zugegriffen). Der Kommentar schreibt ausdrücklich:Headermuss das erste Feld sein, weil die Aufgabenstruktur gleichzeitig von*mut Cellund*mut Headerreferenziert wird📎 tokio/src/runtime/task/core.rs:37-43。
Noch entscheidender ist die Cache-Line-Ausrichtung.Cellträgt eine lange Reihe von#[cfg_attr(..., repr(align(...)))]die je nach Zielarchitektur die Anzahl der Ausrichtungsbytes wählt: x86_64/aarch64/powerpc64 verwenden 128 Bytes, arm/mips/sparc/hexagon verwenden 32 Bytes, m68k verwendet 16 Bytes, s390x verwendet 256 Bytes, der Rest standardmäßig 64 Bytes📎 tokio/src/runtime/task/core.rs:64-125Der Kommentar erklärt, warum x86_64 128 statt 64 verwenden muss: Seit Intel Sandy Bridge lädt der Spatial Prefetcherpaarweise64-Byte-Cache-Lines, daher muss auf 128 Bytes ausgerichtet werden, um False Sharing zu vermeiden📎 tokio/src/runtime/task/core.rs:45-53。
Der Preis dieser Ausrichtungsstrategie ist, dass jede Aufgabe mindestens eine Cache-Line Speicher verschwendet. Aber die Aufgabenstatusbits (state) werden von mehreren Worker-Threads mit hoher Frequenz gelesen und geschrieben – ein Thread setzt beim Poll das RUNNING-Bit, ein anderer Thread liest beim Wecken das NOTIFIED-Bit – wenn die Statusbits zweier Aufgaben in derselben Cache-Line liegen, löst jeder Zustandsübergang ein Hin- und Herspringen der Cache-Line zwischen den Kernen aus (cache line ping-pong), der Leistungsverlust übersteigt bei weitem die Speicherverschwendung. Tokio wählt Raum gegen Zeit.
Headerselbst ist auf 8 Zeigergrößen beschränkt:
#[test]
#[cfg(not(loom))]
fn header_lte_cache_line() {
assert!(std::mem::size_of::<Header>() <= 8 * std::mem::size_of::<*const ()>());
}📎 tokio/src/runtime/task/core.rs:591-593
Dieser Test stellt sicher, dassHeader64 Bytes (8 × 8) nicht überschreitet, sodass es auf Architekturen mit 64-Byte-Cache-Lines vollständig in eine Zeile passt.HeaderDie Felder von umfassen:state: State(atomare Statusbits),queue_next: UnsafeCell<Option<NonNull<Header>>>(Verkettungszeiger der Injektions-Warteschlange),vtable: &'static Vtable(Funktionszeigertabelle),owner_id: UnsafeCell<Option<NonZeroU64>>(ID der zugehörigenOwnedTasksListe),scheduled_at: UnsafeCell<ScheduleLatencyInstant>(Scheduling-Latenzmessung)📎 tokio/src/runtime/task/core.rs:169-198。
Core<T, S>hält das Scheduler-Handlescheduler: Sdie Aufgaben-IDtask_id: Idsowie das Kernstückstage: CoreStage<T> 📎 tokio/src/runtime/task/core.rs:148-165。Stageist eine dreizuständige Enum:
#[repr(C)]
pub(super) enum Stage<T: Future> {
Running(T),
Finished(super::Result<T::Output>),
Consumed,
}📎 tokio/src/runtime/task/core.rs:225-229
Genau das ist der Schlüssel dafür, dass „Future und Output denselben Speicher wiederverwenden": Während die Aufgabe läuft, hältStage::Runningden Future, nach Abschluss wird er an Ort und Stelle durchStage::Finished(output)ersetzt, nach Entnahme durchJoinHandlewird er zuStage::Consumed。#[repr(C)]Der Kommentar verweist auf ein Miri-Issue, das zeigt, dass dieses Layout harte Anforderungen an die Korrektheit von unsafe-Code stellt📎 tokio/src/runtime/task/core.rs:225-229。
Trailerspeichert kalte Daten:owned: linked_list::Pointers<Header>(OwnedTasksVerkettungszeiger),waker: UnsafeCell<Option<Waker>>(Consumer-Waker, der auf die Fertigstellung der Aufgabe wartet),hooks: TaskHarnessScheduleHooks 📎 tokio/src/runtime/task/core.rs:205-213。
Step-by-Step: Von spawn bis zur Einreihung
Wir versetzen uns in ein konkretes Szenario: In einer multi_thread-Runtime führt Worker-Thread Atokio::spawn(async { 42 })。
aus new_taskErster Schritt: Die drei Teile der Aufgabe konstruieren.
fn new_task<T, S>(
task: T,
scheduler: S,
id: Id,
spawned_at: SpawnLocation,
) -> (Task<S>, Notified<S>, JoinHandle<T::Output>)📎 tokio/src/runtime/task/mod.rs:336-346
KopierenRawTask::new::<T, S>Es ruftCellauf, umrawzu allokieren, und leitet dann aus demselbenTaskZeiger drei Referenzen ab:OwnedTasks)、Notified(Owned-Referenz, wird normalerweise sofort inJoinHandle(Benachrichtigungsreferenz, wird dem Scheduler übergeben),📎 tokio/src/runtime/task/mod.rs:347-363. Beachten Sie, dass alle drei dasselberawteilen, wobei jeder eine Referenzzählung hält.
Zweiter Schritt: Allokieren vonCellund Schreiben des Anfangszustands. Cell::newAllokieren der gesamten Struktur auf dem Heap:
let result = Box::new(Cell {
trailer: Trailer::new(scheduler.hooks()),
header: new_header(state, vtable, ...),
core: Core {
scheduler,
stage: CoreStage {
stage: UnsafeCell::new(Stage::Running(future)),
},
task_id,
...
},
});📎 tokio/src/runtime/task/core.rs:261-278
vtablewird vonraw::vtable::<T, S>()generiert und ist eine auf konkreteTundSmonomorphisierte Funktionstabelle📎 tokio/src/runtime/task/core.rs:260. Das Future wird direkt inStage::Runningverschoben, ohne zusätzliches Boxing.
Dritter Schritt: Debug-Assertion zur Layout-Verifikation.Unterdebug_assertions,Cell::newwird diecheck-Funktion aufgerufen, die mitHeader::get_trailer、Header::get_scheduler、Header::get_id_ptrund anderen auf vtable-Offsets basierenden Zeigerarithmetiken einzeln assertiert, dass „die über den Header zurückermittelte Feldadresse" mit „der tatsächlichen Feldadresse" übereinstimmt📎 tokio/src/runtime/task/core.rs:280-321. Dies ist eine Laufzeit-Selbstprüfung der Korrektheit der vtable-Offsets.
Vierter Schritt: Zustellung an den Scheduler.Der Scheduler erhältNotified<S>und ruftSchedule::schedule 📎 tokio/src/runtime/task/mod.rs:315auf. Unter multi_thread geht dies überpush_back_or_overflow, wobei die Aufgabe in die lokale Warteschlange des aktuellen Workers eingereiht wird; bei voller Warteschlange wird in die Injection-Queue überlaufen.
Die folgende Abbildung skizziert den Kontrollfluss und die Verzweigungen vonnew_taskbis zur Einreihung:
flowchart TD
spawn_call["Handle::spawn(future)"] --> new_task["new_task::<T,S>(future, scheduler, id)"]
new_task --> raw_new["RawTask::new::<T,S>"]
raw_new --> cell_new["Cell::new: Box::new(Cell{header, core, trailer})"]
cell_new --> vtable["raw::vtable::<T,S>() 生成函数指针表"]
cell_new --> stage["Stage::Running(future) 移入"]
cell_new --> debug_check{"debug_assertions?"}
debug_check -->|是| check_layout["check(): 断言 trailer/scheduler/id 偏移量"]
debug_check -->|否| skip_check["跳过"]
check_layout --> triple["派生 (Task, Notified, JoinHandle)"]
skip_check --> triple
triple --> owned["Task 存入 OwnedTasks"]
triple --> sched["Notified 交给 Schedule::schedule"]
sched --> push{"本地队列有容量?"}
push -->|是| local_push["push_back_finish: 写入 buffer[tail & MASK]"]
push -->|否| overflow_check{"steal == real?"}
overflow_check -->|否, 有并发窃取| inject_only["overflow.push(task) 仅注入"]
overflow_check -->|是| push_overflow["push_overflow: CAS 认领后半批"]
push_overflow --> cas_ok{"CAS 成功?"}
cas_ok -->|是| inject_batch["overflow.push_batch(后半批 + 当前 task)"]
cas_ok -->|否| retry["返回 Err(task), 重试 push_back_or_overflow"]
retry --> pushDiese Abbildung offenbart einige Schlüsselverzweigungen: Debug-Assertions wirken nur in Debug-Builds; bei voller lokaler Warteschlange wird nicht direkt überlaufen, sondern zuerst geprüft, ob es nebenläufige Stealer gibt (steal != real), und falls ja, wird nur die aktuelle Aufgabe in die Injection-Queue eingereiht, da der von Stealern freigegebene Platz bald verfügbar ist.
Designüberlegung: Warum drei Referenzen statt einer
new_taskgibt drei Referenzen zurück, nicht eine. Dies ist der Kern des Referenzzählungsdesigns:Taskrepräsentiert „die Laufzeit besitzt diese Aufgabe",Notifiedrepräsentiert „diese Aufgabe wurde benachrichtigt, wartet auf Scheduling",JoinHandlerepräsentiert „jemand interessiert sich für ihr Ergebnis". Die drei haben unabhängige Lebensdauern——JoinHandlekann gedroppt werden (Aufgabe läuft weiter, Ergebnis wird verworfen),Notifiedverschwindet nach poll,Taskwird freigegeben, nachdem die Aufgabe abgeschlossen und ausOwnedTasksentfernt wurde. Mit nur einer Referenz ließe sich der Zustand „Aufgabe läuft noch, aber niemand joint" nicht ausdrücken.
UnownedTaskist eine weitere wichtige Verzweigung: Sie hältzweiReferenzzählungen, für Blocking-Aufgaben (nicht inOwnedTasks)📎 tokio/src/runtime/task/mod.rs:286-295。unownedgespeichert). Diemem::forget(task)-Funktion führt übermem::forget(notified)undUnownedTask 📎 tokio/src/runtime/task/mod.rs:388-397die beiden Referenzen inOwnedTaskszusammen. Die Designmotivation für „zwei Referenzen" ist: Blocking-Aufgaben haben keine
3.2 Statusbits: Wie ein usize den gesamten Lebenszyklus einer Aufgabe kodiert
Intuitives Modell
Stellen Sie sich den Aufgabenstatus als einen „Gesundheitsbericht" mit mehreren unabhängigen Kontrollkästchen vor: Wird gerade gepollt, ist abgeschlossen, wurde benachrichtigt, wurde abgebrochen, hat jemand gejoint. Tokio verwendet nicht mehrere boolesche Felder, sondern presst diese Markierungsbits ineinAtomicUsize. So benötigt jeder Statusübergang nur ein CAS statt mehrerer Lockings. Ohne dieses Design würden Aufgabenstatusübergänge zu verschachtelten mehreren Locks werden, was Deadlock-Risiko und Overhead stark erhöhen würde.
Bitfeld-Layout
StateDas Bitfeld von📎 tokio/src/runtime/task/mod.rs:32-53:
RUNNINGist in der Moduldokumentation vollständig definiert: ob die Aufgabe gerade gepollt oder abgebrochen wird. 📎tokio/src/runtime/task/mod.rs:37-38。COMPLETEDieses Bit dient gleichzeitig als Lock der AufgabeRUNNING: Das Future ist vollständig abgeschlossen und gedroppt. Einmal gesetzt, wird es nie gelöscht und nie gleichzeitig mit📎tokio/src/runtime/task/mod.rs:40-41。NOTIFIEDgesetztNotified: ob derzeit ein📎tokio/src/runtime/task/mod.rs:43。CANCELLED-Objekt existiert📎tokio/src/runtime/task/mod.rs:45-46。JOIN_INTEREST: Die Aufgabe sollte so schnell wie möglich abgebrochen werdenJoinHandle📎tokio/src/runtime/task/mod.rs:48。JOIN_WAKER: existiert📎tokio/src/runtime/task/mod.rs:50-51。
: als Zugriffskontrollbit für den join-handle-waker📎 tokio/src/runtime/task/mod.rs:53。
RUNNINGDie restlichen Bits dienen der ReferenzzählungRUNNINGDass das📎 tokio/src/runtime/task/mod.rs:130-133-Bit als Lock dient, verdient Ausführung. Der Safety-Abschnitt der Moduldokumentation stellt fest: Jeder mutierende Zugriff auf das Future muss nach Erlangen des Locks durch Modifikation desRUNNING-Bits erfolgen, um exklusiven Zugriff zu gewährleisten
. Das bedeutet, beim Pollen einer Aufgabe setzt der Thread zuerst per CAS
JOIN_WAKER, und bei Erfolg hat er exklusiven Zugriff auf das Future; bei Fehlschlag pollt ein anderer Thread, und dieser Poll kehrt direkt zurück. Dies vereint „Poll-Mutex" und „Statusübergang" in einer einzigen atomaren Operation und vermeidet ein separates Mutex.wakerZugriffskontrollprotokoll für JOIN_WAKERTrailerDas-Bit ist der raffinierteste Teil der gesamten Zustandsmaschine. Es löst das Problem:DasJoinHandle-Feld (in) wird von zwei Threads nebenläufig zugegriffen——die Laufzeitliest📎 tokio/src/runtime/task/mod.rs:75-120:
1. JOIN_WAKERes beim Abschluss der Aufgabe, um den Joiner zu wecken,
schreibtJoinHandlees beim Pollen, um den Waker zu registrieren. Die Moduldokumentation gibt 7 Regeln an
ist initial 0.JoinHandle2. Wenn 0,
hat exklusiven (mutierenden) Zugriff auf das waker-Feld.COMPLETE3. Wenn 1,
5. JoinHandlehat nur gemeinsamen (nur lesenden) Zugriff.JOIN_WAKER4. Wenn 1 undJOIN_WAKER1 ist, hat die Laufzeit gemeinsamen (nur lesenden) Zugriff auf das waker-Feld.
6. JoinHandleUm waker zu schreiben, muss man: (i) erfolgreichCOMPLETEauf 0 setzen, um exklusiven Zugriff zu erlangen, (ii) waker schreiben, (iii) erfolgreichJOIN_WAKERauf 1 setzen.COMPLETEdarf
nur ändern, wennJOIN_INTEREST0 ist; die Laufzeit darf nur ändern, wennCOMPLETE1 ist.
7. WennCOMPLETE0 ist und📎 tokio/src/runtime/task/mod.rs:110-1201 ist, hat die Laufzeit exklusiven Zugriff auf das waker-Feld (zum Droppen des wakers).
Regel 6 impliziert eine Race-Condition: Schritt (i) oder (iii) kann fehlschlagen. Wenn (i) fehlschlägt, wird das Schreiben des wakers aufgegeben; wenn (iii) fehlschlägt (ein anderer Thread hat in der Zwischenzeit
Taskgesetzt), wird das waker-Feld geleertUnownedTaskden Drop zweimal dekrementiert:
impl<S: 'static> Drop for Task<S> {
fn drop(&mut self) {
if self.header().state.ref_dec() {
self.raw.dealloc();
}
}
}📎 tokio/src/runtime/task/mod.rs:580-586
impl<S: 'static> Drop for UnownedTask<S> {
fn drop(&mut self) {
if self.raw.header().state.ref_dec_twice() {
self.raw.dealloc();
}
}
}📎 tokio/src/runtime/task/mod.rs:590-596
ref_decgibt zurücktruezeigt an, dass dies die letzte Referenz ist, und erst dann wird tatsächlichCellSpeicher freigegeben.ref_dec_twiceistUnownedTaskder direkte Ausdruck dafür, dass zwei Zähler gehalten werden.
Designüberlegung: Warum Statusbits und Referenzzähler ein gemeinsames Atom verwenden
Statusbits und Referenzzähler im selbenAtomicUsizezu platzieren, dient dazu, die beiden Aktionen „Referenzzähler dekrementieren“ und „Statusbit setzen“ ineinem einzigen CASabzuschließen. Die Moduldokumentation erwähnt im Kommentar zuSchedule::releaseausdrücklich: „Das Task-Modul verarbeitet ref-dec und andere Optionseinstellungen in Batches“📎 tokio/src/runtime/task/mod.rs:302-304. Wenn Statusbits und Referenzzähler zwei separate atomare Variablen wären, entstünde zwischen „letzte Referenz freigeben“ und „als abgeschlossen markieren“ ein Fenster, das zusätzliche Synchronisation erfordern würde. Nach der Zusammenführung kannref_decatomar „Zähler dekrementieren + prüfen, ob null erreicht“ ausführen und vermeidet ABA-ähnliche Probleme.
3.3 JoinHandle: Wie Ergebnisse über Task-Grenzen hinweg zurückgegeben werden
Intuitives Modell
JoinHandleist wie der „Abholschein“, den Ihnen ein Restaurant gibt. Wenn der Task (die Küche) fertig ist, wird das Gericht (output) an die Ausgabe (output) gestelltStage::Finishedund dann Ihr Abholsignal (waker) ausgelöst. Sie kommen mit dem Schein zum Abholen; der Schein selbst enthält nicht das Gericht, sondern ist nur ein Zeiger auf die Ausgabe. Wenn Sie den Schein verlieren (dropJoinHandle), wird das Gericht direkt weggeworfen (output wird gedroppt), aber die Küche stellt deshalb nicht die Arbeit ein.
Datenstruktur
JoinHandle<T>ist ebenfalls eine transparente Verpackung vonRawTask:
pub struct JoinHandle<T> {
raw: RawTask,
_p: PhantomData<T>,
}📎 tokio/src/runtime/task/join.rs:163-166
PhantomData<T>markiert den Ausgabetyp.JoinHandle<T>ist erst beiT: SendSend/Sync 📎 tokio/src/runtime/task/join.rs:169-170, was garantiert, dass Nicht-Send-Ausgaben nicht über Threads hinweg verschoben werden.
Schritt für Schritt: await auf einem JoinHandle
JoinHandleimplementiertFuture, dessenpollder Kern der Ergebnisrückgabe ist:
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
ready!(crate::trace::trace_leaf());
let mut ret = Poll::Pending;
let coop = ready!(crate::task::coop::poll_proceed(cx));
unsafe {
self.raw.try_read_output(&mut ret, cx.waker());
}
if ret.is_ready() {
coop.made_progress();
}
ret
}📎 tokio/src/runtime/task/join.rs:327-354
Beachten Sie einige Details:trace_leafdient der Tracing-Instrumentierung;coop::poll_proceedverbraucht das Kooperationsbudget (ausführlich in Kapitel 12);try_read_outputlöscht Generics über die vtable, legt den Rückgabewert auf den Stack und übergibt ihn mit*mut ()an📎 tokio/src/runtime/task/join.rs:327-354. Diese Technik, den Rückgabewert auf den Stack zu legen, ist nötig, weil vtable-Funktionen den Rückgabetyp nicht generisch machen könnenTund nur über einen Rohzeiger zurückschreiben können.
try_read_outputinterne Logik (in raw.rs, in diesem Kapitel nicht als Quellcode bereitgestellt): zuerstCOMPLETE-Bit prüfen; wenn bereits gesetzt, danntake_outputaufrufen, um das Ergebnis ausStage::Finishedzu entnehmen; andernfallscx.waker()im FeldTrailer::wakerregistrieren undPendingzurückgeben. Der Registrierungsprozess folgt genau demJOIN_WAKER-Protokoll aus Abschnitt 3.2.
Eigentumsübertragung des Ergebnisses
Der Abschnitt „Non-Send output“ der Moduldokumentation beschreibt präzise die Eigentumsregeln des Ergebnisses📎 tokio/src/runtime/task/mod.rs:151-170:
- Wenn der Task abgeschlossen ist, wird output in
Stagegelegt, dann wird die Umwandlung „COMPLETE setzen“ ausgeführt und der aktuelleJOIN_INTEREST-Wert gelesen. - Wenn
JOIN_INTEREST0 ist (keinJoinHandle), wird output sofort gedroppt📎tokio/src/runtime/task/mod.rs:157-158。 - Wenn
JOIN_INTEREST1 ist,JoinHandlefür die Bereinigung von output verantwortlich📎tokio/src/runtime/task/mod.rs:160-161。
Für Nicht-Send-output gibt die Dokumentation eine dreistufige Argumentation: output wird auf dem Thread erzeugt, der das Future pollt;JoinHandle<Output>ist ebenfalls Nicht-Send, wenn Output Nicht-Send ist, also liegt es auch auf dem Spawn-Thread; daher wirdJoinHandlebeim Entnehmen oder Droppen von output nicht über Threads hinweg verschoben📎 tokio/src/runtime/task/mod.rs:164-170。
Drop von JoinHandle: zwei Pfade, schnell und langsam
impl<T> Drop for JoinHandle<T> {
fn drop(&mut self) {
if self.raw.state().drop_join_handle_fast().is_ok() {
return;
}
self.raw.drop_join_handle_slow();
}
}📎 tokio/src/runtime/task/join.rs:358-364
drop_join_handle_fastversucht, mit einem CAS „JOIN_INTEREST-Bit löschen + Referenzzähler dekrementieren“ abzuschließen. Wenn das fehlschlägt (z. B. der Task schließt gerade ab und das Statusbit ist belegt), wird der langsame Pfad vondrop_join_handle_slowgenommen. Dies ist das typische Muster „optimistischer schneller Pfad + pessimistischer langsamer Pfad“.
Designüberlegung: Warum JoinHandle output nicht direkt hält
WennJoinHandleoutput direkt halten würde, müsste output beim Abschluss des Tasks auf den Thread verschoben werden, auf demJoinHandleliegt. AberJoinHandlekann auf einen beliebigen Thread verschoben werden (solangeT: Send), während der Erzeugungsthread von output der Poll-Thread ist. Direktes Halten würde eine threadübergreifende Verschiebung bewirken, bei der „output im Poll-Thread erzeugt, aber im Join-Thread gedroppt wird“, was bei Nicht-Send-output direkt das Typsystem verletzt. Tokio entscheidet sich, output inCellzu belassen (Stage::Finished),JoinHandlehält nurCell, das aufRawTaskzeigt, und entnimmt das Ergebnis übertake_outputan Ort und Stelle. So erfolgt das Drop von output auf dem Thread, auf demJoinHandleliegt, aber nur unter der Voraussetzung, dass dieser Thread mit dem Poll-Thread identisch ist (was im Nicht-Send-Fall gilt).
3.4 Lokale Warteschlange: Produzenten-Konsumenten-Struktur für Work-Stealing
Intuitives Modell
Jeder Worker hat eine „private To-do-Liste“ (lokale Warteschlange) mit Kapazität 256. Der Worker selbst entnimmt Aufgaben vomKopf(LIFO, nutzt Cache-Lokalität), andere Worker stehlen Aufgaben vomEnde(FIFO, nehmen die ältesten und wahrscheinlich bereits abgeschlossenen Aufgaben). Ohne lokale Warteschlange würden alle Aufgaben in der globalen Warteschlange liegen, jede Aufgabenentnahme müsste um das globale Lock konkurrieren, und die Mehrkern-Skalierbarkeit würde zusammenbrechen.
Speicherlayout: Trennung von head und tail
pub(crate) struct Inner<T: 'static> {
head: AtomicUnsignedLong,
tail: AtomicUnsignedShort,
buffer: Box<[UnsafeCell<MaybeUninit<task::Notified<T>>>; LOCAL_QUEUE_CAPACITY]>,
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:36-57
headistAtomicUnsignedLong(64 Bit, falls die Plattform u64 unterstützt),tailistAtomicUnsignedShort(32 Bit). Der Kommentar erklärt, warum die Indizes breiter als eigentlich nötig sind: zur ABA-Milderung und zur Unterscheidung zwischen „voll“ und „leer“ des Puffers📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:37-49。
headpackt internzwei UnsignedShort:Das niedrige Bit ist der „echte Kopf“ (real head), das hohe Bit ist die „erste Position, die der Dieb gerade verarbeitet“ (steal head). Wenn beide gleich sind, gibt es keinen aktiven Dieb📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:39-49. Diese Doppelwert-Packung ist der Kerntrick der work-stealing Queue: Der Dieb aktualisiert zuerst per CAS den steal-Wert, um eine Charge von Aufgaben zu „beanspruchen“, und nach Abschluss zieht er den steal-Wert auf den real-Wert nach, was das Ende des Stealens anzeigt.
LOCAL_QUEUE_CAPACITYist unter Nicht-loom 256, unter loom auf 4 reduziert, um mehr Randfälle zu testen📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:62-69。MASK = LOCAL_QUEUE_CAPACITY - 1, verwendet für den Ringpuffer-Index📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:71。
Schritt für Schritt: Die vollständigen Verzweigungen von push_back_or_overflow
Dies ist die komplexeste Funktion der lokalen Queue; wir analysieren sie Zweig für Zweig:
pub(crate) fn push_back_or_overflow<O: Overflow<T>>(
&mut self,
mut task: task::Notified<T>,
overflow: &O,
stats: &mut Stats,
) {
let tail = loop {
let head = self.inner.head.load(Acquire);
let (steal, real) = unpack(head);
let tail = unsafe { self.inner.tail.unsync_load() };
if tail.wrapping_sub(steal) < LOCAL_QUEUE_CAPACITY as UnsignedShort {
break tail;
} else if steal != real {
overflow.push(task);
return;
} else {
match self.push_overflow(task, real, tail, overflow, stats) {
Ok(_) => return,
Err(v) => { task = v; }
}
}
};
self.push_back_finish(task, tail);
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:188-223
Drei Verzweigungen:
1. Kapazität vorhanden(tail - steal < CAPACITY):break tail, nach Verlassen der Schleife wirdpush_back_finishin den Puffer geschrieben.
2. Keine Kapazität, aber gleichzeitige Diebe(steal != real): Der Dieb schafft Platz, also wird nur die aktuelle Aufgabe in die Injektions-Queue geschoben und sofort zurückgekehrt📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:204-208。
3. Keine Kapazität und keine Diebe: Aufruf vonpush_overflow, um die hintere Hälfte der Aufgaben in die Injektions-Queue überlaufen zu lassen📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:209-219. Wenn CAS fehlschlägt (gegen einen gleichzeitigen Dieb verliert),push_overflowgibtErr(task)zurück, Schleife wiederholt.
push_back_finishschreibt die Aufgabe und aktualisiert tail:
fn push_back_finish(&self, task: task::Notified<T>, tail: UnsignedShort) {
let idx = tail as usize & MASK;
self.inner.buffer[idx].with_mut(|ptr| {
unsafe { ptr::write((*ptr).as_mut_ptr(), task); }
});
self.inner.tail.store(tail.wrapping_add(1), Release);
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:226-244
ReleaseDie Reihenfolge garantiert, dass die geschriebene Aufgabe für Diebe sichtbar ist.
push_overflow: Warum die hintere Hälfte überlaufen lassen
const NUM_TASKS_TAKEN: UnsignedShort = (LOCAL_QUEUE_CAPACITY / 2) as UnsignedShort;📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:265
Beim Überlauf werden 128 Aufgaben entnommen. Der Kommentar erklärt ausführlich, warumdie hintere Hälfteund nicht die vordere Hälfte entnommen wird📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:295-306: Beim Entnehmen von Aufgaben aus der Injektions-Queue werden sie immer im vorderen Teil platziert. Wenn eine Aufgabe also im hinteren Teil liegt, kann man sicher sein, dass sie nicht gerade aus der Injektions-Queue entnommen wurde. Dies garantiert, dass „eine aus der Injektions-Queue entnommene Aufgabe nicht sofort wieder in die Injektions-Queue zurückgelegt wird“ (zumindest bevor sie einmal gepollt wurde).
CAS beansprucht die hintere Hälfte:
if self.inner.head.compare_exchange_weak(
pack(head, head), pack(tail, tail), Release, Relaxed
).is_err() {
return Err(task);
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:283-293
Aktualisiertheadvon(head, head)auf(tail, tail), d. h. steal und real werden gleichzeitig auf tail vorgerückt und alle Aufgaben beansprucht. Nach Erfolg wird tail auftail + NUM_TASKS_TAKENzurückgesetzt, was anzeigt, dass die vordere Hälfte in der lokalen Queue verbleibt📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:314-316。
pop und steal_into: Die zwei Pfade zum Abholen von Aufgaben
popist, wenn der Worker selbst eine Aufgabe abholt (vom Kopf, LIFO):
pub(crate) fn pop(&mut self) -> Option<task::Notified<T>> {
let mut head = self.inner.head.load(Acquire);
let idx = loop {
let (steal, real) = unpack(head);
let tail = unsafe { self.inner.tail.unsync_load() };
if real == tail { return None; }
let next_real = real.wrapping_add(1);
let next = if steal == real {
pack(next_real, next_real)
} else {
assert_ne!(steal, next_real);
pack(steal, next_real)
};
let res = self.inner.head.compare_exchange_weak(head, next, AcqRel, Acquire);
match res {
Ok(_) => break real as usize & MASK,
Err(actual) => head = actual,
}
};
Some(self.inner.buffer[idx].with(|ptr| unsafe { ptr::read(ptr).assume_init() }))
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:361-399
Kritische Verzweigung: Wennsteal == real(kein Dieb), werden beide gleichzeitig vorgerückt; andernfalls wird nur real vorgerückt und steal unverändert gelassen📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:377-384。assert_ne!(steal, next_real)stellt sicher, dass real nicht auf die Position von steal vorgerückt wird, da sonst der Beanspruchungszustand des Diebes zerstört würde.
steal_intoist der Steal-Pfad; zuerst wird geprüft, ob die Ziel-Queue genügend Platz hat:
if dst_tail.wrapping_sub(steal) > LOCAL_QUEUE_CAPACITY as UnsignedShort / 2 {
return None;
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:431-435
Wenn die Ziel-Queue mehr als halb voll ist, wird nicht gestohlen, um zu vermeiden, dass nach dem Stehlen sofort wieder überlaufen wird.
steal_into2ist der Kern des Stealens und berechnet die Anzahl der zu stehlenden Aufgaben:
let n = src_tail.wrapping_sub(src_head_real);
let n = n - n / 2;📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:487-488
Die Hälfte stehlen (aufgerundet). Dann per CAS den steal-Wert von head aktualisieren, um zu beanspruchen:
let steal_to = src_head_real.wrapping_add(n);
next_packed = pack(src_head_steal, steal_to);
let res = self.0.head.compare_exchange_weak(prev_packed, next_packed, AcqRel, Acquire);📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:496-506
Beachten Sie, dass hier nur der real-Wert aktualisiert wird (pack(src_head_steal, steal_to)steal bleibt unverändert), wodurch real aufsteal_tovorgerückt wird. Dies zeigt an: „Diese Aufgaben wurden beansprucht, andere Diebe dürfen sie nicht mehr anfassen.“ Nach Abschluss des Stealens wird steal auf real nachgezogen:
loop {
let head = unpack(prev_packed).1;
next_packed = pack(head, head);
let res = self.0.head.compare_exchange_weak(prev_packed, next_packed, AcqRel, Acquire);
match res {
Ok(_) => return n,
Err(actual) => prev_packed = actual,
}
}📎 tokio/src/runtime/scheduler/multi_thread/queue.rs:548-561
Das folgende Sequenzdiagramm stellt die gleichzeitige Interaktion der drei Parteien „Produzent pusht, Konsument poppt, Dieb stiehlt“ dar:
sequenceDiagram
participant P as "Worker A (生产者)"
participant Q as "Local 队列 Inner"
participant C as "Worker A (消费者 pop)"
participant S as "Worker B (窃取者)"
P->>Q: "load head (Acquire)"
P->>Q: "unsync_load tail"
Note over P: "tail - steal < 256?"
P->>Q: "push_back_finish: buffer[idx] = task"
P->>Q: "store tail+1 (Release)"
C->>Q: "load head (Acquire)"
C->>Q: "unsync_load tail"
Note over C: "real == tail? 空则返回 None"
C->>Q: "CAS head: pack(real+1, real+1)"
Q-->>C: "Ok, 读取 buffer[real & MASK]"
S->>Q: "load head (Acquire)"
S->>Q: "load tail (Acquire)"
Note over S: "src_head_steal != src_head_real? 返回 0"
S->>Q: "CAS head: pack(steal, real+n) 认领一半"
Q-->>S: "Ok, 拷贝 n 个任务到 dst"
S->>Q: "CAS head: pack(real+n, real+n) 完成窃取"
Q-->>S: "返回 n"Designüberlegung: Warum die lokale Queue LIFO ist und das Stehlen FIFO
Der Worker selbst holt vom Kopf (LIFO), weil die zuletzt eingefügte Aufgabe am wahrscheinlichsten noch im CPU-Cache liegt und am wahrscheinlichsten eine Aufgabe ist, die „gerade aufgeweckt wurde und deren Daten noch heiß sind“. Der Dieb holt vom Ende (FIFO), weil die älteste Aufgabe am wahrscheinlichsten schon den Großteil der Arbeit erledigt hat und das Stehlen dieser Aufgabe die Last des Opfers am schnellsten reduziert. Diese Kombination aus „LIFO lokal + FIFO stehlen“ ist das klassische Design der work-stealing-Scheduling und vereint Cache-Lokalität mit Lastausgleich.
Damit hat die Aufgabe ihre Verwandlung von Future zu einer planbaren Entität abgeschlossen: Sie hat einen Referenzzähler erhalten, wurde in dasCellSpeicherlayout eingeordnet und erfolgreich an die lokale Queue des Workers oder die globale Injektions-Queue zugestellt. Aber die Aufgabe in die Queue zu legen ist nur der Anfang; was sie wirklich zum Laufen bringt, ist die Scheduling-Schleife des Worker-Threads. Im nächsten Kapitel betreten wir die zweite Hälfte von „Das Leben einer Aufgabe“ und verfolgen, wie der Worker Aufgaben aus der Queue holt,Future::pollaufruft und bei Rückgabe vonPendingdurchWakereine Weckung registriert, was schließlichscheduleerneut in die Queue einreiht – der vollständige Aufrufpfad des geschlossenen Kreises „Wecken → Einreihen → erneut pollen“ sowie die work-stealing-Strategie und die LIFO-Slot-Optimierung werden dort enthüllt.
Kapitel 4: Das Leben einer Aufgabe (Teil 2): Scheduling-Schleife, poll und der geschlossene Kreis des Weckens
Von der Queue zur Ausführung: Das Skelett der Worker-Hauptschleife
Im vorherigen Kapitel haben wir die Aufgabe in dieLocalQueue oder die globale Injektions-Queue geschickt. Aber die Queue ist nur eine „To-do-Liste“; was die Aufgabe wirklich zum Laufen bringt, ist die niemals endende Schleife im Worker-Thread. In diesem Kapitel verfolgen wirContext::run– sie ist das Herz des gesamten Multithread-Schedulers.
Zuerst eine Intuition: Der Worker-Thread ist wie ein Koch, vor sich einen Stapel eigener Bestellungen (run_queue), daneben ein öffentliches Bestellregal (inject). Der Koch schaut zuerst auf die ihm am nächsten liegende Bestellung (lifo_slot), wenn nicht, nimmt er von seinem eigenen Stapel, wenn auch dort nichts ist, greift er sich eine Handvoll vom öffentlichen Regal, und wenn das auch nicht funktioniert, stiehlt er ein paar Blätter vom Stapel eines anderen Kochs. Erst wenn alles leer ist, geht er sich ausruhen, aber beim Ausruhen bleiben seine Ohren gespitzt – sobald eine Bestellung eingeht, wacht er sofort auf.
Ohne diese Schleife würde eine Aufgabe nach dem Einreihen in die Warteschlange für immer dort liegen,Future::pollniemals aufgerufen werden, und die gesamte Laufzeit wäre nur ein Haufen toter Daten.
Speicherlayout und Zustandsfelder des Core
Der gesamte veränderliche Zustand des Workers ist inCoreuntergebracht, der vonBoxauf dem Heap allokiert wird und überAtomicCell<Core>zwischenWorkerund thread-lokalemContextübergeben wird.
CoreDie Schlüsselfelder von📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:113-167:
tick: u32sind folgende: Bei jeder Schleifeniteration inkrementiert, um periodisch Wartung (maintenance) und globale Warteschlangenprüfung auszulösen.lifo_slot: Option<Notified>:LIFO-Slot, dies ist das raffinierteste Design dieses Kapitels. Wenn ein Worker selbst eine Aufgabe einplant, geht sie nicht inrun_queue, sondern wird in diesen Slot gelegt, und beim nächsten Abrufen einer Aufgabewirdbevorzugt von hier genommen.lifo_enabled: bool: Der Schalter für den LIFO-Slot, um Hunger in Ping-Pong-Szenarien zu verhindern.run_queue: queue::Local<Arc<Handle>>: Lokale Warteschlange, die im vorherigen Kapitel analysierteLocal-Struktur.is_searching: bool: Ob der Worker gerade nach stehlbaren Aufgaben sucht.is_shutdown: bool/is_traced: bool: Shutdown- und Tracing-Flags.park: Option<Parker>: Parker, mitOptionumhüllt, um ihn unter dem Borrow-Checker bequem herausnehmen/zurücklegen zu können.global_queue_interval: u32: Wie oft die globale Warteschlange überprüft wird.rand: FastRand: Schneller Zufallszahlengenerator, um den Startpunkt für das Stehlen zufällig zu wählen.
Beachten Sie, dasslifo_sloteinOption<Notified>und keine Warteschlange ist – es speichert nureineAufgabe. Die Design-Motivation wird in den Quellcode-Kommentaren klar erläutert📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:117-121: Vom Worker selbst eingeplante Aufgaben werden in diesem Slot gespeichert, und der Worker prüft ihnrun_queue bevorer
prüft, mit dem Effekt „die zuletzt eingeplante Aufgabe läuft als nächste" (LIFO). Dies dient der Verbesserung der Lokalität, ist besonders effektiv für Message-Passing-Muster und kann die Latenz senken.
Warum kann LIFO die Latenz senken? Betrachten Sie ein typisches Message-Passing-Szenario: Aufgabe A weckt nach der Verarbeitung einer Nachricht Aufgabe B, B weckt nach der Verarbeitung wieder A. Wenn B sofort läuft, nachdem A es geweckt hat, befinden sich die von B benötigten Daten wahrscheinlich noch im CPU-Cache (da A sie gerade berührt hat). Wenn B ans Ende der Warteschlange gestellt wird und darauf wartet, dass Dutzende Aufgaben davor abgearbeitet werden, ist der Cache längst überschrieben.MAX_LIFO_POLLS_PER_TICK = 3Aber LIFO birgt ein Hunger-Risiko. Der Quellcode verwendet📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:263-263, um
zu begrenzen: Pro Tick wird der LIFO-Slot höchstens 3 Mal bevorzugt, danach wird er deaktiviert, damit andere Aufgaben eine Chance zur Ausführung bekommen.
Walkthrough der Hauptschleife: Ein vollständiger Scheduling-ZyklusparkWir versetzen uns in ein konkretes Szenario: Worker 0 ist gerade ausrun_queueaufgewacht,lifo_slotenthält 5 Aufgaben,
enthält 1 Aufgabe, die globale Warteschlange enthält 3 Aufgaben.Context::run 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:570-642Der Einstieg in die Hauptschleife istlifo_enabled. Zuerst wirdblock_in_placezurückgesetzt (da der Core möglicherweise von📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:571-573gestohlen wurde und der Zustand zurückgesetzt werden muss),while !core.is_shutdowndann wird in die
-Schleife eingetreten.
Jede Schleifeniteration erledigt vier Dinge: core.tick()Erster Schritt: Tick und Wartung.📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:587inkrementiert den Zählerself.maintenance(core). Dann prüfttick % event_interval == 0park_yield, und falls ja, wird📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:809-826。
mit 0-Timeout aufgerufen, um I/O und Timer anzutreiben. core.next_task(&self.worker)Zweiter Schritt: Aufgabe abrufen.📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1090-1156ist die zentrale Aufgabenabruflogik
- . Sie hat zwei Pfade:
tick % global_queue_interval == 0Wenn,wird📎tokio/src/runtime/scheduler/multi_thread/worker.rs:1091-1098bevorzugt aus der globalen Warteschlange abgerufen, und wenn das fehlschlägt, wird die lokale - abgerufen. Dies dient dazu, zu verhindern, dass Aufgaben in der globalen Warteschlange verhungern.Andernfallswird📎
tokio/src/runtime/scheduler/multi_thread/worker.rs:1090-1156。
bevorzugt die lokale Aufgabe abgerufennext_local_task. Das lokale Abrufen von Aufgaben wird von📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1158-1160:
fn next_local_task(&mut self) -> Option<Notified> {
self.lifo_slot.take().or_else(|| self.run_queue.pop())
}Kopieren
Zuerst wird der LIFO-Slot abgerufen, dann der Kopf der Warteschlange (LIFO-Pop). Das ist das im vorherigen Kapitel erwähnte „lokale LIFO".Wenn lokal leer ist, aber die globale Warteschlange nicht leer ist, wird der Workerbatchweise📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1110-1154Aufgaben aus der globalen Warteschlange ziehenn. Die Berechnung der Batch-Größemin(inject.len() / remotes.len() + 1, cap)ist wohlüberlegt:cap, wobeimin(remaining_slots, max_capacity / 2)wiederum📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1120-1131annimmt. Die Quellcode-Kommentare erklären, warum auf die Hälfte der Warteschlangenkapazität begrenzt wird: Um sicherzustellen, dass die gezogenen Aufgaben in dervorderen Hälfte
der lokalen Warteschlange landen, sodass diese Aufgaben selbst bei späterem Überlauf nicht zurück in die globale Warteschlange geschoben werden (Überlauf betrifft nur die hintere Hälfte).Dritter Schritt: Aufgabe ausführen.run_task 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:647-796Nachdem
eine Aufgabe erhalten hat, wirdaufgerufen. Dies ist die komplexeste Funktion dieses Kapitels, die wir im nächsten Abschnitt gesondert behandeln.next_taskVierter Schritt: Stehlen oder Parken.NoneWennsteal_work 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1167-1195parkzurückgibt, bedeutet dies, dass weder lokal noch global Arbeit vorhanden ist, undpark_yield 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:613-621。
wird aufgerufen. Bei fehlgeschlagenem Stehlen wird
flowchart TD
start["Context::run 进入循环"] --> tick["core.tick() 自增"]
tick --> maint{"tick % event_interval == 0?"}
maint -->|是| park_yield["park_yield 驱动 I/O 与定时器"]
maint -->|否| next
park_yield --> next["core.next_task()"]
next --> has_task{"取到任务?"}
has_task -->|是| run_task["run_task 执行 poll"]
run_task --> cont{"core 还在?"}
cont -->|是| tick
cont -->|否| ret["return 退出"]
has_task -->|否| steal["core.steal_work()"]
steal --> steal_ok{"窃取成功?"}
steal_ok -->|是| run_task
steal_ok -->|否| defer_check{"defer 非空?"}
defer_check -->|是| py["park_yield"]
defer_check -->|否| pk["park 阻塞等待"]
py --> tick
pk --> tickbetreten. Der gesamte Kontrollfluss ist wie folgt:
run_taskKopierenpollrun_task: Der geschlossene Kreislauf von poll und LIFO-Slot
ist der Ort, an dem die Aufgabe tatsächlichassert_owner 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:648wird, und auch der Schließpunkt des Kreislaufs „Aufwecken → Einreihen → erneut pollen".NotifiedDas Erste, was nach dem Betreten der Funktion geschieht, istTask, wobei
intransition_from_searching 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:652umgewandelt wird, während gleichzeitig per Debug-Assertion sichergestellt wird, dass der aktuelle Thread tatsächlich der Owner dieser Aufgabe ist.
Danach📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:695-795:
coop::budget(|| {
task.run();
let mut lifo_polls = 0;
loop {
let mut core = match self.core.borrow_mut().take() {
Some(core) => core,
None => return ControlFlow::Break(()),
};
let task = match core.lifo_slot.take() {
Some(task) => task,
None => {
self.reset_lifo_enabled(&mut core);
core.stats.end_poll();
return ControlFlow::Continue(core);
}
};
if !coop::has_budget_remaining() {
core.run_queue.push_back_or_overflow(task, ...);
return ControlFlow::Continue(core);
}
lifo_polls += 1;
if lifo_polls >= MAX_LIFO_POLLS_PER_TICK {
core.lifo_enabled = false;
}
let task = self.worker.handle.shared.owned.assert_owner(task);
*self.core.borrow_mut() = Some(core);
task.run();
}
})Dann folgt die entscheidende Budget-Umhüllungtask.run()KopierenFuture::pollDieser Codeabschnitt offenbart den vollständigen geschlossenen Kreislauf des LIFO-Slots:schedule_localführtlifo_slot 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1396-1408aus, und wenn die Aufgabe während des Polls sich selbst oder eine andere Aufgabe aufweckt,lifo_slotlegtdie neue Aufgabe in. Nach der Rückkehr des Polls prüft die Schleife sofort
, und falls eine Aufgabe vorhanden ist, wird weiter ausgeführt –lifo_slotohne zur Hauptschleife zurückzukehren
, wird direkt innerhalb desselben Budgets kontinuierlich gepollt.self.core.borrow_mut().take()Dies ist die Verkörperung von „Aufwecken → Einreihen → erneut pollen" auf dem LIFO-Pfad: Beim Aufwecken wird die Aufgabe inNonegelegt, und nach der Rückkehr des Polls wird sie sofort entnommen und erneut gepollt, wodurch ein enger geschlossener Kreislauf entsteht.📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:716-724Beachten Sie denblock_in_place-Zweig vonControlFlow::Break(()): Wenn der Core gestohlen wurde (z. B. wenn in der AufgabeContext::runaufgerufen wurde), muss der Workerblock_in_placeInteraktionspunkte mit der Scheduler-Schleife.
Aufwachpfad: Wie der Waker das erneute Einreihen auslöst
WennFuture::pollzurückgibtPending, muss die Task einenWakerregistrieren, um beim Bereitwerden des Events aufgeweckt zu werden. TokiosWaker-Implementierung ist extrem schlank – sie ist lediglich ein Rohzeiger auf die TaskHeaderplus eine vtable.
waker_refKonstruiertWakerRef 📎 tokio/src/runtime/task/waker.rs:11-34, umhüllt mitManuallyDropdieWaker, um beim Drop das Dekrementieren des Referenzzählers zu vermeiden. Die vtable ist statisch📎 tokio/src/runtime/task/waker.rs:119-119:
static WAKER_VTABLE: RawWakerVTable =
RawWakerVTable::new(clone_waker, wake_by_val, wake_by_ref, drop_waker);Alle vier Funktionen stellen lediglich den Rohzeiger wieder alsHeaderher und rufen dann die entsprechende Methode vonRawTaskauf📎 tokio/src/runtime/task/waker.rs:70-116. Zum Beispiel ruftwake_by_refletztendlichraw.wake_by_ref() 📎 tokio/src/runtime/task/waker.rs:106-116。
wake_by_refauf. Die Semantik ist: Den Task-Status vonPENDINGnachSCHEDULEDüberführen, und wenn die Überführung erfolgreich ist (d. h. zuvor tatsächlich PENDING war),Schedule::scheduleaufrufen, um die Task erneut einzureihen.
Für den Multithread-Scheduler,scheduledie Implementierung inHandle::schedule_task 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1353-1376:
pub(super) fn schedule_task(&self, task: Notified, is_yield: bool) {
with_current(|maybe_cx| {
if let Some(cx) = maybe_cx {
if self.ptr_eq(&cx.worker.handle) {
if let Some(core) = cx.core.borrow_mut().as_mut() {
self.schedule_local(core, task, is_yield);
return;
}
}
}
self.push_remote_task(task);
self.notify_parked_remote();
});
}Die Logik teilt sich in zwei Zweige:
- Wenn der aktuelle Thread der Worker dieses Schedulers ist und core hält, gehe über
schedule_local📎tokio/src/runtime/scheduler/multi_thread/worker.rs:1385-1417– ablegen im LIFO-Slot oder in der lokalen Queue. - Andernfalls (Aufwecken von einem externen Thread oder core wurde gestohlen), gehe über
push_remote_taskEinfügen in die globale Injektions-Queue undnotify_parked_remoteeinen geparkten Worker aufwecken📎tokio/src/runtime/scheduler/multi_thread/worker.rs:1379-1383。
schedule_localteilt sich intern wiederum in zwei Zweige📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1385-1417: Wenn esyieldist oder LIFO deaktiviert wurde, einfügen amrun_queueEnde; andernfalls ablegen inlifo_slot, und die Task aus dem ursprünglichen Slot an das Ende der Queue verdrängen.
sequenceDiagram
participant Future as "Future::poll"
participant Waker as "Waker(wake_by_ref)"
participant RawTask as "RawTask::wake_by_ref"
participant Handle as "Handle::schedule_task"
participant Core as "Core(schedule_local)"
participant Inject as "InjectQueue"
participant Parker as "Unparker"
Future->>Waker: "返回 Pending, 注册 waker"
Note over Future: "事件就绪(如 epoll)"
Waker->>RawTask: "raw.wake_by_ref()"
RawTask->>RawTask: "state: PENDING -> SCHEDULED"
RawTask->>Handle: "schedule(Notified)"
alt 当前线程是同一 worker 且持有 core
Handle->>Core: "schedule_local: 放入 lifo_slot"
else 外部线程或 core 被偷走
Handle->>Inject: "push_remote_task"
Handle->>Parker: "notify_parked_remote().unpark()"
endpark und unpark: Zustandsmaschine und Atomarität des Aufweckens
Wenn der Worker nichts zu tun hat, muss er parken, aber park/unpark ist die Stelle, die am anfälligsten für Race Conditions ist. Tokio löst dies mit einerAtomicUsize-Zustandsmaschine plusCondvarals Fallback.
InnerDie Felder von📎 tokio/src/runtime/scheduler/multi_thread/park.rs:31-43:state: AtomicUsize、mutex: Mutex<()>、condvar: Condvar、shared: Arc<Shared>. Es gibt vier Zustandskonstanten📎 tokio/src/runtime/scheduler/multi_thread/park.rs:36-45:
EMPTY = 0: nicht geparkt.PARKED_CONDVAR = 1: geparkt auf der condvar.PARKED_DRIVER = 2: geparkt auf dem I/O driver.NOTIFIED = 3: wurde bereits aufgeweckt.
Dies ist eine explizite Zustandsmaschine; wir verwenden sie, um das Zustandsdiagramm zu zeichnen (dies ist die einzige Stelle in diesem Kapitel, die diestateDiagram-v2-Zulassungsbedingung erfüllt – im Quellcode existieren tatsächlich diese vier Zustandskonstanten):
stateDiagram-v2
[*] --> Empty
Empty --> ParkedCondvar : "park_condvar() CAS(EMPTY->PARKED_CONDVAR)"
Empty --> ParkedDriver : "park_driver() CAS(EMPTY->PARKED_DRIVER)"
Empty --> Notified : "unpark() swap(NOTIFIED)"
ParkedCondvar --> Empty : "condvar 唤醒后 CAS(NOTIFIED->EMPTY)"
ParkedCondvar --> Empty : "超时 swap(EMPTY)"
ParkedDriver --> Empty : "driver 返回后 swap(EMPTY)"
Notified --> Empty : "park() CAS(NOTIFIED->EMPTY) 消费通知"
Notified --> Notified : "再次 unpark() swap(NOTIFIED)"unparkDie Implementierung von📎 tokio/src/runtime/scheduler/multi_thread/park.rs:277-290verwendetswapstatt CAS; der Quellcode-Kommentar erklärt den Grund📎 tokio/src/runtime/scheduler/multi_thread/park.rs:277-290: Es muss eine Release-Operation ausgeführt werden, damit der parkende Thread die Schreibvorgänge vor dem unpark beobachten kann, daher muss selbst wenn state bereitsNOTIFIEDist, einmal geschrieben werden.
parkVersucht zunächst, eine vorhandene Benachrichtigung zu konsumieren📎 tokio/src/runtime/scheduler/multi_thread/park.rs:132-149: Wenn CASNOTIFIED -> EMPTYerfolgreich ist, bedeutet dies, dass zuvor bereits aufgeweckt wurde, und es wird direkt zurückgekehrt ohne zu blockieren. Andernfalls wird versucht, das driver-Lock zu erlangen; wenn erhalten, wird auf dem driver geparkt, andernfalls wird die condvar als Fallback verwendet📎 tokio/src/runtime/scheduler/multi_thread/park.rs:143-148。
park_condvarEs gibt eine klassische doppelte Prüfung📎 tokio/src/runtime/scheduler/multi_thread/park.rs:162-180: Zuerst CASEMPTY -> PARKED_CONDVAR, wenn dies fehlschlägt und esNOTIFIEDist, bedeutet dies, dass vor dem Setzen des Zustands bereits aufgeweckt wurde; in diesem Fall mussswap(EMPTY)ausgeführt werden, um die Schreibvorgänge des unpark zu synchronisieren📎 tokio/src/runtime/scheduler/multi_thread/park.rs:167-177. Der Kommentar betont besonders: Selbst wenn bekannt ist, dass esNOTIFIEDist, muss einmal gelesen werden, da unpark möglicherweise nach unserem Lesen vonNOTIFIEDerneut aufgerufen wurde.
unpark_condvarDer Kommentar von📎 tokio/src/runtime/scheduler/multi_thread/park.rs:292-307weist auf die klassische Falle der condvar hin: Zwischen dem Setzen desPARKED-Zustands durch den parkenden Thread und dem tatsächlichenwaitgibt es ein Zeitfenster; wenn in diesem Zeitraum notify aufgerufen wird, wird es ignoriert. Die Lösung ist, dass der parkende Thread zu diesem Zeitpunktmutexhält, und der unparkende Thread zuerstdrop(self.mutex.lock())das Lock erwirbt (und somit darauf wartet, dass der parkende Thread es freigibt), dannnotify_one。
Designüberlegung: Warum der LIFO-Slot ein einzelner Slot und keine Queue ist
Das Einzel-Slot-Design ist eine bewusste Abwägung. Bei einer Queue müsste bei jedem Aufwecken eingereiht und bei jeder Task-Entnahme ausgereiht werden, was teurer wäre; außerdem würde die Queue mehrere Tasks ansammeln und die Lokalitätsannahme „der zuletzt aufgeweckte läuft zuerst" zerstören. Die Semantik des Einzel-Slots ist „sich nur den letzten merken"; verdrängte Tasks gehen in die normale Queue – was genau dem Gesetz des abnehmenden Lokalitätsnutzens entspricht: Die letzte Task ist am heißesten, die zweite weniger, und ab der dritten wird der Nutzen sehr gering.
MAX_LIFO_POLLS_PER_TICK = 3Diese magische Zahl📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:263-263ist ebenfalls ein Erfahrungswert. Der Quellcode-Kommentar besagt: „Ein paar Durchläufe durch den LIFO-Slot scheinen auszureichen, um von der Lokalität zu profitieren; mehr als 3 könnten übergewichten." Dies verhindert, dass ein Ping-Pong-Szenario, bei dem A B aufweckt und B A aufweckt, andere Tasks aushungert.
Ein weiteres bemerkenswertes Design iststeal_workdie „Halbierungs-Such"-Strategie📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1158-1160: Nur wenn weniger als die Hälfte der Worker suchen, versucht ein neuer Worker tatsächlich zu stehlen. Dies vermeidet CAS-Konkurrenz, die entsteht, wenn alle Worker gleichzeitig wild stehlen.transition_to_searchingKoordiniert durchidle.transition_worker_to_searching()wird📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1197-1203。
Das Stehlen beginnt an einem zufälligen Startpunkt📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1172-1174, durchläuft alle remote, überspringt sich selbst📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1179-1182, ruftsteal_intoauf, um einen Diebstahl zu versuchen. Nachdem alles fehlgeschlagen ist, wird auf die globale Queue zurückgegriffen📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1197-1203。
Zusammenfassung dieses Kapitels
Die Worker-HauptschleifeContext::runist das Herz des Schedulers: Nach jedem Tick wird zuerst eine Task geholt (LIFO-Slot → lokale Queue → globale Queue); wenn eine gefunden wird, wirdrun_taskausgeführt, um poll aufzurufen; wenn keine gefunden wird, wird gestohlen; wenn der Diebstahl fehlschlägt, wird geparkt.run_taskDie interne LIFO-Schleife komprimiert „Aufwecken → Einreihen → erneut poll" in dasselbe Budget und bildet so einen geschlossenen Regelkreis mit niedriger Latenz.WakerIst ein Rohzeiger plus statische vtable,wake_by_refLöst durch Zustandsübergangscheduleaus, und je nachdem, ob der aktuelle Thread derselbe Worker ist, wird entschieden, ob die lokale Queue oder die globale Queue verwendet wird.park/unparkVerwendet eine Vier-Zustands-Atommaschine plus condvar als Fallback und löst die klassische Race Condition des verlorenen Aufweckens.
Im nächsten Kapitel verlassen wir den Scheduler und betreten die I/O-Welt: Wie der Reactor epoll-Events inWakerAufwecken übersetzt, sodassAsyncFddiePendingzuReady。
Denkfragen und Selbsttests dieses Kapitels
Q1: Wenn mannext_local_taskso ändert, dass zuerst geholt wirdrun_queueDann nimmlifo_slot, welche Konsequenzen hat das in Szenarien mit intensivem Nachrichtenaustausch?
Referenzanalyse:next_local_taskDie aktuelle Implementierung istself.lifo_slot.take().or_else(|| self.run_queue.pop()) 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:1158-1160, zuerst wird der LIFO-Slot entnommen. Wenn umgekehrt zuerstrun_queueentnommen würde, dann würden Aufgaben, die gerade aufgeweckt wurden und deren Daten noch heiß sind, hinter anderen Aufgaben in der Warteschlange ausgeführt. Im Nachrichtenaustauschmuster A→B→A würde B nach dem Aufwecken nicht sofort laufen, sondern warten, bis andere Aufgaben in der Warteschlange abgearbeitet sind. Zu diesem Zeitpunkt könnten die von A geschriebenen Daten bereits aus dem CPU-Cache verdrängt worden sein, und der Lokalitätsvorteil ginge verloren. Noch schwerwiegender ist, dasslifo_slotAufgaben inrun_queueso lange warten, bis📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:117-121geleert ist, bevor sie ausgeführt werden, was die Latenz erheblich erhöht. Der Quellcode-Kommentar
Q2: park_condvarweist ausdrücklich darauf hin, dass diese Reihenfolge dazu dient, „die Lokalität zu verbessern, vom Nachrichtenaustauschmuster zu profitieren und die Latenz zu senken“.Err(NOTIFIED), welche Probleme entstehen, wenn imself.state.swap(EMPTY, SeqCst)-Zweigreturnentfernt und nur
beibehalten wird?ReferenzanalyseErr(NOTIFIED): Der Quellcode führt imlet old = self.state.swap(EMPTY, SeqCst) 📎 tokio/src/runtime/scheduler/multi_thread/park.rs:167-177-Zweig📎 tokio/src/runtime/scheduler/multi_thread/park.rs:168-173aus. Der Kommentar erklärtNOTIFIED: unpark könnte nach dem Lesen vonreturnnoch einmal aufgerufen worden sein; es muss eine acquire-Operation ausgeführt werden, um mit diesem unpark zu synchronisieren, damit alle davor erfolgten Schreibvorgänge sichtbar werden. Wenn nurNOTIFIEDohne Swap ausgeführt würde, bliebe state beiNOTIFIED -> EMPTYstehen; beim nächsten park würde CAS
Q3: run_taskerfolgreich sein und sofort zurückkehren (eine bereits abgelaufene Benachrichtigung würde konsumiert), aber schlimmer noch: Der release-Schreibvorgang von unpark wäre nicht synchronisiert, und der park-Thread könnte die vor unpark geschriebenen Daten nicht sehen, was zu Problemen mit der Speichersichtbarkeit führt. Dies ist ein typischer doppelter Bug aus „verlorenem Aufwecken + Speicherordnung“.self.core.borrow_mut().take(), warum wirdNonezurückgegeben, wennControlFlow::Break(())zurückgibt, und nichtContinue?
Referenzanalyse:self.core.borrow_mut().take()Die Rückgabe vonNonebedeutet, dass der Core gestohlen wurde📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:716-724. Der einzige Weg, auf dem ein Core gestohlen werden kann, ist, dass eine Aufgabe internblock_in_placeaufruft, wodurch übermaybe_move_runtimeder Core auscx.coreentnommen und an einen neuen Thread📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:473-497übergeben wird. Zu diesem Zeitpunkt besitzt der aktuelle Thread keine Scheduling-Fähigkeit mehr. WennContinue,Context::runzurückgegeben würde, würde die Schleife weiterlaufen und Methoden wiecore.next_task()aufrufen, die einen Core benötigen, aber der Core ist nicht mehr inself.core, was zu einem Panic oder inkonsistentem Zustand führen würde. Die Rückgabe vonBreaklässtContext::rundirektreturn 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:594-597, wodurch die Kontrolle an dierun-Funktion zurückgegeben wird, die die weitere Verarbeitung übernimmt (zum Beispielcx.defer.wake() 📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:564). Der Kommentar erklärt auch📎 tokio/src/runtime/scheduler/multi_thread/worker.rs:719-721: Zu diesem Zeitpunkt darfreset_lifo_enablednicht aufgerufen werden, weil der Core gestohlen wurde und der Dieb ihn oben inContext::runverarbeitet.
Kapitel 5: I/O-Bereitschaftsbenachrichtigung: Wie der Reactor epoll-Ereignisse in Waker-Aufweckungen übersetzt
Im vorherigen Kapitel haben wir die Hauptschleife des Worker-Threads verfolgt: Eine Aufgabe wird gepollt, bei Rückgabe von Pending wird der Waker irgendwo gespeichert, nach Eintritt der Bereitschaft wird der Waker ausgelöst und die Aufgabe erneut in die Warteschlange eingereiht. Aber wo genau ist dieses „irgendwo“? Wie wird der Waker gefunden, wenn ein epoll-Ereignis eintritt? Genau das ist die Frage, die der Reactor beantworten soll. Zuerst ein intuitives Modell: Stellen Sie sich den gesamten I/O-Bereitschaftsbenachrichtigungsmechanismus wie ein Nummernaufrufsystem in einem Restaurant vor – der Gast (die Aufgabe) stellt sich nach der Bestellung nicht ans Fenster und wartet, sondern nimmt einen Summer (Waker) mit zurück an den Platz; wenn die Küche (der Kernel-epoll) das Essen zubereitet hat, findet die Rezeption (der Reactor) anhand der Bestellnummer (Token) den entsprechenden Summer und drückt den Knopf. Ohne dieses System könnte jede Aufgabe nur den Socket pollen, und die CPU würde verbrannt; oder man würde blockierende Threads zum Warten verwenden, ein Thread pro Verbindung, was nicht skalierbar ist. Der Reactor von Tokio besteht aus drei Dateien mit einer dreischichtigen Struktur und strikter Trennung der Zuständigkeiten: driver.rs ist der Ereignisschleifenkern, hält mio::Poll, ist für den Aufruf von poll() zum blockierenden Warten auf Kernel-Ereignisse verantwortlich und übersetzt Ereignisse in Lese-/Schreibvorgänge auf ScheduledIo; registration.rs ist das benutzerseitige Registrierungshandle, das TcpStream intern hält, und bietet APIs wie poll_read_ready / poll_write_ready; scheduled_io.rs ist der Status-Slot jedes fd, speichert Lese-/Schreibbereitschaftsbits und Waker-Listen und ist die Brücke zwischen Ereignissen und Aufgaben. Die Modulzusammensetzung ist in tokio/src/runtime/io/mod.rs:5-16 zu finden: driver exportiert Driver, Handle, ReadyEvent, registration exportiert Registration, scheduled_io exportiert ScheduledIo. Die folgende Abbildung verankert den vollständigen Datenfluss, der in diesem Kapitel verfolgt werden soll: TcpStream → Registration → ScheduledIo → Handle/Driver → Kernel → zurück zu ScheduledIo → Waker. Als Nächstes zerlegen wir dies Schicht für Schicht.
Treiberschicht:DriverundHandleAufgabenteilung
Intuitives Modell
Driveristdie einzige Entität, diemio::Pollbesitzt, und kann nur in einem einzelnen Thread&mutzugegriffen werden – dies ist die Exklusivitätsanforderung der Ereignisschleife. Hingegen istHandleeinklonbarer, threadübergreifend gemeinsam nutzbarer Registrierungseinstiegspunkt, jeder Thread, der einen neuen fd registrieren möchte, tut dies darüber. Ohne diese Aufteilung müsste man entwedermio::Pollsperren (bei jeder Registrierung Konkurrenz), oder alle Registrierungen zurück zum Driver-Thread leiten (Einführung einer Cross-Thread-Message-Queue). Tokio entscheidet sich,Handledirekt einen Klon vonmio::Registryhalten zu lassen, Registrierungsoperationen können parallel ablaufen, nur das tatsächliche Warten auf Events benötigt Exklusivzugriff.
Speicherlayout und Felder
Zuerst schauen wir unsDriverdie Felder von📎 tokio/src/runtime/io/driver.rs:25-38:
signal_ready: boolan: Ob ein Unix-Signal-Event angekommen ist, wird für signal-getrieben verwendet.events: mio::Events: Haupt-Event-Puffer, wird überturnAufrufe hinweg wiederverwendet, um Allokation bei jedem Aufruf zu vermeiden.events_busy: Option<mio::Events>:Dedizierter Puffer für nicht-blockierendes poll, existiert nur, wennmax_io_events_per_busy_tickgesetzt ist.poll: mio::Poll: Kapselung der Kernel-Event-Queue.
Schauen wir uns nunHandle 📎 tokio/src/runtime/io/driver.rs:41-75:
registry: mio::Registry:mio::Poll::registry()den Klon vonregister/deregister。registrations: RegistrationSetan, verwendet fürToken: Menge aller aktiven Registrierungen, verantwortlich für die Zuweisung vonScheduledIo。synced: Mutex<registration_set::Synced>undRegistrationSet: Schützt den Synchronisationszustand vonwaker: mio::Waker: Wird verwendet, um aus beliebigen Threads den Driver aufzuwecken, der inturnblockiert.metrics: IoDriverMetrics: Zählt die Anzahl der fds und der bereiten Events.
Hier gibt es ein entscheidendes Design:events_busyDie Existenz von📎 tokio/src/runtime/io/driver.rs:25-38dient dazu,das Problem zu lösen, dass nicht-blockierendes poll Events verschluckt. Der Kommentar📎 tokio/src/runtime/io/driver.rs:189-190sagt es ganz klar: Wenn Events, die durch nicht-blockierendes poll entnommen wurden, im Hauptpuffer verbleiben, sind sie beim nächsten poll nicht mehr sichtbar; mit einem separaten Puffer bleiben unbehandelte Events in der Kernel-Queue und werden beim nächsten poll erneut zurückgegeben.
Step-by-Step: Eine Ausführung vonturn
turnist die Kernfunktion des Drivers📎 tokio/src/runtime/io/driver.rs:184-261. Angenommen, ein Worker-Thread stellt fest, dass keine Aufgaben auszuführen sind, und ruftpark → turn(handle, None)auf, um blockierend zu warten:
Erster Schritt: Assertion, dass nicht heruntergefahren📎 tokio/src/runtime/io/driver.rs:185, und Freigabe der zu bereinigenden Registrierungen📎 tokio/src/runtime/io/driver.rs:187。release_pending_registrationsPrüftneeds_release(), falls vorhanden, wirdregistrations.release() 📎 tokio/src/runtime/io/driver.rs:336-340。
aufgerufen. Zweiter Schritt: Auswahl des Event-Puffers📎 tokio/src/runtime/io/driver.rs:191-194. Fallsmax_waitnull ist undevents_busyexistiert, wird der busy-Puffer verwendet; andernfalls der Hauptpuffer.
Dritter Schritt: Aufruf vonself.poll.poll(events, max_wait) 📎 tokio/src/runtime/io/driver.rs:198. Hier wird tatsächlich in epoll_wait blockiert. Die Fehlerbehandlung ist sehr zurückhaltend:Interruptedwird direkt ignoriert (Signalunterbrechung ist normal)📎 tokio/src/runtime/io/driver.rs:200, unter WASI wirdInvalidInputebenfalls ignoriert📎 tokio/src/runtime/io/driver.rs:201-205, andere Fehler führen direkt zu panic📎 tokio/src/runtime/io/driver.rs:206。
Vierter Schritt: Iteration über die Events📎 tokio/src/runtime/io/driver.rs:211-233. Für jedesevent:
- falls
token == TOKEN_WAKEUP(Wert 0)📎tokio/src/runtime/io/driver.rs:214, wird nichts getan – dies wird vonunparkverwendet, um die Blockierung zu unterbrechen. - Falls
token == TOKEN_SIGNAL(Wert 1)📎tokio/src/runtime/io/driver.rs:216, wirdsignal_ready = true。 - gesetzt. Andernfalls handelt es sich um ein normales I/O-Event📎
tokio/src/runtime/io/driver.rs:218-231:mio::Readywird in TokiosReadyumgewandelt, mitEXPOSE_IO.from_exposed_addr(token.0)wird das Token zurück in einen*const ScheduledIo-Zeiger umgewandelt, dannset_readiness(Tick::Set, |curr| curr | ready)werden die Bereitschaftsbits akkumuliert, anschließendio.wake(ready)wird die entsprechende Richtung vonWaker。
ausgelöst. Hier istEXPOSE_IOeinPtrExposeDomain<ScheduledIo> 📎 tokio/src/runtime/io/mod.rs:21-22, das den Zeiger alsusize„exponiert“ alsmio::Token. Der Sicherheitskommentar📎 tokio/src/runtime/io/driver.rs:222-225erklärt, warum diese unsafe-Konvertierung sicher ist: Der Zeiger wird nicht freigegeben, bevor er bei mio abgemeldetundder Driver nicht mehr parallel pollt, und der Driver besitzt das Eigentum anArc<ScheduledIo>.
Fünfter Schritt: Verarbeitung der io_uring Completion-Queue (nur Linux + tokio_unstable)📎 tokio/src/runtime/io/driver.rs:235-258, einschließlich der Flush-Schleife bei CQ-Überlauf.
Sechster Schritt: Akkumulation der Metriken📎 tokio/src/runtime/io/driver.rs:265-267。
flowchart TD
start["turn(handle, max_wait)"] --> assert["debug_assert!(!is_shutdown)"]
assert --> release["release_pending_registrations()"]
release --> pick{"max_wait == 0<br/>且 events_busy 存在?"}
pick -->|是| busy["events = events_busy"]
pick -->|否| main["events = events"]
busy --> poll["poll.poll(events, max_wait)"]
main --> poll
poll --> pollres{"poll 返回?"}
pollres -->|"Ok / Interrupted"| iter["遍历 events.iter()"]
pollres -->|"其他 Err"| panic["panic!(unexpected error)"]
iter --> tok{"event.token()?"}
tok -->|"TOKEN_WAKEUP"| skip["忽略,仅用于打断阻塞"]
tok -->|"TOKEN_SIGNAL"| sig["signal_ready = true"]
tok -->|"普通 fd token"| cast["EXPOSE_IO.from_exposed_addr(token.0)"]
cast --> setr["io.set_readiness(Tick::Set, curr | ready)"]
setr --> wake["io.wake(ready)"]
wake --> iter
skip --> iter
sig --> iter
iter --> uring["dispatch_completions() (io-uring)"]
uring --> metrics["metrics.incr_ready_count_by(ready_count)"]Designüberlegung: WarumHandlehalten mussmio::Waker
unpark 📎 tokio/src/runtime/io/driver.rs:280-283aufrufenself.waker.wake(). Diesesmio::Wakerwird beiDriver::newmitTOKEN_WAKEUPregistriert📎 tokio/src/runtime/io/driver.rs:124. Wenn der Driver inpoll.poll()blockiert, wird durch einen anderen Thread, derunparkaufruft, einTOKEN_WAKEUP-Event in epoll eingefügt,pollkehrt sofort zurück, bei der Iteration wird dieses Token direkt übersprungen📎 tokio/src/runtime/io/driver.rs:214-215。
Dieser Mechanismus wird inderegister_sourceverwendet📎 tokio/src/runtime/io/driver.rs:315-334: Nach der Abmeldung einer Source, fallsregistrations.deregistertrue zurückgibt (was bedeutet, dass dies die letzte Referenz ist), wirdunpark(). Warum? Weil der Driver möglicherweise gerade inpollauf das Event dieses fd wartet, und der fd bereits abgemeldet wurde, sodass der Kernel keine Events mehr erzeugen wird; der Driver muss aktiv aufgeweckt werden, damit er die Registrierungsmenge erneut überprüft und möglicherweise die Blockierung beendet. Andernfalls würde der Driver bis zum Timeout vonmax_waitschlafen, was das Herunterfahren verzögert.
Ein weiteres Detail:deregister_sourceruft zuerstself.registry.deregister(source) 📎 tokio/src/runtime/io/driver.rs:322auf, dann wirdregistrations 📎 tokio/src/runtime/io/driver.rs:315-334bereinigt. Der Kommentar📎 tokio/src/runtime/io/driver.rs:320-321sagt „Cleanup ALWAYS happens“ – selbst wenn die Deregistrierung auf OS-Ebene fehlschlägt, muss der interne Zustand bereinigt werden, erst danach wird der OS-Fehler zurückgegeben📎 tokio/src/runtime/io/driver.rs:336-340. Dies ist ein typischesMuster, bei dem Ressourcenbereinigung Vorrang vor Fehlerpropagierung hat.
Registrierungsschicht:RegistrationWieWakerinScheduledIo
gespeichert wird. Intuitives Modell
Registrationistein Vertrag zwischen Task und fd. Es hält zwei Dinge: einscheduler::Handle(verwendet, um bei Bedarf auf die Runtime zuzugreifen), einArc<ScheduledIo>(der Zustandsslot des fd). Wenn eine Taskpoll_read_readyaufruft,RegistrationübergibtWakeranScheduledIozur Verwahrung; wenn der Driver ein Event empfängt, wirdScheduledIoausWakerentnommen und aufgeweckt.
Speicherlayout und Felder
Registrationhat nur zwei Felder📎 tokio/src/runtime/io/registration.rs:46-54:
handle: scheduler::Handle: Runtime-Handle, Kommentar📎tokio/src/runtime/io/registration.rs:46-54sagt „TODO: this can probably be moved into ScheduledIo“, was zeigt, dass der Autor der Meinung ist, dass die Position dieses Feldes optimiert werden kann.shared: Arc<ScheduledIo>: Gemeinsamer Zustand,Arcstellt sicher, dass sowohl Driver als auch Task darauf zugreifen können.
Beachten Sie, dassRegistrationmanuellSendundSync 📎 tokio/src/runtime/io/registration.rs:57-58implementiert. Warum ist unsafe impl erforderlich? Weilscheduler::Handleintern möglicherweise Felder enthält, die nichtSend/Syncsind (zum BeispielRc), aberRegistrationdas Anwendungsszenario erfordert, dass es über Threads hinweg verwendet werden kann. Der Dokumentationskommentar📎 tokio/src/runtime/io/registration.rs:28-33gibt die entscheidende Einschränkung an:Der Aufrufer muss sicherstellen, dass höchstens zwei Tasks gleichzeitig dasselbeRegistrationverwenden, eine liest, eine schreibt. Eine Verletzung dieser Einschränkung ist zwar speichersicher, führt jedoch zu verlorenen Benachrichtigungen und hängenden Tasks.
Step-by-Step:poll_read_readyDie Aufrufkette von
Angenommen, eine Task ist inTcpStream::poll_readstellt fest, dass der Socket keine Daten hat, und muss ein Lese-Interesse registrieren. Die Aufrufkette istTcpStream::poll_read_priv → PollEvented::poll_read → Registration::poll_read_io → poll_io → poll_ready。
poll_readyist der Kern📎 tokio/src/runtime/io/registration.rs:155-171:
Erster Schritt:trace_leaf() 📎 tokio/src/runtime/io/registration.rs:160, verwendet für Tracing-Instrumentierung.
Zweiter Schritt:coop::poll_proceed(cx) 📎 tokio/src/runtime/io/registration.rs:155-171. Dies ist der kooperative Budget-Mechanismus, der in Kapitel 12 behandelt wird. Wenn das Budget erschöpft ist, wird zurückgegebenPendingund ein speziellesWakerregistriert, damit die Aufgabe in der nächsten Runde neu geplant wird.
Dritter Schritt:self.shared.poll_readiness(cx, direction) 📎 tokio/src/runtime/io/registration.rs:155-171. Dies ist der Ort, an dem tatsächlich mitScheduledIointeragiert wird: Der aktuelle Bereitschaftsstatus wird geprüft; wenn bereits bereit, wird sofort zurückgegebenReady; andernfalls wirdcx.waker()in den entsprechenden Richtungs-Slot vonScheduledIogespeichert und zurückgegebenPending。
Vierter Schritt: Prüfen vonev.is_shutdown 📎 tokio/src/runtime/io/registration.rs:155-171. Wenn die Runtime gerade heruntergefahren wird, wird zurückgegebenRUNTIME_SHUTTING_DOWN_ERROR。
Fünfter Schritt:coop.made_progress() 📎 tokio/src/runtime/io/registration.rs:169, Budgetverbrauch markieren und das Bereitschaftsereignis zurückgeben.
poll_iofügt überpoll_readyeine Retry-Schleife hinzu📎 tokio/src/runtime/io/registration.rs:173-192:
loop {
let ev = ready!(self.poll_ready(cx, direction))?;
match f() {
Ok(ret) => return Poll::Ready(Ok(ret)),
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
self.clear_readiness(ev);
}
Err(e) => return Poll::Ready(Err(e)),
}
}Hier zeigt sichReadiness ist ein Hinweis, keine GarantieDie Kernidee:poll_readymeldet lesbar, aber beim tatsächlichenread()kannWouldBlockzurückgegeben werden (z. B. wenn ein anderer Thread die Daten zuerst weggelesen hat). In diesem Fall mussclear_readiness(ev) 📎 tokio/src/runtime/io/registration.rs:187das Bereitschaftsbit gelöscht und die Schleife erneut auf Warten gesetzt werden. Wenn nicht gelöscht wird, gerät die Aufgabe in eine Busy-Schleife von „glaubt lesbar → read schlägt fehl → glaubt wieder lesbar“.
Designüberlegung:try_ioundasync_ioArbeitsteilung
try_io 📎 tokio/src/runtime/io/registration.rs:194-213ist die synchrone Version: Zuerstready_event(interest)das Bereitschaftsbit prüfen; wenn leer, direkt zurückgebenWouldBlock 📎 tokio/src/runtime/io/registration.rs:194-213; andernfallsf()ausführen; wennf()zurückgibtWouldBlock, dann das Bereitschaftsbit löschen📎 tokio/src/runtime/io/registration.rs:207-210. Esregistriert keinen Wakerund eignet sich fürtry_readSzenarien wie „einmal versuchen und gehen“.
async_io 📎 tokio/src/runtime/io/registration.rs:225-245ist die asynchrone Version:readiness(interest).awaitregistriert einen Waker und wartet, dann wird bei der Ausführung vonf(),WouldBlockdas Bereitschaftsbit gelöscht und die Schleife fortgesetzt. Beachten Sie, dass in der Schleife auchcoop::poll_proceed 📎 tokio/src/runtime/io/registration.rs:233aufgerufen wird, um zu verhindern, dass bei vielenWouldBlockRetries das Budget erschöpft wird.
Produktions-Fallstricke:DropWaker-Bereinigung in
Registration::drop 📎 tokio/src/runtime/io/registration.rs:253-262ruftself.shared.clear_wakers()auf. Der Kommentar📎 tokio/src/runtime/io/registration.rs:253-262erklärt den Grund:ScheduledIoDas inWakergespeicherteArc<driver::Inner>kanndriver::Innerhalten, undScheduledIohält wiederumRegistration, was einen Zirkelbezug bildet. Die Waker-Bereinigung ist ein Mittel, um den Zyklus zu durchbrechen. Der Kommentar gibt jedoch zu, dass dies eine „imperfect solution“ ist – wennWakerselbst in
〔Design-Inferenz und Architektur-Abwägung〕clear_wakersDas Verhalten in der Produktion ist: Wenn viele Verbindungen gedroppt werden, aber die Runtime nicht beendet wird, wird der Speicher nicht sofort freigegeben, bis zum nächstenScheduledIooder Runtime-Shutdown. Für Langzeitverbindungsdienste ist dies normalerweise kein Problem; bei Szenarien mit häufiger Erstellung/Zerstörung von Kurzzeitverbindungen muss jedoch der Rückgewinnungszeitpunkt von
beachtet werden.TcpStream::readVonWakerbis
vollständige Kette des Aufweckens
Intuitives ModellTcpStreamNun werden die drei Ebenen verbunden. Der Benutzer ruft auf.read().awaitaufAsyncRead::poll_read → PollEvented::poll_read → Registration::poll_read_io, tatsächlich ausgeführt wirdWaker. Wenn keine Daten ankommen,ScheduledIowird inScheduledIogespeichert; wenn epoll Lesbarkeit meldet, entnimmt der Driver ausWakerdenpoll_readinessund weckt auf, die Aufgabe wird neu geplant, und beim erneuten Poll entdecktReady,read(), dass das Bereitschaftsbit gesetzt ist, und gibt direkt
erfolgreich zurück.
Step-by-Step: Ein vollständiges Lese-Warten。TcpStream::new 📎 tokio/src/net/tcp/stream.rs:166-169Phase eins: Interesse registrierenPollEvented::new(connected)ruftRegistration::new_with_interest_and_handle 📎 tokio/src/runtime/io/registration.rs:73-81auf, letzteres ruft internhandle.driver().io().add_source(io, interest) 📎 tokio/src/runtime/io/registration.rs:73-81。
add_source 📎 tokio/src/runtime/io/driver.rs:288-312auf, und
1. registrations.allocate(&mut synced.lock())erledigt drei Dinge:ScheduledIoweist einentoken 📎 tokio/src/runtime/io/driver.rs:293-294。
2. self.registry.register(source, token, interest.to_mio())zu, holt📎 tokio/src/runtime/io/driver.rs:298registriertbeim Kernel. Bei FehlschlagmussScheduledIoder gerade zugewiesene📎 tokio/src/runtime/io/driver.rs:300-303aus der Menge entfernt werden
3. metrics.incr_fd_count(), sonst Leck.📎 tokio/src/runtime/io/driver.rs:309。
zähltPhase zwei: Auf Bereitschaft wartenTcpStream::poll_read 📎 tokio/src/net/tcp/stream.rs:1492-1498 → poll_read_priv 📎 tokio/src/net/tcp/stream.rs:1451-1458 → PollEvented::poll_read → Registration::poll_read_io 📎 tokio/src/runtime/io/registration.rs:133-139 → poll_io → poll_ready → ScheduledIo::poll_readiness. Die Aufgabe polltWaker. Wenn zu diesem Zeitpunkt nicht bereit,ScheduledIowird in den Lese-Slot vonPending。
gespeichert und zurückgegebenPhase drei: Ereignis trifft einturn. Derpoll.poll()des Drivers holt das Ereignis📎 tokio/src/runtime/io/driver.rs:198ausio.set_readiness(Tick::Set, |curr| curr | ready), und beim Durchlaufen wird für jedes fd-Ereignisio.wake(ready) 📎 tokio/src/runtime/io/driver.rs:228-229。wakeausgeführt undWakerintern derwake()。
der entsprechenden Richtung entnommen und。Waker::wake()aufgerufenpoll_readinessPhase vier: Aufgabe neu planenReady,read()Die Aufgabe wird erneut in die lokale Queue des Workers eingereiht (im vorherigen Kapitel behandelt). Der Worker pollt die Aufgabe erneut,
sequenceDiagram
participant Task as "任务 (worker 线程)"
participant Reg as "Registration"
participant SIO as "ScheduledIo"
participant Drv as "Driver (I/O 线程)"
participant OS as "epoll/kqueue"
Task->>Reg: "poll_read_ready(cx)"
Reg->>SIO: "poll_readiness(cx, Read)"
SIO-->>Reg: "Pending (Waker 已存入读槽位)"
Reg-->>Task: "Poll::Pending"
Note over Task: 任务让出,worker 去跑别的任务
Drv->>OS: "poll.poll(events, max_wait)"
OS-->>Drv: "event(token=fd_ptr, READABLE)"
Drv->>SIO: "set_readiness(Tick::Set, curr | READABLE)"
Drv->>SIO: "wake(READABLE)"
SIO->>Task: "Waker::wake() 重新入队"
Note over Task: worker 再次 poll 该任务
Task->>Reg: "poll_read_ready(cx)"
Reg->>SIO: "poll_readiness(cx, Read)"
SIO-->>Reg: "Ready(ReadyEvent{ready: READABLE})"
Reg-->>Task: "Poll::Ready(Ok(ev))"
Task->>Task: "read() 成功返回数据"erfolgreich zurück.assume_readyKopieren
TcpStream::new_accepted 📎 tokio/src/net/tcp/stream.rs:174-181Wichtiger Zweig:acceptOptimierungnew_acceptedist eine bemerkenswerte Optimierung.assume_ready(Ready::READABLE | Ready::WRITABLE) 📎 tokio/src/net/tcp/stream.rs:174-181。
assume_readyDer von📎 tokio/src/runtime/io/registration.rs:103-105zurückgegebene Socket ist natürlich beschreibbar und hält normalerweise bereits die ersten Bytes der Gegenseite. Wenn man auf das erste Ereignis des Drivers wartet, könnte dieses Ereignis unter hoher Last hinter allen Ereignissen bereits aufgebauter Verbindungen stehen und Verzögerungen verursachen. Daher ruftWouldBlockdirektWouldBlock,poll_ioauf. Der Kommentarsagt: „A wrong guess costs one, which clears the readiness again.“ – Der Preis für eine falsche Vermutung ist nur eine
Schleife, die das Bereitschaftsbit löscht und erneut wartet. Dies ist ein Design von
.DriverDesignüberlegung: Warum der I/O-Treiber vom Scheduler entkoppelt istDriver〔Design-Inferenz und Architektur-Abwägung〕block_onAus der Quellcode-Struktur geht hervor, dassHandleund Worker-Threads getrennt sind:
1. wird an einer speziellen Stelle der Runtime platziert (normalerweise:HandleThread oder dedizierter I/O-Thread), während Worker-Threads nurmio::Registryhalten. Diese Entkopplung bringt mehrere Vorteile:
2. Lock-freie Registrierunghält einenepoll_wait-Klon, jeder Worker kann parallel neue fds registrieren, ohne zum Driver-Thread zurückzukehren.
3. Zentralisiertes Ereignis-Warten: Nur ein Thread blockiert aufScheduledIo, wodurch das Thundering-Herd-Problem vermieden wird, bei dem mehrere Threads gleichzeitig denselben epoll-fd pollen.Waker::wake(),wake()Kurzer Aufweckpfad
: Nach Erhalt eines Ereignisses operiert der Driver direkt aufScheduledIound ruftset_readinessauf; intern wird die Aufgabe in die Worker-Queue geschoben, ohne threadübergreifende Nachrichtenübermittlung.poll_readinessDer Preis ist, dass
konkurrierende Zugriffe behandeln muss (is_shutdownundRUNTIME_SHUTTING_DOWN_ERROR
poll_readykönnen gleichzeitig auftreten), was durch atomare Operationen und interne Locks gelöst wird.ev.is_shutdown 📎 tokio/src/runtime/io/registration.rs:155-171Produktions-Fallstricke:gone() 📎 tokio/src/runtime/io/registration.rs:265-267undRUNTIME_SHUTTING_DOWN_ERROR。
, wenn wahr, zurückgebenshutdown 📎 tokio/src/runtime/io/driver.rs:174-182durchläuft alle registrierten und ruft aufio.shutdown(), setztis_shutdownund weckt alle Wartenden auf. Wenn dieses Flag nicht geprüft wird, könnte eine Task noch versuchen, den Socket zu lesen, nachdem die Runtime bereits das Scheduling eingestellt hat, was zu undefiniertem Verhalten oder Hängen führen kann. In Produktionsumgebungen, wenn SieRUNTIME_SHUTTING_DOWN_ERRORsehen, bedeutet das normalerweise, dass eine Task noch läuft, nachdem die Runtime gedroppt wurde – prüfen Sie, ob esspawnTasks gibt, die nicht korrekt gejoint wurden.
Eine weitere Falle istderegister_sourcevonunpark 📎 tokio/src/runtime/io/driver.rs:328. Wenn der Driver blockiert inpollwartet und zu diesem Zeitpunkt der letzteRegistrationgedroppt wird,unparkweckt den Driver auf. Aber wenn der Driver nicht blockiert ist (z. B. gerade andere Events verarbeitet),unparklässt nur den nächstenturnsofort📎 tokio/src/runtime/io/driver.rs:280-283zurückgeben. Diese Semantik ist in den Dokumentationskommentaren vonHandle::unparkbeschrieben.
Designüberlegung: Die drei entscheidenden Kompromisse des Reactors
Kompromiss eins:Tokenverwendet Zeiger statt Indizes。EXPOSE_IO.from_exposed_addr(token.0) 📎 tokio/src/runtime/io/driver.rs:220behandeltmio::Tokendirekt als*const ScheduledIoAdresse. Dies vermeidet die Pflege einerToken → ScheduledIoMapping-Tabelle, die Suche ist O(1) und lockfrei. Der Preis ist, dass die Sicherheit von striktem Lebenszeitmanagement abhängt: Der Zeiger darf erst freigegeben werden, nachdem die Registrierung aufgehoben wurde und der Driver nicht mehr pollt📎 tokio/src/runtime/io/driver.rs:222-225。
Kompromiss zwei: Zwei getrennte Waker-Slots für Lesen und Schreiben。RegistrationDokumentation📎 tokio/src/runtime/io/registration.rs:24-26sagt „A registration instance represents two separate readiness streams“ – Lesen und Schreiben haben jeweils einen unabhängigenWakerSlot. Dies erlaubt, dass Lese- und Schreib-Tasks desselben Sockets sich separat registrieren, ohne sich gegenseitig zu stören. Aberpoll_read_readyKommentar📎 tokio/src/net/tcp/stream.rs:549-552erinnert daran: Mehrfache Aufrufe vonpoll_read_ready/poll_read/poll_peekbehalten nur den letztenWaker– die Lese-Richtung hat nur einen Slot.
Kompromiss drei:events_busyunabhängiger Puffer. Test📎 tokio/src/runtime/io/driver.rs:364-386verifiziert dieses Verhalten:Driver::new(16, Some(2))Erstellt einen Driver mit busy-Kapazität 2, registriert 5 lesbare Sources, nicht-blockierendesturnnimmt nur 2 Events📎 tokio/src/runtime/io/driver.rs:375-376, die restlichen 3 bleiben in der Kernel-Queue, beim nächsten blockierendenturnwerden📎 tokio/src/runtime/io/driver.rs:379-380geholt. Dies verhindert, dass ein nicht-blockierender Poll alle Events auf einmal verschlingt und nachfolgende Polls aushungert.
Zusammenfassung dieses Kapitels
Dieses Kapitel verfolgte die vollständige Reactor-Kette hinterTcpStream::read:
- Treiber-Schicht:
Driverexklusivmio::Poll,turnblockierendes Warten auf Events, mitEXPOSE_IOwirdTokenzurück inScheduledIoZeiger umgewandelt, Aufruf vonset_readiness+wakelöstWaker。Handleaus, bietet einen thread-übergreifenden Registrierungseingang,unparkdient zum Unterbrechen der Blockierung. - Registrierungs-Schicht:
RegistrationhältArc<ScheduledIo>,poll_readyprüft Ready-Bits oder speichert inWaker,poll_iomitWouldBlockRetry-Schleife behandelt False Positives,try_io/async_iobedient jeweils synchrone und asynchrone Szenarien. - Zustands-Schicht:
ScheduledIoist der Zustands-Slot des fd, speichert Lese-/Schreib-Ready-Bits und zweiWakerSlots, ist die einzige Brücke zwischen Events und Tasks.
Kapitel-Überlegungen und Selbsttest
Q1: Wenn man inpoll_ioWouldBlockZweigself.clear_readiness(ev)löscht, in welchem Szenario führt das zu einer Busy-Loop der Task? Warum?
Referenzanalyse:poll_ioSchleife📎 tokio/src/runtime/io/registration.rs:173-192inf()ruftWouldBlockauf, wennclear_readiness(ev) 📎 tokio/src/runtime/io/registration.rs:187。evzurückgibtpoll_readyistReadyEventRückgabe vonclear_readiness, enthält aktuelle Ready-Bits.ScheduledIolöscht diese Bits aus
.poll_ready → poll_readinessWenn nicht bereinigt, beim nächsten Schleifenaufruf vonScheduledIobleiben die alten „lesbar“-Bits inpoll_readinesserhalten,Readygibt sofortf()zurück (da Ready-Bits nicht leer), dannread()führt erneutWouldBlockaus, wenn der Socket tatsächlich keine Daten hat, gibt wiederPendingzurück, Schleife geht weiter. Da die Ready-Bits nie gelöscht werden, erreicht diese Schleife nie
, die Task verbraucht dauerhaft CPU durch Polling.RegistrationAuslöseszenario: Mehrere Tasks teilen sich die Lese-Richtung desselben Sockets (obwohl📎 tokio/src/runtime/io/registration.rs:28-33Dokumentationtry_readsagt maximal zwei Tasks, aber die Lese-Richtung hat nur einen Slot), oderpoll_readundread()gemischt verwendet werden. Häufiger: Nachdem epoll Lesbarkeit meldet, hat ein anderer Thread die Daten zuerst weggelesen, dieWouldBlockder aktuellen Task gibt
Q2: add_sourcezurück, dann müssen die Ready-Bits gelöscht werden, sonst wird endlos wiederholt.registry.registerWarum wird beiregistrations.removeFehlschlag
aufgerufen? Was passiert, wenn nicht aufgerufen?:add_source 📎 tokio/src/runtime/io/driver.rs:288-312Referenzanalyseregistrations.allocatezuerstScheduledIo 📎 tokio/src/runtime/io/driver.rs:293allokiertregistry.register, dann📎 tokio/src/runtime/io/driver.rs:298registriert beim KernelScheduledIo. Wenn die Registrierung fehlschlägt,RegistrationSetwurde bereits allokiert, aber kein fd ist damit verknüpft; wenn nicht entfernt, bleibt es für immer in
.📎 tokio/src/runtime/io/driver.rs:296-297Kommentarscheduled_io from the registrations set if registering the source with the OS fails. Otherwise it will leak the scheduled_iosagt explizit: „we should remove the
remove.“ – das ist ein Speicherleck.📎 tokio/src/runtime/io/driver.rs:300-303Aufruf vonScheduledIoist in einen unsafe-Block gehüllt, weilRegistrationSetTeil vonRegistrationSetist, die Entfernung muss sicherstellen, dass keine anderen Referenzen existieren. Konsequenzen des Lecks:Tokenwächst kontinuierlich,allocateSpeicher wird verschwendet, kann schließlich zu
Q3: deregister_sourceFehlschlag oder Speichererschöpfung führen. In Szenarien mit häufiger Erstellung/Zerstörung von Verbindungen (z. B. Kurzverbindungs-Server), wenn die Registrierungsfehlerrate hoch ist (z. B. fd-Erschöpfung), beschleunigt das Leck die Ressourcenerschöpfung.unpark()Inregistrations.deregister, warum wird
nur aufgerufen, wenn:deregister_source 📎 tokio/src/runtime/io/driver.rs:315-334true zurückgibt? Was wäre das Problem bei bedingungslosem Aufruf?registry.deregister(source)Referenzanalyse📎 tokio/src/runtime/io/driver.rs:322Logik ist: zuerstregistrations.deregisterderegistriert beim Kernel📎 tokio/src/runtime/io/driver.rs:315-334, dannunpark() 📎 tokio/src/runtime/io/driver.rs:328。
registrations.deregisterbereinigt internen ZustandScheduledIo, wenn true zurückgegeben wird, dannpolltrue zurückgeben bedeutet, dies ist die letzte Referenz,unparkwird tatsächlich entfernt. Zu diesem Zeitpunkt könnte der Driver blockiert inmio::Wakerauf Events dieses fd warten, aber der fd ist bereits deregistriert, der Kernel wird keine Events mehr erzeugen.TOKEN_WAKEUPdurch📎 tokio/src/runtime/io/driver.rs:280-283wird einpollEvent in epoll eingefügt
, lässtunparksofort zurückkehren, der Driver prüft die Registrierungsmenge erneut und könnte die Blockierung beenden.ScheduledIoWenn bedingungslosTcpStreamaufgerufen wird: Jede Deregistrierung einer nicht-letzten Referenz weckt den Driver auf, verursacht unnötige Wakeups. In Szenarien, wo viele Verbindungen dasselbesplit后读写两半),每次 drop 一个半都会唤醒 driver,增加 CPU 开销。更严重
本章我们拆解了 Reactor 如何把 epoll 事件翻译成 Waker 唤醒:从 TcpStream 的 poll_read_ready 出发,经过 Registration 的注册与查询,落到 ScheduledIo 的就绪位与 Waker 槽位,再由 Driver 在事件循环中根据 Token 定位并触发唤醒。关键设计包括:Token 即指针实现 O(1) 查找,读写双 Waker 槽位支持并发读写分离,events_busy 独立缓冲区防止事件饥饿,assume_ready 乐观猜测优化 accept 场景。至此,I/O 就绪通知的闭环已经完整。但异步运行时还需要处理另一类「就绪」——时间。下一章我们将剖析 tokio::time::sleep 与 timeout 的实现:定时器如何被插入时间轮、时间轮如何按到期时间分级、driver 如何计算下一次 park 的超时并触发到期任务。你会看到「时间也是一种 I/O 事件」这一统一抽象,以及 start_paused 与 test clock 如何让时间在测试中可控。