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}