Coverage Report

Created: 2026-08-13 06:46

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/rust/registry/src/index.crates.io-1949cf8c6b5b557f/tokio-1.53.1/src/runtime/builder.rs
Line
Count
Source
1
#![cfg_attr(loom, allow(unused_imports))]
2
3
use crate::runtime::handle::Handle;
4
use crate::runtime::{
5
    blocking, driver, Callback, HistogramBuilder, Runtime, TaskCallback, TimerFlavor,
6
};
7
#[cfg(tokio_unstable)]
8
use crate::runtime::{metrics::HistogramConfiguration, TaskMeta};
9
10
use crate::runtime::{LocalOptions, LocalRuntime};
11
use crate::util::rand::{RngSeed, RngSeedGenerator};
12
13
use crate::runtime::blocking::BlockingPool;
14
use crate::runtime::scheduler::CurrentThread;
15
use std::fmt;
16
use std::io;
17
use std::thread::ThreadId;
18
use std::time::Duration;
19
20
/// Builds Tokio Runtime with custom configuration values.
21
///
22
/// Methods can be chained in order to set the configuration values. The
23
/// Runtime is constructed by calling [`build`].
24
///
25
/// New instances of `Builder` are obtained via [`Builder::new_multi_thread`]
26
/// or [`Builder::new_current_thread`].
27
///
28
/// See function level documentation for details on the various configuration
29
/// settings.
30
///
31
/// [`build`]: method@Self::build
32
/// [`Builder::new_multi_thread`]: method@Self::new_multi_thread
33
/// [`Builder::new_current_thread`]: method@Self::new_current_thread
34
///
35
/// # Examples
36
///
37
/// ```
38
/// # #[cfg(not(target_family = "wasm"))]
39
/// # {
40
/// use tokio::runtime::Builder;
41
///
42
/// fn main() {
43
///     // build runtime
44
///     let runtime = Builder::new_multi_thread()
45
///         .worker_threads(4)
46
///         .thread_name("my-custom-name")
47
///         .thread_stack_size(3 * 1024 * 1024)
48
///         .build()
49
///         .unwrap();
50
///
51
///     // use runtime ...
52
/// }
53
/// # }
54
/// ```
55
pub struct Builder {
56
    /// Runtime type
57
    kind: Kind,
58
59
    /// Name of the runtime.
60
    name: Option<String>,
61
62
    /// Whether or not to enable the I/O driver
63
    enable_io: bool,
64
    nevents: usize,
65
66
    /// Whether or not to enable the time driver
67
    enable_time: bool,
68
69
    /// Whether or not the clock should start paused.
70
    start_paused: bool,
71
72
    /// The number of worker threads, used by Runtime.
73
    ///
74
    /// Only used when not using the current-thread executor.
75
    worker_threads: Option<usize>,
76
77
    /// Cap on thread usage.
78
    max_blocking_threads: usize,
79
80
    /// Name fn used for threads spawned by the runtime.
81
    pub(super) thread_name: ThreadNameFn,
82
83
    /// Stack size used for threads spawned by the runtime.
84
    pub(super) thread_stack_size: Option<usize>,
85
86
    /// Callback to run after each thread starts.
87
    pub(super) after_start: Option<Callback>,
88
89
    /// To run before each worker thread stops
90
    pub(super) before_stop: Option<Callback>,
91
92
    /// To run before each worker thread is parked.
93
    pub(super) before_park: Option<Callback>,
94
95
    /// To run after each thread is unparked.
96
    pub(super) after_unpark: Option<Callback>,
97
98
    /// To run before each task is spawned.
99
    pub(super) before_spawn: Option<TaskCallback>,
100
101
    /// To run before each poll
102
    #[cfg(tokio_unstable)]
103
    pub(super) before_poll: Option<TaskCallback>,
104
105
    /// To run after each poll
106
    #[cfg(tokio_unstable)]
107
    pub(super) after_poll: Option<TaskCallback>,
108
109
    /// To run after each task is terminated.
110
    pub(super) after_termination: Option<TaskCallback>,
111
112
    /// Customizable keep alive timeout for `BlockingPool`
113
    pub(super) keep_alive: Option<Duration>,
114
115
    /// How many ticks before pulling a task from the global/remote queue?
116
    ///
117
    /// When `None`, the value is unspecified and behavior details are left to
118
    /// the scheduler. Each scheduler flavor could choose to either pick its own
119
    /// default value or use some other strategy to decide when to poll from the
120
    /// global queue. For example, the multi-threaded scheduler uses a
121
    /// self-tuning strategy based on mean task poll times.
122
    pub(super) global_queue_interval: Option<u32>,
123
124
    /// How many ticks before yielding to the driver for timer and I/O events?
125
    pub(super) event_interval: u32,
126
127
    /// When true, the multi-threade scheduler LIFO slot should not be used.
128
    ///
129
    /// This option should only be exposed as unstable.
130
    pub(super) disable_lifo_slot: bool,
131
132
    /// Specify a random number generator seed to provide deterministic results
133
    pub(super) seed_generator: RngSeedGenerator,
134
135
    /// When true, enables task poll count histogram instrumentation.
136
    pub(super) metrics_poll_count_histogram_enable: bool,
137
138
    /// Configures the task poll count histogram
139
    pub(super) metrics_poll_count_histogram: HistogramBuilder,
140
141
    /// When true, enables task schedule latency instrumentation.
142
    pub(super) metrics_schedule_latency_histogram_enabled: bool,
143
144
    /// Configures the task schedule latency histogram.
145
    pub(super) metrics_schedule_latency_histogram: HistogramBuilder,
146
147
    #[cfg(tokio_unstable)]
148
    pub(super) unhandled_panic: UnhandledPanic,
149
150
    timer_flavor: TimerFlavor,
151
152
    /// Whether or not to enable eager hand-off for the I/O and time drivers (in
153
    /// `tokio_unstable`).
154
    enable_eager_driver_handoff: bool,
155
}
156
157
cfg_unstable! {
158
    /// How the runtime should respond to unhandled panics.
159
    ///
160
    /// Instances of `UnhandledPanic` are passed to `Builder::unhandled_panic`
161
    /// to configure the runtime behavior when a spawned task panics.
162
    ///
163
    /// See [`Builder::unhandled_panic`] for more details.
164
    #[derive(Debug, Clone)]
165
    #[non_exhaustive]
166
    pub enum UnhandledPanic {
167
        /// The runtime should ignore panics on spawned tasks.
168
        ///
169
        /// The panic is forwarded to the task's [`JoinHandle`] and all spawned
170
        /// tasks continue running normally.
171
        ///
172
        /// This is the default behavior.
173
        ///
174
        /// # Examples
175
        ///
176
        /// ```
177
        /// # #[cfg(not(target_family = "wasm"))]
178
        /// # {
179
        /// use tokio::runtime::{self, UnhandledPanic};
180
        ///
181
        /// # pub fn main() {
182
        /// let rt = runtime::Builder::new_current_thread()
183
        ///     .unhandled_panic(UnhandledPanic::Ignore)
184
        ///     .build()
185
        ///     .unwrap();
186
        ///
187
        /// let task1 = rt.spawn(async { panic!("boom"); });
188
        /// let task2 = rt.spawn(async {
189
        ///     // This task completes normally
190
        ///     "done"
191
        /// });
192
        ///
193
        /// rt.block_on(async {
194
        ///     // The panic on the first task is forwarded to the `JoinHandle`
195
        ///     assert!(task1.await.is_err());
196
        ///
197
        ///     // The second task completes normally
198
        ///     assert!(task2.await.is_ok());
199
        /// })
200
        /// # }
201
        /// # }
202
        /// ```
203
        ///
204
        /// [`JoinHandle`]: struct@crate::task::JoinHandle
205
        Ignore,
206
207
        /// The runtime should immediately shutdown if a spawned task panics.
208
        ///
209
        /// The runtime will immediately shutdown even if the panicked task's
210
        /// [`JoinHandle`] is still available. All further spawned tasks will be
211
        /// immediately dropped and call to [`Runtime::block_on`] will panic.
212
        ///
213
        /// # Examples
214
        ///
215
        /// ```should_panic
216
        /// use tokio::runtime::{self, UnhandledPanic};
217
        ///
218
        /// # pub fn main() {
219
        /// let rt = runtime::Builder::new_current_thread()
220
        ///     .unhandled_panic(UnhandledPanic::ShutdownRuntime)
221
        ///     .build()
222
        ///     .unwrap();
223
        ///
224
        /// rt.spawn(async { panic!("boom"); });
225
        /// rt.spawn(async {
226
        ///     // This task never completes.
227
        /// });
228
        ///
229
        /// rt.block_on(async {
230
        ///     // Do some work
231
        /// # loop { tokio::task::yield_now().await; }
232
        /// })
233
        /// # }
234
        /// ```
235
        ///
236
        /// [`JoinHandle`]: struct@crate::task::JoinHandle
237
        ShutdownRuntime,
238
    }
239
}
240
241
pub(crate) type ThreadNameFn = std::sync::Arc<dyn Fn() -> String + Send + Sync + 'static>;
242
243
#[derive(Clone, Copy)]
244
pub(crate) enum Kind {
245
    CurrentThread,
246
    #[cfg(feature = "rt-multi-thread")]
247
    MultiThread,
248
}
249
250
impl Builder {
251
    /// Returns a new builder with the current thread scheduler selected.
252
    ///
253
    /// Configuration methods can be chained on the return value.
254
    ///
255
    /// To spawn non-`Send` tasks on the resulting runtime, combine it with a
256
    /// [`LocalSet`], or call [`build_local`] to create a [`LocalRuntime`].
257
    ///
258
    /// [`LocalSet`]: crate::task::LocalSet
259
    /// [`LocalRuntime`]: crate::runtime::LocalRuntime
260
    /// [`build_local`]: crate::runtime::Builder::build_local
261
0
    pub fn new_current_thread() -> Builder {
262
        #[cfg(loom)]
263
        const EVENT_INTERVAL: u32 = 4;
264
        // The number `61` is fairly arbitrary. I believe this value was copied from golang.
265
        #[cfg(not(loom))]
266
        const EVENT_INTERVAL: u32 = 61;
267
268
0
        Builder::new(Kind::CurrentThread, EVENT_INTERVAL)
269
0
    }
270
271
    /// Returns a new builder with the multi thread scheduler selected.
272
    ///
273
    /// Configuration methods can be chained on the return value.
274
    #[cfg(feature = "rt-multi-thread")]
275
    #[cfg_attr(docsrs, doc(cfg(feature = "rt-multi-thread")))]
276
18.6k
    pub fn new_multi_thread() -> Builder {
277
        // The number `61` is fairly arbitrary. I believe this value was copied from golang.
278
18.6k
        Builder::new(Kind::MultiThread, 61)
279
18.6k
    }
280
281
    /// Returns a new runtime builder initialized with default configuration
282
    /// values.
283
    ///
284
    /// Configuration methods can be chained on the return value.
285
18.6k
    pub(crate) fn new(kind: Kind, event_interval: u32) -> Builder {
286
        Builder {
287
18.6k
            kind,
288
289
            // Default runtime name
290
18.6k
            name: None,
291
292
            // I/O defaults to "off"
293
            enable_io: false,
294
            nevents: 1024,
295
296
            // Time defaults to "off"
297
            enable_time: false,
298
299
            // The clock starts not-paused
300
            start_paused: false,
301
302
            // Read from environment variable first in multi-threaded mode.
303
            // Default to lazy auto-detection (one thread per CPU core)
304
18.6k
            worker_threads: None,
305
306
            max_blocking_threads: 512,
307
308
            // Default thread name
309
597k
            thread_name: std::sync::Arc::new(|| "tokio-rt-worker".into()),
310
311
            // Do not set a stack size by default
312
18.6k
            thread_stack_size: None,
313
314
            // No worker thread callbacks
315
18.6k
            after_start: None,
316
18.6k
            before_stop: None,
317
18.6k
            before_park: None,
318
18.6k
            after_unpark: None,
319
320
18.6k
            before_spawn: None,
321
18.6k
            after_termination: None,
322
323
            #[cfg(tokio_unstable)]
324
            before_poll: None,
325
            #[cfg(tokio_unstable)]
326
            after_poll: None,
327
328
18.6k
            keep_alive: None,
329
330
            // Defaults for these values depend on the scheduler kind, so we get them
331
            // as parameters.
332
18.6k
            global_queue_interval: None,
333
18.6k
            event_interval,
334
335
18.6k
            seed_generator: RngSeedGenerator::new(RngSeed::new()),
336
337
            #[cfg(tokio_unstable)]
338
            unhandled_panic: UnhandledPanic::Ignore,
339
340
            metrics_poll_count_histogram_enable: false,
341
342
18.6k
            metrics_poll_count_histogram: HistogramBuilder::default(),
343
344
            metrics_schedule_latency_histogram_enabled: false,
345
346
18.6k
            metrics_schedule_latency_histogram: HistogramBuilder::default(),
347
348
            disable_lifo_slot: false,
349
350
18.6k
            timer_flavor: TimerFlavor::Traditional,
351
352
            // Eager driver handoff is disabled by default.
353
            enable_eager_driver_handoff: false,
354
        }
355
18.6k
    }
356
357
    /// Enables both I/O and time drivers.
358
    ///
359
    /// Doing this is a shorthand for calling `enable_io` and `enable_time`
360
    /// individually. If additional components are added to Tokio in the future,
361
    /// `enable_all` will include these future components.
362
    ///
363
    /// # Examples
364
    ///
365
    /// ```
366
    /// # #[cfg(not(target_family = "wasm"))]
367
    /// # {
368
    /// use tokio::runtime;
369
    ///
370
    /// let rt = runtime::Builder::new_multi_thread()
371
    ///     .enable_all()
372
    ///     .build()
373
    ///     .unwrap();
374
    /// # }
375
    /// ```
376
18.6k
    pub fn enable_all(&mut self) -> &mut Self {
377
        #[cfg(any(
378
            feature = "net",
379
            all(unix, feature = "process"),
380
            all(unix, feature = "signal")
381
        ))]
382
18.6k
        self.enable_io();
383
384
        #[cfg(all(
385
            tokio_unstable,
386
            feature = "io-uring",
387
            feature = "rt",
388
            feature = "fs",
389
            target_os = "linux",
390
        ))]
391
        self.enable_io_uring();
392
393
        #[cfg(feature = "time")]
394
18.6k
        self.enable_time();
395
396
18.6k
        self
397
18.6k
    }
398
399
    /// Enables the alternative timer implementation, which is disabled by default.
400
    ///
401
    /// The alternative timer implementation is an unstable feature that may
402
    /// provide better performance on multi-threaded runtimes with a large number
403
    /// of worker threads.
404
    ///
405
    /// This option only applies to multi-threaded runtimes. Attempting to use
406
    /// this option with any other runtime type will have no effect.
407
    ///
408
    /// [Click here to share your experience with the alternative timer](https://github.com/tokio-rs/tokio/issues/7745)
409
    ///
410
    /// # Examples
411
    ///
412
    /// ```
413
    /// # #[cfg(not(target_family = "wasm"))]
414
    /// # {
415
    /// use tokio::runtime;
416
    ///
417
    /// let rt = runtime::Builder::new_multi_thread()
418
    ///   .enable_alt_timer()
419
    ///   .build()
420
    ///   .unwrap();
421
    /// # }
422
    /// ```
423
    #[cfg(all(tokio_unstable, feature = "time", feature = "rt-multi-thread"))]
424
    #[cfg_attr(
425
        docsrs,
426
        doc(cfg(all(tokio_unstable, feature = "time", feature = "rt-multi-thread")))
427
    )]
428
    pub fn enable_alt_timer(&mut self) -> &mut Self {
429
        self.enable_time();
430
        self.timer_flavor = TimerFlavor::Alternative;
431
        self
432
    }
433
434
    /// Enable eager hand-off of the I/O and time drivers for multi-threaded
435
    /// runtimes, which is disabled by default.
436
    ///
437
    /// When this option is enabled, a worker thread which has parked on the I/O
438
    /// or time driver will notify another worker thread once it is preparing to
439
    /// begin polling a task from the run queue, so that the notified worker can
440
    /// begin polling the I/O or time driver. This can reduce the latency with
441
    /// which I/O and timer notifications are processed, especially when some
442
    /// tasks have polls that take a long time to complete. In addition, it can
443
    /// reduce the risk of a deadlock which may occur when a task blocks the
444
    /// worker thread which is holding the I/O or time driver until some other
445
    /// task, which is waiting for a notification from *that* driver, unblocks
446
    /// it.
447
    ///
448
    /// This option is disabled by default, as enabling it may potentially
449
    /// increase contention due to extra synchronization in cross-driver
450
    /// wakeups.
451
    ///
452
    /// This option only applies to multi-threaded runtimes. Attempting to use
453
    /// this option with any other runtime type will have no effect.
454
    ///
455
    /// **Note**: This is an [unstable API][unstable]. Eager driver hand-off is
456
    /// an experimental feature whose behavior may be removed or changed in 1.x
457
    /// releases. See [the documentation on unstable features][unstable] for
458
    /// details.
459
    ///
460
    /// [unstable]: crate#unstable-features
461
    #[cfg(all(tokio_unstable, feature = "rt-multi-thread"))]
462
    #[cfg_attr(docsrs, doc(cfg(all(tokio_unstable, feature = "rt-multi-thread"))))]
463
    pub fn enable_eager_driver_handoff(&mut self) -> &mut Self {
464
        self.enable_eager_driver_handoff = true;
465
        self
466
    }
467
468
    /// Sets the number of worker threads the `Runtime` will use.
469
    ///
470
    /// This can be any number above 0 though it is advised to keep this value
471
    /// on the smaller side.
472
    ///
473
    /// This will override the value read from environment variable `TOKIO_WORKER_THREADS`.
474
    ///
475
    /// # Default
476
    ///
477
    /// The default value is the number of cores available to the system.
478
    ///
479
    /// When using the `current_thread` runtime this method has no effect.
480
    ///
481
    /// # Examples
482
    ///
483
    /// ## Multi threaded runtime with 4 threads
484
    ///
485
    /// ```
486
    /// # #[cfg(not(target_family = "wasm"))]
487
    /// # {
488
    /// use tokio::runtime;
489
    ///
490
    /// // This will spawn a work-stealing runtime with 4 worker threads.
491
    /// let rt = runtime::Builder::new_multi_thread()
492
    ///     .worker_threads(4)
493
    ///     .build()
494
    ///     .unwrap();
495
    ///
496
    /// rt.spawn(async move {});
497
    /// # }
498
    /// ```
499
    ///
500
    /// ## Current thread runtime (will only run on the current thread via `Runtime::block_on`)
501
    ///
502
    /// ```
503
    /// use tokio::runtime;
504
    ///
505
    /// // Create a runtime that _must_ be driven from a call
506
    /// // to `Runtime::block_on`.
507
    /// let rt = runtime::Builder::new_current_thread()
508
    ///     .build()
509
    ///     .unwrap();
510
    ///
511
    /// // This will run the runtime and future on the current thread
512
    /// rt.block_on(async move {});
513
    /// ```
514
    ///
515
    /// # Panics
516
    ///
517
    /// This will panic if `val` is not larger than `0`.
518
    #[track_caller]
519
0
    pub fn worker_threads(&mut self, val: usize) -> &mut Self {
520
0
        assert!(val > 0, "Worker threads cannot be set to 0");
521
0
        self.worker_threads = Some(val);
522
0
        self
523
0
    }
524
525
    /// Specifies the limit for additional threads spawned by the Runtime.
526
    ///
527
    /// These threads are used for blocking operations like tasks spawned
528
    /// through [`spawn_blocking`], this includes but is not limited to:
529
    /// - [`fs`] operations
530
    /// - dns resolution through [`ToSocketAddrs`]
531
    /// - writing to [`Stdout`] or [`Stderr`]
532
    /// - reading from [`Stdin`]
533
    ///
534
    /// Unlike the [`worker_threads`], they are not always active and will exit
535
    /// if left idle for too long. You can change this timeout duration with [`thread_keep_alive`].
536
    ///
537
    /// It's recommended to not set this limit too low in order to avoid hanging on operations
538
    /// requiring [`spawn_blocking`].
539
    ///
540
    /// The default value is 512.
541
    ///
542
    /// # Queue Behavior
543
    ///
544
    /// When a blocking task is submitted, it will be inserted into a queue. If available, one of
545
    /// the idle threads will be notified to run the task. Otherwise, if the threshold set by this
546
    /// method has not been reached, a new thread will be spawned. If no idle thread is available
547
    /// and no more threads are allowed to be spawned, the task will remain in the queue until one
548
    /// of the busy threads pick it up. Note that since the queue does not apply any backpressure,
549
    /// it could potentially grow unbounded.
550
    ///
551
    /// # Panics
552
    ///
553
    /// This will panic if `val` is not larger than `0`.
554
    ///
555
    /// # Upgrading from 0.x
556
    ///
557
    /// In old versions `max_threads` limited both blocking and worker threads, but the
558
    /// current `max_blocking_threads` does not include async worker threads in the count.
559
    ///
560
    /// [`spawn_blocking`]: fn@crate::task::spawn_blocking
561
    /// [`fs`]: mod@crate::fs
562
    /// [`ToSocketAddrs`]: trait@crate::net::ToSocketAddrs
563
    /// [`Stdout`]: struct@crate::io::Stdout
564
    /// [`Stdin`]: struct@crate::io::Stdin
565
    /// [`Stderr`]: struct@crate::io::Stderr
566
    /// [`worker_threads`]: Self::worker_threads
567
    /// [`thread_keep_alive`]: Self::thread_keep_alive
568
    #[track_caller]
569
    #[cfg_attr(docsrs, doc(alias = "max_threads"))]
570
0
    pub fn max_blocking_threads(&mut self, val: usize) -> &mut Self {
571
0
        assert!(val > 0, "Max blocking threads cannot be set to 0");
572
0
        self.max_blocking_threads = val;
573
0
        self
574
0
    }
575
576
    /// Sets name of threads spawned by the `Runtime`'s thread pool.
577
    ///
578
    /// The default name is "tokio-rt-worker".
579
    ///
580
    /// # Examples
581
    ///
582
    /// ```
583
    /// # #[cfg(not(target_family = "wasm"))]
584
    /// # {
585
    /// # use tokio::runtime;
586
    ///
587
    /// # pub fn main() {
588
    /// let rt = runtime::Builder::new_multi_thread()
589
    ///     .thread_name("my-pool")
590
    ///     .build();
591
    /// # }
592
    /// # }
593
    /// ```
594
0
    pub fn thread_name(&mut self, val: impl Into<String>) -> &mut Self {
595
0
        let val = val.into();
596
0
        self.thread_name = std::sync::Arc::new(move || val.clone());
597
0
        self
598
0
    }
599
600
    /// Sets the name of the runtime.
601
    ///
602
    /// # Examples
603
    ///
604
    /// ```
605
    /// # #[cfg(not(target_family = "wasm"))]
606
    /// # {
607
    /// # use tokio::runtime;
608
    ///
609
    /// # pub fn main() {
610
    /// let rt = runtime::Builder::new_multi_thread()
611
    ///     .name("my-runtime")
612
    ///     .build();
613
    /// # }
614
    /// # }
615
    /// ```
616
    /// # Panics
617
    ///
618
    /// This function will panic if an empty value is passed as an argument.
619
    ///
620
    #[track_caller]
621
0
    pub fn name(&mut self, val: impl Into<String>) -> &mut Self {
622
0
        let val = val.into();
623
0
        assert!(!val.trim().is_empty(), "runtime name shouldn't be empty");
624
0
        self.name = Some(val);
625
0
        self
626
0
    }
627
628
    /// Sets a function used to generate the name of threads spawned by the `Runtime`'s thread pool.
629
    ///
630
    /// The default name fn is `|| "tokio-rt-worker".into()`.
631
    ///
632
    /// # Examples
633
    ///
634
    /// ```
635
    /// # #[cfg(not(target_family = "wasm"))]
636
    /// # {
637
    /// # use tokio::runtime;
638
    /// # use std::sync::atomic::{AtomicUsize, Ordering};
639
    /// # pub fn main() {
640
    /// let rt = runtime::Builder::new_multi_thread()
641
    ///     .thread_name_fn(|| {
642
    ///        static ATOMIC_ID: AtomicUsize = AtomicUsize::new(0);
643
    ///        let id = ATOMIC_ID.fetch_add(1, Ordering::SeqCst);
644
    ///        format!("my-pool-{}", id)
645
    ///     })
646
    ///     .build();
647
    /// # }
648
    /// # }
649
    /// ```
650
0
    pub fn thread_name_fn<F>(&mut self, f: F) -> &mut Self
651
0
    where
652
0
        F: Fn() -> String + Send + Sync + 'static,
653
    {
654
0
        self.thread_name = std::sync::Arc::new(f);
655
0
        self
656
0
    }
657
658
    /// Sets the stack size (in bytes) for worker threads.
659
    ///
660
    /// The actual stack size may be greater than this value if the platform
661
    /// specifies minimal stack size.
662
    ///
663
    /// The default stack size for spawned threads is 2 MiB, though this
664
    /// particular stack size is subject to change in the future.
665
    ///
666
    /// # Examples
667
    ///
668
    /// ```
669
    /// # #[cfg(not(target_family = "wasm"))]
670
    /// # {
671
    /// # use tokio::runtime;
672
    ///
673
    /// # pub fn main() {
674
    /// let rt = runtime::Builder::new_multi_thread()
675
    ///     .thread_stack_size(32 * 1024)
676
    ///     .build();
677
    /// # }
678
    /// # }
679
    /// ```
680
0
    pub fn thread_stack_size(&mut self, val: usize) -> &mut Self {
681
0
        self.thread_stack_size = Some(val);
682
0
        self
683
0
    }
684
685
    /// Executes function `f` after each thread is started but before it starts
686
    /// doing work.
687
    ///
688
    /// This is intended for bookkeeping and monitoring use cases.
689
    ///
690
    /// # Examples
691
    ///
692
    /// ```
693
    /// # #[cfg(not(target_family = "wasm"))]
694
    /// # {
695
    /// # use tokio::runtime;
696
    /// # pub fn main() {
697
    /// let runtime = runtime::Builder::new_multi_thread()
698
    ///     .on_thread_start(|| {
699
    ///         println!("thread started");
700
    ///     })
701
    ///     .build();
702
    /// # }
703
    /// # }
704
    /// ```
705
    #[cfg(not(loom))]
706
0
    pub fn on_thread_start<F>(&mut self, f: F) -> &mut Self
707
0
    where
708
0
        F: Fn() + Send + Sync + 'static,
709
    {
710
0
        self.after_start = Some(std::sync::Arc::new(f));
711
0
        self
712
0
    }
713
714
    /// Executes function `f` before each thread stops.
715
    ///
716
    /// This is intended for bookkeeping and monitoring use cases.
717
    ///
718
    /// # Examples
719
    ///
720
    /// ```
721
    /// # #[cfg(not(target_family = "wasm"))]
722
    /// {
723
    /// # use tokio::runtime;
724
    /// # pub fn main() {
725
    /// let runtime = runtime::Builder::new_multi_thread()
726
    ///     .on_thread_stop(|| {
727
    ///         println!("thread stopping");
728
    ///     })
729
    ///     .build();
730
    /// # }
731
    /// # }
732
    /// ```
733
    #[cfg(not(loom))]
734
0
    pub fn on_thread_stop<F>(&mut self, f: F) -> &mut Self
735
0
    where
736
0
        F: Fn() + Send + Sync + 'static,
737
    {
738
0
        self.before_stop = Some(std::sync::Arc::new(f));
739
0
        self
740
0
    }
741
742
    /// Executes function `f` just before a thread is parked (goes idle).
743
    /// `f` is called within the Tokio context, so functions like [`tokio::spawn`](crate::spawn)
744
    /// can be called, and may result in this thread being unparked immediately.
745
    ///
746
    /// This can be used to start work only when the executor is idle, or for bookkeeping
747
    /// and monitoring purposes.
748
    ///
749
    /// Note: There can only be one park callback for a runtime; calling this function
750
    /// more than once replaces the last callback defined, rather than adding to it.
751
    ///
752
    /// # Examples
753
    ///
754
    /// ## Multithreaded executor
755
    /// ```
756
    /// # #[cfg(not(target_family = "wasm"))]
757
    /// # {
758
    /// # use std::sync::Arc;
759
    /// # use std::sync::atomic::{AtomicBool, Ordering};
760
    /// # use tokio::runtime;
761
    /// # use tokio::sync::Barrier;
762
    /// # pub fn main() {
763
    /// let once = AtomicBool::new(true);
764
    /// let barrier = Arc::new(Barrier::new(2));
765
    ///
766
    /// let runtime = runtime::Builder::new_multi_thread()
767
    ///     .worker_threads(1)
768
    ///     .on_thread_park({
769
    ///         let barrier = barrier.clone();
770
    ///         move || {
771
    ///             let barrier = barrier.clone();
772
    ///             if once.swap(false, Ordering::Relaxed) {
773
    ///                 tokio::spawn(async move { barrier.wait().await; });
774
    ///            }
775
    ///         }
776
    ///     })
777
    ///     .build()
778
    ///     .unwrap();
779
    ///
780
    /// runtime.block_on(async {
781
    ///    barrier.wait().await;
782
    /// })
783
    /// # }
784
    /// # }
785
    /// ```
786
    /// ## Current thread executor
787
    /// ```
788
    /// # use std::sync::Arc;
789
    /// # use std::sync::atomic::{AtomicBool, Ordering};
790
    /// # use tokio::runtime;
791
    /// # use tokio::sync::Barrier;
792
    /// # pub fn main() {
793
    /// let once = AtomicBool::new(true);
794
    /// let barrier = Arc::new(Barrier::new(2));
795
    ///
796
    /// let runtime = runtime::Builder::new_current_thread()
797
    ///     .on_thread_park({
798
    ///         let barrier = barrier.clone();
799
    ///         move || {
800
    ///             let barrier = barrier.clone();
801
    ///             if once.swap(false, Ordering::Relaxed) {
802
    ///                 tokio::spawn(async move { barrier.wait().await; });
803
    ///            }
804
    ///         }
805
    ///     })
806
    ///     .build()
807
    ///     .unwrap();
808
    ///
809
    /// runtime.block_on(async {
810
    ///    barrier.wait().await;
811
    /// })
812
    /// # }
813
    /// ```
814
    #[cfg(not(loom))]
815
0
    pub fn on_thread_park<F>(&mut self, f: F) -> &mut Self
816
0
    where
817
0
        F: Fn() + Send + Sync + 'static,
818
    {
819
0
        self.before_park = Some(std::sync::Arc::new(f));
820
0
        self
821
0
    }
822
823
    /// Executes function `f` just after a thread unparks (starts executing tasks).
824
    ///
825
    /// This is intended for bookkeeping and monitoring use cases; note that work
826
    /// in this callback will increase latencies when the application has allowed one or
827
    /// more runtime threads to go idle.
828
    ///
829
    /// Note: There can only be one unpark callback for a runtime; calling this function
830
    /// more than once replaces the last callback defined, rather than adding to it.
831
    ///
832
    /// # Examples
833
    ///
834
    /// ```
835
    /// # #[cfg(not(target_family = "wasm"))]
836
    /// # {
837
    /// # use tokio::runtime;
838
    /// # pub fn main() {
839
    /// let runtime = runtime::Builder::new_multi_thread()
840
    ///     .on_thread_unpark(|| {
841
    ///         println!("thread unparking");
842
    ///     })
843
    ///     .build();
844
    ///
845
    /// runtime.unwrap().block_on(async {
846
    ///    tokio::task::yield_now().await;
847
    ///    println!("Hello from Tokio!");
848
    /// })
849
    /// # }
850
    /// # }
851
    /// ```
852
    #[cfg(not(loom))]
853
0
    pub fn on_thread_unpark<F>(&mut self, f: F) -> &mut Self
854
0
    where
855
0
        F: Fn() + Send + Sync + 'static,
856
    {
857
0
        self.after_unpark = Some(std::sync::Arc::new(f));
858
0
        self
859
0
    }
860
861
    /// Executes function `f` just before a task is spawned.
862
    ///
863
    /// `f` is called within the Tokio context, so functions like
864
    /// [`tokio::spawn`](crate::spawn) can be called, and may result in this callback being
865
    /// invoked immediately.
866
    ///
867
    /// This can be used for bookkeeping or monitoring purposes.
868
    ///
869
    /// Note: There can only be one spawn callback for a runtime; calling this function more
870
    /// than once replaces the last callback defined, rather than adding to it.
871
    ///
872
    /// This *does not* support [`LocalSet`](crate::task::LocalSet) at this time.
873
    ///
874
    /// **Note**: This is an [unstable API][unstable]. The public API of this type
875
    /// may break in 1.x releases. See [the documentation on unstable
876
    /// features][unstable] for details.
877
    ///
878
    /// [unstable]: crate#unstable-features
879
    ///
880
    /// # Examples
881
    ///
882
    /// ```
883
    /// # use tokio::runtime;
884
    /// # pub fn main() {
885
    /// let runtime = runtime::Builder::new_current_thread()
886
    ///     .on_task_spawn(|_| {
887
    ///         println!("spawning task");
888
    ///     })
889
    ///     .build()
890
    ///     .unwrap();
891
    ///
892
    /// runtime.block_on(async {
893
    ///     tokio::task::spawn(std::future::ready(()));
894
    ///
895
    ///     for _ in 0..64 {
896
    ///         tokio::task::yield_now().await;
897
    ///     }
898
    /// })
899
    /// # }
900
    /// ```
901
    #[cfg(all(not(loom), tokio_unstable))]
902
    #[cfg_attr(docsrs, doc(cfg(tokio_unstable)))]
903
    pub fn on_task_spawn<F>(&mut self, f: F) -> &mut Self
904
    where
905
        F: Fn(&TaskMeta<'_>) + Send + Sync + 'static,
906
    {
907
        self.before_spawn = Some(std::sync::Arc::new(f));
908
        self
909
    }
910
911
    /// Executes function `f` just before a task is polled
912
    ///
913
    /// `f` is called within the Tokio context, so functions like
914
    /// [`tokio::spawn`](crate::spawn) can be called, and may result in this callback being
915
    /// invoked immediately.
916
    ///
917
    /// **Note**: This is an [unstable API][unstable]. The public API of this type
918
    /// may break in 1.x releases. See [the documentation on unstable
919
    /// features][unstable] for details.
920
    ///
921
    /// [unstable]: crate#unstable-features
922
    ///
923
    /// # Examples
924
    ///
925
    /// ```
926
    /// # #[cfg(not(target_family = "wasm"))]
927
    /// # {
928
    /// # use std::sync::{atomic::AtomicUsize, Arc};
929
    /// # use tokio::task::yield_now;
930
    /// # pub fn main() {
931
    /// let poll_start_counter = Arc::new(AtomicUsize::new(0));
932
    /// let poll_start = poll_start_counter.clone();
933
    /// let rt = tokio::runtime::Builder::new_multi_thread()
934
    ///     .enable_all()
935
    ///     .on_before_task_poll(move |meta| {
936
    ///         println!("task {} is about to be polled", meta.id())
937
    ///     })
938
    ///     .build()
939
    ///     .unwrap();
940
    /// let task = rt.spawn(async {
941
    ///     yield_now().await;
942
    /// });
943
    /// let _ = rt.block_on(task);
944
    ///
945
    /// # }
946
    /// # }
947
    /// ```
948
    #[cfg(tokio_unstable)]
949
    #[cfg_attr(docsrs, doc(cfg(tokio_unstable)))]
950
    pub fn on_before_task_poll<F>(&mut self, f: F) -> &mut Self
951
    where
952
        F: Fn(&TaskMeta<'_>) + Send + Sync + 'static,
953
    {
954
        self.before_poll = Some(std::sync::Arc::new(f));
955
        self
956
    }
957
958
    /// Executes function `f` just after a task is polled
959
    ///
960
    /// `f` is called within the Tokio context, so functions like
961
    /// [`tokio::spawn`](crate::spawn) can be called, and may result in this callback being
962
    /// invoked immediately.
963
    ///
964
    /// **Note**: This is an [unstable API][unstable]. The public API of this type
965
    /// may break in 1.x releases. See [the documentation on unstable
966
    /// features][unstable] for details.
967
    ///
968
    /// [unstable]: crate#unstable-features
969
    ///
970
    /// # Examples
971
    ///
972
    /// ```
973
    /// # #[cfg(not(target_family = "wasm"))]
974
    /// # {
975
    /// # use std::sync::{atomic::AtomicUsize, Arc};
976
    /// # use tokio::task::yield_now;
977
    /// # pub fn main() {
978
    /// let poll_stop_counter = Arc::new(AtomicUsize::new(0));
979
    /// let poll_stop = poll_stop_counter.clone();
980
    /// let rt = tokio::runtime::Builder::new_multi_thread()
981
    ///     .enable_all()
982
    ///     .on_after_task_poll(move |meta| {
983
    ///         println!("task {} completed polling", meta.id());
984
    ///     })
985
    ///     .build()
986
    ///     .unwrap();
987
    /// let task = rt.spawn(async {
988
    ///     yield_now().await;
989
    /// });
990
    /// let _ = rt.block_on(task);
991
    ///
992
    /// # }
993
    /// # }
994
    /// ```
995
    #[cfg(tokio_unstable)]
996
    #[cfg_attr(docsrs, doc(cfg(tokio_unstable)))]
997
    pub fn on_after_task_poll<F>(&mut self, f: F) -> &mut Self
998
    where
999
        F: Fn(&TaskMeta<'_>) + Send + Sync + 'static,
1000
    {
1001
        self.after_poll = Some(std::sync::Arc::new(f));
1002
        self
1003
    }
1004
1005
    /// Executes function `f` just after a task is terminated.
1006
    ///
1007
    /// `f` is called within the Tokio context, so functions like
1008
    /// [`tokio::spawn`](crate::spawn) can be called.
1009
    ///
1010
    /// This can be used for bookkeeping or monitoring purposes.
1011
    ///
1012
    /// Note: There can only be one task termination callback for a runtime; calling this
1013
    /// function more than once replaces the last callback defined, rather than adding to it.
1014
    ///
1015
    /// This *does not* support [`LocalSet`](crate::task::LocalSet) at this time.
1016
    ///
1017
    /// **Note**: This is an [unstable API][unstable]. The public API of this type
1018
    /// may break in 1.x releases. See [the documentation on unstable
1019
    /// features][unstable] for details.
1020
    ///
1021
    /// [unstable]: crate#unstable-features
1022
    ///
1023
    /// # Examples
1024
    ///
1025
    /// ```
1026
    /// # use tokio::runtime;
1027
    /// # pub fn main() {
1028
    /// let runtime = runtime::Builder::new_current_thread()
1029
    ///     .on_task_terminate(|_| {
1030
    ///         println!("killing task");
1031
    ///     })
1032
    ///     .build()
1033
    ///     .unwrap();
1034
    ///
1035
    /// runtime.block_on(async {
1036
    ///     tokio::task::spawn(std::future::ready(()));
1037
    ///
1038
    ///     for _ in 0..64 {
1039
    ///         tokio::task::yield_now().await;
1040
    ///     }
1041
    /// })
1042
    /// # }
1043
    /// ```
1044
    #[cfg(all(not(loom), tokio_unstable))]
1045
    #[cfg_attr(docsrs, doc(cfg(tokio_unstable)))]
1046
    pub fn on_task_terminate<F>(&mut self, f: F) -> &mut Self
1047
    where
1048
        F: Fn(&TaskMeta<'_>) + Send + Sync + 'static,
1049
    {
1050
        self.after_termination = Some(std::sync::Arc::new(f));
1051
        self
1052
    }
1053
1054
    /// Creates the configured `Runtime`.
1055
    ///
1056
    /// The returned `Runtime` instance is ready to spawn tasks.
1057
    ///
1058
    /// # Examples
1059
    ///
1060
    /// ```
1061
    /// # #[cfg(not(target_family = "wasm"))]
1062
    /// # {
1063
    /// use tokio::runtime::Builder;
1064
    ///
1065
    /// let rt  = Builder::new_multi_thread().build().unwrap();
1066
    ///
1067
    /// rt.block_on(async {
1068
    ///     println!("Hello from the Tokio runtime");
1069
    /// });
1070
    /// # }
1071
    /// ```
1072
18.6k
    pub fn build(&mut self) -> io::Result<Runtime> {
1073
18.6k
        match &self.kind {
1074
0
            Kind::CurrentThread => self.build_current_thread_runtime(),
1075
            #[cfg(feature = "rt-multi-thread")]
1076
18.6k
            Kind::MultiThread => self.build_threaded_runtime(),
1077
        }
1078
18.6k
    }
1079
1080
    /// Creates the configured [`LocalRuntime`].
1081
    ///
1082
    /// The returned [`LocalRuntime`] instance is ready to spawn tasks.
1083
    ///
1084
    /// # Panics
1085
    ///
1086
    /// This will panic if the runtime is configured with [`new_multi_thread()`].
1087
    ///
1088
    /// [`new_multi_thread()`]: Builder::new_multi_thread
1089
    ///
1090
    /// # Examples
1091
    ///
1092
    /// ```
1093
    /// use tokio::runtime::{Builder, LocalOptions};
1094
    ///
1095
    /// let rt = Builder::new_current_thread()
1096
    ///     .build_local(LocalOptions::default())
1097
    ///     .unwrap();
1098
    ///
1099
    /// rt.spawn_local(async {
1100
    ///     println!("Hello from the Tokio runtime");
1101
    /// });
1102
    /// ```
1103
    #[allow(unused_variables, unreachable_patterns)]
1104
0
    pub fn build_local(&mut self, options: LocalOptions) -> io::Result<LocalRuntime> {
1105
0
        match &self.kind {
1106
0
            Kind::CurrentThread => self.build_current_thread_local_runtime(),
1107
            #[cfg(feature = "rt-multi-thread")]
1108
0
            Kind::MultiThread => panic!("multi_thread is not supported for LocalRuntime"),
1109
        }
1110
0
    }
1111
1112
18.6k
    fn get_cfg(&self) -> driver::Cfg {
1113
        driver::Cfg {
1114
18.6k
            enable_pause_time: match self.kind {
1115
0
                Kind::CurrentThread => true,
1116
                #[cfg(feature = "rt-multi-thread")]
1117
18.6k
                Kind::MultiThread => false,
1118
            },
1119
18.6k
            enable_io: self.enable_io,
1120
18.6k
            enable_time: self.enable_time,
1121
18.6k
            start_paused: self.start_paused,
1122
18.6k
            nevents: self.nevents,
1123
18.6k
            timer_flavor: self.timer_flavor,
1124
        }
1125
18.6k
    }
1126
1127
    /// Sets a custom timeout for a thread in the blocking pool.
1128
    ///
1129
    /// By default, the timeout for a thread is set to 10 seconds. This can
1130
    /// be overridden using `.thread_keep_alive()`.
1131
    ///
1132
    /// # Example
1133
    ///
1134
    /// ```
1135
    /// # #[cfg(not(target_family = "wasm"))]
1136
    /// # {
1137
    /// # use tokio::runtime;
1138
    /// # use std::time::Duration;
1139
    /// # pub fn main() {
1140
    /// let rt = runtime::Builder::new_multi_thread()
1141
    ///     .thread_keep_alive(Duration::from_millis(100))
1142
    ///     .build();
1143
    /// # }
1144
    /// # }
1145
    /// ```
1146
0
    pub fn thread_keep_alive(&mut self, duration: Duration) -> &mut Self {
1147
0
        self.keep_alive = Some(duration);
1148
0
        self
1149
0
    }
1150
1151
    /// Sets the number of scheduler ticks after which the scheduler will poll the global
1152
    /// task queue.
1153
    ///
1154
    /// A scheduler "tick" roughly corresponds to one `poll` invocation on a task.
1155
    ///
1156
    /// By default the global queue interval is 31 for the current-thread scheduler. Please see
1157
    /// [the module documentation] for the default behavior of the multi-thread scheduler.
1158
    ///
1159
    /// Schedulers have a local queue of already-claimed tasks, and a global queue of incoming
1160
    /// tasks. Setting the interval to a smaller value increases the fairness of the scheduler,
1161
    /// at the cost of more synchronization overhead. That can be beneficial for prioritizing
1162
    /// getting started on new work, especially if tasks frequently yield rather than complete
1163
    /// or await on further I/O. Setting the interval to `1` will prioritize the global queue and
1164
    /// tasks from the local queue will be executed only if the global queue is empty.
1165
    /// Conversely, a higher value prioritizes existing work, and is a good choice when most
1166
    /// tasks quickly complete polling.
1167
    ///
1168
    /// [the module documentation]: crate::runtime#multi-threaded-runtime-behavior-at-the-time-of-writing
1169
    ///
1170
    /// # Panics
1171
    ///
1172
    /// This function will panic if 0 is passed as an argument.
1173
    ///
1174
    /// # Examples
1175
    ///
1176
    /// ```
1177
    /// # #[cfg(not(target_family = "wasm"))]
1178
    /// # {
1179
    /// # use tokio::runtime;
1180
    /// # pub fn main() {
1181
    /// let rt = runtime::Builder::new_multi_thread()
1182
    ///     .global_queue_interval(31)
1183
    ///     .build();
1184
    /// # }
1185
    /// # }
1186
    /// ```
1187
    #[track_caller]
1188
0
    pub fn global_queue_interval(&mut self, val: u32) -> &mut Self {
1189
0
        assert!(val > 0, "global_queue_interval must be greater than 0");
1190
0
        self.global_queue_interval = Some(val);
1191
0
        self
1192
0
    }
1193
1194
    /// Sets the number of scheduler ticks after which the scheduler will poll for
1195
    /// external events (timers, I/O, and so on).
1196
    ///
1197
    /// A scheduler "tick" roughly corresponds to one `poll` invocation on a task.
1198
    ///
1199
    /// By default, the event interval is `61` for all scheduler types.
1200
    ///
1201
    /// Setting the event interval determines the effective "priority" of delivering
1202
    /// these external events (which may wake up additional tasks), compared to
1203
    /// executing tasks that are currently ready to run. A smaller value is useful
1204
    /// when tasks frequently spend a long time in polling, or infrequently yield,
1205
    /// which can result in overly long delays picking up I/O events. Conversely,
1206
    /// picking up new events requires extra synchronization and syscall overhead,
1207
    /// so if tasks generally complete their polling quickly, a higher event interval
1208
    /// will minimize that overhead while still keeping the scheduler responsive to
1209
    /// events.
1210
    ///
1211
    /// # Panics
1212
    ///
1213
    /// This function will panic if 0 is passed as an argument.
1214
    ///
1215
    /// # Examples
1216
    ///
1217
    /// ```
1218
    /// # #[cfg(not(target_family = "wasm"))]
1219
    /// # {
1220
    /// # use tokio::runtime;
1221
    /// # pub fn main() {
1222
    /// let rt = runtime::Builder::new_multi_thread()
1223
    ///     .event_interval(31)
1224
    ///     .build();
1225
    /// # }
1226
    /// # }
1227
    /// ```
1228
    #[track_caller]
1229
0
    pub fn event_interval(&mut self, val: u32) -> &mut Self {
1230
0
        assert!(val > 0, "event_interval must be greater than 0");
1231
0
        self.event_interval = val;
1232
0
        self
1233
0
    }
1234
1235
    cfg_unstable! {
1236
        /// Configure how the runtime responds to an unhandled panic on a
1237
        /// spawned task.
1238
        ///
1239
        /// By default, an unhandled panic (i.e. a panic not caught by
1240
        /// [`std::panic::catch_unwind`]) has no impact on the runtime's
1241
        /// execution. The panic's error value is forwarded to the task's
1242
        /// [`JoinHandle`] and all other spawned tasks continue running.
1243
        ///
1244
        /// The `unhandled_panic` option enables configuring this behavior.
1245
        ///
1246
        /// * `UnhandledPanic::Ignore` is the default behavior. Panics on
1247
        ///   spawned tasks have no impact on the runtime's execution.
1248
        /// * `UnhandledPanic::ShutdownRuntime` will force the runtime to
1249
        ///   shutdown immediately when a spawned task panics even if that
1250
        ///   task's `JoinHandle` has not been dropped. All other spawned tasks
1251
        ///   will immediately terminate and further calls to
1252
        ///   [`Runtime::block_on`] will panic.
1253
        ///
1254
        /// # Panics
1255
        /// This method panics if called with [`UnhandledPanic::ShutdownRuntime`]
1256
        /// on a runtime other than the current thread runtime.
1257
        ///
1258
        /// # Unstable
1259
        ///
1260
        /// This option is currently unstable and its implementation is
1261
        /// incomplete. The API may change or be removed in the future. See
1262
        /// issue [tokio-rs/tokio#4516] for more details.
1263
        ///
1264
        /// # Examples
1265
        ///
1266
        /// The following demonstrates a runtime configured to shutdown on
1267
        /// panic. The first spawned task panics and results in the runtime
1268
        /// shutting down. The second spawned task never has a chance to
1269
        /// execute. The call to `block_on` will panic due to the runtime being
1270
        /// forcibly shutdown.
1271
        ///
1272
        /// ```should_panic
1273
        /// use tokio::runtime::{self, UnhandledPanic};
1274
        ///
1275
        /// # pub fn main() {
1276
        /// let rt = runtime::Builder::new_current_thread()
1277
        ///     .unhandled_panic(UnhandledPanic::ShutdownRuntime)
1278
        ///     .build()
1279
        ///     .unwrap();
1280
        ///
1281
        /// rt.spawn(async { panic!("boom"); });
1282
        /// rt.spawn(async {
1283
        ///     // This task never completes.
1284
        /// });
1285
        ///
1286
        /// rt.block_on(async {
1287
        ///     // Do some work
1288
        /// # loop { tokio::task::yield_now().await; }
1289
        /// })
1290
        /// # }
1291
        /// ```
1292
        ///
1293
        /// [`JoinHandle`]: struct@crate::task::JoinHandle
1294
        /// [tokio-rs/tokio#4516]: https://github.com/tokio-rs/tokio/issues/4516
1295
        pub fn unhandled_panic(&mut self, behavior: UnhandledPanic) -> &mut Self {
1296
            if !matches!(self.kind, Kind::CurrentThread) && matches!(behavior, UnhandledPanic::ShutdownRuntime) {
1297
                panic!("UnhandledPanic::ShutdownRuntime is only supported in current thread runtime");
1298
            }
1299
1300
            self.unhandled_panic = behavior;
1301
            self
1302
        }
1303
1304
        /// Disables the LIFO task scheduler heuristic.
1305
        ///
1306
        /// The multi-threaded scheduler includes a heuristic for optimizing
1307
        /// message-passing patterns. This heuristic results in the **last**
1308
        /// scheduled task being polled first.
1309
        ///
1310
        /// To implement this heuristic, each worker thread has a slot which
1311
        /// holds the task that should be polled next. However, this slot cannot
1312
        /// be stolen by other worker threads, which can result in lower total
1313
        /// throughput when tasks tend to have longer poll times.
1314
        ///
1315
        /// This configuration option will disable this heuristic resulting in
1316
        /// all scheduled tasks being pushed into the worker-local queue, which
1317
        /// is stealable.
1318
        ///
1319
        /// Consider trying this option when the task "scheduled" time is high
1320
        /// but the runtime is underutilized. Use [tokio-rs/tokio-metrics] to
1321
        /// collect this data.
1322
        ///
1323
        /// # Unstable
1324
        ///
1325
        /// This configuration option is considered a workaround for the LIFO
1326
        /// slot not being stealable. When the slot becomes stealable, we will
1327
        /// revisit whether or not this option is necessary. See
1328
        /// issue [tokio-rs/tokio#4941].
1329
        ///
1330
        /// # Examples
1331
        ///
1332
        /// ```
1333
        /// # #[cfg(not(target_family = "wasm"))]
1334
        /// # {
1335
        /// use tokio::runtime;
1336
        ///
1337
        /// let rt = runtime::Builder::new_multi_thread()
1338
        ///     .disable_lifo_slot()
1339
        ///     .build()
1340
        ///     .unwrap();
1341
        /// # }
1342
        /// ```
1343
        ///
1344
        /// [tokio-rs/tokio-metrics]: https://github.com/tokio-rs/tokio-metrics
1345
        /// [tokio-rs/tokio#4941]: https://github.com/tokio-rs/tokio/issues/4941
1346
        pub fn disable_lifo_slot(&mut self) -> &mut Self {
1347
            self.disable_lifo_slot = true;
1348
            self
1349
        }
1350
1351
        /// Specifies the random number generation seed to use within all
1352
        /// threads associated with the runtime being built.
1353
        ///
1354
        /// This option is intended to make certain parts of the runtime
1355
        /// deterministic (e.g. the [`tokio::select!`] macro). In the case of
1356
        /// [`tokio::select!`] it will ensure that the order that branches are
1357
        /// polled is deterministic.
1358
        ///
1359
        /// In addition to the code specifying `rng_seed` and interacting with
1360
        /// the runtime, the internals of Tokio and the Rust compiler may affect
1361
        /// the sequences of random numbers. In order to ensure repeatable
1362
        /// results, the version of Tokio, the versions of all other
1363
        /// dependencies that interact with Tokio, and the Rust compiler version
1364
        /// should also all remain constant.
1365
        ///
1366
        /// # Examples
1367
        ///
1368
        /// ```
1369
        /// # use tokio::runtime::{self, RngSeed};
1370
        /// # pub fn main() {
1371
        /// let seed = RngSeed::from_bytes(b"place your seed here");
1372
        /// let rt = runtime::Builder::new_current_thread()
1373
        ///     .rng_seed(seed)
1374
        ///     .build();
1375
        /// # }
1376
        /// ```
1377
        ///
1378
        /// [`tokio::select!`]: crate::select
1379
        pub fn rng_seed(&mut self, seed: RngSeed) -> &mut Self {
1380
            self.seed_generator = RngSeedGenerator::new(seed);
1381
            self
1382
        }
1383
    }
1384
1385
    cfg_unstable_metrics! {
1386
        /// Enables tracking the distribution of task poll times.
1387
        ///
1388
        /// Task poll times are not instrumented by default as doing so requires
1389
        /// calling [`Instant::now()`] twice per task poll, which could add
1390
        /// measurable overhead. Use the [`Handle::metrics()`] to access the
1391
        /// metrics data.
1392
        ///
1393
        /// The histogram uses fixed bucket sizes. In other words, the histogram
1394
        /// buckets are not dynamic based on input values. Use the
1395
        /// `metrics_poll_time_histogram` builder methods to configure the
1396
        /// histogram details.
1397
        ///
1398
        /// By default, a linear histogram with 10 buckets each 100 microseconds wide will be used.
1399
        /// This has an extremely low memory footprint, but may not provide enough granularity. For
1400
        /// better granularity with low memory usage, use [`metrics_poll_time_histogram_configuration()`]
1401
        /// to select [`LogHistogram`] instead.
1402
        ///
1403
        /// # Examples
1404
        ///
1405
        /// ```
1406
        /// # #[cfg(not(target_family = "wasm"))]
1407
        /// # {
1408
        /// use tokio::runtime;
1409
        ///
1410
        /// let rt = runtime::Builder::new_multi_thread()
1411
        ///     .enable_metrics_poll_time_histogram()
1412
        ///     .build()
1413
        ///     .unwrap();
1414
        /// # // Test default values here
1415
        /// # fn us(n: u64) -> std::time::Duration { std::time::Duration::from_micros(n) }
1416
        /// # let m = rt.handle().metrics();
1417
        /// # assert_eq!(m.poll_time_histogram_num_buckets(), 10);
1418
        /// # assert_eq!(m.poll_time_histogram_bucket_range(0), us(0)..us(100));
1419
        /// # assert_eq!(m.poll_time_histogram_bucket_range(1), us(100)..us(200));
1420
        /// # }
1421
        /// ```
1422
        ///
1423
        /// [`Handle::metrics()`]: crate::runtime::Handle::metrics
1424
        /// [`Instant::now()`]: std::time::Instant::now
1425
        /// [`LogHistogram`]: crate::runtime::LogHistogram
1426
        /// [`metrics_poll_time_histogram_configuration()`]: Builder::metrics_poll_time_histogram_configuration
1427
        pub fn enable_metrics_poll_time_histogram(&mut self) -> &mut Self {
1428
            self.metrics_poll_count_histogram_enable = true;
1429
            self
1430
        }
1431
1432
        /// Deprecated. Use [`enable_metrics_poll_time_histogram()`] instead.
1433
        ///
1434
        /// [`enable_metrics_poll_time_histogram()`]: Builder::enable_metrics_poll_time_histogram
1435
        #[deprecated(note = "`poll_count_histogram` related methods have been renamed `poll_time_histogram` to better reflect their functionality.")]
1436
        #[doc(hidden)]
1437
        pub fn enable_metrics_poll_count_histogram(&mut self) -> &mut Self {
1438
            self.enable_metrics_poll_time_histogram()
1439
        }
1440
1441
        /// Sets the histogram scale for tracking the distribution of task poll
1442
        /// times.
1443
        ///
1444
        /// Tracking the distribution of task poll times can be done using a
1445
        /// linear or log scale. When using linear scale, each histogram bucket
1446
        /// will represent the same range of poll times. When using log scale,
1447
        /// each histogram bucket will cover a range twice as big as the
1448
        /// previous bucket.
1449
        ///
1450
        /// **Default:** linear scale.
1451
        ///
1452
        /// # Examples
1453
        ///
1454
        /// ```
1455
        /// # #[cfg(not(target_family = "wasm"))]
1456
        /// # {
1457
        /// use tokio::runtime::{self, HistogramScale};
1458
        ///
1459
        /// # #[allow(deprecated)]
1460
        /// let rt = runtime::Builder::new_multi_thread()
1461
        ///     .enable_metrics_poll_time_histogram()
1462
        ///     .metrics_poll_count_histogram_scale(HistogramScale::Log)
1463
        ///     .build()
1464
        ///     .unwrap();
1465
        /// # }
1466
        /// ```
1467
        #[deprecated(note = "use `metrics_poll_time_histogram_configuration`")]
1468
        pub fn metrics_poll_count_histogram_scale(&mut self, histogram_scale: crate::runtime::HistogramScale) -> &mut Self {
1469
            self.metrics_poll_count_histogram.legacy_mut(|b|b.scale = histogram_scale);
1470
            self
1471
        }
1472
1473
        /// Configure the histogram for tracking poll times
1474
        ///
1475
        /// By default, a linear histogram with 10 buckets each 100 microseconds wide will be used.
1476
        /// This has an extremely low memory footprint, but may not provide enough granularity. For
1477
        /// better granularity with low memory usage, use [`LogHistogram`] instead.
1478
        ///
1479
        /// # Examples
1480
        /// Configure a [`LogHistogram`] with [default configuration]:
1481
        /// ```
1482
        /// # #[cfg(not(target_family = "wasm"))]
1483
        /// # {
1484
        /// use tokio::runtime;
1485
        /// use tokio::runtime::{HistogramConfiguration, LogHistogram};
1486
        ///
1487
        /// let rt = runtime::Builder::new_multi_thread()
1488
        ///     .enable_metrics_poll_time_histogram()
1489
        ///     .metrics_poll_time_histogram_configuration(
1490
        ///         HistogramConfiguration::log(LogHistogram::default())
1491
        ///     )
1492
        ///     .build()
1493
        ///     .unwrap();
1494
        /// # }
1495
        /// ```
1496
        ///
1497
        /// Configure a linear histogram with 100 buckets, each 10μs wide
1498
        /// ```
1499
        /// # #[cfg(not(target_family = "wasm"))]
1500
        /// # {
1501
        /// use tokio::runtime;
1502
        /// use std::time::Duration;
1503
        /// use tokio::runtime::HistogramConfiguration;
1504
        ///
1505
        /// let rt = runtime::Builder::new_multi_thread()
1506
        ///     .enable_metrics_poll_time_histogram()
1507
        ///     .metrics_poll_time_histogram_configuration(
1508
        ///         HistogramConfiguration::linear(Duration::from_micros(10), 100)
1509
        ///     )
1510
        ///     .build()
1511
        ///     .unwrap();
1512
        /// # }
1513
        /// ```
1514
        ///
1515
        /// Configure a [`LogHistogram`] with the following settings:
1516
        /// - Measure times from 100ns to 120s
1517
        /// - Max error of 0.1
1518
        /// - No more than 1024 buckets
1519
        /// ```
1520
        /// # #[cfg(not(target_family = "wasm"))]
1521
        /// # {
1522
        /// use std::time::Duration;
1523
        /// use tokio::runtime;
1524
        /// use tokio::runtime::{HistogramConfiguration, LogHistogram};
1525
        ///
1526
        /// let rt = runtime::Builder::new_multi_thread()
1527
        ///     .enable_metrics_poll_time_histogram()
1528
        ///     .metrics_poll_time_histogram_configuration(
1529
        ///         HistogramConfiguration::log(LogHistogram::builder()
1530
        ///             .max_value(Duration::from_secs(120))
1531
        ///             .min_value(Duration::from_nanos(100))
1532
        ///             .max_error(0.1)
1533
        ///             .max_buckets(1024)
1534
        ///             .expect("configuration uses 488 buckets")
1535
        ///         )
1536
        ///     )
1537
        ///     .build()
1538
        ///     .unwrap();
1539
        /// # }
1540
        /// ```
1541
        ///
1542
        /// When migrating from the legacy histogram ([`HistogramScale::Log`]) and wanting
1543
        /// to match the previous behavior, use `precision_exact(0)`. This creates a histogram
1544
        /// where each bucket is twice the size of the previous bucket.
1545
        /// ```rust
1546
        /// use std::time::Duration;
1547
        /// use tokio::runtime::{HistogramConfiguration, LogHistogram};
1548
        /// let rt = tokio::runtime::Builder::new_current_thread()
1549
        ///     .enable_all()
1550
        ///     .enable_metrics_poll_time_histogram()
1551
        ///     .metrics_poll_time_histogram_configuration(HistogramConfiguration::log(
1552
        ///         LogHistogram::builder()
1553
        ///             .min_value(Duration::from_micros(20))
1554
        ///             .max_value(Duration::from_millis(4))
1555
        ///             // Set `precision_exact` to `0` to match `HistogramScale::Log`
1556
        ///             .precision_exact(0)
1557
        ///             .max_buckets(10)
1558
        ///             .unwrap(),
1559
        ///     ))
1560
        ///     .build()
1561
        ///     .unwrap();
1562
        /// ```
1563
        ///
1564
        /// [`LogHistogram`]: crate::runtime::LogHistogram
1565
        /// [default configuration]: crate::runtime::LogHistogramBuilder
1566
        /// [`HistogramScale::Log`]: crate::runtime::HistogramScale::Log
1567
        pub fn metrics_poll_time_histogram_configuration(&mut self, configuration: HistogramConfiguration) -> &mut Self {
1568
            self.metrics_poll_count_histogram.histogram_type = configuration.inner;
1569
            self
1570
        }
1571
1572
        /// Sets the histogram resolution for tracking the distribution of task
1573
        /// poll times.
1574
        ///
1575
        /// The resolution is the histogram's first bucket's range. When using a
1576
        /// linear histogram scale, each bucket will cover the same range. When
1577
        /// using a log scale, each bucket will cover a range twice as big as
1578
        /// the previous bucket. In the log case, the resolution represents the
1579
        /// smallest bucket range.
1580
        ///
1581
        /// Note that, when using log scale, the resolution is rounded up to the
1582
        /// nearest power of 2 in nanoseconds.
1583
        ///
1584
        /// **Default:** 100 microseconds.
1585
        ///
1586
        /// # Examples
1587
        ///
1588
        /// ```
1589
        /// # #[cfg(not(target_family = "wasm"))]
1590
        /// # {
1591
        /// use tokio::runtime;
1592
        /// use std::time::Duration;
1593
        ///
1594
        /// # #[allow(deprecated)]
1595
        /// let rt = runtime::Builder::new_multi_thread()
1596
        ///     .enable_metrics_poll_time_histogram()
1597
        ///     .metrics_poll_count_histogram_resolution(Duration::from_micros(100))
1598
        ///     .build()
1599
        ///     .unwrap();
1600
        /// # }
1601
        /// ```
1602
        #[deprecated(note = "use `metrics_poll_time_histogram_configuration`")]
1603
        pub fn metrics_poll_count_histogram_resolution(&mut self, resolution: Duration) -> &mut Self {
1604
            assert!(resolution > Duration::from_secs(0));
1605
            // Sanity check the argument and also make the cast below safe.
1606
            assert!(resolution <= Duration::from_secs(1));
1607
1608
            let resolution = resolution.as_nanos() as u64;
1609
1610
            self.metrics_poll_count_histogram.legacy_mut(|b|b.resolution = resolution);
1611
            self
1612
        }
1613
1614
        /// Sets the number of buckets for the histogram tracking the
1615
        /// distribution of task poll times.
1616
        ///
1617
        /// The last bucket tracks all greater values that fall out of other
1618
        /// ranges. So, configuring the histogram using a linear scale,
1619
        /// resolution of 50ms, and 10 buckets, the 10th bucket will track task
1620
        /// polls that take more than 450ms to complete.
1621
        ///
1622
        /// **Default:** 10
1623
        ///
1624
        /// # Examples
1625
        ///
1626
        /// ```
1627
        /// # #[cfg(not(target_family = "wasm"))]
1628
        /// # {
1629
        /// use tokio::runtime;
1630
        ///
1631
        /// # #[allow(deprecated)]
1632
        /// let rt = runtime::Builder::new_multi_thread()
1633
        ///     .enable_metrics_poll_time_histogram()
1634
        ///     .metrics_poll_count_histogram_buckets(15)
1635
        ///     .build()
1636
        ///     .unwrap();
1637
        /// # }
1638
        /// ```
1639
        #[deprecated(note = "use `metrics_poll_time_histogram_configuration`")]
1640
        pub fn metrics_poll_count_histogram_buckets(&mut self, buckets: usize) -> &mut Self {
1641
            self.metrics_poll_count_histogram.legacy_mut(|b|b.num_buckets = buckets);
1642
            self
1643
        }
1644
    }
1645
1646
0
    fn build_current_thread_runtime(&mut self) -> io::Result<Runtime> {
1647
        use crate::runtime::runtime::Scheduler;
1648
1649
0
        let (scheduler, handle, blocking_pool) =
1650
0
            self.build_current_thread_runtime_components(None)?;
1651
1652
0
        Ok(Runtime::from_parts(
1653
0
            Scheduler::CurrentThread(scheduler),
1654
0
            handle,
1655
0
            blocking_pool,
1656
0
        ))
1657
0
    }
1658
1659
0
    fn build_current_thread_local_runtime(&mut self) -> io::Result<LocalRuntime> {
1660
        use crate::runtime::local_runtime::LocalRuntimeScheduler;
1661
1662
0
        let tid = std::thread::current().id();
1663
1664
0
        let (scheduler, handle, blocking_pool) =
1665
0
            self.build_current_thread_runtime_components(Some(tid))?;
1666
1667
0
        Ok(LocalRuntime::from_parts(
1668
0
            LocalRuntimeScheduler::CurrentThread(scheduler),
1669
0
            handle,
1670
0
            blocking_pool,
1671
0
        ))
1672
0
    }
1673
1674
0
    fn build_current_thread_runtime_components(
1675
0
        &mut self,
1676
0
        local_tid: Option<ThreadId>,
1677
0
    ) -> io::Result<(CurrentThread, Handle, BlockingPool)> {
1678
        use crate::runtime::scheduler;
1679
        use crate::runtime::Config;
1680
1681
0
        let mut cfg = self.get_cfg();
1682
0
        cfg.timer_flavor = TimerFlavor::Traditional;
1683
0
        let (driver, driver_handle) = driver::Driver::new(cfg)?;
1684
1685
        // Blocking pool
1686
0
        let blocking_pool = blocking::create_blocking_pool(self, self.max_blocking_threads);
1687
0
        let blocking_spawner = blocking_pool.spawner().clone();
1688
1689
        // Generate a rng seed for this runtime.
1690
0
        let seed_generator_1 = self.seed_generator.next_generator();
1691
0
        let seed_generator_2 = self.seed_generator.next_generator();
1692
1693
        // And now put a single-threaded scheduler on top of the timer. When
1694
        // there are no futures ready to do something, it'll let the timer or
1695
        // the reactor to generate some new stimuli for the futures to continue
1696
        // in their life.
1697
0
        let (scheduler, handle) = CurrentThread::new(
1698
0
            driver,
1699
0
            driver_handle,
1700
0
            blocking_spawner,
1701
0
            seed_generator_2,
1702
0
            Config {
1703
0
                before_park: self.before_park.clone(),
1704
0
                after_unpark: self.after_unpark.clone(),
1705
0
                before_spawn: self.before_spawn.clone(),
1706
0
                #[cfg(tokio_unstable)]
1707
0
                before_poll: self.before_poll.clone(),
1708
0
                #[cfg(tokio_unstable)]
1709
0
                after_poll: self.after_poll.clone(),
1710
0
                after_termination: self.after_termination.clone(),
1711
0
                global_queue_interval: self.global_queue_interval,
1712
0
                event_interval: self.event_interval,
1713
0
                #[cfg(tokio_unstable)]
1714
0
                unhandled_panic: self.unhandled_panic.clone(),
1715
0
                disable_lifo_slot: self.disable_lifo_slot,
1716
0
                // This setting never makes sense for a current thread runtime,
1717
0
                // as it only configures how the I/O driver is stolen across
1718
0
                // workers.
1719
0
                enable_eager_driver_handoff: false,
1720
0
                seed_generator: seed_generator_1,
1721
0
                metrics_poll_count_histogram: self.metrics_poll_count_histogram_builder(),
1722
0
                metrics_schedule_latency_histogram: self
1723
0
                    .metrics_schedule_latency_histogram_builder(),
1724
0
            },
1725
0
            local_tid,
1726
0
            self.name.clone(),
1727
0
        );
1728
1729
0
        let handle = Handle {
1730
0
            inner: scheduler::Handle::CurrentThread(handle),
1731
0
        };
1732
1733
0
        Ok((scheduler, handle, blocking_pool))
1734
0
    }
1735
1736
18.6k
    fn metrics_poll_count_histogram_builder(&self) -> Option<HistogramBuilder> {
1737
18.6k
        if self.metrics_poll_count_histogram_enable {
1738
0
            Some(self.metrics_poll_count_histogram.clone())
1739
        } else {
1740
18.6k
            None
1741
        }
1742
18.6k
    }
1743
1744
18.6k
    fn metrics_schedule_latency_histogram_builder(&self) -> Option<HistogramBuilder> {
1745
18.6k
        if self.metrics_schedule_latency_histogram_enabled {
1746
0
            Some(self.metrics_schedule_latency_histogram.clone())
1747
        } else {
1748
18.6k
            None
1749
        }
1750
18.6k
    }
1751
}
1752
1753
cfg_io_driver! {
1754
    impl Builder {
1755
        /// Enables the I/O driver.
1756
        ///
1757
        /// Doing this enables using net, process, signal, and some I/O types on
1758
        /// the runtime.
1759
        ///
1760
        /// # Examples
1761
        ///
1762
        /// ```
1763
        /// use tokio::runtime;
1764
        ///
1765
        /// let rt = runtime::Builder::new_multi_thread()
1766
        ///     .enable_io()
1767
        ///     .build()
1768
        ///     .unwrap();
1769
        /// ```
1770
18.6k
        pub fn enable_io(&mut self) -> &mut Self {
1771
18.6k
            self.enable_io = true;
1772
18.6k
            self
1773
18.6k
        }
1774
1775
        /// Enables the I/O driver and configures the max number of events to be
1776
        /// processed per tick.
1777
        ///
1778
        /// # Examples
1779
        ///
1780
        /// ```
1781
        /// use tokio::runtime;
1782
        ///
1783
        /// let rt = runtime::Builder::new_current_thread()
1784
        ///     .enable_io()
1785
        ///     .max_io_events_per_tick(1024)
1786
        ///     .build()
1787
        ///     .unwrap();
1788
        /// ```
1789
0
        pub fn max_io_events_per_tick(&mut self, capacity: usize) -> &mut Self {
1790
0
            self.nevents = capacity;
1791
0
            self
1792
0
        }
1793
    }
1794
}
1795
1796
cfg_time! {
1797
    impl Builder {
1798
        /// Enables the time driver.
1799
        ///
1800
        /// Doing this enables using `tokio::time` on the runtime.
1801
        ///
1802
        /// # Examples
1803
        ///
1804
        /// ```
1805
        /// # #[cfg(not(target_family = "wasm"))]
1806
        /// # {
1807
        /// use tokio::runtime;
1808
        ///
1809
        /// let rt = runtime::Builder::new_multi_thread()
1810
        ///     .enable_time()
1811
        ///     .build()
1812
        ///     .unwrap();
1813
        /// # }
1814
        /// ```
1815
18.6k
        pub fn enable_time(&mut self) -> &mut Self {
1816
18.6k
            self.enable_time = true;
1817
18.6k
            self
1818
18.6k
        }
1819
    }
1820
}
1821
1822
cfg_io_uring! {
1823
    impl Builder {
1824
        /// Enables the tokio's io_uring driver.
1825
        ///
1826
        /// Doing this enables using io_uring operations on the runtime.
1827
        ///
1828
        /// # Examples
1829
        ///
1830
        /// ```
1831
        /// use tokio::runtime;
1832
        ///
1833
        /// let rt = runtime::Builder::new_multi_thread()
1834
        ///     .enable_io_uring()
1835
        ///     .build()
1836
        ///     .unwrap();
1837
        /// ```
1838
        #[cfg_attr(docsrs, doc(cfg(feature = "io-uring")))]
1839
        pub fn enable_io_uring(&mut self) -> &mut Self {
1840
            // Currently, the uring flag is equivalent to `enable_io`.
1841
            self.enable_io = true;
1842
            self
1843
        }
1844
    }
1845
}
1846
1847
cfg_test_util! {
1848
    impl Builder {
1849
        /// Controls if the runtime's clock starts paused or advancing.
1850
        ///
1851
        /// Pausing time requires the current-thread runtime; construction of
1852
        /// the runtime will panic otherwise.
1853
        ///
1854
        /// # Examples
1855
        ///
1856
        /// ```
1857
        /// use tokio::runtime;
1858
        ///
1859
        /// let rt = runtime::Builder::new_current_thread()
1860
        ///     .enable_time()
1861
        ///     .start_paused(true)
1862
        ///     .build()
1863
        ///     .unwrap();
1864
        /// ```
1865
0
        pub fn start_paused(&mut self, start_paused: bool) -> &mut Self {
1866
0
            self.start_paused = start_paused;
1867
0
            self
1868
0
        }
1869
    }
1870
}
1871
1872
cfg_schedule_latency! {
1873
    impl Builder {
1874
        /// Enables tracking the distribution of task schedule latencies. Task
1875
        /// schedule latency is the time between when a task is scheduled for
1876
        /// execution and when it is polled.
1877
        ///
1878
        /// **This feature is only supported on 64-bit targets.**
1879
        ///
1880
        /// Task schedule latencies are not instrumented by default as doing
1881
        /// so requires calling [`Instant::now()`] when a task is scheduled
1882
        /// and when it is polled, which could add measurable overhead. Use
1883
        /// the [`Handle::metrics()`] to access the metrics data.
1884
        ///
1885
        /// By default, a linear histogram with 10 buckets each 100 microseconds wide will be used.
1886
        /// This has an extremely low memory footprint, but may not provide enough granularity. For
1887
        /// better granularity with low memory usage, use [`metrics_schedule_latency_histogram_configuration()`]
1888
        /// to select [`LogHistogram`] instead.
1889
        ///
1890
        /// # Examples
1891
        ///
1892
        /// ```
1893
        /// # #[cfg(not(target_family = "wasm"))]
1894
        /// # {
1895
        /// use tokio::runtime;
1896
        ///
1897
        /// let rt = runtime::Builder::new_multi_thread()
1898
        ///     .enable_metrics_schedule_latency_histogram()
1899
        ///     .build()
1900
        ///     .unwrap();
1901
        /// # // Test default values here
1902
        /// # fn us(n: u64) -> std::time::Duration { std::time::Duration::from_micros(n) }
1903
        /// # let m = rt.handle().metrics();
1904
        /// # assert_eq!(m.schedule_latency_histogram_num_buckets(), 10);
1905
        /// # assert_eq!(m.schedule_latency_histogram_bucket_range(0), us(0)..us(100));
1906
        /// # assert_eq!(m.schedule_latency_histogram_bucket_range(1), us(100)..us(200));
1907
        /// # }
1908
        /// ```
1909
        ///
1910
        /// [`Handle::metrics()`]: crate::runtime::Handle::metrics
1911
        /// [`Instant::now()`]: std::time::Instant::now
1912
        /// [`LogHistogram`]: crate::runtime::LogHistogram
1913
        /// [`metrics_schedule_latency_histogram_configuration()`]: Builder::metrics_schedule_latency_histogram_configuration
1914
        pub fn enable_metrics_schedule_latency_histogram(&mut self) -> &mut Self {
1915
            self.metrics_schedule_latency_histogram_enabled = true;
1916
            self
1917
        }
1918
1919
        /// Configure the histogram for tracking task schedule latencies.
1920
        ///
1921
        /// Tracking of task schedule latencies must be enabled with
1922
        /// [`enable_metrics_schedule_latency_histogram()`] for this function
1923
        /// to have any effect.
1924
        ///
1925
        /// By default, a linear histogram with 10 buckets each 100 microseconds wide will be used.
1926
        /// This has an extremely low memory footprint, but may not provide enough granularity. For
1927
        /// better granularity with low memory usage, use [`LogHistogram`] instead.
1928
        ///
1929
        /// # Examples
1930
        /// Configure a [`LogHistogram`] with [default configuration]:
1931
        /// ```
1932
        /// # #[cfg(not(target_family = "wasm"))]
1933
        /// # {
1934
        /// use tokio::runtime;
1935
        /// use tokio::runtime::{HistogramConfiguration, LogHistogram};
1936
        ///
1937
        /// let rt = runtime::Builder::new_multi_thread()
1938
        ///     .enable_metrics_schedule_latency_histogram()
1939
        ///     .metrics_schedule_latency_histogram_configuration(
1940
        ///         HistogramConfiguration::log(LogHistogram::default())
1941
        ///     )
1942
        ///     .build()
1943
        ///     .unwrap();
1944
        /// # }
1945
        /// ```
1946
        ///
1947
        /// Configure a linear histogram with 100 buckets, each 10μs wide
1948
        /// ```
1949
        /// # #[cfg(not(target_family = "wasm"))]
1950
        /// # {
1951
        /// use tokio::runtime;
1952
        /// use std::time::Duration;
1953
        /// use tokio::runtime::HistogramConfiguration;
1954
        ///
1955
        /// let rt = runtime::Builder::new_multi_thread()
1956
        ///     .enable_metrics_schedule_latency_histogram()
1957
        ///     .metrics_schedule_latency_histogram_configuration(
1958
        ///         HistogramConfiguration::linear(Duration::from_micros(10), 100)
1959
        ///     )
1960
        ///     .build()
1961
        ///     .unwrap();
1962
        /// # }
1963
        /// ```
1964
        ///
1965
        /// Configure a [`LogHistogram`] with the following settings:
1966
        /// - Measure times from 100ns to 120s
1967
        /// - Max error of 0.1
1968
        /// - No more than 1024 buckets
1969
        /// ```
1970
        /// # #[cfg(not(target_family = "wasm"))]
1971
        /// # {
1972
        /// use std::time::Duration;
1973
        /// use tokio::runtime;
1974
        /// use tokio::runtime::{HistogramConfiguration, LogHistogram};
1975
        ///
1976
        /// let rt = runtime::Builder::new_multi_thread()
1977
        ///     .enable_metrics_schedule_latency_histogram()
1978
        ///     .metrics_schedule_latency_histogram_configuration(
1979
        ///         HistogramConfiguration::log(LogHistogram::builder()
1980
        ///             .max_value(Duration::from_secs(120))
1981
        ///             .min_value(Duration::from_nanos(100))
1982
        ///             .max_error(0.1)
1983
        ///             .max_buckets(1024)
1984
        ///             .expect("configuration uses 488 buckets")
1985
        ///         )
1986
        ///     )
1987
        ///     .build()
1988
        ///     .unwrap();
1989
        /// # }
1990
        /// ```
1991
        ///
1992
        /// [`LogHistogram`]: crate::runtime::LogHistogram
1993
        /// [`enable_metrics_schedule_latency_histogram()`]: Builder::enable_metrics_schedule_latency_histogram
1994
        pub fn metrics_schedule_latency_histogram_configuration(&mut self, configuration: HistogramConfiguration) -> &mut Self {
1995
            self.metrics_schedule_latency_histogram.histogram_type = configuration.inner;
1996
            self
1997
        }
1998
    }
1999
}
2000
2001
cfg_rt_multi_thread! {
2002
    impl Builder {
2003
18.6k
        fn build_threaded_runtime(&mut self) -> io::Result<Runtime> {
2004
            use crate::loom::sys::num_cpus;
2005
            use crate::runtime::{Config, runtime::Scheduler};
2006
            use crate::runtime::scheduler::{self, MultiThread};
2007
2008
18.6k
            let worker_threads = self.worker_threads.unwrap_or_else(num_cpus);
2009
2010
18.6k
            let (driver, driver_handle) = driver::Driver::new(self.get_cfg())?;
2011
2012
            // Create the blocking pool
2013
18.6k
            let blocking_pool =
2014
18.6k
                blocking::create_blocking_pool(self, self.max_blocking_threads + worker_threads);
2015
18.6k
            let blocking_spawner = blocking_pool.spawner().clone();
2016
2017
            // Generate a rng seed for this runtime.
2018
18.6k
            let seed_generator_1 = self.seed_generator.next_generator();
2019
18.6k
            let seed_generator_2 = self.seed_generator.next_generator();
2020
2021
18.6k
            let (scheduler, handle, launch) = MultiThread::new(
2022
18.6k
                worker_threads,
2023
18.6k
                driver,
2024
18.6k
                driver_handle,
2025
18.6k
                blocking_spawner,
2026
18.6k
                seed_generator_2,
2027
18.6k
                Config {
2028
18.6k
                    before_park: self.before_park.clone(),
2029
18.6k
                    after_unpark: self.after_unpark.clone(),
2030
18.6k
                    before_spawn: self.before_spawn.clone(),
2031
18.6k
                    #[cfg(tokio_unstable)]
2032
18.6k
                    before_poll: self.before_poll.clone(),
2033
18.6k
                    #[cfg(tokio_unstable)]
2034
18.6k
                    after_poll: self.after_poll.clone(),
2035
18.6k
                    after_termination: self.after_termination.clone(),
2036
18.6k
                    global_queue_interval: self.global_queue_interval,
2037
18.6k
                    event_interval: self.event_interval,
2038
18.6k
                    #[cfg(tokio_unstable)]
2039
18.6k
                    unhandled_panic: self.unhandled_panic.clone(),
2040
18.6k
                    disable_lifo_slot: self.disable_lifo_slot,
2041
18.6k
                    enable_eager_driver_handoff: self.enable_eager_driver_handoff,
2042
18.6k
                    seed_generator: seed_generator_1,
2043
18.6k
                    metrics_poll_count_histogram: self.metrics_poll_count_histogram_builder(),
2044
18.6k
                    metrics_schedule_latency_histogram: self.metrics_schedule_latency_histogram_builder(),
2045
18.6k
                },
2046
18.6k
                self.timer_flavor,
2047
18.6k
                self.name.clone(),
2048
18.6k
            );
2049
2050
18.6k
            let handle = Handle { inner: scheduler::Handle::MultiThread(handle) };
2051
2052
            // Spawn the thread pool workers
2053
18.6k
            let _enter = handle.enter();
2054
18.6k
            launch.launch();
2055
2056
18.6k
            Ok(Runtime::from_parts(Scheduler::MultiThread(scheduler), handle, blocking_pool))
2057
18.6k
        }
2058
    }
2059
}
2060
2061
impl fmt::Debug for Builder {
2062
0
    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
2063
0
        let mut debug = fmt.debug_struct("Builder");
2064
2065
0
        if let Some(name) = &self.name {
2066
0
            debug.field("name", name);
2067
0
        }
2068
2069
0
        debug
2070
0
            .field("worker_threads", &self.worker_threads)
2071
0
            .field("max_blocking_threads", &self.max_blocking_threads)
2072
0
            .field(
2073
0
                "thread_name",
2074
0
                &"<dyn Fn() -> String + Send + Sync + 'static>",
2075
0
            )
2076
0
            .field("thread_stack_size", &self.thread_stack_size)
2077
0
            .field("after_start", &self.after_start.as_ref().map(|_| "..."))
2078
0
            .field("before_stop", &self.before_stop.as_ref().map(|_| "..."))
2079
0
            .field("before_park", &self.before_park.as_ref().map(|_| "..."))
2080
0
            .field("after_unpark", &self.after_unpark.as_ref().map(|_| "..."))
2081
0
            .field(
2082
0
                "enable_eager_driver_handoff",
2083
0
                &self.enable_eager_driver_handoff,
2084
            );
2085
2086
0
        if self.name.is_none() {
2087
0
            debug.finish_non_exhaustive()
2088
        } else {
2089
0
            debug.finish()
2090
        }
2091
0
    }
2092
}