Coverage Report

Created: 2026-09-04 06:11

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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;