Skip to main content

zbox/
file.rs

1use std::fmt::{self, Debug};
2use std::io::{self, Error as IoError, ErrorKind, Read, Seek, SeekFrom, Write};
3
4use super::{Error, Result};
5use fs::fnode::{
6    Fnode, Metadata, Reader as FnodeReader, Version, Writer as FnodeWriter,
7};
8use fs::Handle;
9use trans::{TxHandle, TxMgr};
10
11/// A reader for a specific vesion of file content.
12///
13/// This reader can be obtained by [`version_reader`] method, and it
14/// implements [`Read`] trait.
15///
16/// [`version_reader`]: struct.File.html#method.version_reader
17/// [`Read`]: https://doc.rust-lang.org/std/io/trait.Read.html
18#[derive(Debug)]
19pub struct VersionReader {
20    handle: Handle,
21    rdr: FnodeReader,
22}
23
24impl VersionReader {
25    fn new(handle: &Handle, ver: usize) -> Result<Self> {
26        let rdr = FnodeReader::new(handle.fnode.clone(), ver, &handle.store)?;
27        Ok(VersionReader {
28            handle: handle.clone(),
29            rdr,
30        })
31    }
32
33    /// Returns the content version associated with this reader.
34    pub fn version(&self) -> Result<Version> {
35        let fnode = self.handle.fnode.read().unwrap();
36        fnode
37            .ver(self.rdr.version_num())
38            .cloned()
39            .ok_or(Error::NoVersion)
40    }
41}
42
43impl Read for VersionReader {
44    #[inline]
45    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
46        self.rdr.read(buf)
47    }
48}
49
50impl Seek for VersionReader {
51    #[inline]
52    fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
53        self.rdr.seek(pos)
54    }
55}
56
57/// A reference to an opened file in the repository.
58///
59/// An instance of a `File` can be read and/or written depending on what options
60/// it was opened with. Files also implement [`Seek`] to alter the logical
61/// cursor that the file contains internally.
62///
63/// Files are automatically closed when they go out of scope.
64///
65/// As ZboxFS internally cached file content, it is no need to use buffered
66/// reader, such as [`BufReader<R>`].
67///
68/// # Examples
69///
70/// Create a new file and write bytes to it:
71///
72/// ```
73/// use std::io::prelude::*;
74/// # use zbox::{init_env, Result, RepoOpener};
75///
76/// # fn foo() -> Result<()> {
77/// # init_env();
78/// # let mut repo = RepoOpener::new().create(true).open("mem://foo", "pwd")?;
79/// let mut file = repo.create_file("/foo.txt")?;
80/// file.write_all(b"Hello, world!")?;
81/// file.finish()?;
82/// # Ok(())
83/// # }
84/// # foo().unwrap();
85/// ```
86///
87/// Read the content of a file into a [`String`]:
88///
89/// ```
90/// # use zbox::{init_env, Result, RepoOpener};
91/// use std::io::prelude::*;
92/// # use zbox::OpenOptions;
93///
94/// # fn foo() -> Result<()> {
95/// # init_env();
96/// # let mut repo = RepoOpener::new().create(true).open("mem://foo", "pwd")?;
97/// # {
98/// #     let mut file = OpenOptions::new()
99/// #         .create(true)
100/// #         .open(&mut repo, "/foo.txt")?;
101/// #     file.write_all(b"Hello, world!")?;
102/// #     file.finish()?;
103/// # }
104/// let mut file = repo.open_file("/foo.txt")?;
105/// let mut content = String::new();
106/// file.read_to_string(&mut content)?;
107/// assert_eq!(content, "Hello, world!");
108/// # Ok(())
109/// # }
110/// # foo().unwrap();
111/// ```
112///
113/// # Versioning
114///
115/// `File` contents support up to 255 revision versions. [`Version`] is
116/// immutable once it is created.
117///
118/// By default, the maximum number of versions of a `File` is `1`, which is
119/// configurable by [`version_limit`] on both `Repo` and `File` level. File
120/// level option takes precedence.
121///
122/// After reaching this limit, the oldest [`Version`] will be automatically
123/// deleted after adding a new one.
124///
125/// Version number starts from `1` and continuously increases by 1.
126///
127/// # Writing
128///
129/// File content is cached internally for deduplication and will be handled
130/// automatically, thus calling [`flush`] is **not** recommended.
131///
132/// `File` can be sent to multiple threads, but only one thread can modify it at
133/// a time, which is similar to a `RwLock`.
134///
135/// `File` is multi-versioned, each time updating its content will create a new
136/// permanent [`Version`]. There are two ways of writing data to a file:
137///
138/// - **Multi-part Write**
139///
140///   This is done by updating `File` using [`Write`] trait multiple times.
141///   After all writing operations, [`finish`] must be called to create a new
142///   version. Unless [`finish`] was successfully returned, no data will be
143///   written to the file.
144///
145///   Internally, a transaction is created when writing to the file first time
146///   and calling [`finish`] will commit that transaction. If any errors
147///   happened during [`write`], that transaction will be aborted. Thus, you
148///   should not call [`finish`] after any failed [`write`].
149///
150///   Because transactions is thread local, multi-part write should be done in
151///   one transaction.
152///
153///   ## Examples
154///
155///   ```
156///   # use zbox::{init_env, Result, RepoOpener};
157///   use std::io::prelude::*;
158///   use std::io::SeekFrom;
159///   # use zbox::OpenOptions;
160///
161///   # fn foo() -> Result<()> {
162///   # init_env();
163///   # let mut repo = RepoOpener::new().create(true).open("mem://foo", "pwd")?;
164///   let mut file = OpenOptions::new()
165///       .create(true)
166///       .open(&mut repo, "/foo.txt")?;
167///   file.write_all(b"foo ")?;
168///   file.write_all(b"bar")?;
169///   file.finish()?;
170///
171///   let mut content = String::new();
172///   file.seek(SeekFrom::Start(0))?;
173///   file.read_to_string(&mut content)?;
174///   assert_eq!(content, "foo bar");
175///
176///   # Ok(())
177///   # }
178///   # foo().unwrap();
179///   ```
180///
181/// - **Single-part Write**
182///
183///   This can be done by calling [`write_once`], which will call [`finish`]
184///   internally to create a new version. Unless this method was successfully
185///   returned, no data will be written to the file.
186///
187///   ## Examples
188///
189///   ```
190///   # #![allow(unused_mut, unused_variables)]
191///   # use zbox::{init_env, Result, RepoOpener};
192///   use std::io::{Read, Seek, SeekFrom};
193///   # use zbox::OpenOptions;
194///
195///   # fn foo() -> Result<()> {
196///   # init_env();
197///   # let mut repo = RepoOpener::new().create(true).open("mem://foo", "pwd")?;
198///   let mut file = OpenOptions::new()
199///       .create(true)
200///       .open(&mut repo, "/foo.txt")?;
201///   file.write_once(b"foo bar")?;
202///
203///   let mut content = String::new();
204///   file.seek(SeekFrom::Start(0))?;
205///   file.read_to_string(&mut content)?;
206///   assert_eq!(content, "foo bar");
207///
208///   # Ok(())
209///   # }
210///   # foo().unwrap();
211///   ```
212///
213/// To gurantee atomicity, ZboxFS uses transaction when updating file so the
214/// data either be wholly persisted or nothing has been written.
215///
216/// - For multi-part write, the transaction begins in the first-time [`write`]
217///   and will be committed in [`finish`]. Any failure in [`write`] will abort
218///   the transaction, thus [`finish`] should not be called after that. If error
219///   happened during [`finish`], the transaction will also be aborted.
220/// - For single-part write, [`write_once`] itself is transactional. The
221///   transaction begins and will be committed inside this method.
222///
223/// Keep in mind of those characteristics, especially when writing a large
224/// amount of data to file, because any uncomitted transactions will abort
225/// and data in those transactions won't be persisted.
226///
227/// # Reading
228///
229/// As `File` can contain multiple versions, [`Read`] operation can be
230/// associated with different versions. By default, reading on `File` object is
231/// always bound to the latest version. To read a specific version, a
232/// [`VersionReader`], which supports [`Read`] trait as well, can be used.
233///
234/// ## Examples
235///
236/// Read the file content while it is in writing, notice that reading is always
237/// bound to latest content version.
238///
239/// ```
240/// use std::io::prelude::*;
241/// use std::io::SeekFrom;
242/// # use zbox::{init_env, Result, RepoOpener};
243/// # use zbox::OpenOptions;
244///
245/// # fn foo() -> Result<()> {
246/// # init_env();
247/// # let mut repo = RepoOpener::new().create(true).open("mem://foo", "pwd")?;
248/// // create a file and write data to it
249/// let mut file = OpenOptions::new().create(true).open(&mut repo, "/foo.txt")?;
250/// file.write_once(&[1, 2, 3, 4, 5, 6])?;
251///
252/// // read the first 2 bytes
253/// let mut buf = [0; 2];
254/// file.seek(SeekFrom::Start(0))?;
255/// file.read_exact(&mut buf)?;
256/// assert_eq!(&buf[..], &[1, 2]);
257///
258/// // create a new version, now the file content is [1, 2, 7, 8, 5, 6]
259/// file.write_once(&[7, 8])?;
260///
261/// // notice that reading is on the latest version
262/// file.seek(SeekFrom::Current(-2))?;
263/// file.read_exact(&mut buf)?;
264/// assert_eq!(&buf[..], &[7, 8]);
265///
266/// # Ok(())
267/// # }
268/// # foo().unwrap();
269/// ```
270///
271/// Read multiple versions using [`VersionReader`].
272///
273/// ```
274/// use std::io::prelude::*;
275/// # use zbox::{init_env, Result, RepoOpener};
276/// # use zbox::OpenOptions;
277///
278/// # fn foo() -> Result<()> {
279/// # init_env();
280/// # let mut repo = RepoOpener::new().create(true).open("mem://foo", "pwd")?;
281/// // create a file and write 2 versions
282/// let mut file = OpenOptions::new()
283///     .version_limit(4)
284///     .create(true)
285///     .open(&mut repo, "/foo.txt")?;
286/// file.write_once(b"foo")?;
287/// file.write_once(b"bar")?;
288///
289/// // get latest version number
290/// let curr_ver = file.curr_version()?;
291///
292/// // create a version reader and read latest version of content
293/// let mut rdr = file.version_reader(curr_ver)?;
294/// let mut content = String::new();
295/// rdr.read_to_string(&mut content)?;
296/// assert_eq!(content, "foobar");
297///
298/// // create another version reader and read previous version of content
299/// let mut rdr = file.version_reader(curr_ver - 1)?;
300/// let mut content = String::new();
301/// rdr.read_to_string(&mut content)?;
302/// assert_eq!(content, "foo");
303///
304/// # Ok(())
305/// # }
306/// # foo().unwrap();
307/// ```
308///
309/// [`Seek`]: https://doc.rust-lang.org/std/io/trait.Seek.html
310/// [`BufReader<R>`]: https://doc.rust-lang.org/std/io/struct.BufReader.html
311/// [`flush`]: https://doc.rust-lang.org/std/io/trait.Write.html#tymethod.flush
312/// [`String`]: https://doc.rust-lang.org/std/string/struct.String.html
313/// [`Read`]: https://doc.rust-lang.org/std/io/trait.Read.html
314/// [`Write`]: https://doc.rust-lang.org/std/io/trait.Write.html
315/// [`Version`]: struct.Version.html
316/// [`VersionReader`]: struct.VersionReader.html
317/// [`version_limit`]: struct.OpenOptions.html#method.version_limit
318/// [`finish`]: struct.File.html#method.finish
319/// [`write_once`]: struct.File.html#method.write_once
320pub struct File {
321    handle: Handle,
322    pos: SeekFrom, // must always be SeekFrom::Start
323    rdr: Option<FnodeReader>,
324    wtr: Option<FnodeWriter>,
325    tx_handle: Option<TxHandle>,
326    can_read: bool,
327    can_write: bool,
328}
329
330impl File {
331    pub(super) fn new(
332        handle: Handle,
333        pos: SeekFrom,
334        can_read: bool,
335        can_write: bool,
336    ) -> Self {
337        File {
338            handle,
339            pos,
340            rdr: None,
341            wtr: None,
342            tx_handle: None,
343            can_read,
344            can_write,
345        }
346    }
347
348    /// Check if file system is closed
349    fn check_closed(&self) -> Result<()> {
350        let shutter = self.handle.shutter.read().unwrap();
351        if shutter.is_closed() {
352            return Err(Error::RepoClosed);
353        }
354        Ok(())
355    }
356
357    /// Queries metadata about the file.
358    pub fn metadata(&self) -> Result<Metadata> {
359        self.check_closed()?;
360        let fnode = self.handle.fnode.read().unwrap();
361        Ok(fnode.metadata())
362    }
363
364    /// Returns a list of all the file content versions.
365    pub fn history(&self) -> Result<Vec<Version>> {
366        self.check_closed()?;
367        let fnode = self.handle.fnode.read().unwrap();
368        Ok(fnode.history())
369    }
370
371    /// Returns the current content version number.
372    pub fn curr_version(&self) -> Result<usize> {
373        self.check_closed()?;
374        let fnode = self.handle.fnode.read().unwrap();
375        Ok(fnode.curr_ver_num())
376    }
377
378    /// Returns content byte size of the current version.
379    fn curr_len(&self) -> usize {
380        let fnode = self.handle.fnode.read().unwrap();
381        fnode.curr_len()
382    }
383
384    /// Get a reader of the specified version.
385    ///
386    /// The returned reader implements [`Read`] trait. To get the version
387    /// number, first call [`history`] to get the list of all versions and
388    /// then choose the version number from it.
389    ///
390    /// [`Read`]: https://doc.rust-lang.org/std/io/trait.Read.html
391    /// [`history`]: struct.File.html#method.history
392    pub fn version_reader(&self, ver_num: usize) -> Result<VersionReader> {
393        self.check_closed()?;
394        if !self.can_read {
395            return Err(Error::CannotRead);
396        }
397        VersionReader::new(&self.handle, ver_num)
398    }
399
400    // calculate the seek position from the start based on file current size
401    fn seek_pos(&self, pos: SeekFrom) -> SeekFrom {
402        let curr_len = self.curr_len();
403        let pos: i64 = match pos {
404            SeekFrom::Start(p) => p as i64,
405            SeekFrom::End(p) => curr_len as i64 + p,
406            SeekFrom::Current(p) => match self.pos {
407                SeekFrom::Start(q) => p + q as i64,
408                SeekFrom::End(_) => unreachable!(),
409                SeekFrom::Current(_) => unreachable!(),
410            },
411        };
412        SeekFrom::Start(pos as u64)
413    }
414
415    fn begin_write(&mut self) -> Result<()> {
416        if !self.can_write {
417            return Err(Error::CannotWrite);
418        }
419
420        if self.wtr.is_some() {
421            return Err(Error::NotFinish);
422        }
423
424        assert!(self.tx_handle.is_none());
425
426        // append zeros if current position is beyond EOF
427        let curr_len = self.curr_len();
428        match self.pos {
429            SeekFrom::Start(pos) => {
430                let pos = pos as usize;
431                if pos > curr_len {
432                    // append zeros by setting file length
433                    self.set_len(pos)?;
434
435                    // then seek to new EOF
436                    self.pos = self.seek_pos(SeekFrom::End(0));
437                }
438            }
439            _ => unreachable!(),
440        }
441
442        // begin write
443        let txmgr = self.handle.txmgr.upgrade().ok_or(Error::RepoClosed)?;
444        let tx_handle = TxMgr::begin_trans(&txmgr)?;
445        tx_handle.run(|| {
446            let mut wtr =
447                FnodeWriter::new(self.handle.clone(), tx_handle.txid)?;
448            wtr.seek(self.seek_pos(self.pos))?;
449            self.wtr = Some(wtr);
450            Ok(())
451        })?;
452        self.tx_handle = Some(tx_handle);
453
454        Ok(())
455    }
456
457    // re-create reader on latest version
458    fn renew_reader(&mut self) -> Result<()> {
459        let mut rdr = FnodeReader::new_current(
460            self.handle.fnode.clone(),
461            &self.handle.store,
462        )?;
463        rdr.seek(self.pos)?;
464        self.rdr = Some(rdr);
465        Ok(())
466    }
467
468    /// Complete multi-part write to file and create a new version.
469    ///
470    /// This method will try to commit the transaction internally, no data will
471    /// be persisted if it failed. Do not call this method if any previous
472    /// [`write`] failed.
473    ///
474    /// # Errors
475    ///
476    /// Calling this method without writing data before will return
477    /// [`Error::NotWrite`] error.
478    ///
479    /// [`Write`]: https://doc.rust-lang.org/std/io/trait.Write.html
480    /// [`Error::NotWrite`]: enum.Error.html
481    pub fn finish(&mut self) -> Result<()> {
482        self.check_closed()?;
483
484        match self.wtr.take() {
485            Some(wtr) => {
486                let tx_handle = self.tx_handle.take().unwrap();
487                let mut end_pos = 0;
488
489                tx_handle.run_all_exclusive(|| {
490                    end_pos = wtr.finish()?;
491                    Ok(())
492                })?;
493
494                // set position
495                self.pos = SeekFrom::Start(end_pos as u64);
496            }
497            None => return Err(Error::NotWrite),
498        }
499
500        // re-create reader if there is an existing reader
501        if self.rdr.is_some() {
502            self.renew_reader()?;
503        }
504
505        Ok(())
506    }
507
508    /// Single-part write to file and create a new version.
509    ///
510    /// This method provides a convenient way of combining [`Write`] and
511    /// [`finish`].
512    ///
513    /// This method is atomic.
514    ///
515    /// [`Write`]: https://doc.rust-lang.org/std/io/trait.Write.html
516    /// [`finish`]: struct.File.html#method.finish
517    pub fn write_once(&mut self, buf: &[u8]) -> Result<()> {
518        self.check_closed()?;
519        match self.wtr {
520            Some(_) => Err(Error::NotFinish),
521            None => {
522                self.begin_write()?;
523                match self.wtr {
524                    Some(ref mut wtr) => match self.tx_handle {
525                        Some(ref tx_handle) => {
526                            tx_handle.run(|| {
527                                wtr.write_all(buf)?;
528                                Ok(())
529                            })?;
530                        }
531                        None => unreachable!(),
532                    },
533                    None => unreachable!(),
534                }
535                self.finish()
536            }
537        }
538    }
539
540    /// Truncates or extends the underlying file, create a new version of
541    /// content which size to become `size`.
542    ///
543    /// If the size is less than the current content size, then the new
544    /// content will be shrunk. If it is greater than the current content size,
545    /// then the content will be extended to `size` and have all of the
546    /// intermediate data filled in with 0s.
547    ///
548    /// This method is atomic.
549    ///
550    /// # Errors
551    ///
552    /// This method will return an error if the file is not opened for writing
553    /// or not finished writing.
554    pub fn set_len(&mut self, len: usize) -> Result<()> {
555        self.check_closed()?;
556        if self.wtr.is_some() {
557            return Err(Error::NotFinish);
558        }
559
560        if !self.can_write {
561            return Err(Error::CannotWrite);
562        }
563
564        let txmgr = self.handle.txmgr.upgrade().ok_or(Error::RepoClosed)?;
565        let tx_handle = TxMgr::begin_trans(&txmgr)?;
566        tx_handle.run_all_exclusive(|| {
567            Fnode::set_len(self.handle.clone(), len, tx_handle.txid)
568        })?;
569
570        // re-create reader if there is an existing reader
571        if self.rdr.is_some() {
572            self.renew_reader()?;
573        }
574
575        Ok(())
576    }
577}
578
579impl Read for File {
580    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
581        map_io_err!(self.check_closed())?;
582        if !self.can_read {
583            return Err(IoError::new(
584                ErrorKind::Other,
585                Error::CannotRead.to_string(),
586            ));
587        }
588
589        // if reader is not created yet, create a new reader and seek to
590        // the current file position
591        if self.rdr.is_none() {
592            map_io_err!(self.renew_reader())?;
593        }
594
595        match self.rdr {
596            Some(ref mut rdr) => {
597                let read = rdr.read(buf)?;
598                let new_pos = rdr.seek(SeekFrom::Current(0)).unwrap();
599                self.pos = SeekFrom::Start(new_pos);
600                Ok(read)
601            }
602            None => unreachable!(),
603        }
604    }
605}
606
607impl Write for File {
608    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
609        map_io_err!(self.check_closed())?;
610        if self.wtr.is_none() {
611            map_io_err!(self.begin_write())?;
612        }
613
614        let mut ret = 0;
615        map_io_err!(match self.wtr {
616            Some(ref mut wtr) => match self.tx_handle {
617                Some(ref tx_handle) => tx_handle
618                    .run(|| {
619                        ret = wtr.write(buf)?;
620                        Ok(())
621                    })
622                    .map(|_| ret),
623                None => unreachable!(),
624            },
625            None => unreachable!(),
626        }
627        .or_else(|err| {
628            // when write failed the tx has been aborted, so we need to clean up
629            // writer and tx handle here
630            self.wtr.take();
631            self.tx_handle.take();
632            Err(err)
633        }))
634    }
635
636    fn flush(&mut self) -> io::Result<()> {
637        map_io_err!(self.check_closed())?;
638        match self.wtr {
639            Some(ref mut wtr) => match self.tx_handle {
640                Some(ref tx_handle) => {
641                    map_io_err!(tx_handle.run(|| {
642                        wtr.flush()?;
643                        Ok(())
644                    }))?;
645                    Ok(())
646                }
647                None => unreachable!(),
648            },
649            None => Err(IoError::new(
650                ErrorKind::PermissionDenied,
651                Error::CannotWrite.to_string(),
652            )),
653        }
654    }
655}
656
657impl Seek for File {
658    fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
659        map_io_err!(self.check_closed())?;
660        if self.wtr.is_some() {
661            return Err(IoError::new(
662                ErrorKind::Other,
663                Error::NotFinish.to_string(),
664            ));
665        }
666
667        self.pos = match self.rdr {
668            Some(ref mut rdr) => SeekFrom::Start(rdr.seek(pos)?),
669            None => self.seek_pos(pos),
670        };
671
672        match self.pos {
673            SeekFrom::Start(pos) => Ok(pos),
674            _ => unreachable!(),
675        }
676    }
677}
678
679impl Debug for File {
680    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
681        f.debug_struct("File")
682            .field("pos", &self.pos)
683            .field("rdr", &self.rdr)
684            .field("wtr", &self.wtr)
685            .field("can_read", &self.can_read)
686            .field("can_write", &self.can_write)
687            .finish()
688    }
689}