852 lines
77 KiB
HTML
852 lines
77 KiB
HTML
<!DOCTYPE html><html lang="en"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width, initial-scale=1.0"><meta name="generator" content="rustdoc"><meta name="description" content="Source of the Rust file `/home/newkirk/.cargo/registry/src/index.crates.io-1949cf8c6b5b557f/tokio-1.49.0/src/task/join_set.rs`."><title>join_set.rs - source</title><script>if(window.location.protocol!=="file:")document.head.insertAdjacentHTML("beforeend","SourceSerif4-Regular-6b053e98.ttf.woff2,FiraSans-Italic-81dc35de.woff2,FiraSans-Regular-0fe48ade.woff2,FiraSans-MediumItalic-ccf7e434.woff2,FiraSans-Medium-e1aa3f0a.woff2,SourceCodePro-Regular-8badfe75.ttf.woff2,SourceCodePro-Semibold-aa29a496.ttf.woff2".split(",").map(f=>`<link rel="preload" as="font" type="font/woff2"href="../../../static.files/${f}">`).join(""))</script><link rel="stylesheet" href="../../../static.files/normalize-9960930a.css"><link rel="stylesheet" href="../../../static.files/rustdoc-ca0dd0c4.css"><meta name="rustdoc-vars" data-root-path="../../../" data-static-root-path="../../../static.files/" data-current-crate="tokio" data-themes="" data-resource-suffix="" data-rustdoc-version="1.93.1 (01f6ddf75 2026-02-11) (Arch Linux rust 1:1.93.1-1)" data-channel="1.93.1" data-search-js="search-9e2438ea.js" data-stringdex-js="stringdex-a3946164.js" data-settings-js="settings-c38705f0.js" ><script src="../../../static.files/storage-e2aeef58.js"></script><script defer src="../../../static.files/src-script-813739b1.js"></script><script defer src="../../../src-files.js"></script><script defer src="../../../static.files/main-a410ff4d.js"></script><noscript><link rel="stylesheet" href="../../../static.files/noscript-263c88ec.css"></noscript><link rel="alternate icon" type="image/png" href="../../../static.files/favicon-32x32-eab170b8.png"><link rel="icon" type="image/svg+xml" href="../../../static.files/favicon-044be391.svg"></head><body class="rustdoc src"><!--[if lte IE 11]><div class="warning">This old browser is unsupported and will most likely display funky things.</div><![endif]--><nav class="sidebar"><div class="src-sidebar-title"><h2>Files</h2></div></nav><div class="sidebar-resizer" title="Drag to resize sidebar"></div><main><section id="main-content" class="content"><div class="main-heading"><h1><div class="sub-heading">tokio/task/</div>join_set.rs</h1><rustdoc-toolbar></rustdoc-toolbar></div><div class="example-wrap digits-3"><pre class="rust"><code><a href=#1 id=1 data-nosnippet>1</a><span class="doccomment">//! A collection of tasks spawned on a Tokio runtime.
|
|
<a href=#2 id=2 data-nosnippet>2</a>//!
|
|
<a href=#3 id=3 data-nosnippet>3</a>//! This module provides the [`JoinSet`] type, a collection which stores a set
|
|
<a href=#4 id=4 data-nosnippet>4</a>//! of spawned tasks and allows asynchronously awaiting the output of those
|
|
<a href=#5 id=5 data-nosnippet>5</a>//! tasks as they complete. See the documentation for the [`JoinSet`] type for
|
|
<a href=#6 id=6 data-nosnippet>6</a>//! details.
|
|
<a href=#7 id=7 data-nosnippet>7</a></span><span class="kw">use </span>std::future::Future;
|
|
<a href=#8 id=8 data-nosnippet>8</a><span class="kw">use </span>std::pin::Pin;
|
|
<a href=#9 id=9 data-nosnippet>9</a><span class="kw">use </span>std::task::{Context, Poll};
|
|
<a href=#10 id=10 data-nosnippet>10</a><span class="kw">use </span>std::{fmt, panic};
|
|
<a href=#11 id=11 data-nosnippet>11</a>
|
|
<a href=#12 id=12 data-nosnippet>12</a><span class="kw">use </span><span class="kw">crate</span>::runtime::Handle;
|
|
<a href=#13 id=13 data-nosnippet>13</a><span class="kw">use </span><span class="kw">crate</span>::task::Id;
|
|
<a href=#14 id=14 data-nosnippet>14</a><span class="kw">use </span><span class="kw">crate</span>::task::{unconstrained, AbortHandle, JoinError, JoinHandle, LocalSet};
|
|
<a href=#15 id=15 data-nosnippet>15</a><span class="kw">use </span><span class="kw">crate</span>::util::IdleNotifiedSet;
|
|
<a href=#16 id=16 data-nosnippet>16</a>
|
|
<a href=#17 id=17 data-nosnippet>17</a><span class="doccomment">/// A collection of tasks spawned on a Tokio runtime.
|
|
<a href=#18 id=18 data-nosnippet>18</a>///
|
|
<a href=#19 id=19 data-nosnippet>19</a>/// A `JoinSet` can be used to await the completion of some or all of the tasks
|
|
<a href=#20 id=20 data-nosnippet>20</a>/// in the set. The set is not ordered, and the tasks will be returned in the
|
|
<a href=#21 id=21 data-nosnippet>21</a>/// order they complete.
|
|
<a href=#22 id=22 data-nosnippet>22</a>///
|
|
<a href=#23 id=23 data-nosnippet>23</a>/// All of the tasks must have the same return type `T`.
|
|
<a href=#24 id=24 data-nosnippet>24</a>///
|
|
<a href=#25 id=25 data-nosnippet>25</a>/// When the `JoinSet` is dropped, all tasks in the `JoinSet` are immediately aborted.
|
|
<a href=#26 id=26 data-nosnippet>26</a>///
|
|
<a href=#27 id=27 data-nosnippet>27</a>/// # Examples
|
|
<a href=#28 id=28 data-nosnippet>28</a>///
|
|
<a href=#29 id=29 data-nosnippet>29</a>/// Spawn multiple tasks and wait for them.
|
|
<a href=#30 id=30 data-nosnippet>30</a>///
|
|
<a href=#31 id=31 data-nosnippet>31</a>/// ```
|
|
<a href=#32 id=32 data-nosnippet>32</a>/// use tokio::task::JoinSet;
|
|
<a href=#33 id=33 data-nosnippet>33</a>///
|
|
<a href=#34 id=34 data-nosnippet>34</a>/// # #[tokio::main(flavor = "current_thread")]
|
|
<a href=#35 id=35 data-nosnippet>35</a>/// # async fn main() {
|
|
<a href=#36 id=36 data-nosnippet>36</a>/// let mut set = JoinSet::new();
|
|
<a href=#37 id=37 data-nosnippet>37</a>///
|
|
<a href=#38 id=38 data-nosnippet>38</a>/// for i in 0..10 {
|
|
<a href=#39 id=39 data-nosnippet>39</a>/// set.spawn(async move { i });
|
|
<a href=#40 id=40 data-nosnippet>40</a>/// }
|
|
<a href=#41 id=41 data-nosnippet>41</a>///
|
|
<a href=#42 id=42 data-nosnippet>42</a>/// let mut seen = [false; 10];
|
|
<a href=#43 id=43 data-nosnippet>43</a>/// while let Some(res) = set.join_next().await {
|
|
<a href=#44 id=44 data-nosnippet>44</a>/// let idx = res.unwrap();
|
|
<a href=#45 id=45 data-nosnippet>45</a>/// seen[idx] = true;
|
|
<a href=#46 id=46 data-nosnippet>46</a>/// }
|
|
<a href=#47 id=47 data-nosnippet>47</a>///
|
|
<a href=#48 id=48 data-nosnippet>48</a>/// for i in 0..10 {
|
|
<a href=#49 id=49 data-nosnippet>49</a>/// assert!(seen[i]);
|
|
<a href=#50 id=50 data-nosnippet>50</a>/// }
|
|
<a href=#51 id=51 data-nosnippet>51</a>/// # }
|
|
<a href=#52 id=52 data-nosnippet>52</a>/// ```
|
|
<a href=#53 id=53 data-nosnippet>53</a>///
|
|
<a href=#54 id=54 data-nosnippet>54</a>/// # Task ID guarantees
|
|
<a href=#55 id=55 data-nosnippet>55</a>///
|
|
<a href=#56 id=56 data-nosnippet>56</a>/// While a task is tracked in a `JoinSet`, that task's ID is unique relative
|
|
<a href=#57 id=57 data-nosnippet>57</a>/// to all other running tasks in Tokio. For this purpose, tracking a task in a
|
|
<a href=#58 id=58 data-nosnippet>58</a>/// `JoinSet` is equivalent to holding a [`JoinHandle`] to it. See the [task ID]
|
|
<a href=#59 id=59 data-nosnippet>59</a>/// documentation for more info.
|
|
<a href=#60 id=60 data-nosnippet>60</a>///
|
|
<a href=#61 id=61 data-nosnippet>61</a>/// [`JoinHandle`]: crate::task::JoinHandle
|
|
<a href=#62 id=62 data-nosnippet>62</a>/// [task ID]: crate::task::Id
|
|
<a href=#63 id=63 data-nosnippet>63</a></span><span class="attr">#[cfg_attr(docsrs, doc(cfg(feature = <span class="string">"rt"</span>)))]
|
|
<a href=#64 id=64 data-nosnippet>64</a></span><span class="kw">pub struct </span>JoinSet<T> {
|
|
<a href=#65 id=65 data-nosnippet>65</a> inner: IdleNotifiedSet<JoinHandle<T>>,
|
|
<a href=#66 id=66 data-nosnippet>66</a>}
|
|
<a href=#67 id=67 data-nosnippet>67</a>
|
|
<a href=#68 id=68 data-nosnippet>68</a><span class="doccomment">/// A variant of [`task::Builder`] that spawns tasks on a [`JoinSet`] rather
|
|
<a href=#69 id=69 data-nosnippet>69</a>/// than on the current default runtime.
|
|
<a href=#70 id=70 data-nosnippet>70</a>///
|
|
<a href=#71 id=71 data-nosnippet>71</a>/// [`task::Builder`]: crate::task::Builder
|
|
<a href=#72 id=72 data-nosnippet>72</a></span><span class="attr">#[cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))]
|
|
<a href=#73 id=73 data-nosnippet>73</a>#[cfg_attr(docsrs, doc(cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))))]
|
|
<a href=#74 id=74 data-nosnippet>74</a>#[must_use = <span class="string">"builders do nothing unless used to spawn a task"</span>]
|
|
<a href=#75 id=75 data-nosnippet>75</a></span><span class="kw">pub struct </span>Builder<<span class="lifetime">'a</span>, T> {
|
|
<a href=#76 id=76 data-nosnippet>76</a> joinset: <span class="kw-2">&</span><span class="lifetime">'a </span><span class="kw-2">mut </span>JoinSet<T>,
|
|
<a href=#77 id=77 data-nosnippet>77</a> builder: <span class="kw">super</span>::Builder<<span class="lifetime">'a</span>>,
|
|
<a href=#78 id=78 data-nosnippet>78</a>}
|
|
<a href=#79 id=79 data-nosnippet>79</a>
|
|
<a href=#80 id=80 data-nosnippet>80</a><span class="kw">impl</span><T> JoinSet<T> {
|
|
<a href=#81 id=81 data-nosnippet>81</a> <span class="doccomment">/// Create a new `JoinSet`.
|
|
<a href=#82 id=82 data-nosnippet>82</a> </span><span class="kw">pub fn </span>new() -> <span class="self">Self </span>{
|
|
<a href=#83 id=83 data-nosnippet>83</a> <span class="self">Self </span>{
|
|
<a href=#84 id=84 data-nosnippet>84</a> inner: IdleNotifiedSet::new(),
|
|
<a href=#85 id=85 data-nosnippet>85</a> }
|
|
<a href=#86 id=86 data-nosnippet>86</a> }
|
|
<a href=#87 id=87 data-nosnippet>87</a>
|
|
<a href=#88 id=88 data-nosnippet>88</a> <span class="doccomment">/// Returns the number of tasks currently in the `JoinSet`.
|
|
<a href=#89 id=89 data-nosnippet>89</a> </span><span class="kw">pub fn </span>len(<span class="kw-2">&</span><span class="self">self</span>) -> usize {
|
|
<a href=#90 id=90 data-nosnippet>90</a> <span class="self">self</span>.inner.len()
|
|
<a href=#91 id=91 data-nosnippet>91</a> }
|
|
<a href=#92 id=92 data-nosnippet>92</a>
|
|
<a href=#93 id=93 data-nosnippet>93</a> <span class="doccomment">/// Returns whether the `JoinSet` is empty.
|
|
<a href=#94 id=94 data-nosnippet>94</a> </span><span class="kw">pub fn </span>is_empty(<span class="kw-2">&</span><span class="self">self</span>) -> bool {
|
|
<a href=#95 id=95 data-nosnippet>95</a> <span class="self">self</span>.inner.is_empty()
|
|
<a href=#96 id=96 data-nosnippet>96</a> }
|
|
<a href=#97 id=97 data-nosnippet>97</a>}
|
|
<a href=#98 id=98 data-nosnippet>98</a>
|
|
<a href=#99 id=99 data-nosnippet>99</a><span class="kw">impl</span><T: <span class="lifetime">'static</span>> JoinSet<T> {
|
|
<a href=#100 id=100 data-nosnippet>100</a> <span class="doccomment">/// Returns a [`Builder`] that can be used to configure a task prior to
|
|
<a href=#101 id=101 data-nosnippet>101</a> /// spawning it on this `JoinSet`.
|
|
<a href=#102 id=102 data-nosnippet>102</a> ///
|
|
<a href=#103 id=103 data-nosnippet>103</a> /// # Examples
|
|
<a href=#104 id=104 data-nosnippet>104</a> ///
|
|
<a href=#105 id=105 data-nosnippet>105</a> /// ```
|
|
<a href=#106 id=106 data-nosnippet>106</a> /// use tokio::task::JoinSet;
|
|
<a href=#107 id=107 data-nosnippet>107</a> ///
|
|
<a href=#108 id=108 data-nosnippet>108</a> /// #[tokio::main]
|
|
<a href=#109 id=109 data-nosnippet>109</a> /// async fn main() -> std::io::Result<()> {
|
|
<a href=#110 id=110 data-nosnippet>110</a> /// let mut set = JoinSet::new();
|
|
<a href=#111 id=111 data-nosnippet>111</a> ///
|
|
<a href=#112 id=112 data-nosnippet>112</a> /// // Use the builder to configure a task's name before spawning it.
|
|
<a href=#113 id=113 data-nosnippet>113</a> /// set.build_task()
|
|
<a href=#114 id=114 data-nosnippet>114</a> /// .name("my_task")
|
|
<a href=#115 id=115 data-nosnippet>115</a> /// .spawn(async { /* ... */ })?;
|
|
<a href=#116 id=116 data-nosnippet>116</a> ///
|
|
<a href=#117 id=117 data-nosnippet>117</a> /// Ok(())
|
|
<a href=#118 id=118 data-nosnippet>118</a> /// }
|
|
<a href=#119 id=119 data-nosnippet>119</a> /// ```
|
|
<a href=#120 id=120 data-nosnippet>120</a> </span><span class="attr">#[cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))]
|
|
<a href=#121 id=121 data-nosnippet>121</a> #[cfg_attr(docsrs, doc(cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))))]
|
|
<a href=#122 id=122 data-nosnippet>122</a> </span><span class="kw">pub fn </span>build_task(<span class="kw-2">&mut </span><span class="self">self</span>) -> Builder<<span class="lifetime">'_</span>, T> {
|
|
<a href=#123 id=123 data-nosnippet>123</a> Builder {
|
|
<a href=#124 id=124 data-nosnippet>124</a> builder: <span class="kw">super</span>::Builder::new(),
|
|
<a href=#125 id=125 data-nosnippet>125</a> joinset: <span class="self">self</span>,
|
|
<a href=#126 id=126 data-nosnippet>126</a> }
|
|
<a href=#127 id=127 data-nosnippet>127</a> }
|
|
<a href=#128 id=128 data-nosnippet>128</a>
|
|
<a href=#129 id=129 data-nosnippet>129</a> <span class="doccomment">/// Spawn the provided task on the `JoinSet`, returning an [`AbortHandle`]
|
|
<a href=#130 id=130 data-nosnippet>130</a> /// that can be used to remotely cancel the task.
|
|
<a href=#131 id=131 data-nosnippet>131</a> ///
|
|
<a href=#132 id=132 data-nosnippet>132</a> /// The provided future will start running in the background immediately
|
|
<a href=#133 id=133 data-nosnippet>133</a> /// when this method is called, even if you don't await anything on this
|
|
<a href=#134 id=134 data-nosnippet>134</a> /// `JoinSet`.
|
|
<a href=#135 id=135 data-nosnippet>135</a> ///
|
|
<a href=#136 id=136 data-nosnippet>136</a> /// # Panics
|
|
<a href=#137 id=137 data-nosnippet>137</a> ///
|
|
<a href=#138 id=138 data-nosnippet>138</a> /// This method panics if called outside of a Tokio runtime.
|
|
<a href=#139 id=139 data-nosnippet>139</a> ///
|
|
<a href=#140 id=140 data-nosnippet>140</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#141 id=141 data-nosnippet>141</a> </span><span class="attr">#[track_caller]
|
|
<a href=#142 id=142 data-nosnippet>142</a> </span><span class="kw">pub fn </span>spawn<F>(<span class="kw-2">&mut </span><span class="self">self</span>, task: F) -> AbortHandle
|
|
<a href=#143 id=143 data-nosnippet>143</a> <span class="kw">where
|
|
<a href=#144 id=144 data-nosnippet>144</a> </span>F: Future<Output = T>,
|
|
<a href=#145 id=145 data-nosnippet>145</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#146 id=146 data-nosnippet>146</a> T: Send,
|
|
<a href=#147 id=147 data-nosnippet>147</a> {
|
|
<a href=#148 id=148 data-nosnippet>148</a> <span class="self">self</span>.insert(<span class="kw">crate</span>::spawn(task))
|
|
<a href=#149 id=149 data-nosnippet>149</a> }
|
|
<a href=#150 id=150 data-nosnippet>150</a>
|
|
<a href=#151 id=151 data-nosnippet>151</a> <span class="doccomment">/// Spawn the provided task on the provided runtime and store it in this
|
|
<a href=#152 id=152 data-nosnippet>152</a> /// `JoinSet` returning an [`AbortHandle`] that can be used to remotely
|
|
<a href=#153 id=153 data-nosnippet>153</a> /// cancel the task.
|
|
<a href=#154 id=154 data-nosnippet>154</a> ///
|
|
<a href=#155 id=155 data-nosnippet>155</a> /// The provided future will start running in the background immediately
|
|
<a href=#156 id=156 data-nosnippet>156</a> /// when this method is called, even if you don't await anything on this
|
|
<a href=#157 id=157 data-nosnippet>157</a> /// `JoinSet`.
|
|
<a href=#158 id=158 data-nosnippet>158</a> ///
|
|
<a href=#159 id=159 data-nosnippet>159</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#160 id=160 data-nosnippet>160</a> </span><span class="attr">#[track_caller]
|
|
<a href=#161 id=161 data-nosnippet>161</a> </span><span class="kw">pub fn </span>spawn_on<F>(<span class="kw-2">&mut </span><span class="self">self</span>, task: F, handle: <span class="kw-2">&</span>Handle) -> AbortHandle
|
|
<a href=#162 id=162 data-nosnippet>162</a> <span class="kw">where
|
|
<a href=#163 id=163 data-nosnippet>163</a> </span>F: Future<Output = T>,
|
|
<a href=#164 id=164 data-nosnippet>164</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#165 id=165 data-nosnippet>165</a> T: Send,
|
|
<a href=#166 id=166 data-nosnippet>166</a> {
|
|
<a href=#167 id=167 data-nosnippet>167</a> <span class="self">self</span>.insert(handle.spawn(task))
|
|
<a href=#168 id=168 data-nosnippet>168</a> }
|
|
<a href=#169 id=169 data-nosnippet>169</a>
|
|
<a href=#170 id=170 data-nosnippet>170</a> <span class="doccomment">/// Spawn the provided task on the current [`LocalSet`] or [`LocalRuntime`]
|
|
<a href=#171 id=171 data-nosnippet>171</a> /// and store it in this `JoinSet`, returning an [`AbortHandle`] that can
|
|
<a href=#172 id=172 data-nosnippet>172</a> /// be used to remotely cancel the task.
|
|
<a href=#173 id=173 data-nosnippet>173</a> ///
|
|
<a href=#174 id=174 data-nosnippet>174</a> /// The provided future will start running in the background immediately
|
|
<a href=#175 id=175 data-nosnippet>175</a> /// when this method is called, even if you don't await anything on this
|
|
<a href=#176 id=176 data-nosnippet>176</a> /// `JoinSet`.
|
|
<a href=#177 id=177 data-nosnippet>177</a> ///
|
|
<a href=#178 id=178 data-nosnippet>178</a> /// # Panics
|
|
<a href=#179 id=179 data-nosnippet>179</a> ///
|
|
<a href=#180 id=180 data-nosnippet>180</a> /// This method panics if it is called outside of a `LocalSet` or `LocalRuntime`.
|
|
<a href=#181 id=181 data-nosnippet>181</a> ///
|
|
<a href=#182 id=182 data-nosnippet>182</a> /// [`LocalSet`]: crate::task::LocalSet
|
|
<a href=#183 id=183 data-nosnippet>183</a> /// [`LocalRuntime`]: crate::runtime::LocalRuntime
|
|
<a href=#184 id=184 data-nosnippet>184</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#185 id=185 data-nosnippet>185</a> </span><span class="attr">#[track_caller]
|
|
<a href=#186 id=186 data-nosnippet>186</a> </span><span class="kw">pub fn </span>spawn_local<F>(<span class="kw-2">&mut </span><span class="self">self</span>, task: F) -> AbortHandle
|
|
<a href=#187 id=187 data-nosnippet>187</a> <span class="kw">where
|
|
<a href=#188 id=188 data-nosnippet>188</a> </span>F: Future<Output = T>,
|
|
<a href=#189 id=189 data-nosnippet>189</a> F: <span class="lifetime">'static</span>,
|
|
<a href=#190 id=190 data-nosnippet>190</a> {
|
|
<a href=#191 id=191 data-nosnippet>191</a> <span class="self">self</span>.insert(<span class="kw">crate</span>::task::spawn_local(task))
|
|
<a href=#192 id=192 data-nosnippet>192</a> }
|
|
<a href=#193 id=193 data-nosnippet>193</a>
|
|
<a href=#194 id=194 data-nosnippet>194</a> <span class="doccomment">/// Spawn the provided task on the provided [`LocalSet`] and store it in
|
|
<a href=#195 id=195 data-nosnippet>195</a> /// this `JoinSet`, returning an [`AbortHandle`] that can be used to
|
|
<a href=#196 id=196 data-nosnippet>196</a> /// remotely cancel the task.
|
|
<a href=#197 id=197 data-nosnippet>197</a> ///
|
|
<a href=#198 id=198 data-nosnippet>198</a> /// Unlike the [`spawn_local`] method, this method may be used to spawn local
|
|
<a href=#199 id=199 data-nosnippet>199</a> /// tasks on a `LocalSet` that is _not_ currently running. The provided
|
|
<a href=#200 id=200 data-nosnippet>200</a> /// future will start running whenever the `LocalSet` is next started.
|
|
<a href=#201 id=201 data-nosnippet>201</a> ///
|
|
<a href=#202 id=202 data-nosnippet>202</a> /// [`LocalSet`]: crate::task::LocalSet
|
|
<a href=#203 id=203 data-nosnippet>203</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#204 id=204 data-nosnippet>204</a> /// [`spawn_local`]: Self::spawn_local
|
|
<a href=#205 id=205 data-nosnippet>205</a> </span><span class="attr">#[track_caller]
|
|
<a href=#206 id=206 data-nosnippet>206</a> </span><span class="kw">pub fn </span>spawn_local_on<F>(<span class="kw-2">&mut </span><span class="self">self</span>, task: F, local_set: <span class="kw-2">&</span>LocalSet) -> AbortHandle
|
|
<a href=#207 id=207 data-nosnippet>207</a> <span class="kw">where
|
|
<a href=#208 id=208 data-nosnippet>208</a> </span>F: Future<Output = T>,
|
|
<a href=#209 id=209 data-nosnippet>209</a> F: <span class="lifetime">'static</span>,
|
|
<a href=#210 id=210 data-nosnippet>210</a> {
|
|
<a href=#211 id=211 data-nosnippet>211</a> <span class="self">self</span>.insert(local_set.spawn_local(task))
|
|
<a href=#212 id=212 data-nosnippet>212</a> }
|
|
<a href=#213 id=213 data-nosnippet>213</a>
|
|
<a href=#214 id=214 data-nosnippet>214</a> <span class="doccomment">/// Spawn the blocking code on the blocking threadpool and store
|
|
<a href=#215 id=215 data-nosnippet>215</a> /// it in this `JoinSet`, returning an [`AbortHandle`] that can be
|
|
<a href=#216 id=216 data-nosnippet>216</a> /// used to remotely cancel the task.
|
|
<a href=#217 id=217 data-nosnippet>217</a> ///
|
|
<a href=#218 id=218 data-nosnippet>218</a> /// # Examples
|
|
<a href=#219 id=219 data-nosnippet>219</a> ///
|
|
<a href=#220 id=220 data-nosnippet>220</a> /// Spawn multiple blocking tasks and wait for them.
|
|
<a href=#221 id=221 data-nosnippet>221</a> ///
|
|
<a href=#222 id=222 data-nosnippet>222</a> /// ```
|
|
<a href=#223 id=223 data-nosnippet>223</a> /// # #[cfg(not(target_family = "wasm"))]
|
|
<a href=#224 id=224 data-nosnippet>224</a> /// # {
|
|
<a href=#225 id=225 data-nosnippet>225</a> /// use tokio::task::JoinSet;
|
|
<a href=#226 id=226 data-nosnippet>226</a> ///
|
|
<a href=#227 id=227 data-nosnippet>227</a> /// #[tokio::main]
|
|
<a href=#228 id=228 data-nosnippet>228</a> /// async fn main() {
|
|
<a href=#229 id=229 data-nosnippet>229</a> /// let mut set = JoinSet::new();
|
|
<a href=#230 id=230 data-nosnippet>230</a> ///
|
|
<a href=#231 id=231 data-nosnippet>231</a> /// for i in 0..10 {
|
|
<a href=#232 id=232 data-nosnippet>232</a> /// set.spawn_blocking(move || { i });
|
|
<a href=#233 id=233 data-nosnippet>233</a> /// }
|
|
<a href=#234 id=234 data-nosnippet>234</a> ///
|
|
<a href=#235 id=235 data-nosnippet>235</a> /// let mut seen = [false; 10];
|
|
<a href=#236 id=236 data-nosnippet>236</a> /// while let Some(res) = set.join_next().await {
|
|
<a href=#237 id=237 data-nosnippet>237</a> /// let idx = res.unwrap();
|
|
<a href=#238 id=238 data-nosnippet>238</a> /// seen[idx] = true;
|
|
<a href=#239 id=239 data-nosnippet>239</a> /// }
|
|
<a href=#240 id=240 data-nosnippet>240</a> ///
|
|
<a href=#241 id=241 data-nosnippet>241</a> /// for i in 0..10 {
|
|
<a href=#242 id=242 data-nosnippet>242</a> /// assert!(seen[i]);
|
|
<a href=#243 id=243 data-nosnippet>243</a> /// }
|
|
<a href=#244 id=244 data-nosnippet>244</a> /// }
|
|
<a href=#245 id=245 data-nosnippet>245</a> /// # }
|
|
<a href=#246 id=246 data-nosnippet>246</a> /// ```
|
|
<a href=#247 id=247 data-nosnippet>247</a> ///
|
|
<a href=#248 id=248 data-nosnippet>248</a> /// # Panics
|
|
<a href=#249 id=249 data-nosnippet>249</a> ///
|
|
<a href=#250 id=250 data-nosnippet>250</a> /// This method panics if called outside of a Tokio runtime.
|
|
<a href=#251 id=251 data-nosnippet>251</a> ///
|
|
<a href=#252 id=252 data-nosnippet>252</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#253 id=253 data-nosnippet>253</a> </span><span class="attr">#[track_caller]
|
|
<a href=#254 id=254 data-nosnippet>254</a> </span><span class="kw">pub fn </span>spawn_blocking<F>(<span class="kw-2">&mut </span><span class="self">self</span>, f: F) -> AbortHandle
|
|
<a href=#255 id=255 data-nosnippet>255</a> <span class="kw">where
|
|
<a href=#256 id=256 data-nosnippet>256</a> </span>F: FnOnce() -> T,
|
|
<a href=#257 id=257 data-nosnippet>257</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#258 id=258 data-nosnippet>258</a> T: Send,
|
|
<a href=#259 id=259 data-nosnippet>259</a> {
|
|
<a href=#260 id=260 data-nosnippet>260</a> <span class="self">self</span>.insert(<span class="kw">crate</span>::runtime::spawn_blocking(f))
|
|
<a href=#261 id=261 data-nosnippet>261</a> }
|
|
<a href=#262 id=262 data-nosnippet>262</a>
|
|
<a href=#263 id=263 data-nosnippet>263</a> <span class="doccomment">/// Spawn the blocking code on the blocking threadpool of the
|
|
<a href=#264 id=264 data-nosnippet>264</a> /// provided runtime and store it in this `JoinSet`, returning an
|
|
<a href=#265 id=265 data-nosnippet>265</a> /// [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#266 id=266 data-nosnippet>266</a> ///
|
|
<a href=#267 id=267 data-nosnippet>267</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#268 id=268 data-nosnippet>268</a> </span><span class="attr">#[track_caller]
|
|
<a href=#269 id=269 data-nosnippet>269</a> </span><span class="kw">pub fn </span>spawn_blocking_on<F>(<span class="kw-2">&mut </span><span class="self">self</span>, f: F, handle: <span class="kw-2">&</span>Handle) -> AbortHandle
|
|
<a href=#270 id=270 data-nosnippet>270</a> <span class="kw">where
|
|
<a href=#271 id=271 data-nosnippet>271</a> </span>F: FnOnce() -> T,
|
|
<a href=#272 id=272 data-nosnippet>272</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#273 id=273 data-nosnippet>273</a> T: Send,
|
|
<a href=#274 id=274 data-nosnippet>274</a> {
|
|
<a href=#275 id=275 data-nosnippet>275</a> <span class="self">self</span>.insert(handle.spawn_blocking(f))
|
|
<a href=#276 id=276 data-nosnippet>276</a> }
|
|
<a href=#277 id=277 data-nosnippet>277</a>
|
|
<a href=#278 id=278 data-nosnippet>278</a> <span class="kw">fn </span>insert(<span class="kw-2">&mut </span><span class="self">self</span>, jh: JoinHandle<T>) -> AbortHandle {
|
|
<a href=#279 id=279 data-nosnippet>279</a> <span class="kw">let </span>abort = jh.abort_handle();
|
|
<a href=#280 id=280 data-nosnippet>280</a> <span class="kw">let </span><span class="kw-2">mut </span>entry = <span class="self">self</span>.inner.insert_idle(jh);
|
|
<a href=#281 id=281 data-nosnippet>281</a>
|
|
<a href=#282 id=282 data-nosnippet>282</a> <span class="comment">// Set the waker that is notified when the task completes.
|
|
<a href=#283 id=283 data-nosnippet>283</a> </span>entry.with_value_and_context(|jh, ctx| jh.set_join_waker(ctx.waker()));
|
|
<a href=#284 id=284 data-nosnippet>284</a> abort
|
|
<a href=#285 id=285 data-nosnippet>285</a> }
|
|
<a href=#286 id=286 data-nosnippet>286</a>
|
|
<a href=#287 id=287 data-nosnippet>287</a> <span class="doccomment">/// Waits until one of the tasks in the set completes and returns its output.
|
|
<a href=#288 id=288 data-nosnippet>288</a> ///
|
|
<a href=#289 id=289 data-nosnippet>289</a> /// Returns `None` if the set is empty.
|
|
<a href=#290 id=290 data-nosnippet>290</a> ///
|
|
<a href=#291 id=291 data-nosnippet>291</a> /// # Cancel Safety
|
|
<a href=#292 id=292 data-nosnippet>292</a> ///
|
|
<a href=#293 id=293 data-nosnippet>293</a> /// This method is cancel safe. If `join_next` is used as the event in a `tokio::select!`
|
|
<a href=#294 id=294 data-nosnippet>294</a> /// statement and some other branch completes first, it is guaranteed that no tasks were
|
|
<a href=#295 id=295 data-nosnippet>295</a> /// removed from this `JoinSet`.
|
|
<a href=#296 id=296 data-nosnippet>296</a> </span><span class="kw">pub async fn </span>join_next(<span class="kw-2">&mut </span><span class="self">self</span>) -> <span class="prelude-ty">Option</span><<span class="prelude-ty">Result</span><T, JoinError>> {
|
|
<a href=#297 id=297 data-nosnippet>297</a> std::future::poll_fn(|cx| <span class="self">self</span>.poll_join_next(cx)).<span class="kw">await
|
|
<a href=#298 id=298 data-nosnippet>298</a> </span>}
|
|
<a href=#299 id=299 data-nosnippet>299</a>
|
|
<a href=#300 id=300 data-nosnippet>300</a> <span class="doccomment">/// Waits until one of the tasks in the set completes and returns its
|
|
<a href=#301 id=301 data-nosnippet>301</a> /// output, along with the [task ID] of the completed task.
|
|
<a href=#302 id=302 data-nosnippet>302</a> ///
|
|
<a href=#303 id=303 data-nosnippet>303</a> /// Returns `None` if the set is empty.
|
|
<a href=#304 id=304 data-nosnippet>304</a> ///
|
|
<a href=#305 id=305 data-nosnippet>305</a> /// When this method returns an error, then the id of the task that failed can be accessed
|
|
<a href=#306 id=306 data-nosnippet>306</a> /// using the [`JoinError::id`] method.
|
|
<a href=#307 id=307 data-nosnippet>307</a> ///
|
|
<a href=#308 id=308 data-nosnippet>308</a> /// # Cancel Safety
|
|
<a href=#309 id=309 data-nosnippet>309</a> ///
|
|
<a href=#310 id=310 data-nosnippet>310</a> /// This method is cancel safe. If `join_next_with_id` is used as the event in a `tokio::select!`
|
|
<a href=#311 id=311 data-nosnippet>311</a> /// statement and some other branch completes first, it is guaranteed that no tasks were
|
|
<a href=#312 id=312 data-nosnippet>312</a> /// removed from this `JoinSet`.
|
|
<a href=#313 id=313 data-nosnippet>313</a> ///
|
|
<a href=#314 id=314 data-nosnippet>314</a> /// [task ID]: crate::task::Id
|
|
<a href=#315 id=315 data-nosnippet>315</a> /// [`JoinError::id`]: fn@crate::task::JoinError::id
|
|
<a href=#316 id=316 data-nosnippet>316</a> </span><span class="kw">pub async fn </span>join_next_with_id(<span class="kw-2">&mut </span><span class="self">self</span>) -> <span class="prelude-ty">Option</span><<span class="prelude-ty">Result</span><(Id, T), JoinError>> {
|
|
<a href=#317 id=317 data-nosnippet>317</a> std::future::poll_fn(|cx| <span class="self">self</span>.poll_join_next_with_id(cx)).<span class="kw">await
|
|
<a href=#318 id=318 data-nosnippet>318</a> </span>}
|
|
<a href=#319 id=319 data-nosnippet>319</a>
|
|
<a href=#320 id=320 data-nosnippet>320</a> <span class="doccomment">/// Tries to join one of the tasks in the set that has completed and return its output.
|
|
<a href=#321 id=321 data-nosnippet>321</a> ///
|
|
<a href=#322 id=322 data-nosnippet>322</a> /// Returns `None` if there are no completed tasks, or if the set is empty.
|
|
<a href=#323 id=323 data-nosnippet>323</a> </span><span class="kw">pub fn </span>try_join_next(<span class="kw-2">&mut </span><span class="self">self</span>) -> <span class="prelude-ty">Option</span><<span class="prelude-ty">Result</span><T, JoinError>> {
|
|
<a href=#324 id=324 data-nosnippet>324</a> <span class="comment">// Loop over all notified `JoinHandle`s to find one that's ready, or until none are left.
|
|
<a href=#325 id=325 data-nosnippet>325</a> </span><span class="kw">loop </span>{
|
|
<a href=#326 id=326 data-nosnippet>326</a> <span class="kw">let </span><span class="kw-2">mut </span>entry = <span class="self">self</span>.inner.try_pop_notified()<span class="question-mark">?</span>;
|
|
<a href=#327 id=327 data-nosnippet>327</a>
|
|
<a href=#328 id=328 data-nosnippet>328</a> <span class="kw">let </span>res = entry.with_value_and_context(|jh, ctx| {
|
|
<a href=#329 id=329 data-nosnippet>329</a> <span class="comment">// Since this function is not async and cannot be forced to yield, we should
|
|
<a href=#330 id=330 data-nosnippet>330</a> // disable budgeting when we want to check for the `JoinHandle` readiness.
|
|
<a href=#331 id=331 data-nosnippet>331</a> </span>Pin::new(<span class="kw-2">&mut </span>unconstrained(jh)).poll(ctx)
|
|
<a href=#332 id=332 data-nosnippet>332</a> });
|
|
<a href=#333 id=333 data-nosnippet>333</a>
|
|
<a href=#334 id=334 data-nosnippet>334</a> <span class="kw">if let </span>Poll::Ready(res) = res {
|
|
<a href=#335 id=335 data-nosnippet>335</a> <span class="kw">let </span>_entry = entry.remove();
|
|
<a href=#336 id=336 data-nosnippet>336</a>
|
|
<a href=#337 id=337 data-nosnippet>337</a> <span class="kw">return </span><span class="prelude-val">Some</span>(res);
|
|
<a href=#338 id=338 data-nosnippet>338</a> }
|
|
<a href=#339 id=339 data-nosnippet>339</a> }
|
|
<a href=#340 id=340 data-nosnippet>340</a> }
|
|
<a href=#341 id=341 data-nosnippet>341</a>
|
|
<a href=#342 id=342 data-nosnippet>342</a> <span class="doccomment">/// Tries to join one of the tasks in the set that has completed and return its output,
|
|
<a href=#343 id=343 data-nosnippet>343</a> /// along with the [task ID] of the completed task.
|
|
<a href=#344 id=344 data-nosnippet>344</a> ///
|
|
<a href=#345 id=345 data-nosnippet>345</a> /// Returns `None` if there are no completed tasks, or if the set is empty.
|
|
<a href=#346 id=346 data-nosnippet>346</a> ///
|
|
<a href=#347 id=347 data-nosnippet>347</a> /// When this method returns an error, then the id of the task that failed can be accessed
|
|
<a href=#348 id=348 data-nosnippet>348</a> /// using the [`JoinError::id`] method.
|
|
<a href=#349 id=349 data-nosnippet>349</a> ///
|
|
<a href=#350 id=350 data-nosnippet>350</a> /// [task ID]: crate::task::Id
|
|
<a href=#351 id=351 data-nosnippet>351</a> /// [`JoinError::id`]: fn@crate::task::JoinError::id
|
|
<a href=#352 id=352 data-nosnippet>352</a> </span><span class="kw">pub fn </span>try_join_next_with_id(<span class="kw-2">&mut </span><span class="self">self</span>) -> <span class="prelude-ty">Option</span><<span class="prelude-ty">Result</span><(Id, T), JoinError>> {
|
|
<a href=#353 id=353 data-nosnippet>353</a> <span class="comment">// Loop over all notified `JoinHandle`s to find one that's ready, or until none are left.
|
|
<a href=#354 id=354 data-nosnippet>354</a> </span><span class="kw">loop </span>{
|
|
<a href=#355 id=355 data-nosnippet>355</a> <span class="kw">let </span><span class="kw-2">mut </span>entry = <span class="self">self</span>.inner.try_pop_notified()<span class="question-mark">?</span>;
|
|
<a href=#356 id=356 data-nosnippet>356</a>
|
|
<a href=#357 id=357 data-nosnippet>357</a> <span class="kw">let </span>res = entry.with_value_and_context(|jh, ctx| {
|
|
<a href=#358 id=358 data-nosnippet>358</a> <span class="comment">// Since this function is not async and cannot be forced to yield, we should
|
|
<a href=#359 id=359 data-nosnippet>359</a> // disable budgeting when we want to check for the `JoinHandle` readiness.
|
|
<a href=#360 id=360 data-nosnippet>360</a> </span>Pin::new(<span class="kw-2">&mut </span>unconstrained(jh)).poll(ctx)
|
|
<a href=#361 id=361 data-nosnippet>361</a> });
|
|
<a href=#362 id=362 data-nosnippet>362</a>
|
|
<a href=#363 id=363 data-nosnippet>363</a> <span class="kw">if let </span>Poll::Ready(res) = res {
|
|
<a href=#364 id=364 data-nosnippet>364</a> <span class="kw">let </span>entry = entry.remove();
|
|
<a href=#365 id=365 data-nosnippet>365</a>
|
|
<a href=#366 id=366 data-nosnippet>366</a> <span class="kw">return </span><span class="prelude-val">Some</span>(res.map(|output| (entry.id(), output)));
|
|
<a href=#367 id=367 data-nosnippet>367</a> }
|
|
<a href=#368 id=368 data-nosnippet>368</a> }
|
|
<a href=#369 id=369 data-nosnippet>369</a> }
|
|
<a href=#370 id=370 data-nosnippet>370</a>
|
|
<a href=#371 id=371 data-nosnippet>371</a> <span class="doccomment">/// Aborts all tasks and waits for them to finish shutting down.
|
|
<a href=#372 id=372 data-nosnippet>372</a> ///
|
|
<a href=#373 id=373 data-nosnippet>373</a> /// Calling this method is equivalent to calling [`abort_all`] and then calling [`join_next`] in
|
|
<a href=#374 id=374 data-nosnippet>374</a> /// a loop until it returns `None`.
|
|
<a href=#375 id=375 data-nosnippet>375</a> ///
|
|
<a href=#376 id=376 data-nosnippet>376</a> /// This method ignores any panics in the tasks shutting down. When this call returns, the
|
|
<a href=#377 id=377 data-nosnippet>377</a> /// `JoinSet` will be empty.
|
|
<a href=#378 id=378 data-nosnippet>378</a> ///
|
|
<a href=#379 id=379 data-nosnippet>379</a> /// [`abort_all`]: fn@Self::abort_all
|
|
<a href=#380 id=380 data-nosnippet>380</a> /// [`join_next`]: fn@Self::join_next
|
|
<a href=#381 id=381 data-nosnippet>381</a> </span><span class="kw">pub async fn </span>shutdown(<span class="kw-2">&mut </span><span class="self">self</span>) {
|
|
<a href=#382 id=382 data-nosnippet>382</a> <span class="self">self</span>.abort_all();
|
|
<a href=#383 id=383 data-nosnippet>383</a> <span class="kw">while </span><span class="self">self</span>.join_next().<span class="kw">await</span>.is_some() {}
|
|
<a href=#384 id=384 data-nosnippet>384</a> }
|
|
<a href=#385 id=385 data-nosnippet>385</a>
|
|
<a href=#386 id=386 data-nosnippet>386</a> <span class="doccomment">/// Awaits the completion of all tasks in this `JoinSet`, returning a vector of their results.
|
|
<a href=#387 id=387 data-nosnippet>387</a> ///
|
|
<a href=#388 id=388 data-nosnippet>388</a> /// The results will be stored in the order they completed not the order they were spawned.
|
|
<a href=#389 id=389 data-nosnippet>389</a> /// This is a convenience method that is equivalent to calling [`join_next`] in
|
|
<a href=#390 id=390 data-nosnippet>390</a> /// a loop. If any tasks on the `JoinSet` fail with an [`JoinError`], then this call
|
|
<a href=#391 id=391 data-nosnippet>391</a> /// to `join_all` will panic and all remaining tasks on the `JoinSet` are
|
|
<a href=#392 id=392 data-nosnippet>392</a> /// cancelled. To handle errors in any other way, manually call [`join_next`]
|
|
<a href=#393 id=393 data-nosnippet>393</a> /// in a loop.
|
|
<a href=#394 id=394 data-nosnippet>394</a> ///
|
|
<a href=#395 id=395 data-nosnippet>395</a> /// # Examples
|
|
<a href=#396 id=396 data-nosnippet>396</a> ///
|
|
<a href=#397 id=397 data-nosnippet>397</a> /// Spawn multiple tasks and `join_all` them.
|
|
<a href=#398 id=398 data-nosnippet>398</a> ///
|
|
<a href=#399 id=399 data-nosnippet>399</a> /// ```
|
|
<a href=#400 id=400 data-nosnippet>400</a> /// use tokio::task::JoinSet;
|
|
<a href=#401 id=401 data-nosnippet>401</a> /// use std::time::Duration;
|
|
<a href=#402 id=402 data-nosnippet>402</a> ///
|
|
<a href=#403 id=403 data-nosnippet>403</a> /// # #[tokio::main(flavor = "current_thread")]
|
|
<a href=#404 id=404 data-nosnippet>404</a> /// # async fn main() {
|
|
<a href=#405 id=405 data-nosnippet>405</a> /// let mut set = JoinSet::new();
|
|
<a href=#406 id=406 data-nosnippet>406</a> ///
|
|
<a href=#407 id=407 data-nosnippet>407</a> /// for i in 0..3 {
|
|
<a href=#408 id=408 data-nosnippet>408</a> /// set.spawn(async move {
|
|
<a href=#409 id=409 data-nosnippet>409</a> /// tokio::time::sleep(Duration::from_secs(3 - i)).await;
|
|
<a href=#410 id=410 data-nosnippet>410</a> /// i
|
|
<a href=#411 id=411 data-nosnippet>411</a> /// });
|
|
<a href=#412 id=412 data-nosnippet>412</a> /// }
|
|
<a href=#413 id=413 data-nosnippet>413</a> ///
|
|
<a href=#414 id=414 data-nosnippet>414</a> /// let output = set.join_all().await;
|
|
<a href=#415 id=415 data-nosnippet>415</a> /// assert_eq!(output, vec![2, 1, 0]);
|
|
<a href=#416 id=416 data-nosnippet>416</a> /// # }
|
|
<a href=#417 id=417 data-nosnippet>417</a> /// ```
|
|
<a href=#418 id=418 data-nosnippet>418</a> ///
|
|
<a href=#419 id=419 data-nosnippet>419</a> /// Equivalent implementation of `join_all`, using [`join_next`] and loop.
|
|
<a href=#420 id=420 data-nosnippet>420</a> ///
|
|
<a href=#421 id=421 data-nosnippet>421</a> /// ```
|
|
<a href=#422 id=422 data-nosnippet>422</a> /// use tokio::task::JoinSet;
|
|
<a href=#423 id=423 data-nosnippet>423</a> /// use std::panic;
|
|
<a href=#424 id=424 data-nosnippet>424</a> ///
|
|
<a href=#425 id=425 data-nosnippet>425</a> /// # #[tokio::main(flavor = "current_thread")]
|
|
<a href=#426 id=426 data-nosnippet>426</a> /// # async fn main() {
|
|
<a href=#427 id=427 data-nosnippet>427</a> /// let mut set = JoinSet::new();
|
|
<a href=#428 id=428 data-nosnippet>428</a> ///
|
|
<a href=#429 id=429 data-nosnippet>429</a> /// for i in 0..3 {
|
|
<a href=#430 id=430 data-nosnippet>430</a> /// set.spawn(async move {i});
|
|
<a href=#431 id=431 data-nosnippet>431</a> /// }
|
|
<a href=#432 id=432 data-nosnippet>432</a> ///
|
|
<a href=#433 id=433 data-nosnippet>433</a> /// let mut output = Vec::new();
|
|
<a href=#434 id=434 data-nosnippet>434</a> /// while let Some(res) = set.join_next().await{
|
|
<a href=#435 id=435 data-nosnippet>435</a> /// match res {
|
|
<a href=#436 id=436 data-nosnippet>436</a> /// Ok(t) => output.push(t),
|
|
<a href=#437 id=437 data-nosnippet>437</a> /// Err(err) if err.is_panic() => panic::resume_unwind(err.into_panic()),
|
|
<a href=#438 id=438 data-nosnippet>438</a> /// Err(err) => panic!("{err}"),
|
|
<a href=#439 id=439 data-nosnippet>439</a> /// }
|
|
<a href=#440 id=440 data-nosnippet>440</a> /// }
|
|
<a href=#441 id=441 data-nosnippet>441</a> /// assert_eq!(output.len(),3);
|
|
<a href=#442 id=442 data-nosnippet>442</a> /// # }
|
|
<a href=#443 id=443 data-nosnippet>443</a> /// ```
|
|
<a href=#444 id=444 data-nosnippet>444</a> /// [`join_next`]: fn@Self::join_next
|
|
<a href=#445 id=445 data-nosnippet>445</a> /// [`JoinError::id`]: fn@crate::task::JoinError::id
|
|
<a href=#446 id=446 data-nosnippet>446</a> </span><span class="kw">pub async fn </span>join_all(<span class="kw-2">mut </span><span class="self">self</span>) -> Vec<T> {
|
|
<a href=#447 id=447 data-nosnippet>447</a> <span class="kw">let </span><span class="kw-2">mut </span>output = Vec::with_capacity(<span class="self">self</span>.len());
|
|
<a href=#448 id=448 data-nosnippet>448</a>
|
|
<a href=#449 id=449 data-nosnippet>449</a> <span class="kw">while let </span><span class="prelude-val">Some</span>(res) = <span class="self">self</span>.join_next().<span class="kw">await </span>{
|
|
<a href=#450 id=450 data-nosnippet>450</a> <span class="kw">match </span>res {
|
|
<a href=#451 id=451 data-nosnippet>451</a> <span class="prelude-val">Ok</span>(t) => output.push(t),
|
|
<a href=#452 id=452 data-nosnippet>452</a> <span class="prelude-val">Err</span>(err) <span class="kw">if </span>err.is_panic() => panic::resume_unwind(err.into_panic()),
|
|
<a href=#453 id=453 data-nosnippet>453</a> <span class="prelude-val">Err</span>(err) => <span class="macro">panic!</span>(<span class="string">"{err}"</span>),
|
|
<a href=#454 id=454 data-nosnippet>454</a> }
|
|
<a href=#455 id=455 data-nosnippet>455</a> }
|
|
<a href=#456 id=456 data-nosnippet>456</a> output
|
|
<a href=#457 id=457 data-nosnippet>457</a> }
|
|
<a href=#458 id=458 data-nosnippet>458</a>
|
|
<a href=#459 id=459 data-nosnippet>459</a> <span class="doccomment">/// Aborts all tasks on this `JoinSet`.
|
|
<a href=#460 id=460 data-nosnippet>460</a> ///
|
|
<a href=#461 id=461 data-nosnippet>461</a> /// This does not remove the tasks from the `JoinSet`. To wait for the tasks to complete
|
|
<a href=#462 id=462 data-nosnippet>462</a> /// cancellation, you should call `join_next` in a loop until the `JoinSet` is empty.
|
|
<a href=#463 id=463 data-nosnippet>463</a> </span><span class="kw">pub fn </span>abort_all(<span class="kw-2">&mut </span><span class="self">self</span>) {
|
|
<a href=#464 id=464 data-nosnippet>464</a> <span class="self">self</span>.inner.for_each(|jh| jh.abort());
|
|
<a href=#465 id=465 data-nosnippet>465</a> }
|
|
<a href=#466 id=466 data-nosnippet>466</a>
|
|
<a href=#467 id=467 data-nosnippet>467</a> <span class="doccomment">/// Removes all tasks from this `JoinSet` without aborting them.
|
|
<a href=#468 id=468 data-nosnippet>468</a> ///
|
|
<a href=#469 id=469 data-nosnippet>469</a> /// The tasks removed by this call will continue to run in the background even if the `JoinSet`
|
|
<a href=#470 id=470 data-nosnippet>470</a> /// is dropped.
|
|
<a href=#471 id=471 data-nosnippet>471</a> </span><span class="kw">pub fn </span>detach_all(<span class="kw-2">&mut </span><span class="self">self</span>) {
|
|
<a href=#472 id=472 data-nosnippet>472</a> <span class="self">self</span>.inner.drain(drop);
|
|
<a href=#473 id=473 data-nosnippet>473</a> }
|
|
<a href=#474 id=474 data-nosnippet>474</a>
|
|
<a href=#475 id=475 data-nosnippet>475</a> <span class="doccomment">/// Polls for one of the tasks in the set to complete.
|
|
<a href=#476 id=476 data-nosnippet>476</a> ///
|
|
<a href=#477 id=477 data-nosnippet>477</a> /// If this returns `Poll::Ready(Some(_))`, then the task that completed is removed from the set.
|
|
<a href=#478 id=478 data-nosnippet>478</a> ///
|
|
<a href=#479 id=479 data-nosnippet>479</a> /// When the method returns `Poll::Pending`, the `Waker` in the provided `Context` is scheduled
|
|
<a href=#480 id=480 data-nosnippet>480</a> /// to receive a wakeup when a task in the `JoinSet` completes. Note that on multiple calls to
|
|
<a href=#481 id=481 data-nosnippet>481</a> /// `poll_join_next`, only the `Waker` from the `Context` passed to the most recent call is
|
|
<a href=#482 id=482 data-nosnippet>482</a> /// scheduled to receive a wakeup.
|
|
<a href=#483 id=483 data-nosnippet>483</a> ///
|
|
<a href=#484 id=484 data-nosnippet>484</a> /// # Returns
|
|
<a href=#485 id=485 data-nosnippet>485</a> ///
|
|
<a href=#486 id=486 data-nosnippet>486</a> /// This function returns:
|
|
<a href=#487 id=487 data-nosnippet>487</a> ///
|
|
<a href=#488 id=488 data-nosnippet>488</a> /// * `Poll::Pending` if the `JoinSet` is not empty but there is no task whose output is
|
|
<a href=#489 id=489 data-nosnippet>489</a> /// available right now.
|
|
<a href=#490 id=490 data-nosnippet>490</a> /// * `Poll::Ready(Some(Ok(value)))` if one of the tasks in this `JoinSet` has completed.
|
|
<a href=#491 id=491 data-nosnippet>491</a> /// The `value` is the return value of one of the tasks that completed.
|
|
<a href=#492 id=492 data-nosnippet>492</a> /// * `Poll::Ready(Some(Err(err)))` if one of the tasks in this `JoinSet` has panicked or been
|
|
<a href=#493 id=493 data-nosnippet>493</a> /// aborted. The `err` is the `JoinError` from the panicked/aborted task.
|
|
<a href=#494 id=494 data-nosnippet>494</a> /// * `Poll::Ready(None)` if the `JoinSet` is empty.
|
|
<a href=#495 id=495 data-nosnippet>495</a> ///
|
|
<a href=#496 id=496 data-nosnippet>496</a> /// Note that this method may return `Poll::Pending` even if one of the tasks has completed.
|
|
<a href=#497 id=497 data-nosnippet>497</a> /// This can happen if the [coop budget] is reached.
|
|
<a href=#498 id=498 data-nosnippet>498</a> ///
|
|
<a href=#499 id=499 data-nosnippet>499</a> /// [coop budget]: crate::task::coop#cooperative-scheduling
|
|
<a href=#500 id=500 data-nosnippet>500</a> </span><span class="kw">pub fn </span>poll_join_next(<span class="kw-2">&mut </span><span class="self">self</span>, cx: <span class="kw-2">&mut </span>Context<<span class="lifetime">'_</span>>) -> Poll<<span class="prelude-ty">Option</span><<span class="prelude-ty">Result</span><T, JoinError>>> {
|
|
<a href=#501 id=501 data-nosnippet>501</a> <span class="comment">// The call to `pop_notified` moves the entry to the `idle` list. It is moved back to
|
|
<a href=#502 id=502 data-nosnippet>502</a> // the `notified` list if the waker is notified in the `poll` call below.
|
|
<a href=#503 id=503 data-nosnippet>503</a> </span><span class="kw">let </span><span class="kw-2">mut </span>entry = <span class="kw">match </span><span class="self">self</span>.inner.pop_notified(cx.waker()) {
|
|
<a href=#504 id=504 data-nosnippet>504</a> <span class="prelude-val">Some</span>(entry) => entry,
|
|
<a href=#505 id=505 data-nosnippet>505</a> <span class="prelude-val">None </span>=> {
|
|
<a href=#506 id=506 data-nosnippet>506</a> <span class="kw">if </span><span class="self">self</span>.is_empty() {
|
|
<a href=#507 id=507 data-nosnippet>507</a> <span class="kw">return </span>Poll::Ready(<span class="prelude-val">None</span>);
|
|
<a href=#508 id=508 data-nosnippet>508</a> } <span class="kw">else </span>{
|
|
<a href=#509 id=509 data-nosnippet>509</a> <span class="comment">// The waker was set by `pop_notified`.
|
|
<a href=#510 id=510 data-nosnippet>510</a> </span><span class="kw">return </span>Poll::Pending;
|
|
<a href=#511 id=511 data-nosnippet>511</a> }
|
|
<a href=#512 id=512 data-nosnippet>512</a> }
|
|
<a href=#513 id=513 data-nosnippet>513</a> };
|
|
<a href=#514 id=514 data-nosnippet>514</a>
|
|
<a href=#515 id=515 data-nosnippet>515</a> <span class="kw">let </span>res = entry.with_value_and_context(|jh, ctx| Pin::new(jh).poll(ctx));
|
|
<a href=#516 id=516 data-nosnippet>516</a>
|
|
<a href=#517 id=517 data-nosnippet>517</a> <span class="kw">if let </span>Poll::Ready(res) = res {
|
|
<a href=#518 id=518 data-nosnippet>518</a> <span class="kw">let </span>_entry = entry.remove();
|
|
<a href=#519 id=519 data-nosnippet>519</a> Poll::Ready(<span class="prelude-val">Some</span>(res))
|
|
<a href=#520 id=520 data-nosnippet>520</a> } <span class="kw">else </span>{
|
|
<a href=#521 id=521 data-nosnippet>521</a> <span class="comment">// A JoinHandle generally won't emit a wakeup without being ready unless
|
|
<a href=#522 id=522 data-nosnippet>522</a> // the coop limit has been reached. We yield to the executor in this
|
|
<a href=#523 id=523 data-nosnippet>523</a> // case.
|
|
<a href=#524 id=524 data-nosnippet>524</a> </span>cx.waker().wake_by_ref();
|
|
<a href=#525 id=525 data-nosnippet>525</a> Poll::Pending
|
|
<a href=#526 id=526 data-nosnippet>526</a> }
|
|
<a href=#527 id=527 data-nosnippet>527</a> }
|
|
<a href=#528 id=528 data-nosnippet>528</a>
|
|
<a href=#529 id=529 data-nosnippet>529</a> <span class="doccomment">/// Polls for one of the tasks in the set to complete.
|
|
<a href=#530 id=530 data-nosnippet>530</a> ///
|
|
<a href=#531 id=531 data-nosnippet>531</a> /// If this returns `Poll::Ready(Some(_))`, then the task that completed is removed from the set.
|
|
<a href=#532 id=532 data-nosnippet>532</a> ///
|
|
<a href=#533 id=533 data-nosnippet>533</a> /// When the method returns `Poll::Pending`, the `Waker` in the provided `Context` is scheduled
|
|
<a href=#534 id=534 data-nosnippet>534</a> /// to receive a wakeup when a task in the `JoinSet` completes. Note that on multiple calls to
|
|
<a href=#535 id=535 data-nosnippet>535</a> /// `poll_join_next`, only the `Waker` from the `Context` passed to the most recent call is
|
|
<a href=#536 id=536 data-nosnippet>536</a> /// scheduled to receive a wakeup.
|
|
<a href=#537 id=537 data-nosnippet>537</a> ///
|
|
<a href=#538 id=538 data-nosnippet>538</a> /// # Returns
|
|
<a href=#539 id=539 data-nosnippet>539</a> ///
|
|
<a href=#540 id=540 data-nosnippet>540</a> /// This function returns:
|
|
<a href=#541 id=541 data-nosnippet>541</a> ///
|
|
<a href=#542 id=542 data-nosnippet>542</a> /// * `Poll::Pending` if the `JoinSet` is not empty but there is no task whose output is
|
|
<a href=#543 id=543 data-nosnippet>543</a> /// available right now.
|
|
<a href=#544 id=544 data-nosnippet>544</a> /// * `Poll::Ready(Some(Ok((id, value))))` if one of the tasks in this `JoinSet` has completed.
|
|
<a href=#545 id=545 data-nosnippet>545</a> /// The `value` is the return value of one of the tasks that completed, and
|
|
<a href=#546 id=546 data-nosnippet>546</a> /// `id` is the [task ID] of that task.
|
|
<a href=#547 id=547 data-nosnippet>547</a> /// * `Poll::Ready(Some(Err(err)))` if one of the tasks in this `JoinSet` has panicked or been
|
|
<a href=#548 id=548 data-nosnippet>548</a> /// aborted. The `err` is the `JoinError` from the panicked/aborted task.
|
|
<a href=#549 id=549 data-nosnippet>549</a> /// * `Poll::Ready(None)` if the `JoinSet` is empty.
|
|
<a href=#550 id=550 data-nosnippet>550</a> ///
|
|
<a href=#551 id=551 data-nosnippet>551</a> /// Note that this method may return `Poll::Pending` even if one of the tasks has completed.
|
|
<a href=#552 id=552 data-nosnippet>552</a> /// This can happen if the [coop budget] is reached.
|
|
<a href=#553 id=553 data-nosnippet>553</a> ///
|
|
<a href=#554 id=554 data-nosnippet>554</a> /// [coop budget]: crate::task::coop#cooperative-scheduling
|
|
<a href=#555 id=555 data-nosnippet>555</a> /// [task ID]: crate::task::Id
|
|
<a href=#556 id=556 data-nosnippet>556</a> </span><span class="kw">pub fn </span>poll_join_next_with_id(
|
|
<a href=#557 id=557 data-nosnippet>557</a> <span class="kw-2">&mut </span><span class="self">self</span>,
|
|
<a href=#558 id=558 data-nosnippet>558</a> cx: <span class="kw-2">&mut </span>Context<<span class="lifetime">'_</span>>,
|
|
<a href=#559 id=559 data-nosnippet>559</a> ) -> Poll<<span class="prelude-ty">Option</span><<span class="prelude-ty">Result</span><(Id, T), JoinError>>> {
|
|
<a href=#560 id=560 data-nosnippet>560</a> <span class="comment">// The call to `pop_notified` moves the entry to the `idle` list. It is moved back to
|
|
<a href=#561 id=561 data-nosnippet>561</a> // the `notified` list if the waker is notified in the `poll` call below.
|
|
<a href=#562 id=562 data-nosnippet>562</a> </span><span class="kw">let </span><span class="kw-2">mut </span>entry = <span class="kw">match </span><span class="self">self</span>.inner.pop_notified(cx.waker()) {
|
|
<a href=#563 id=563 data-nosnippet>563</a> <span class="prelude-val">Some</span>(entry) => entry,
|
|
<a href=#564 id=564 data-nosnippet>564</a> <span class="prelude-val">None </span>=> {
|
|
<a href=#565 id=565 data-nosnippet>565</a> <span class="kw">if </span><span class="self">self</span>.is_empty() {
|
|
<a href=#566 id=566 data-nosnippet>566</a> <span class="kw">return </span>Poll::Ready(<span class="prelude-val">None</span>);
|
|
<a href=#567 id=567 data-nosnippet>567</a> } <span class="kw">else </span>{
|
|
<a href=#568 id=568 data-nosnippet>568</a> <span class="comment">// The waker was set by `pop_notified`.
|
|
<a href=#569 id=569 data-nosnippet>569</a> </span><span class="kw">return </span>Poll::Pending;
|
|
<a href=#570 id=570 data-nosnippet>570</a> }
|
|
<a href=#571 id=571 data-nosnippet>571</a> }
|
|
<a href=#572 id=572 data-nosnippet>572</a> };
|
|
<a href=#573 id=573 data-nosnippet>573</a>
|
|
<a href=#574 id=574 data-nosnippet>574</a> <span class="kw">let </span>res = entry.with_value_and_context(|jh, ctx| Pin::new(jh).poll(ctx));
|
|
<a href=#575 id=575 data-nosnippet>575</a>
|
|
<a href=#576 id=576 data-nosnippet>576</a> <span class="kw">if let </span>Poll::Ready(res) = res {
|
|
<a href=#577 id=577 data-nosnippet>577</a> <span class="kw">let </span>entry = entry.remove();
|
|
<a href=#578 id=578 data-nosnippet>578</a> <span class="comment">// If the task succeeded, add the task ID to the output. Otherwise, the
|
|
<a href=#579 id=579 data-nosnippet>579</a> // `JoinError` will already have the task's ID.
|
|
<a href=#580 id=580 data-nosnippet>580</a> </span>Poll::Ready(<span class="prelude-val">Some</span>(res.map(|output| (entry.id(), output))))
|
|
<a href=#581 id=581 data-nosnippet>581</a> } <span class="kw">else </span>{
|
|
<a href=#582 id=582 data-nosnippet>582</a> <span class="comment">// A JoinHandle generally won't emit a wakeup without being ready unless
|
|
<a href=#583 id=583 data-nosnippet>583</a> // the coop limit has been reached. We yield to the executor in this
|
|
<a href=#584 id=584 data-nosnippet>584</a> // case.
|
|
<a href=#585 id=585 data-nosnippet>585</a> </span>cx.waker().wake_by_ref();
|
|
<a href=#586 id=586 data-nosnippet>586</a> Poll::Pending
|
|
<a href=#587 id=587 data-nosnippet>587</a> }
|
|
<a href=#588 id=588 data-nosnippet>588</a> }
|
|
<a href=#589 id=589 data-nosnippet>589</a>}
|
|
<a href=#590 id=590 data-nosnippet>590</a>
|
|
<a href=#591 id=591 data-nosnippet>591</a><span class="kw">impl</span><T> Drop <span class="kw">for </span>JoinSet<T> {
|
|
<a href=#592 id=592 data-nosnippet>592</a> <span class="kw">fn </span>drop(<span class="kw-2">&mut </span><span class="self">self</span>) {
|
|
<a href=#593 id=593 data-nosnippet>593</a> <span class="self">self</span>.inner.drain(|join_handle| join_handle.abort());
|
|
<a href=#594 id=594 data-nosnippet>594</a> }
|
|
<a href=#595 id=595 data-nosnippet>595</a>}
|
|
<a href=#596 id=596 data-nosnippet>596</a>
|
|
<a href=#597 id=597 data-nosnippet>597</a><span class="kw">impl</span><T> fmt::Debug <span class="kw">for </span>JoinSet<T> {
|
|
<a href=#598 id=598 data-nosnippet>598</a> <span class="kw">fn </span>fmt(<span class="kw-2">&</span><span class="self">self</span>, f: <span class="kw-2">&mut </span>fmt::Formatter<<span class="lifetime">'_</span>>) -> fmt::Result {
|
|
<a href=#599 id=599 data-nosnippet>599</a> f.debug_struct(<span class="string">"JoinSet"</span>).field(<span class="string">"len"</span>, <span class="kw-2">&</span><span class="self">self</span>.len()).finish()
|
|
<a href=#600 id=600 data-nosnippet>600</a> }
|
|
<a href=#601 id=601 data-nosnippet>601</a>}
|
|
<a href=#602 id=602 data-nosnippet>602</a>
|
|
<a href=#603 id=603 data-nosnippet>603</a><span class="kw">impl</span><T> Default <span class="kw">for </span>JoinSet<T> {
|
|
<a href=#604 id=604 data-nosnippet>604</a> <span class="kw">fn </span>default() -> <span class="self">Self </span>{
|
|
<a href=#605 id=605 data-nosnippet>605</a> <span class="self">Self</span>::new()
|
|
<a href=#606 id=606 data-nosnippet>606</a> }
|
|
<a href=#607 id=607 data-nosnippet>607</a>}
|
|
<a href=#608 id=608 data-nosnippet>608</a>
|
|
<a href=#609 id=609 data-nosnippet>609</a><span class="doccomment">/// Collect an iterator of futures into a [`JoinSet`].
|
|
<a href=#610 id=610 data-nosnippet>610</a>///
|
|
<a href=#611 id=611 data-nosnippet>611</a>/// This is equivalent to calling [`JoinSet::spawn`] on each element of the iterator.
|
|
<a href=#612 id=612 data-nosnippet>612</a>///
|
|
<a href=#613 id=613 data-nosnippet>613</a>/// # Examples
|
|
<a href=#614 id=614 data-nosnippet>614</a>///
|
|
<a href=#615 id=615 data-nosnippet>615</a>/// The main example from [`JoinSet`]'s documentation can also be written using [`collect`]:
|
|
<a href=#616 id=616 data-nosnippet>616</a>///
|
|
<a href=#617 id=617 data-nosnippet>617</a>/// ```
|
|
<a href=#618 id=618 data-nosnippet>618</a>/// use tokio::task::JoinSet;
|
|
<a href=#619 id=619 data-nosnippet>619</a>///
|
|
<a href=#620 id=620 data-nosnippet>620</a>/// # #[tokio::main(flavor = "current_thread")]
|
|
<a href=#621 id=621 data-nosnippet>621</a>/// # async fn main() {
|
|
<a href=#622 id=622 data-nosnippet>622</a>/// let mut set: JoinSet<_> = (0..10).map(|i| async move { i }).collect();
|
|
<a href=#623 id=623 data-nosnippet>623</a>///
|
|
<a href=#624 id=624 data-nosnippet>624</a>/// let mut seen = [false; 10];
|
|
<a href=#625 id=625 data-nosnippet>625</a>/// while let Some(res) = set.join_next().await {
|
|
<a href=#626 id=626 data-nosnippet>626</a>/// let idx = res.unwrap();
|
|
<a href=#627 id=627 data-nosnippet>627</a>/// seen[idx] = true;
|
|
<a href=#628 id=628 data-nosnippet>628</a>/// }
|
|
<a href=#629 id=629 data-nosnippet>629</a>///
|
|
<a href=#630 id=630 data-nosnippet>630</a>/// for i in 0..10 {
|
|
<a href=#631 id=631 data-nosnippet>631</a>/// assert!(seen[i]);
|
|
<a href=#632 id=632 data-nosnippet>632</a>/// }
|
|
<a href=#633 id=633 data-nosnippet>633</a>/// # }
|
|
<a href=#634 id=634 data-nosnippet>634</a>/// ```
|
|
<a href=#635 id=635 data-nosnippet>635</a>///
|
|
<a href=#636 id=636 data-nosnippet>636</a>/// [`collect`]: std::iter::Iterator::collect
|
|
<a href=#637 id=637 data-nosnippet>637</a></span><span class="kw">impl</span><T, F> std::iter::FromIterator<F> <span class="kw">for </span>JoinSet<T>
|
|
<a href=#638 id=638 data-nosnippet>638</a><span class="kw">where
|
|
<a href=#639 id=639 data-nosnippet>639</a> </span>F: Future<Output = T>,
|
|
<a href=#640 id=640 data-nosnippet>640</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#641 id=641 data-nosnippet>641</a> T: Send + <span class="lifetime">'static</span>,
|
|
<a href=#642 id=642 data-nosnippet>642</a>{
|
|
<a href=#643 id=643 data-nosnippet>643</a> <span class="kw">fn </span>from_iter<I: IntoIterator<Item = F>>(iter: I) -> <span class="self">Self </span>{
|
|
<a href=#644 id=644 data-nosnippet>644</a> <span class="kw">let </span><span class="kw-2">mut </span>set = <span class="self">Self</span>::new();
|
|
<a href=#645 id=645 data-nosnippet>645</a> iter.into_iter().for_each(|task| {
|
|
<a href=#646 id=646 data-nosnippet>646</a> set.spawn(task);
|
|
<a href=#647 id=647 data-nosnippet>647</a> });
|
|
<a href=#648 id=648 data-nosnippet>648</a> set
|
|
<a href=#649 id=649 data-nosnippet>649</a> }
|
|
<a href=#650 id=650 data-nosnippet>650</a>}
|
|
<a href=#651 id=651 data-nosnippet>651</a>
|
|
<a href=#652 id=652 data-nosnippet>652</a><span class="doccomment">/// Extend a [`JoinSet`] with futures from an iterator.
|
|
<a href=#653 id=653 data-nosnippet>653</a>///
|
|
<a href=#654 id=654 data-nosnippet>654</a>/// This is equivalent to calling [`JoinSet::spawn`] on each element of the iterator.
|
|
<a href=#655 id=655 data-nosnippet>655</a>///
|
|
<a href=#656 id=656 data-nosnippet>656</a>/// # Examples
|
|
<a href=#657 id=657 data-nosnippet>657</a>///
|
|
<a href=#658 id=658 data-nosnippet>658</a>/// ```
|
|
<a href=#659 id=659 data-nosnippet>659</a>/// # #[cfg(not(target_family = "wasm"))]
|
|
<a href=#660 id=660 data-nosnippet>660</a>/// # {
|
|
<a href=#661 id=661 data-nosnippet>661</a>/// use tokio::task::JoinSet;
|
|
<a href=#662 id=662 data-nosnippet>662</a>///
|
|
<a href=#663 id=663 data-nosnippet>663</a>/// #[tokio::main]
|
|
<a href=#664 id=664 data-nosnippet>664</a>/// async fn main() {
|
|
<a href=#665 id=665 data-nosnippet>665</a>/// let mut set: JoinSet<_> = (0..5).map(|i| async move { i }).collect();
|
|
<a href=#666 id=666 data-nosnippet>666</a>///
|
|
<a href=#667 id=667 data-nosnippet>667</a>/// set.extend((5..10).map(|i| async move { i }));
|
|
<a href=#668 id=668 data-nosnippet>668</a>///
|
|
<a href=#669 id=669 data-nosnippet>669</a>/// let mut seen = [false; 10];
|
|
<a href=#670 id=670 data-nosnippet>670</a>/// while let Some(res) = set.join_next().await {
|
|
<a href=#671 id=671 data-nosnippet>671</a>/// let idx = res.unwrap();
|
|
<a href=#672 id=672 data-nosnippet>672</a>/// seen[idx] = true;
|
|
<a href=#673 id=673 data-nosnippet>673</a>/// }
|
|
<a href=#674 id=674 data-nosnippet>674</a>///
|
|
<a href=#675 id=675 data-nosnippet>675</a>/// for i in 0..10 {
|
|
<a href=#676 id=676 data-nosnippet>676</a>/// assert!(seen[i]);
|
|
<a href=#677 id=677 data-nosnippet>677</a>/// }
|
|
<a href=#678 id=678 data-nosnippet>678</a>/// }
|
|
<a href=#679 id=679 data-nosnippet>679</a>/// # }
|
|
<a href=#680 id=680 data-nosnippet>680</a>/// ```
|
|
<a href=#681 id=681 data-nosnippet>681</a></span><span class="kw">impl</span><T, F> std::iter::Extend<F> <span class="kw">for </span>JoinSet<T>
|
|
<a href=#682 id=682 data-nosnippet>682</a><span class="kw">where
|
|
<a href=#683 id=683 data-nosnippet>683</a> </span>F: Future<Output = T>,
|
|
<a href=#684 id=684 data-nosnippet>684</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#685 id=685 data-nosnippet>685</a> T: Send + <span class="lifetime">'static</span>,
|
|
<a href=#686 id=686 data-nosnippet>686</a>{
|
|
<a href=#687 id=687 data-nosnippet>687</a> <span class="kw">fn </span>extend<I>(<span class="kw-2">&mut </span><span class="self">self</span>, iter: I)
|
|
<a href=#688 id=688 data-nosnippet>688</a> <span class="kw">where
|
|
<a href=#689 id=689 data-nosnippet>689</a> </span>I: IntoIterator<Item = F>,
|
|
<a href=#690 id=690 data-nosnippet>690</a> {
|
|
<a href=#691 id=691 data-nosnippet>691</a> iter.into_iter().for_each(|task| {
|
|
<a href=#692 id=692 data-nosnippet>692</a> <span class="self">self</span>.spawn(task);
|
|
<a href=#693 id=693 data-nosnippet>693</a> });
|
|
<a href=#694 id=694 data-nosnippet>694</a> }
|
|
<a href=#695 id=695 data-nosnippet>695</a>}
|
|
<a href=#696 id=696 data-nosnippet>696</a>
|
|
<a href=#697 id=697 data-nosnippet>697</a><span class="comment">// === impl Builder ===
|
|
<a href=#698 id=698 data-nosnippet>698</a>
|
|
<a href=#699 id=699 data-nosnippet>699</a></span><span class="attr">#[cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))]
|
|
<a href=#700 id=700 data-nosnippet>700</a>#[cfg_attr(docsrs, doc(cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))))]
|
|
<a href=#701 id=701 data-nosnippet>701</a></span><span class="kw">impl</span><<span class="lifetime">'a</span>, T: <span class="lifetime">'static</span>> Builder<<span class="lifetime">'a</span>, T> {
|
|
<a href=#702 id=702 data-nosnippet>702</a> <span class="doccomment">/// Assigns a name to the task which will be spawned.
|
|
<a href=#703 id=703 data-nosnippet>703</a> </span><span class="kw">pub fn </span>name(<span class="self">self</span>, name: <span class="kw-2">&</span><span class="lifetime">'a </span>str) -> <span class="self">Self </span>{
|
|
<a href=#704 id=704 data-nosnippet>704</a> <span class="kw">let </span>builder = <span class="self">self</span>.builder.name(name);
|
|
<a href=#705 id=705 data-nosnippet>705</a> <span class="self">Self </span>{ builder, ..<span class="self">self </span>}
|
|
<a href=#706 id=706 data-nosnippet>706</a> }
|
|
<a href=#707 id=707 data-nosnippet>707</a>
|
|
<a href=#708 id=708 data-nosnippet>708</a> <span class="doccomment">/// Spawn the provided task with this builder's settings and store it in the
|
|
<a href=#709 id=709 data-nosnippet>709</a> /// [`JoinSet`], returning an [`AbortHandle`] that can be used to remotely
|
|
<a href=#710 id=710 data-nosnippet>710</a> /// cancel the task.
|
|
<a href=#711 id=711 data-nosnippet>711</a> ///
|
|
<a href=#712 id=712 data-nosnippet>712</a> /// # Returns
|
|
<a href=#713 id=713 data-nosnippet>713</a> ///
|
|
<a href=#714 id=714 data-nosnippet>714</a> /// An [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#715 id=715 data-nosnippet>715</a> ///
|
|
<a href=#716 id=716 data-nosnippet>716</a> /// # Panics
|
|
<a href=#717 id=717 data-nosnippet>717</a> ///
|
|
<a href=#718 id=718 data-nosnippet>718</a> /// This method panics if called outside of a Tokio runtime.
|
|
<a href=#719 id=719 data-nosnippet>719</a> ///
|
|
<a href=#720 id=720 data-nosnippet>720</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#721 id=721 data-nosnippet>721</a> </span><span class="attr">#[track_caller]
|
|
<a href=#722 id=722 data-nosnippet>722</a> </span><span class="kw">pub fn </span>spawn<F>(<span class="self">self</span>, future: F) -> std::io::Result<AbortHandle>
|
|
<a href=#723 id=723 data-nosnippet>723</a> <span class="kw">where
|
|
<a href=#724 id=724 data-nosnippet>724</a> </span>F: Future<Output = T>,
|
|
<a href=#725 id=725 data-nosnippet>725</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#726 id=726 data-nosnippet>726</a> T: Send,
|
|
<a href=#727 id=727 data-nosnippet>727</a> {
|
|
<a href=#728 id=728 data-nosnippet>728</a> <span class="prelude-val">Ok</span>(<span class="self">self</span>.joinset.insert(<span class="self">self</span>.builder.spawn(future)<span class="question-mark">?</span>))
|
|
<a href=#729 id=729 data-nosnippet>729</a> }
|
|
<a href=#730 id=730 data-nosnippet>730</a>
|
|
<a href=#731 id=731 data-nosnippet>731</a> <span class="doccomment">/// Spawn the provided task on the provided [runtime handle] with this
|
|
<a href=#732 id=732 data-nosnippet>732</a> /// builder's settings, and store it in the [`JoinSet`].
|
|
<a href=#733 id=733 data-nosnippet>733</a> ///
|
|
<a href=#734 id=734 data-nosnippet>734</a> /// # Returns
|
|
<a href=#735 id=735 data-nosnippet>735</a> ///
|
|
<a href=#736 id=736 data-nosnippet>736</a> /// An [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#737 id=737 data-nosnippet>737</a> ///
|
|
<a href=#738 id=738 data-nosnippet>738</a> ///
|
|
<a href=#739 id=739 data-nosnippet>739</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#740 id=740 data-nosnippet>740</a> /// [runtime handle]: crate::runtime::Handle
|
|
<a href=#741 id=741 data-nosnippet>741</a> </span><span class="attr">#[track_caller]
|
|
<a href=#742 id=742 data-nosnippet>742</a> </span><span class="kw">pub fn </span>spawn_on<F>(<span class="self">self</span>, future: F, handle: <span class="kw-2">&</span>Handle) -> std::io::Result<AbortHandle>
|
|
<a href=#743 id=743 data-nosnippet>743</a> <span class="kw">where
|
|
<a href=#744 id=744 data-nosnippet>744</a> </span>F: Future<Output = T>,
|
|
<a href=#745 id=745 data-nosnippet>745</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#746 id=746 data-nosnippet>746</a> T: Send,
|
|
<a href=#747 id=747 data-nosnippet>747</a> {
|
|
<a href=#748 id=748 data-nosnippet>748</a> <span class="prelude-val">Ok</span>(<span class="self">self</span>.joinset.insert(<span class="self">self</span>.builder.spawn_on(future, handle)<span class="question-mark">?</span>))
|
|
<a href=#749 id=749 data-nosnippet>749</a> }
|
|
<a href=#750 id=750 data-nosnippet>750</a>
|
|
<a href=#751 id=751 data-nosnippet>751</a> <span class="doccomment">/// Spawn the blocking code on the blocking threadpool with this builder's
|
|
<a href=#752 id=752 data-nosnippet>752</a> /// settings, and store it in the [`JoinSet`].
|
|
<a href=#753 id=753 data-nosnippet>753</a> ///
|
|
<a href=#754 id=754 data-nosnippet>754</a> /// # Returns
|
|
<a href=#755 id=755 data-nosnippet>755</a> ///
|
|
<a href=#756 id=756 data-nosnippet>756</a> /// An [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#757 id=757 data-nosnippet>757</a> ///
|
|
<a href=#758 id=758 data-nosnippet>758</a> /// # Panics
|
|
<a href=#759 id=759 data-nosnippet>759</a> ///
|
|
<a href=#760 id=760 data-nosnippet>760</a> /// This method panics if called outside of a Tokio runtime.
|
|
<a href=#761 id=761 data-nosnippet>761</a> ///
|
|
<a href=#762 id=762 data-nosnippet>762</a> /// [`JoinSet`]: crate::task::JoinSet
|
|
<a href=#763 id=763 data-nosnippet>763</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#764 id=764 data-nosnippet>764</a> </span><span class="attr">#[track_caller]
|
|
<a href=#765 id=765 data-nosnippet>765</a> </span><span class="kw">pub fn </span>spawn_blocking<F>(<span class="self">self</span>, f: F) -> std::io::Result<AbortHandle>
|
|
<a href=#766 id=766 data-nosnippet>766</a> <span class="kw">where
|
|
<a href=#767 id=767 data-nosnippet>767</a> </span>F: FnOnce() -> T,
|
|
<a href=#768 id=768 data-nosnippet>768</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#769 id=769 data-nosnippet>769</a> T: Send,
|
|
<a href=#770 id=770 data-nosnippet>770</a> {
|
|
<a href=#771 id=771 data-nosnippet>771</a> <span class="prelude-val">Ok</span>(<span class="self">self</span>.joinset.insert(<span class="self">self</span>.builder.spawn_blocking(f)<span class="question-mark">?</span>))
|
|
<a href=#772 id=772 data-nosnippet>772</a> }
|
|
<a href=#773 id=773 data-nosnippet>773</a>
|
|
<a href=#774 id=774 data-nosnippet>774</a> <span class="doccomment">/// Spawn the blocking code on the blocking threadpool of the provided
|
|
<a href=#775 id=775 data-nosnippet>775</a> /// runtime handle with this builder's settings, and store it in the
|
|
<a href=#776 id=776 data-nosnippet>776</a> /// [`JoinSet`].
|
|
<a href=#777 id=777 data-nosnippet>777</a> ///
|
|
<a href=#778 id=778 data-nosnippet>778</a> /// # Returns
|
|
<a href=#779 id=779 data-nosnippet>779</a> ///
|
|
<a href=#780 id=780 data-nosnippet>780</a> /// An [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#781 id=781 data-nosnippet>781</a> ///
|
|
<a href=#782 id=782 data-nosnippet>782</a> /// [`JoinSet`]: crate::task::JoinSet
|
|
<a href=#783 id=783 data-nosnippet>783</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#784 id=784 data-nosnippet>784</a> </span><span class="attr">#[track_caller]
|
|
<a href=#785 id=785 data-nosnippet>785</a> </span><span class="kw">pub fn </span>spawn_blocking_on<F>(<span class="self">self</span>, f: F, handle: <span class="kw-2">&</span>Handle) -> std::io::Result<AbortHandle>
|
|
<a href=#786 id=786 data-nosnippet>786</a> <span class="kw">where
|
|
<a href=#787 id=787 data-nosnippet>787</a> </span>F: FnOnce() -> T,
|
|
<a href=#788 id=788 data-nosnippet>788</a> F: Send + <span class="lifetime">'static</span>,
|
|
<a href=#789 id=789 data-nosnippet>789</a> T: Send,
|
|
<a href=#790 id=790 data-nosnippet>790</a> {
|
|
<a href=#791 id=791 data-nosnippet>791</a> <span class="prelude-val">Ok</span>(<span class="self">self
|
|
<a href=#792 id=792 data-nosnippet>792</a> </span>.joinset
|
|
<a href=#793 id=793 data-nosnippet>793</a> .insert(<span class="self">self</span>.builder.spawn_blocking_on(f, handle)<span class="question-mark">?</span>))
|
|
<a href=#794 id=794 data-nosnippet>794</a> }
|
|
<a href=#795 id=795 data-nosnippet>795</a>
|
|
<a href=#796 id=796 data-nosnippet>796</a> <span class="doccomment">/// Spawn the provided task on the current [`LocalSet`] or [`LocalRuntime`]
|
|
<a href=#797 id=797 data-nosnippet>797</a> /// with this builder's settings, and store it in the [`JoinSet`].
|
|
<a href=#798 id=798 data-nosnippet>798</a> ///
|
|
<a href=#799 id=799 data-nosnippet>799</a> /// # Returns
|
|
<a href=#800 id=800 data-nosnippet>800</a> ///
|
|
<a href=#801 id=801 data-nosnippet>801</a> /// An [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#802 id=802 data-nosnippet>802</a> ///
|
|
<a href=#803 id=803 data-nosnippet>803</a> /// # Panics
|
|
<a href=#804 id=804 data-nosnippet>804</a> ///
|
|
<a href=#805 id=805 data-nosnippet>805</a> /// This method panics if it is called outside of a `LocalSet` or `LocalRuntime`.
|
|
<a href=#806 id=806 data-nosnippet>806</a> ///
|
|
<a href=#807 id=807 data-nosnippet>807</a> /// [`LocalSet`]: crate::task::LocalSet
|
|
<a href=#808 id=808 data-nosnippet>808</a> /// [`LocalRuntime`]: crate::runtime::LocalRuntime
|
|
<a href=#809 id=809 data-nosnippet>809</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#810 id=810 data-nosnippet>810</a> </span><span class="attr">#[track_caller]
|
|
<a href=#811 id=811 data-nosnippet>811</a> </span><span class="kw">pub fn </span>spawn_local<F>(<span class="self">self</span>, future: F) -> std::io::Result<AbortHandle>
|
|
<a href=#812 id=812 data-nosnippet>812</a> <span class="kw">where
|
|
<a href=#813 id=813 data-nosnippet>813</a> </span>F: Future<Output = T>,
|
|
<a href=#814 id=814 data-nosnippet>814</a> F: <span class="lifetime">'static</span>,
|
|
<a href=#815 id=815 data-nosnippet>815</a> {
|
|
<a href=#816 id=816 data-nosnippet>816</a> <span class="prelude-val">Ok</span>(<span class="self">self</span>.joinset.insert(<span class="self">self</span>.builder.spawn_local(future)<span class="question-mark">?</span>))
|
|
<a href=#817 id=817 data-nosnippet>817</a> }
|
|
<a href=#818 id=818 data-nosnippet>818</a>
|
|
<a href=#819 id=819 data-nosnippet>819</a> <span class="doccomment">/// Spawn the provided task on the provided [`LocalSet`] with this builder's
|
|
<a href=#820 id=820 data-nosnippet>820</a> /// settings, and store it in the [`JoinSet`].
|
|
<a href=#821 id=821 data-nosnippet>821</a> ///
|
|
<a href=#822 id=822 data-nosnippet>822</a> /// # Returns
|
|
<a href=#823 id=823 data-nosnippet>823</a> ///
|
|
<a href=#824 id=824 data-nosnippet>824</a> /// An [`AbortHandle`] that can be used to remotely cancel the task.
|
|
<a href=#825 id=825 data-nosnippet>825</a> ///
|
|
<a href=#826 id=826 data-nosnippet>826</a> /// [`LocalSet`]: crate::task::LocalSet
|
|
<a href=#827 id=827 data-nosnippet>827</a> /// [`AbortHandle`]: crate::task::AbortHandle
|
|
<a href=#828 id=828 data-nosnippet>828</a> </span><span class="attr">#[track_caller]
|
|
<a href=#829 id=829 data-nosnippet>829</a> </span><span class="kw">pub fn </span>spawn_local_on<F>(<span class="self">self</span>, future: F, local_set: <span class="kw-2">&</span>LocalSet) -> std::io::Result<AbortHandle>
|
|
<a href=#830 id=830 data-nosnippet>830</a> <span class="kw">where
|
|
<a href=#831 id=831 data-nosnippet>831</a> </span>F: Future<Output = T>,
|
|
<a href=#832 id=832 data-nosnippet>832</a> F: <span class="lifetime">'static</span>,
|
|
<a href=#833 id=833 data-nosnippet>833</a> {
|
|
<a href=#834 id=834 data-nosnippet>834</a> <span class="prelude-val">Ok</span>(<span class="self">self
|
|
<a href=#835 id=835 data-nosnippet>835</a> </span>.joinset
|
|
<a href=#836 id=836 data-nosnippet>836</a> .insert(<span class="self">self</span>.builder.spawn_local_on(future, local_set)<span class="question-mark">?</span>))
|
|
<a href=#837 id=837 data-nosnippet>837</a> }
|
|
<a href=#838 id=838 data-nosnippet>838</a>}
|
|
<a href=#839 id=839 data-nosnippet>839</a>
|
|
<a href=#840 id=840 data-nosnippet>840</a><span class="comment">// Manual `Debug` impl so that `Builder` is `Debug` regardless of whether `T` is
|
|
<a href=#841 id=841 data-nosnippet>841</a>// `Debug`.
|
|
<a href=#842 id=842 data-nosnippet>842</a></span><span class="attr">#[cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))]
|
|
<a href=#843 id=843 data-nosnippet>843</a>#[cfg_attr(docsrs, doc(cfg(all(tokio_unstable, feature = <span class="string">"tracing"</span>))))]
|
|
<a href=#844 id=844 data-nosnippet>844</a></span><span class="kw">impl</span><<span class="lifetime">'a</span>, T> fmt::Debug <span class="kw">for </span>Builder<<span class="lifetime">'a</span>, T> {
|
|
<a href=#845 id=845 data-nosnippet>845</a> <span class="kw">fn </span>fmt(<span class="kw-2">&</span><span class="self">self</span>, f: <span class="kw-2">&mut </span>fmt::Formatter<<span class="lifetime">'_</span>>) -> fmt::Result {
|
|
<a href=#846 id=846 data-nosnippet>846</a> f.debug_struct(<span class="string">"join_set::Builder"</span>)
|
|
<a href=#847 id=847 data-nosnippet>847</a> .field(<span class="string">"joinset"</span>, <span class="kw-2">&</span><span class="self">self</span>.joinset)
|
|
<a href=#848 id=848 data-nosnippet>848</a> .field(<span class="string">"builder"</span>, <span class="kw-2">&</span><span class="self">self</span>.builder)
|
|
<a href=#849 id=849 data-nosnippet>849</a> .finish()
|
|
<a href=#850 id=850 data-nosnippet>850</a> }
|
|
<a href=#851 id=851 data-nosnippet>851</a>}
|
|
</code></pre></div></section></main></body></html> |