/rust/registry/src/index.crates.io-1949cf8c6b5b557f/tokio-1.53.1/src/fs/file.rs
Line | Count | Source |
1 | | //! Types for working with [`File`]. |
2 | | //! |
3 | | //! [`File`]: File |
4 | | |
5 | | use crate::fs::{asyncify, OpenOptions}; |
6 | | use crate::io::blocking::{Buf, DEFAULT_MAX_BUF_SIZE}; |
7 | | use crate::io::{AsyncRead, AsyncSeek, AsyncWrite, ReadBuf}; |
8 | | use crate::sync::Mutex; |
9 | | |
10 | | use std::cmp; |
11 | | use std::fmt; |
12 | | use std::fs::{Metadata, Permissions}; |
13 | | use std::future::Future; |
14 | | use std::io::{self, Seek, SeekFrom}; |
15 | | use std::path::Path; |
16 | | use std::pin::Pin; |
17 | | use std::sync::Arc; |
18 | | use std::task::{ready, Context, Poll}; |
19 | | |
20 | | #[cfg(test)] |
21 | | use super::mocks::JoinHandle; |
22 | | #[cfg(test)] |
23 | | use super::mocks::MockFile as StdFile; |
24 | | #[cfg(test)] |
25 | | use super::mocks::{spawn_blocking, spawn_mandatory_blocking}; |
26 | | #[cfg(not(test))] |
27 | | use crate::blocking::JoinHandle; |
28 | | #[cfg(not(test))] |
29 | | use crate::blocking::{spawn_blocking, spawn_mandatory_blocking}; |
30 | | #[cfg(not(test))] |
31 | | use std::fs::File as StdFile; |
32 | | |
33 | | cfg_io_uring! { |
34 | | #[cfg(not(test))] |
35 | | use crate::spawn; |
36 | | } |
37 | | |
38 | | /// A reference to an open file on the filesystem. |
39 | | /// |
40 | | /// This is a specialized version of [`std::fs::File`] for usage from the |
41 | | /// Tokio runtime. |
42 | | /// |
43 | | /// An instance of a `File` can be read and/or written depending on what options |
44 | | /// it was opened with. Files also implement [`AsyncSeek`] to alter the logical |
45 | | /// cursor that the file contains internally. |
46 | | /// |
47 | | /// A file will not be closed immediately when it goes out of scope if there |
48 | | /// are any IO operations that have not yet completed. To ensure that a file is |
49 | | /// closed immediately when it is dropped, you should call [`flush`] before |
50 | | /// dropping it. Note that this does not ensure that the file has been fully |
51 | | /// written to disk; the operating system might keep the changes around in an |
52 | | /// in-memory buffer. See the [`sync_all`] method for telling the OS to write |
53 | | /// the data to disk. |
54 | | /// |
55 | | /// Reading and writing to a `File` is usually done using the convenience |
56 | | /// methods found on the [`AsyncReadExt`] and [`AsyncWriteExt`] traits. |
57 | | /// |
58 | | /// [`AsyncSeek`]: trait@crate::io::AsyncSeek |
59 | | /// [`flush`]: fn@crate::io::AsyncWriteExt::flush |
60 | | /// [`sync_all`]: fn@crate::fs::File::sync_all |
61 | | /// [`AsyncReadExt`]: trait@crate::io::AsyncReadExt |
62 | | /// [`AsyncWriteExt`]: trait@crate::io::AsyncWriteExt |
63 | | /// |
64 | | /// # Examples |
65 | | /// |
66 | | /// Create a new file and asynchronously write bytes to it: |
67 | | /// |
68 | | /// ```no_run |
69 | | /// use tokio::fs::File; |
70 | | /// use tokio::io::AsyncWriteExt; // for write_all() |
71 | | /// |
72 | | /// # async fn dox() -> std::io::Result<()> { |
73 | | /// let mut file = File::create("foo.txt").await?; |
74 | | /// file.write_all(b"hello, world!").await?; |
75 | | /// # Ok(()) |
76 | | /// # } |
77 | | /// ``` |
78 | | /// |
79 | | /// Read the contents of a file into a buffer: |
80 | | /// |
81 | | /// ```no_run |
82 | | /// use tokio::fs::File; |
83 | | /// use tokio::io::AsyncReadExt; // for read_to_end() |
84 | | /// |
85 | | /// # async fn dox() -> std::io::Result<()> { |
86 | | /// let mut file = File::open("foo.txt").await?; |
87 | | /// |
88 | | /// let mut contents = vec![]; |
89 | | /// file.read_to_end(&mut contents).await?; |
90 | | /// |
91 | | /// println!("len = {}", contents.len()); |
92 | | /// # Ok(()) |
93 | | /// # } |
94 | | /// ``` |
95 | | pub struct File { |
96 | | std: Arc<StdFile>, |
97 | | inner: Mutex<Inner>, |
98 | | max_buf_size: usize, |
99 | | } |
100 | | |
101 | | struct Inner { |
102 | | state: State, |
103 | | |
104 | | /// Errors from writes/flushes are returned in write/flush calls. If a write |
105 | | /// error is observed while performing a read, it is saved until the next |
106 | | /// write / flush call. |
107 | | last_write_err: Option<io::ErrorKind>, |
108 | | |
109 | | pos: u64, |
110 | | } |
111 | | |
112 | | #[derive(Debug)] |
113 | | enum State { |
114 | | Idle(Option<Buf>), |
115 | | Busy(JoinHandle<(Operation, Buf)>), |
116 | | } |
117 | | |
118 | | #[derive(Debug)] |
119 | | enum Operation { |
120 | | Read(io::Result<usize>), |
121 | | Write(io::Result<()>), |
122 | | Seek(io::Result<u64>), |
123 | | } |
124 | | |
125 | | impl File { |
126 | | /// Attempts to open a file in read-only mode. |
127 | | /// |
128 | | /// See [`OpenOptions`] for more details. |
129 | | /// |
130 | | /// # Errors |
131 | | /// |
132 | | /// This function will return an error if called from outside of the Tokio |
133 | | /// runtime or if path does not already exist. Other errors may also be |
134 | | /// returned according to `OpenOptions::open`. |
135 | | /// |
136 | | /// # Examples |
137 | | /// |
138 | | /// ```no_run |
139 | | /// use tokio::fs::File; |
140 | | /// use tokio::io::AsyncReadExt; |
141 | | /// |
142 | | /// # async fn dox() -> std::io::Result<()> { |
143 | | /// let mut file = File::open("foo.txt").await?; |
144 | | /// |
145 | | /// let mut contents = vec![]; |
146 | | /// file.read_to_end(&mut contents).await?; |
147 | | /// |
148 | | /// println!("len = {}", contents.len()); |
149 | | /// # Ok(()) |
150 | | /// # } |
151 | | /// ``` |
152 | | /// |
153 | | /// The [`read_to_end`] method is defined on the [`AsyncReadExt`] trait. |
154 | | /// |
155 | | /// [`read_to_end`]: fn@crate::io::AsyncReadExt::read_to_end |
156 | | /// [`AsyncReadExt`]: trait@crate::io::AsyncReadExt |
157 | 0 | pub async fn open(path: impl AsRef<Path>) -> io::Result<File> { |
158 | 0 | Self::options().read(true).open(path).await |
159 | 0 | } |
160 | | |
161 | | /// Opens a file in write-only mode. |
162 | | /// |
163 | | /// This function will create a file if it does not exist, and will truncate |
164 | | /// it if it does. |
165 | | /// |
166 | | /// See [`OpenOptions`] for more details. |
167 | | /// |
168 | | /// # Errors |
169 | | /// |
170 | | /// Results in an error if called from outside of the Tokio runtime or if |
171 | | /// the underlying [`create`] call results in an error. |
172 | | /// |
173 | | /// [`create`]: std::fs::File::create |
174 | | /// |
175 | | /// # Examples |
176 | | /// |
177 | | /// ```no_run |
178 | | /// use tokio::fs::File; |
179 | | /// use tokio::io::AsyncWriteExt; |
180 | | /// |
181 | | /// # async fn dox() -> std::io::Result<()> { |
182 | | /// let mut file = File::create("foo.txt").await?; |
183 | | /// file.write_all(b"hello, world!").await?; |
184 | | /// # Ok(()) |
185 | | /// # } |
186 | | /// ``` |
187 | | /// |
188 | | /// The [`write_all`] method is defined on the [`AsyncWriteExt`] trait. |
189 | | /// |
190 | | /// [`write_all`]: fn@crate::io::AsyncWriteExt::write_all |
191 | | /// [`AsyncWriteExt`]: trait@crate::io::AsyncWriteExt |
192 | 0 | pub async fn create(path: impl AsRef<Path>) -> io::Result<File> { |
193 | 0 | Self::options() |
194 | 0 | .write(true) |
195 | 0 | .create(true) |
196 | 0 | .truncate(true) |
197 | 0 | .open(path) |
198 | 0 | .await |
199 | 0 | } |
200 | | |
201 | | /// Opens a file in read-write mode. |
202 | | /// |
203 | | /// This function will create a file if it does not exist, or return an error |
204 | | /// if it does. This way, if the call succeeds, the file returned is guaranteed |
205 | | /// to be new. |
206 | | /// |
207 | | /// This option is useful because it is atomic. Otherwise between checking |
208 | | /// whether a file exists and creating a new one, the file may have been |
209 | | /// created by another process (a TOCTOU race condition / attack). |
210 | | /// |
211 | | /// This can also be written using `File::options().read(true).write(true).create_new(true).open(...)`. |
212 | | /// |
213 | | /// See [`OpenOptions`] for more details. |
214 | | /// |
215 | | /// # Examples |
216 | | /// |
217 | | /// ```no_run |
218 | | /// use tokio::fs::File; |
219 | | /// use tokio::io::AsyncWriteExt; |
220 | | /// |
221 | | /// # async fn dox() -> std::io::Result<()> { |
222 | | /// let mut file = File::create_new("foo.txt").await?; |
223 | | /// file.write_all(b"hello, world!").await?; |
224 | | /// # Ok(()) |
225 | | /// # } |
226 | | /// ``` |
227 | | /// |
228 | | /// The [`write_all`] method is defined on the [`AsyncWriteExt`] trait. |
229 | | /// |
230 | | /// [`write_all`]: fn@crate::io::AsyncWriteExt::write_all |
231 | | /// [`AsyncWriteExt`]: trait@crate::io::AsyncWriteExt |
232 | 0 | pub async fn create_new<P: AsRef<Path>>(path: P) -> std::io::Result<File> { |
233 | 0 | Self::options() |
234 | 0 | .read(true) |
235 | 0 | .write(true) |
236 | 0 | .create_new(true) |
237 | 0 | .open(path) |
238 | 0 | .await |
239 | 0 | } |
240 | | |
241 | | /// Returns a new [`OpenOptions`] object. |
242 | | /// |
243 | | /// This function returns a new `OpenOptions` object that you can use to |
244 | | /// open or create a file with specific options if `open()` or `create()` |
245 | | /// are not appropriate. |
246 | | /// |
247 | | /// It is equivalent to `OpenOptions::new()`, but allows you to write more |
248 | | /// readable code. Instead of |
249 | | /// `OpenOptions::new().append(true).open("example.log")`, |
250 | | /// you can write `File::options().append(true).open("example.log")`. This |
251 | | /// also avoids the need to import `OpenOptions`. |
252 | | /// |
253 | | /// See the [`OpenOptions::new`] function for more details. |
254 | | /// |
255 | | /// # Examples |
256 | | /// |
257 | | /// ```no_run |
258 | | /// use tokio::fs::File; |
259 | | /// use tokio::io::AsyncWriteExt; |
260 | | /// |
261 | | /// # async fn dox() -> std::io::Result<()> { |
262 | | /// let mut f = File::options().append(true).open("example.log").await?; |
263 | | /// f.write_all(b"new line\n").await?; |
264 | | /// # Ok(()) |
265 | | /// # } |
266 | | /// ``` |
267 | | #[must_use] |
268 | 0 | pub fn options() -> OpenOptions { |
269 | 0 | OpenOptions::new() |
270 | 0 | } |
271 | | |
272 | | /// Converts a [`std::fs::File`] to a [`tokio::fs::File`](File). |
273 | | /// |
274 | | /// # Examples |
275 | | /// |
276 | | /// ```no_run |
277 | | /// // This line could block. It is not recommended to do this on the Tokio |
278 | | /// // runtime. |
279 | | /// let std_file = std::fs::File::open("foo.txt").unwrap(); |
280 | | /// let file = tokio::fs::File::from_std(std_file); |
281 | | /// ``` |
282 | 0 | pub fn from_std(std: StdFile) -> File { |
283 | 0 | File { |
284 | 0 | std: Arc::new(std), |
285 | 0 | inner: Mutex::new(Inner { |
286 | 0 | state: State::Idle(Some(Buf::with_capacity(0))), |
287 | 0 | last_write_err: None, |
288 | 0 | pos: 0, |
289 | 0 | }), |
290 | 0 | max_buf_size: DEFAULT_MAX_BUF_SIZE, |
291 | 0 | } |
292 | 0 | } |
293 | | |
294 | | /// Attempts to sync all OS-internal metadata to disk. |
295 | | /// |
296 | | /// This function will attempt to ensure that all in-core data reaches the |
297 | | /// filesystem before returning. |
298 | | /// |
299 | | /// # Examples |
300 | | /// |
301 | | /// ```no_run |
302 | | /// use tokio::fs::File; |
303 | | /// use tokio::io::AsyncWriteExt; |
304 | | /// |
305 | | /// # async fn dox() -> std::io::Result<()> { |
306 | | /// let mut file = File::create("foo.txt").await?; |
307 | | /// file.write_all(b"hello, world!").await?; |
308 | | /// file.sync_all().await?; |
309 | | /// # Ok(()) |
310 | | /// # } |
311 | | /// ``` |
312 | | /// |
313 | | /// The [`write_all`] method is defined on the [`AsyncWriteExt`] trait. |
314 | | /// |
315 | | /// [`write_all`]: fn@crate::io::AsyncWriteExt::write_all |
316 | | /// [`AsyncWriteExt`]: trait@crate::io::AsyncWriteExt |
317 | 0 | pub async fn sync_all(&self) -> io::Result<()> { |
318 | 0 | let mut inner = self.inner.lock().await; |
319 | 0 | inner.complete_inflight().await; |
320 | | |
321 | 0 | let std = self.std.clone(); |
322 | 0 | asyncify(move || std.sync_all()).await |
323 | 0 | } |
324 | | |
325 | | /// This function is similar to `sync_all`, except that it may not |
326 | | /// synchronize file metadata to the filesystem. |
327 | | /// |
328 | | /// This is intended for use cases that must synchronize content, but don't |
329 | | /// need the metadata on disk. The goal of this method is to reduce disk |
330 | | /// operations. |
331 | | /// |
332 | | /// Note that some platforms may simply implement this in terms of `sync_all`. |
333 | | /// |
334 | | /// # Examples |
335 | | /// |
336 | | /// ```no_run |
337 | | /// use tokio::fs::File; |
338 | | /// use tokio::io::AsyncWriteExt; |
339 | | /// |
340 | | /// # async fn dox() -> std::io::Result<()> { |
341 | | /// let mut file = File::create("foo.txt").await?; |
342 | | /// file.write_all(b"hello, world!").await?; |
343 | | /// file.sync_data().await?; |
344 | | /// # Ok(()) |
345 | | /// # } |
346 | | /// ``` |
347 | | /// |
348 | | /// The [`write_all`] method is defined on the [`AsyncWriteExt`] trait. |
349 | | /// |
350 | | /// [`write_all`]: fn@crate::io::AsyncWriteExt::write_all |
351 | | /// [`AsyncWriteExt`]: trait@crate::io::AsyncWriteExt |
352 | 0 | pub async fn sync_data(&self) -> io::Result<()> { |
353 | 0 | let mut inner = self.inner.lock().await; |
354 | 0 | inner.complete_inflight().await; |
355 | | |
356 | 0 | let std = self.std.clone(); |
357 | 0 | asyncify(move || std.sync_data()).await |
358 | 0 | } |
359 | | |
360 | | /// Truncates or extends the underlying file, updating the size of this file to become size. |
361 | | /// |
362 | | /// If the size is less than the current file's size, then the file will be |
363 | | /// shrunk. If it is greater than the current file's size, then the file |
364 | | /// will be extended to size and have all of the intermediate data filled in |
365 | | /// with 0s. |
366 | | /// |
367 | | /// # Errors |
368 | | /// |
369 | | /// This function will return an error if the file is not opened for |
370 | | /// writing. |
371 | | /// |
372 | | /// # Examples |
373 | | /// |
374 | | /// ```no_run |
375 | | /// use tokio::fs::File; |
376 | | /// use tokio::io::AsyncWriteExt; |
377 | | /// |
378 | | /// # async fn dox() -> std::io::Result<()> { |
379 | | /// let mut file = File::create("foo.txt").await?; |
380 | | /// file.write_all(b"hello, world!").await?; |
381 | | /// file.set_len(10).await?; |
382 | | /// # Ok(()) |
383 | | /// # } |
384 | | /// ``` |
385 | | /// |
386 | | /// The [`write_all`] method is defined on the [`AsyncWriteExt`] trait. |
387 | | /// |
388 | | /// [`write_all`]: fn@crate::io::AsyncWriteExt::write_all |
389 | | /// [`AsyncWriteExt`]: trait@crate::io::AsyncWriteExt |
390 | 0 | pub async fn set_len(&self, size: u64) -> io::Result<()> { |
391 | 0 | let mut inner = self.inner.lock().await; |
392 | 0 | inner.complete_inflight().await; |
393 | | |
394 | 0 | let mut buf = match inner.state { |
395 | 0 | State::Idle(ref mut buf_cell) => buf_cell.take().unwrap(), |
396 | 0 | _ => unreachable!(), |
397 | | }; |
398 | | |
399 | 0 | let seek = if !buf.is_empty() { |
400 | 0 | Some(SeekFrom::Current(buf.discard_read())) |
401 | | } else { |
402 | 0 | None |
403 | | }; |
404 | | |
405 | 0 | let std = self.std.clone(); |
406 | | |
407 | 0 | inner.state = State::Busy(spawn_blocking(move || { |
408 | 0 | let res = if let Some(seek) = seek { |
409 | 0 | (&*std).seek(seek).and_then(|_| std.set_len(size)) |
410 | | } else { |
411 | 0 | std.set_len(size) |
412 | | } |
413 | 0 | .map(|()| 0); // the value is discarded later |
414 | | |
415 | | // Return the result as a seek |
416 | 0 | (Operation::Seek(res), buf) |
417 | 0 | })); |
418 | | |
419 | 0 | let (op, buf) = match inner.state { |
420 | 0 | State::Idle(_) => unreachable!(), |
421 | 0 | State::Busy(ref mut rx) => rx.await?, |
422 | | }; |
423 | | |
424 | 0 | inner.state = State::Idle(Some(buf)); |
425 | | |
426 | 0 | match op { |
427 | 0 | Operation::Seek(res) => res.map(|pos| { |
428 | 0 | inner.pos = pos; |
429 | 0 | }), |
430 | 0 | _ => unreachable!(), |
431 | | } |
432 | 0 | } |
433 | | |
434 | | /// Queries metadata about the underlying file. |
435 | | /// |
436 | | /// # Examples |
437 | | /// |
438 | | /// ```no_run |
439 | | /// use tokio::fs::File; |
440 | | /// |
441 | | /// # async fn dox() -> std::io::Result<()> { |
442 | | /// let file = File::open("foo.txt").await?; |
443 | | /// let metadata = file.metadata().await?; |
444 | | /// |
445 | | /// println!("{:?}", metadata); |
446 | | /// # Ok(()) |
447 | | /// # } |
448 | | /// ``` |
449 | 0 | pub async fn metadata(&self) -> io::Result<Metadata> { |
450 | 0 | let std = self.std.clone(); |
451 | 0 | asyncify(move || std.metadata()).await |
452 | 0 | } |
453 | | |
454 | | /// Creates a new `File` instance that shares the same underlying file handle |
455 | | /// as the existing `File` instance. Reads, writes, and seeks will affect both |
456 | | /// File instances simultaneously. |
457 | | /// |
458 | | /// # Examples |
459 | | /// |
460 | | /// ```no_run |
461 | | /// use tokio::fs::File; |
462 | | /// |
463 | | /// # async fn dox() -> std::io::Result<()> { |
464 | | /// let file = File::open("foo.txt").await?; |
465 | | /// let file_clone = file.try_clone().await?; |
466 | | /// # Ok(()) |
467 | | /// # } |
468 | | /// ``` |
469 | 0 | pub async fn try_clone(&self) -> io::Result<File> { |
470 | 0 | self.inner.lock().await.complete_inflight().await; |
471 | 0 | let std = self.std.clone(); |
472 | 0 | let std_file = asyncify(move || std.try_clone()).await?; |
473 | 0 | let mut file = File::from_std(std_file); |
474 | 0 | file.set_max_buf_size(self.max_buf_size); |
475 | 0 | Ok(file) |
476 | 0 | } |
477 | | |
478 | | /// Destructures `File` into a [`std::fs::File`]. This function is |
479 | | /// async to allow any in-flight operations to complete. |
480 | | /// |
481 | | /// Use `File::try_into_std` to attempt conversion immediately. |
482 | | /// |
483 | | /// # Examples |
484 | | /// |
485 | | /// ```no_run |
486 | | /// use tokio::fs::File; |
487 | | /// |
488 | | /// # async fn dox() -> std::io::Result<()> { |
489 | | /// let tokio_file = File::open("foo.txt").await?; |
490 | | /// let std_file = tokio_file.into_std().await; |
491 | | /// # Ok(()) |
492 | | /// # } |
493 | | /// ``` |
494 | 0 | pub async fn into_std(mut self) -> StdFile { |
495 | 0 | self.inner.get_mut().complete_inflight().await; |
496 | 0 | Arc::try_unwrap(self.std).expect("Arc::try_unwrap failed") |
497 | 0 | } |
498 | | |
499 | | /// Tries to immediately destructure `File` into a [`std::fs::File`]. |
500 | | /// |
501 | | /// # Errors |
502 | | /// |
503 | | /// This function will return an error containing the file if some |
504 | | /// operation is in-flight. |
505 | | /// |
506 | | /// # Examples |
507 | | /// |
508 | | /// ```no_run |
509 | | /// use tokio::fs::File; |
510 | | /// |
511 | | /// # async fn dox() -> std::io::Result<()> { |
512 | | /// let tokio_file = File::open("foo.txt").await?; |
513 | | /// let std_file = tokio_file.try_into_std().unwrap(); |
514 | | /// # Ok(()) |
515 | | /// # } |
516 | | /// ``` |
517 | | #[allow(clippy::result_large_err)] |
518 | 0 | pub fn try_into_std(mut self) -> Result<StdFile, Self> { |
519 | 0 | match Arc::try_unwrap(self.std) { |
520 | 0 | Ok(file) => Ok(file), |
521 | 0 | Err(std_file_arc) => { |
522 | 0 | self.std = std_file_arc; |
523 | 0 | Err(self) |
524 | | } |
525 | | } |
526 | 0 | } |
527 | | |
528 | | /// Changes the permissions on the underlying file. |
529 | | /// |
530 | | /// # Platform-specific behavior |
531 | | /// |
532 | | /// This function currently corresponds to the `fchmod` function on Unix and |
533 | | /// the `SetFileInformationByHandle` function on Windows. Note that, this |
534 | | /// [may change in the future][changes]. |
535 | | /// |
536 | | /// [changes]: https://doc.rust-lang.org/std/io/index.html#platform-specific-behavior |
537 | | /// |
538 | | /// # Errors |
539 | | /// |
540 | | /// This function will return an error if the user lacks permission change |
541 | | /// attributes on the underlying file. It may also return an error in other |
542 | | /// os-specific unspecified cases. |
543 | | /// |
544 | | /// # Examples |
545 | | /// |
546 | | /// ```no_run |
547 | | /// use tokio::fs::File; |
548 | | /// |
549 | | /// # async fn dox() -> std::io::Result<()> { |
550 | | /// let file = File::open("foo.txt").await?; |
551 | | /// let mut perms = file.metadata().await?.permissions(); |
552 | | /// perms.set_readonly(true); |
553 | | /// file.set_permissions(perms).await?; |
554 | | /// # Ok(()) |
555 | | /// # } |
556 | | /// ``` |
557 | 0 | pub async fn set_permissions(&self, perm: Permissions) -> io::Result<()> { |
558 | 0 | let std = self.std.clone(); |
559 | 0 | asyncify(move || std.set_permissions(perm)).await |
560 | 0 | } |
561 | | |
562 | | /// Set the maximum buffer size for the underlying [`AsyncRead`] / [`AsyncWrite`] operation. |
563 | | /// |
564 | | /// Although Tokio uses a sensible default value for this buffer size, this function would be |
565 | | /// useful for changing that default depending on the situation. |
566 | | /// |
567 | | /// # Examples |
568 | | /// |
569 | | /// ```no_run |
570 | | /// use tokio::fs::File; |
571 | | /// use tokio::io::AsyncWriteExt; |
572 | | /// |
573 | | /// # async fn dox() -> std::io::Result<()> { |
574 | | /// let mut file = File::open("foo.txt").await?; |
575 | | /// |
576 | | /// // Set maximum buffer size to 8 MiB |
577 | | /// file.set_max_buf_size(8 * 1024 * 1024); |
578 | | /// |
579 | | /// let mut buf = vec![1; 1024 * 1024 * 1024]; |
580 | | /// |
581 | | /// // Write the 1 GiB buffer in chunks up to 8 MiB each. |
582 | | /// file.write_all(&mut buf).await?; |
583 | | /// # Ok(()) |
584 | | /// # } |
585 | | /// ``` |
586 | 0 | pub fn set_max_buf_size(&mut self, max_buf_size: usize) { |
587 | 0 | self.max_buf_size = max_buf_size; |
588 | 0 | } |
589 | | |
590 | | /// Get the maximum buffer size for the underlying [`AsyncRead`] / [`AsyncWrite`] operation. |
591 | 0 | pub fn max_buf_size(&self) -> usize { |
592 | 0 | self.max_buf_size |
593 | 0 | } |
594 | | } |
595 | | |
596 | | impl AsyncRead for File { |
597 | 0 | fn poll_read( |
598 | 0 | self: Pin<&mut Self>, |
599 | 0 | cx: &mut Context<'_>, |
600 | 0 | dst: &mut ReadBuf<'_>, |
601 | 0 | ) -> Poll<io::Result<()>> { |
602 | 0 | ready!(crate::trace::trace_leaf()); |
603 | | |
604 | 0 | let me = self.get_mut(); |
605 | 0 | let inner = me.inner.get_mut(); |
606 | | |
607 | | loop { |
608 | 0 | match inner.state { |
609 | 0 | State::Idle(ref mut buf_cell) => { |
610 | 0 | let mut buf = buf_cell.take().unwrap(); |
611 | | |
612 | 0 | if !buf.is_empty() || dst.remaining() == 0 { |
613 | 0 | buf.copy_to(dst); |
614 | 0 | *buf_cell = Some(buf); |
615 | 0 | return Poll::Ready(Ok(())); |
616 | 0 | } |
617 | | |
618 | 0 | let std = me.std.clone(); |
619 | | |
620 | 0 | let max_buf_size = cmp::min(dst.remaining(), me.max_buf_size); |
621 | 0 | inner.state = State::Busy(Inner::poll_read_inner(std, buf, max_buf_size)?); |
622 | | } |
623 | 0 | State::Busy(ref mut rx) => { |
624 | 0 | let (op, mut buf) = ready!(Pin::new(rx).poll(cx))?; |
625 | | |
626 | 0 | match op { |
627 | | Operation::Read(Ok(_)) => { |
628 | 0 | buf.copy_to(dst); |
629 | 0 | inner.state = State::Idle(Some(buf)); |
630 | 0 | return Poll::Ready(Ok(())); |
631 | | } |
632 | 0 | Operation::Read(Err(e)) => { |
633 | 0 | assert!(buf.is_empty()); |
634 | | |
635 | 0 | inner.state = State::Idle(Some(buf)); |
636 | 0 | return Poll::Ready(Err(e)); |
637 | | } |
638 | | Operation::Write(Ok(())) => { |
639 | 0 | assert!(buf.is_empty()); |
640 | 0 | inner.state = State::Idle(Some(buf)); |
641 | 0 | continue; |
642 | | } |
643 | 0 | Operation::Write(Err(e)) => { |
644 | 0 | assert!(inner.last_write_err.is_none()); |
645 | 0 | inner.last_write_err = Some(e.kind()); |
646 | 0 | inner.state = State::Idle(Some(buf)); |
647 | | } |
648 | 0 | Operation::Seek(result) => { |
649 | 0 | assert!(buf.is_empty()); |
650 | 0 | inner.state = State::Idle(Some(buf)); |
651 | 0 | if let Ok(pos) = result { |
652 | 0 | inner.pos = pos; |
653 | 0 | } |
654 | 0 | continue; |
655 | | } |
656 | | } |
657 | | } |
658 | | } |
659 | | } |
660 | 0 | } |
661 | | } |
662 | | |
663 | | impl AsyncSeek for File { |
664 | 0 | fn start_seek(self: Pin<&mut Self>, mut pos: SeekFrom) -> io::Result<()> { |
665 | 0 | let me = self.get_mut(); |
666 | 0 | let inner = me.inner.get_mut(); |
667 | | |
668 | 0 | match inner.state { |
669 | 0 | State::Busy(_) => Err(io::Error::new( |
670 | 0 | io::ErrorKind::Other, |
671 | 0 | "other file operation is pending, call poll_complete before start_seek", |
672 | 0 | )), |
673 | 0 | State::Idle(ref mut buf_cell) => { |
674 | 0 | let mut buf = buf_cell.take().unwrap(); |
675 | | |
676 | | // Factor in any unread data from the buf |
677 | 0 | if !buf.is_empty() { |
678 | 0 | let n = buf.discard_read(); |
679 | | |
680 | 0 | if let SeekFrom::Current(ref mut offset) = pos { |
681 | 0 | *offset += n; |
682 | 0 | } |
683 | 0 | } |
684 | | |
685 | 0 | let std = me.std.clone(); |
686 | | |
687 | 0 | inner.state = State::Busy(spawn_blocking(move || { |
688 | 0 | let res = (&*std).seek(pos); |
689 | 0 | (Operation::Seek(res), buf) |
690 | 0 | })); |
691 | 0 | Ok(()) |
692 | | } |
693 | | } |
694 | 0 | } |
695 | | |
696 | 0 | fn poll_complete(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<u64>> { |
697 | 0 | ready!(crate::trace::trace_leaf()); |
698 | 0 | let inner = self.inner.get_mut(); |
699 | | |
700 | | loop { |
701 | 0 | match inner.state { |
702 | 0 | State::Idle(_) => return Poll::Ready(Ok(inner.pos)), |
703 | 0 | State::Busy(ref mut rx) => { |
704 | 0 | let (op, buf) = ready!(Pin::new(rx).poll(cx))?; |
705 | 0 | inner.state = State::Idle(Some(buf)); |
706 | | |
707 | 0 | match op { |
708 | 0 | Operation::Read(_) => {} |
709 | 0 | Operation::Write(Err(e)) => { |
710 | 0 | assert!(inner.last_write_err.is_none()); |
711 | 0 | inner.last_write_err = Some(e.kind()); |
712 | | } |
713 | 0 | Operation::Write(_) => {} |
714 | 0 | Operation::Seek(res) => { |
715 | 0 | if let Ok(pos) = res { |
716 | 0 | inner.pos = pos; |
717 | 0 | } |
718 | 0 | return Poll::Ready(res); |
719 | | } |
720 | | } |
721 | | } |
722 | | } |
723 | | } |
724 | 0 | } |
725 | | } |
726 | | |
727 | | impl AsyncWrite for File { |
728 | 0 | fn poll_write( |
729 | 0 | self: Pin<&mut Self>, |
730 | 0 | cx: &mut Context<'_>, |
731 | 0 | src: &[u8], |
732 | 0 | ) -> Poll<io::Result<usize>> { |
733 | 0 | ready!(crate::trace::trace_leaf()); |
734 | 0 | let me = self.get_mut(); |
735 | 0 | let inner = me.inner.get_mut(); |
736 | | |
737 | 0 | if let Some(e) = inner.last_write_err.take() { |
738 | 0 | return Poll::Ready(Err(e.into())); |
739 | 0 | } |
740 | | |
741 | | loop { |
742 | 0 | match inner.state { |
743 | 0 | State::Idle(ref mut buf_cell) => { |
744 | 0 | let mut buf = buf_cell.take().unwrap(); |
745 | | |
746 | 0 | let seek = if !buf.is_empty() { |
747 | 0 | Some(SeekFrom::Current(buf.discard_read())) |
748 | | } else { |
749 | 0 | None |
750 | | }; |
751 | | |
752 | 0 | let n = buf.copy_from(src, me.max_buf_size); |
753 | 0 | let std = me.std.clone(); |
754 | | |
755 | 0 | let blocking_task_join_handle = spawn_mandatory_blocking(move || { |
756 | 0 | let res = if let Some(seek) = seek { |
757 | 0 | (&*std).seek(seek).and_then(|_| buf.write_to(&mut &*std)) |
758 | | } else { |
759 | 0 | buf.write_to(&mut &*std) |
760 | | }; |
761 | | |
762 | 0 | (Operation::Write(res), buf) |
763 | 0 | }) |
764 | 0 | .ok_or_else(|| { |
765 | 0 | io::Error::new(io::ErrorKind::Other, "background task failed") |
766 | 0 | })?; |
767 | | |
768 | 0 | inner.state = State::Busy(blocking_task_join_handle); |
769 | | |
770 | 0 | return Poll::Ready(Ok(n)); |
771 | | } |
772 | 0 | State::Busy(ref mut rx) => { |
773 | 0 | let (op, buf) = ready!(Pin::new(rx).poll(cx))?; |
774 | 0 | inner.state = State::Idle(Some(buf)); |
775 | | |
776 | 0 | match op { |
777 | | Operation::Read(_) => { |
778 | | // We don't care about the result here. The fact |
779 | | // that the cursor has advanced will be reflected in |
780 | | // the next iteration of the loop |
781 | 0 | continue; |
782 | | } |
783 | 0 | Operation::Write(res) => { |
784 | | // If the previous write was successful, continue. |
785 | | // Otherwise, error. |
786 | 0 | res?; |
787 | 0 | continue; |
788 | | } |
789 | | Operation::Seek(_) => { |
790 | | // Ignore the seek |
791 | 0 | continue; |
792 | | } |
793 | | } |
794 | | } |
795 | | } |
796 | | } |
797 | 0 | } |
798 | | |
799 | 0 | fn poll_write_vectored( |
800 | 0 | self: Pin<&mut Self>, |
801 | 0 | cx: &mut Context<'_>, |
802 | 0 | bufs: &[io::IoSlice<'_>], |
803 | 0 | ) -> Poll<Result<usize, io::Error>> { |
804 | 0 | ready!(crate::trace::trace_leaf()); |
805 | 0 | let me = self.get_mut(); |
806 | 0 | let inner = me.inner.get_mut(); |
807 | | |
808 | 0 | if let Some(e) = inner.last_write_err.take() { |
809 | 0 | return Poll::Ready(Err(e.into())); |
810 | 0 | } |
811 | | |
812 | | loop { |
813 | 0 | match inner.state { |
814 | 0 | State::Idle(ref mut buf_cell) => { |
815 | 0 | let mut buf = buf_cell.take().unwrap(); |
816 | | |
817 | 0 | let seek = if !buf.is_empty() { |
818 | 0 | Some(SeekFrom::Current(buf.discard_read())) |
819 | | } else { |
820 | 0 | None |
821 | | }; |
822 | | |
823 | 0 | let n = buf.copy_from_bufs(bufs, me.max_buf_size); |
824 | 0 | let std = me.std.clone(); |
825 | | |
826 | 0 | let blocking_task_join_handle = spawn_mandatory_blocking(move || { |
827 | 0 | let res = if let Some(seek) = seek { |
828 | 0 | (&*std).seek(seek).and_then(|_| buf.write_to(&mut &*std)) |
829 | | } else { |
830 | 0 | buf.write_to(&mut &*std) |
831 | | }; |
832 | | |
833 | 0 | (Operation::Write(res), buf) |
834 | 0 | }) |
835 | 0 | .ok_or_else(|| { |
836 | 0 | io::Error::new(io::ErrorKind::Other, "background task failed") |
837 | 0 | })?; |
838 | | |
839 | 0 | inner.state = State::Busy(blocking_task_join_handle); |
840 | | |
841 | 0 | return Poll::Ready(Ok(n)); |
842 | | } |
843 | 0 | State::Busy(ref mut rx) => { |
844 | 0 | let (op, buf) = ready!(Pin::new(rx).poll(cx))?; |
845 | 0 | inner.state = State::Idle(Some(buf)); |
846 | | |
847 | 0 | match op { |
848 | | Operation::Read(_) => { |
849 | | // We don't care about the result here. The fact |
850 | | // that the cursor has advanced will be reflected in |
851 | | // the next iteration of the loop |
852 | 0 | continue; |
853 | | } |
854 | 0 | Operation::Write(res) => { |
855 | | // If the previous write was successful, continue. |
856 | | // Otherwise, error. |
857 | 0 | res?; |
858 | 0 | continue; |
859 | | } |
860 | | Operation::Seek(_) => { |
861 | | // Ignore the seek |
862 | 0 | continue; |
863 | | } |
864 | | } |
865 | | } |
866 | | } |
867 | | } |
868 | 0 | } |
869 | | |
870 | 0 | fn is_write_vectored(&self) -> bool { |
871 | 0 | true |
872 | 0 | } |
873 | | |
874 | 0 | fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> { |
875 | 0 | ready!(crate::trace::trace_leaf()); |
876 | 0 | let inner = self.inner.get_mut(); |
877 | 0 | inner.poll_flush(cx) |
878 | 0 | } |
879 | | |
880 | 0 | fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> { |
881 | 0 | ready!(crate::trace::trace_leaf()); |
882 | 0 | self.poll_flush(cx) |
883 | 0 | } |
884 | | } |
885 | | |
886 | | impl From<StdFile> for File { |
887 | 0 | fn from(std: StdFile) -> Self { |
888 | 0 | Self::from_std(std) |
889 | 0 | } |
890 | | } |
891 | | |
892 | | impl fmt::Debug for File { |
893 | 0 | fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result { |
894 | 0 | fmt.debug_struct("tokio::fs::File") |
895 | 0 | .field("std", &self.std) |
896 | 0 | .finish() |
897 | 0 | } |
898 | | } |
899 | | |
900 | | #[cfg(unix)] |
901 | | impl From<std::os::fd::OwnedFd> for File { |
902 | 0 | fn from(fd: std::os::fd::OwnedFd) -> Self { |
903 | 0 | Self::from_std(StdFile::from(fd)) |
904 | 0 | } |
905 | | } |
906 | | |
907 | | #[cfg(unix)] |
908 | | impl std::os::unix::io::AsRawFd for File { |
909 | 0 | fn as_raw_fd(&self) -> std::os::unix::io::RawFd { |
910 | 0 | self.std.as_raw_fd() |
911 | 0 | } |
912 | | } |
913 | | |
914 | | #[cfg(unix)] |
915 | | impl std::os::unix::io::AsFd for File { |
916 | 0 | fn as_fd(&self) -> std::os::unix::io::BorrowedFd<'_> { |
917 | | unsafe { |
918 | 0 | std::os::unix::io::BorrowedFd::borrow_raw(std::os::unix::io::AsRawFd::as_raw_fd(self)) |
919 | | } |
920 | 0 | } |
921 | | } |
922 | | |
923 | | #[cfg(unix)] |
924 | | impl std::os::unix::io::FromRawFd for File { |
925 | 0 | unsafe fn from_raw_fd(fd: std::os::unix::io::RawFd) -> Self { |
926 | | // Safety: exactly the same safety contract as |
927 | | // `std::os::unix::io::FromRawFd::from_raw_fd`. |
928 | 0 | unsafe { StdFile::from_raw_fd(fd).into() } |
929 | 0 | } |
930 | | } |
931 | | |
932 | | cfg_windows! { |
933 | | use crate::os::windows::io::{AsRawHandle, FromRawHandle, RawHandle, AsHandle, BorrowedHandle, OwnedHandle}; |
934 | | |
935 | | impl From<OwnedHandle> for File { |
936 | | fn from(handle: OwnedHandle) -> Self { |
937 | | Self::from_std(StdFile::from(handle)) |
938 | | } |
939 | | } |
940 | | |
941 | | impl AsRawHandle for File { |
942 | | fn as_raw_handle(&self) -> RawHandle { |
943 | | self.std.as_raw_handle() |
944 | | } |
945 | | } |
946 | | |
947 | | impl AsHandle for File { |
948 | | fn as_handle(&self) -> BorrowedHandle<'_> { |
949 | | unsafe { |
950 | | BorrowedHandle::borrow_raw( |
951 | | AsRawHandle::as_raw_handle(self), |
952 | | ) |
953 | | } |
954 | | } |
955 | | } |
956 | | |
957 | | impl FromRawHandle for File { |
958 | | unsafe fn from_raw_handle(handle: RawHandle) -> Self { |
959 | | // Safety: exactly the same safety contract as |
960 | | // `FromRawHandle::from_raw_handle`. |
961 | | unsafe { StdFile::from_raw_handle(handle).into() } |
962 | | } |
963 | | } |
964 | | } |
965 | | |
966 | | impl Inner { |
967 | 0 | fn poll_read_inner( |
968 | 0 | std: Arc<StdFile>, |
969 | 0 | buf: Buf, |
970 | 0 | max_buf_size: usize, |
971 | 0 | ) -> io::Result<JoinHandle<(Operation, Buf)>> { |
972 | | // Unit tests use `MockFile` and the mock `spawn_blocking` infrastructure, |
973 | | // which can't drive real io_uring operations. The io_uring read path |
974 | | // is tested through integration tests in `tests/fs_uring_file_read.rs`. |
975 | | #[cfg(all( |
976 | | not(test), |
977 | | tokio_unstable, |
978 | | feature = "io-uring", |
979 | | feature = "rt", |
980 | | feature = "fs", |
981 | | target_os = "linux", |
982 | | ))] |
983 | | { |
984 | | if let Ok(handle) = crate::runtime::Handle::try_current() { |
985 | | let driver_handle = handle.inner.driver().io(); |
986 | | |
987 | | if driver_handle.is_uring_ready(io_uring::opcode::Read::CODE) { |
988 | | // Fast path: uring already initialized and Read supported. |
989 | | let fd: crate::io::uring::utils::ArcFd = std; |
990 | | return Ok(spawn(Self::uring_read(fd, buf, max_buf_size))); |
991 | | } |
992 | | |
993 | | if !driver_handle.is_uring_probed() { |
994 | | // Not yet probed: lazy init inside an async task so |
995 | | // `File::from_std()` can still benefit from io-uring. |
996 | | return Ok(spawn(Self::lazy_init_read(std, buf, max_buf_size))); |
997 | | } |
998 | | // Probed but unsupported: fall through to spawn_blocking. |
999 | | } |
1000 | | } |
1001 | | |
1002 | | // Fallback: spawn_blocking |
1003 | 0 | let join = Self::spawn_blocking_read(buf, std, max_buf_size); |
1004 | 0 | Ok(join) |
1005 | 0 | } |
1006 | | |
1007 | | /// Perform an io-uring read with interrupt retry. |
1008 | | #[cfg(all( |
1009 | | not(test), |
1010 | | tokio_unstable, |
1011 | | feature = "io-uring", |
1012 | | feature = "rt", |
1013 | | feature = "fs", |
1014 | | target_os = "linux", |
1015 | | ))] |
1016 | | async fn uring_read( |
1017 | | mut fd: crate::io::uring::utils::ArcFd, |
1018 | | mut buf: Buf, |
1019 | | max_buf_size: usize, |
1020 | | ) -> (Operation, Buf) { |
1021 | | use crate::runtime::driver::op::Op; |
1022 | | |
1023 | | loop { |
1024 | | let (res, r_fd, r_buf) = |
1025 | | // u64::MAX to use and advance the file position |
1026 | | Op::read_at(fd, buf, max_buf_size, u64::MAX).await; |
1027 | | match res { |
1028 | | Err(e) if e.kind() == io::ErrorKind::Interrupted => { |
1029 | | buf = r_buf; |
1030 | | fd = r_fd; |
1031 | | continue; |
1032 | | } |
1033 | | Err(e) => break (Operation::Read(Err(e)), r_buf), |
1034 | | Ok(n) => break (Operation::Read(Ok(n as usize)), r_buf), |
1035 | | } |
1036 | | } |
1037 | | } |
1038 | | |
1039 | | /// Attempt lazy io-uring initialization, then read via uring or fall back |
1040 | | /// to a blocking read. Covers the `File::from_std()` path where |
1041 | | /// `check_and_init()` hasn't been called yet. |
1042 | | #[cfg(all( |
1043 | | not(test), |
1044 | | tokio_unstable, |
1045 | | feature = "io-uring", |
1046 | | feature = "rt", |
1047 | | feature = "fs", |
1048 | | target_os = "linux", |
1049 | | ))] |
1050 | | async fn lazy_init_read(std: Arc<StdFile>, buf: Buf, max_buf_size: usize) -> (Operation, Buf) { |
1051 | | let handle = crate::runtime::Handle::current(); |
1052 | | let driver_handle = handle.inner.driver().io(); |
1053 | | if driver_handle |
1054 | | .check_and_init(io_uring::opcode::Read::CODE) |
1055 | | .await |
1056 | | .unwrap_or(false) |
1057 | | { |
1058 | | let fd: crate::io::uring::utils::ArcFd = std; |
1059 | | Self::uring_read(fd, buf, max_buf_size).await |
1060 | | } else { |
1061 | | match Self::spawn_blocking_read(buf, std, max_buf_size).await { |
1062 | | Ok(result) => result, |
1063 | | Err(e) => ( |
1064 | | Operation::Read(Err(io::Error::new(io::ErrorKind::Other, e))), |
1065 | | Buf::with_capacity(0), |
1066 | | ), |
1067 | | } |
1068 | | } |
1069 | | } |
1070 | | |
1071 | 0 | fn spawn_blocking_read( |
1072 | 0 | buf: Buf, |
1073 | 0 | std: Arc<StdFile>, |
1074 | 0 | max_buf_size: usize, |
1075 | 0 | ) -> JoinHandle<(Operation, Buf)> { |
1076 | 0 | spawn_blocking(move || { |
1077 | 0 | let mut buf = buf; |
1078 | | // SAFETY: the `Read` implementation of `std` does not |
1079 | | // read from the buffer it is borrowing and correctly |
1080 | | // reports the length of the data written into the buffer. |
1081 | 0 | let res = unsafe { buf.read_from(&mut &*std, max_buf_size) }; |
1082 | 0 | (Operation::Read(res), buf) |
1083 | 0 | }) |
1084 | 0 | } |
1085 | | |
1086 | 0 | async fn complete_inflight(&mut self) { |
1087 | | use std::future::poll_fn; |
1088 | | |
1089 | 0 | poll_fn(|cx| self.poll_complete_inflight(cx)).await; |
1090 | 0 | } |
1091 | | |
1092 | 0 | fn poll_complete_inflight(&mut self, cx: &mut Context<'_>) -> Poll<()> { |
1093 | 0 | ready!(crate::trace::trace_leaf()); |
1094 | 0 | match self.poll_flush(cx) { |
1095 | 0 | Poll::Ready(Err(e)) => { |
1096 | 0 | self.last_write_err = Some(e.kind()); |
1097 | 0 | Poll::Ready(()) |
1098 | | } |
1099 | 0 | Poll::Ready(Ok(())) => Poll::Ready(()), |
1100 | 0 | Poll::Pending => Poll::Pending, |
1101 | | } |
1102 | 0 | } |
1103 | | |
1104 | 0 | fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> { |
1105 | 0 | if let Some(e) = self.last_write_err.take() { |
1106 | 0 | return Poll::Ready(Err(e.into())); |
1107 | 0 | } |
1108 | | |
1109 | 0 | let (op, buf) = match self.state { |
1110 | 0 | State::Idle(_) => return Poll::Ready(Ok(())), |
1111 | 0 | State::Busy(ref mut rx) => ready!(Pin::new(rx).poll(cx))?, |
1112 | | }; |
1113 | | |
1114 | | // The buffer is not used here |
1115 | 0 | self.state = State::Idle(Some(buf)); |
1116 | | |
1117 | 0 | match op { |
1118 | 0 | Operation::Read(_) => Poll::Ready(Ok(())), |
1119 | 0 | Operation::Write(res) => Poll::Ready(res), |
1120 | 0 | Operation::Seek(_) => Poll::Ready(Ok(())), |
1121 | | } |
1122 | 0 | } |
1123 | | } |
1124 | | |
1125 | | #[cfg(test)] |
1126 | | mod tests; |