Skip to main content

atom_file/
lib.rs

1//! [`AtomicFile`] provides buffered concurrent access to files with async atomic commit.
2//!
3//! [`BasicAtomicFile`] is a non-async alternative.
4//!
5//! [`MultiFileStorage`] is the recommended backing storage for AtomicFile.
6//!
7//! [`FastFileStorage`] is the recommended temporary storage for AtomicFile.
8//!
9//!# Features
10//!
11//! This crate supports the following cargo features:
12//! - `pstd` : Use pstd crate for `BTreeMap` (allocated in `GTemp`).
13//! - `unsafe-optim` : Enable unsafe optimisations in release mode.
14
15#![deny(missing_docs)]
16
17use rustc_hash::FxHashMap as HashMap;
18use std::cell::Cell;
19use std::cmp::min;
20use std::sync::{Mutex, RwLock};
21
22pub use std::sync::Arc;
23
24#[cfg(not(feature = "pstd"))]
25pub use std::{vec as pvec, vec::Vec as PVec};
26
27/// Vec type for Data. pvec is macro.
28#[cfg(feature = "pstd")]
29pub type PVec<T> = VecA<T, Perm>;
30
31#[cfg(feature = "pstd")]
32pub use pstd::veca as pvec;
33
34#[cfg(feature = "pstd")]
35use pstd::{
36    VecA,
37    collections::{BTreeMapA, btree_map::CustomTuning},
38    localalloc::GTemp, localalloc::Perm,
39    veca as gvec,
40};
41
42#[cfg(not(feature = "pstd"))]
43use std::{collections::BTreeMap, vec as gvec, vec::Vec as GVec};
44
45#[cfg(feature = "pstd")]
46type BTreeMap<K, V> = BTreeMapA<K, V, CustomTuning<GTemp>>;
47
48#[cfg(feature = "pstd")]
49type GVec<T> = VecA<T, GTemp>;
50
51/// ```Arc<Vec<u8>>```
52pub type Data = Arc<PVec<u8>>;
53
54/// Based on [BasicAtomicFile] which makes sure that updates are all-or-nothing.
55/// Performs commit asyncronously.
56///
57/// #Example
58///
59/// ```
60/// use atom_file::{AtomicFile,DummyFile,MemFile,BasicStorage};
61/// let mut af = AtomicFile::new(MemFile::new(), DummyFile::new());
62/// af.write( 0, &[1,2,3,4] );
63/// af.commit(4);
64/// af.wait_complete();
65/// ```
66///
67/// Atomic file has two maps of writes. On commit, the latest batch of writes are sent to be written to underlying
68/// storage, and are also applied to the second map in the "CommitFile". The CommitFile map is reset when all
69/// the updates to underlying storage have been applied.
70pub struct AtomicFile {
71    /// New updates are written here.
72    map: WMap,
73    /// Underlying file, with previous updates mapped.
74    cf: Arc<RwLock<CommitFile>>,
75    /// File size.
76    size: u64,
77    /// For sending update maps to be saved.
78    tx: std::sync::mpsc::Sender<(u64, WMap)>,
79    /// Held by update process while it is active.
80    busy: Arc<Mutex<()>>,
81    /// Limit on size of CommitFile map.
82    map_lim: usize,
83}
84
85impl AtomicFile {
86    /// Construct AtomicFile with default limits. stg is the main underlying storage, upd is temporary storage for updates during commit.
87    pub fn new(stg: Box<dyn Storage>, upd: Box<dyn BasicStorage>) -> Box<Self> {
88        Self::new_with_limits(stg, upd, &Limits::default())
89    }
90
91    /// Construct Atomic file with specified limits.
92    pub fn new_with_limits(
93        stg: Box<dyn Storage>,
94        upd: Box<dyn BasicStorage>,
95        lim: &Limits,
96    ) -> Box<Self> {
97        let size = stg.size();
98        let mut baf = BasicAtomicFile::new(stg.clone(), upd, lim);
99
100        let (tx, rx) = std::sync::mpsc::channel::<(u64, WMap)>();
101        let cf = Arc::new(RwLock::new(CommitFile::new(stg, lim.rbuf_mem)));
102        let busy = Arc::new(Mutex::new(())); // Lock held while async save thread is active.
103
104        // Start the thread which does save asyncronously.
105        let (cf1, busy1) = (cf.clone(), busy.clone());
106
107        std::thread::spawn(move || {
108            // Loop that recieves a map of updates and applies it to BasicAtomicFile.
109            while let Ok((size, map)) = rx.recv() {
110                let _lock = busy1.lock();
111                baf.map = map;
112                baf.commit(size);
113                cf1.write().unwrap().done_one();
114            }
115        });
116        Box::new(Self {
117            map: WMap::default(),
118            cf,
119            size,
120            tx,
121            busy,
122            map_lim: lim.map_lim,
123        })
124    }
125}
126
127impl Storage for AtomicFile {
128    fn clone(&self) -> Box<dyn Storage> {
129        panic!()
130    }
131}
132
133impl BasicStorage for AtomicFile {
134    fn commit(&mut self, size: u64) {
135        self.size = size;
136        if self.map.is_empty() {
137            return;
138        }
139        if self.cf.read().unwrap().map.len() > self.map_lim {
140            self.wait_complete();
141        }
142        let map = std::mem::take(&mut self.map);
143        let stop =
144        {
145            let cf = &mut *self.cf.write().unwrap();
146            if cf.stop { true }
147            else
148            {
149                cf.todo += 1;
150                // Apply map of updates to CommitFile.
151                map.to_storage(cf);
152                // Send map of updates to thread to be written to underlying storage.
153                self.tx.send((size, map)).unwrap();
154                false
155            }
156        };
157        if stop {
158            // Program is terminating, loop forever.
159            loop { std::thread::sleep(std::time::Duration::from_millis(100)); }
160        }
161    }
162
163    fn size(&self) -> u64 {
164        self.size
165    }
166
167    fn read(&self, start: u64, data: &mut [u8]) {
168        self.map.read(start, data, &*self.cf.read().unwrap());
169    }
170
171    fn write_data(&mut self, start: u64, data: Data, off: usize, len: usize) {
172        self.map.write(start, data, off, len);
173    }
174
175    fn write(&mut self, start: u64, data: &[u8]) {
176        let len = data.len();
177        let d = Arc::new(PVec::from(data));
178        self.write_data(start, d, 0, len);
179    }
180
181    fn wait_complete(&self) {
182       while self.cf.read().unwrap().todo != 0 {
183           std::thread::yield_now();
184           let _x = self.busy.lock();
185       }
186    }   
187
188    fn shutdown(&mut self) {
189        self.cf.write().unwrap().stop = true; // Prevents new commits from being added.
190        self.wait_complete();
191    }       
192}
193
194struct CommitFile {
195    /// Buffered underlying storage.
196    stg: ReadBufStg<256>,
197    /// Map of committed updates.
198    map: WMap,
199    /// Number of outstanding unsaved commits.
200    todo: usize,
201    /// Flag to prevent new commits starting
202    stop: bool,
203}
204
205impl CommitFile {
206    fn new(stg: Box<dyn Storage>, buf_mem: usize) -> Self {
207        Self {
208            stg: ReadBufStg::<256>::new(stg, 50, buf_mem / 256),
209            map: WMap::default(),
210            todo: 0,
211            stop: false,
212        }
213    }
214
215    fn done_one(&mut self) {
216        self.todo -= 1;
217        if self.todo == 0 {
218            self.map = WMap::default();
219            self.stg.reset();
220        }
221    }
222}
223
224impl BasicStorage for CommitFile {
225    fn commit(&mut self, _size: u64) {
226        panic!()
227    }
228
229    fn size(&self) -> u64 {
230        panic!()
231    }
232
233    fn read(&self, start: u64, data: &mut [u8]) {
234        self.map.read(start, data, &self.stg);
235    }
236
237    fn write_data(&mut self, start: u64, data: Data, off: usize, len: usize) {
238        self.map.write(start, data, off, len);
239    }
240
241    fn write(&mut self, _start: u64, _data: &[u8]) {
242        panic!()
243    }
244}
245
246/// Storage interface - BasicStorage is some kind of "file" storage.
247///
248/// read and write methods take a start which is a byte offset in the underlying file.
249pub trait BasicStorage: Send {
250    /// Get the size of the underlying storage.
251    /// Note : this is valid initially and after a commit but is not defined after write is called.
252    fn size(&self) -> u64;
253
254    /// Read data.
255    fn read(&self, start: u64, data: &mut [u8]);
256
257    /// Write byte slice to storage.
258    fn write(&mut self, start: u64, data: &[u8]);
259
260    /// Write byte Vec.
261    fn write_vec(&mut self, start: u64, data: PVec<u8>) {
262        let len = data.len();
263        let d = Arc::new(data);
264        self.write_data(start, d, 0, len);
265    }
266
267    /// Write Data slice.
268    fn write_data(&mut self, start: u64, data: Data, off: usize, len: usize) {
269        self.write(start, &data[off..off + len]);
270    }
271
272    /// Finish write transaction, size is new size of underlying storage.
273    fn commit(&mut self, size: u64);
274
275    /// Write u64.
276    fn write_u64(&mut self, start: u64, value: u64) {
277        self.write(start, &value.to_le_bytes());
278    }
279
280    /// Read u64.
281    fn read_u64(&self, start: u64) -> u64 {
282        let mut bytes = [0; 8];
283        self.read(start, &mut bytes);
284        u64::from_le_bytes(bytes)
285    }
286
287    /// Wait until current writes are complete.
288    fn wait_complete(&self){}
289
290    /// Called on program termination.
291    fn shutdown(&mut self){}
292}
293
294/// BasicStorage with Sync and clone.
295pub trait Storage: BasicStorage + Sync {
296    /// Clone.
297    fn clone(&self) -> Box<dyn Storage>;
298}
299
300/// Simple implementation of [Storage] using `Arc<Mutex<Vec<u8>>`.
301#[derive(Default)]
302pub struct MemFile {
303    v: Arc<Mutex<Vec<u8>>>,
304}
305
306impl MemFile {
307    /// Get a new (boxed) MemFile.
308    pub fn new() -> Box<Self> {
309        Box::default()
310    }
311}
312
313impl Storage for MemFile {
314    fn clone(&self) -> Box<dyn Storage> {
315        Box::new(Self { v: self.v.clone() })
316    }
317}
318
319impl BasicStorage for MemFile {
320    fn size(&self) -> u64 {
321        let v = self.v.lock().unwrap();
322        v.len() as u64
323    }
324
325    fn read(&self, off: u64, bytes: &mut [u8]) {
326        let off = off as usize;
327        let len = bytes.len();
328        let mut v = self.v.lock().unwrap();
329        if off + len > v.len() {
330            v.resize(off + len, 0);
331        }
332        bytes.copy_from_slice(&v[off..off + len]);
333    }
334
335    fn write(&mut self, off: u64, bytes: &[u8]) {
336        let off = off as usize;
337        let len = bytes.len();
338        let mut v = self.v.lock().unwrap();
339        if off + len > v.len() {
340            v.resize(off + len, 0);
341        }
342        v[off..off + len].copy_from_slice(bytes);
343    }
344
345    fn commit(&mut self, size: u64) {
346        let mut v = self.v.lock().unwrap();
347        v.resize(size as usize, 0);
348    }
349}
350
351use std::{fs, fs::OpenOptions, io::Read, io::Seek, io::SeekFrom, io::Write};
352
353struct FileInner {
354    f: fs::File,
355}
356
357impl FileInner {
358    /// Construct from filename.
359    pub fn new(filename: &str) -> Self {
360        Self {
361            f: OpenOptions::new()
362                .read(true)
363                .write(true)
364                .create(true)
365                .truncate(false)
366                .open(filename)
367                .unwrap(),
368        }
369    }
370
371    fn size(&mut self) -> u64 {
372        self.f.seek(SeekFrom::End(0)).unwrap()
373    }
374
375    fn read(&mut self, off: u64, bytes: &mut [u8]) {
376        self.f.seek(SeekFrom::Start(off)).unwrap();
377        let _ = self.f.read(bytes).unwrap();
378    }
379
380    fn write(&mut self, off: u64, bytes: &[u8]) {
381        // The list of operating systems which auto-zero is likely more than this...research is todo.
382        #[cfg(not(any(target_os = "windows", target_os = "linux")))]
383        {
384            let size = self.f.seek(SeekFrom::End(0)).unwrap();
385            if off > size {
386                self.f.set_len(off).unwrap();
387            }
388        }
389        self.f.seek(SeekFrom::Start(off)).unwrap();
390        let _ = self.f.write(bytes).unwrap();
391    }
392
393    fn commit(&mut self, size: u64) {
394        self.f.set_len(size).unwrap();
395        self.f.sync_all().unwrap();
396    }
397}
398
399/// For atomic upd file, if not unix or windows.
400pub struct UpdFileStorage {
401    file: Cell<Option<FileInner>>,
402}
403
404impl UpdFileStorage {
405    /// Construct from filename.
406    pub fn new(filename: &str) -> Box<Self> {
407        Box::new(Self {
408            file: Cell::new(Some(FileInner::new(filename))),
409        })
410    }
411}
412
413impl BasicStorage for UpdFileStorage {
414    fn size(&self) -> u64 {
415        let mut f = self.file.take().unwrap();
416        let result = f.size();
417        self.file.set(Some(f));
418        result
419    }
420    fn read(&self, off: u64, bytes: &mut [u8]) {
421        let mut f = self.file.take().unwrap();
422        f.read(off, bytes);
423        self.file.set(Some(f));
424    }
425
426    fn write(&mut self, off: u64, bytes: &[u8]) {
427        let mut f = self.file.take().unwrap();
428        f.write(off, bytes);
429        self.file.set(Some(f));
430    }
431
432    fn commit(&mut self, size: u64) {
433        let mut f = self.file.take().unwrap();
434        f.commit(size);
435        self.file.set(Some(f));
436    }
437}
438
439/// Simple implementation of [Storage] using [`std::fs::File`].
440pub struct SimpleFileStorage {
441    file: Arc<Mutex<FileInner>>,
442}
443
444impl SimpleFileStorage {
445    /// Construct from filename.
446    pub fn new(filename: &str) -> Box<Self> {
447        Box::new(Self {
448            file: Arc::new(Mutex::new(FileInner::new(filename))),
449        })
450    }
451}
452
453impl Storage for SimpleFileStorage {
454    fn clone(&self) -> Box<dyn Storage> {
455        Box::new(Self {
456            file: self.file.clone(),
457        })
458    }
459}
460
461impl BasicStorage for SimpleFileStorage {
462    fn size(&self) -> u64 {
463        self.file.lock().unwrap().size()
464    }
465
466    fn read(&self, off: u64, bytes: &mut [u8]) {
467        self.file.lock().unwrap().read(off, bytes);
468    }
469
470    fn write(&mut self, off: u64, bytes: &[u8]) {
471        self.file.lock().unwrap().write(off, bytes);
472    }
473
474    fn commit(&mut self, size: u64) {
475        self.file.lock().unwrap().commit(size);
476    }
477}
478
479/// Alternative to SimpleFileStorage that uses multiple [SimpleFileStorage]s to allow parallel reads by different threads.
480pub struct AnyFileStorage {
481    filename: String,
482    files: Arc<Mutex<Vec<FileInner>>>,
483}
484
485impl AnyFileStorage {
486    /// Create new.
487    pub fn new(filename: &str) -> Box<Self> {
488        Box::new(Self {
489            filename: filename.to_owned(),
490            files: Arc::new(Mutex::new(Vec::new())),
491        })
492    }
493
494    fn get_file(&self) -> FileInner {
495        match self.files.lock().unwrap().pop() {
496            Some(f) => f,
497            _ => FileInner::new(&self.filename),
498        }
499    }
500
501    fn put_file(&self, f: FileInner) {
502        self.files.lock().unwrap().push(f);
503    }
504}
505
506impl Storage for AnyFileStorage {
507    fn clone(&self) -> Box<dyn Storage> {
508        Box::new(Self {
509            filename: self.filename.clone(),
510            files: self.files.clone(),
511        })
512    }
513}
514
515impl BasicStorage for AnyFileStorage {
516    fn size(&self) -> u64 {
517        let mut f = self.get_file();
518        let result = f.size();
519        self.put_file(f);
520        result
521    }
522
523    fn read(&self, off: u64, bytes: &mut [u8]) {
524        let mut f = self.get_file();
525        f.read(off, bytes);
526        self.put_file(f);
527    }
528
529    fn write(&mut self, off: u64, bytes: &[u8]) {
530        let mut f = self.get_file();
531        f.write(off, bytes);
532        self.put_file(f);
533    }
534
535    fn commit(&mut self, size: u64) {
536        let mut f = self.get_file();
537        f.commit(size);
538        self.put_file(f);
539    }
540}
541
542/// Dummy Stg that can be used for Atomic upd file if "reliable" atomic commits are not required.
543pub struct DummyFile {}
544impl DummyFile {
545    /// Construct.
546    pub fn new() -> Box<Self> {
547        Box::new(Self {})
548    }
549}
550
551impl Storage for DummyFile {
552    fn clone(&self) -> Box<dyn Storage> {
553        Self::new()
554    }
555}
556
557impl BasicStorage for DummyFile {
558    fn size(&self) -> u64 {
559        0
560    }
561
562    fn read(&self, _off: u64, _bytes: &mut [u8]) {}
563
564    fn write(&mut self, _off: u64, _bytes: &[u8]) {}
565
566    fn commit(&mut self, _size: u64) {}
567}
568
569/// Memory configuration limits for [`AtomicFile`].
570#[non_exhaustive]
571pub struct Limits {
572    /// Limit on size of commit write map, default is 5000.
573    pub map_lim: usize,
574    /// Memory for buffering small reads, default is 0x200000 ( 2MB ).
575    pub rbuf_mem: usize,
576    /// Memory for buffering writes to main storage, default is 0x100000 (1MB).
577    pub swbuf: usize,
578    /// Memory for buffering writes to temporary storage, default is 0x100000 (1MB).
579    pub uwbuf: usize,
580}
581
582impl Default for Limits {
583    fn default() -> Self {
584        Self {
585            map_lim: 5000,
586            rbuf_mem: 0x200000,
587            swbuf: 0x100000,
588            uwbuf: 0x100000,
589        }
590    }
591}
592
593/// Write Buffer.
594struct WriteBuffer {
595    /// Current write index into buf.
596    ix: usize,
597    /// Current file position.
598    pos: u64,
599    /// Underlying storage.
600    pub stg: Box<dyn BasicStorage>,
601    /// Buffer.
602    buf: Vec<u8>,
603}
604
605impl WriteBuffer {
606    /// Construct.
607    pub fn new(stg: Box<dyn BasicStorage>, buf_size: usize) -> Self {
608        Self {
609            ix: 0,
610            pos: u64::MAX,
611            stg,
612            buf: vec![0; buf_size],
613        }
614    }
615
616    /// Write data to specified offset,
617    pub fn write(&mut self, off: u64, data: &[u8]) {
618        if self.pos + self.ix as u64 != off {
619            self.flush(off);
620        }
621        let mut done: usize = 0;
622        let mut todo: usize = data.len();
623        while todo > 0 {
624            let mut n: usize = self.buf.len() - self.ix;
625            if n == 0 {
626                self.flush(off + done as u64);
627                n = self.buf.len();
628            }
629            if n > todo {
630                n = todo;
631            }
632            self.buf[self.ix..self.ix + n].copy_from_slice(&data[done..done + n]);
633            todo -= n;
634            done += n;
635            self.ix += n;
636        }
637    }
638
639    fn flush(&mut self, new_pos: u64) {
640        if self.ix > 0 {
641            self.stg.write(self.pos, &self.buf[0..self.ix]);
642        }
643        self.ix = 0;
644        self.pos = new_pos;
645    }
646
647    /// Commit.
648    pub fn commit(&mut self, size: u64) {
649        self.flush(u64::MAX);
650        self.stg.commit(size);
651    }
652
653    /// Write u64.
654    pub fn write_u64(&mut self, start: u64, value: u64) {
655        self.write(start, &value.to_le_bytes());
656    }
657}
658
659/// ReadBufStg buffers small (up to limit) reads to the underlying storage using multiple buffers. Only supported functions are read and reset.
660///
661/// See implementation of AtomicFile for how this is used in conjunction with WMap.
662///
663/// N is buffer size.
664struct ReadBufStg<const N: usize> {
665    /// Underlying storage.
666    stg: Box<dyn Storage>,
667    /// Buffers.
668    buf: Mutex<ReadBuffer<N>>,
669    /// Read size that is considered small.
670    limit: usize,
671}
672
673impl<const N: usize> Drop for ReadBufStg<N> {
674    fn drop(&mut self) {
675        self.reset();
676    }
677}
678
679impl<const N: usize> ReadBufStg<N> {
680    /// limit is the size of a read that is considered "small", max_buf is the maximum number of buffers used.
681    pub fn new(stg: Box<dyn Storage>, limit: usize, max_buf: usize) -> Self {
682        Self {
683            stg,
684            buf: Mutex::new(ReadBuffer::<N>::new(max_buf)),
685            limit,
686        }
687    }
688
689    /// Clears the buffers.
690    fn reset(&mut self) {
691        self.buf.lock().unwrap().reset();
692    }
693}
694
695impl<const N: usize> BasicStorage for ReadBufStg<N> {
696    /// Read data from storage.
697    fn read(&self, start: u64, data: &mut [u8]) {
698        if data.len() <= self.limit {
699            self.buf.lock().unwrap().read(&*self.stg, start, data);
700        } else {
701            self.stg.read(start, data);
702        }
703    }
704
705    /// Panics.
706    fn size(&self) -> u64 {
707        panic!()
708    }
709
710    /// Panics.
711    fn write(&mut self, _start: u64, _data: &[u8]) {
712        panic!();
713    }
714
715    /// Panics.
716    fn commit(&mut self, _size: u64) {
717        panic!();
718    }
719}
720
721struct ReadBuffer<const N: usize> {
722    /// Maps sector mumbers cached buffers.
723    map: HashMap<u64, Box<[u8; N]>>,
724    /// Maximum number of buffers.
725    max_buf: usize,
726}
727
728impl<const N: usize> ReadBuffer<N> {
729    fn new(max_buf: usize) -> Self {
730        Self {
731            map: HashMap::default(),
732            max_buf,
733        }
734    }
735
736    fn reset(&mut self) {
737        self.map.clear();
738    }
739
740    fn read(&mut self, stg: &dyn BasicStorage, off: u64, data: &mut [u8]) {
741        let mut done = 0;
742        while done < data.len() {
743            let off = off + done as u64;
744            let sector = off / N as u64;
745            let disp = (off % N as u64) as usize;
746            let amount = min(data.len() - done, N - disp);
747
748            let p = self.map.entry(sector).or_insert_with(|| {
749                let mut p: Box<[u8; N]> = vec![0; N].try_into().unwrap();
750                stg.read(sector * N as u64, &mut *p);
751                p
752            });
753            data[done..done + amount].copy_from_slice(&p[disp..disp + amount]);
754            done += amount;
755        }
756        if self.map.len() >= self.max_buf {
757            self.reset();
758        }
759    }
760}
761
762#[derive(Default)]
763/// Slice of Data to be written to storage.
764struct DataSlice {
765    /// Slice data.
766    pub data: Data,
767    /// Start of slice.
768    pub off: usize,
769    /// Length of slice.
770    pub len: usize,
771}
772
773impl DataSlice {
774    /// Get reference to the whole slice.
775    pub fn all(&self) -> &[u8] {
776        &self.data[self.off..self.off + self.len]
777    }
778    /// Get reference to part of slice.
779    pub fn part(&self, off: usize, len: usize) -> &[u8] {
780        &self.data[self.off + off..self.off + off + len]
781    }
782    /// Trim specified amount from start of slice.
783    pub fn trim(&mut self, trim: usize) {
784        self.off += trim;
785        self.len -= trim;
786    }
787    /// Take the data.
788    #[allow(dead_code)]
789    pub fn take(&mut self) -> Data {
790        std::mem::take(&mut self.data)
791    }
792}
793
794#[derive(Default)]
795/// Updateable store based on some underlying storage.
796struct WMap {
797    /// Map of writes. Key is the end of the slice.
798    map: BTreeMap<u64, DataSlice>,
799}
800
801impl WMap {
802    /// Is the map empty?
803    pub fn is_empty(&self) -> bool {
804        self.map.is_empty()
805    }
806
807    /// Number of key-value pairs in the map.
808    pub fn len(&self) -> usize {
809        self.map.len()
810    }
811
812    /// Take the map and convert it to a Vec.
813    pub fn convert_to_vec(&mut self) -> GVec<(u64, DataSlice)> {
814        let map = std::mem::take(&mut self.map);
815        let mut result = GVec::with_capacity(map.len());
816        for (end, v) in map {
817            let start = end - v.len as u64;
818            result.push((start, v));
819        }
820        result
821    }
822
823    /// Write the map into storage.
824    pub fn to_storage(&self, stg: &mut dyn BasicStorage) {
825        for (end, v) in self.map.iter() {
826            let start = end - v.len as u64;
827            stg.write_data(start, v.data.clone(), v.off, v.len);
828        }
829    }
830
831    #[cfg(not(feature = "pstd"))]
832    /// Write to storage, existing writes which overlap with new write need to be trimmed or removed.
833    pub fn write(&mut self, start: u64, data: Data, off: usize, len: usize) {
834        if len != 0 {
835            let (mut insert, mut remove) = (Vec::new(), Vec::new());
836            let end = start + len as u64;
837            for (ee, v) in self.map.range_mut(start + 1..) {
838                let ee = *ee;
839                let es = ee - v.len as u64; // Existing write Start.
840                if es >= end {
841                    // Existing write starts after end of new write, nothing to do.
842                    break;
843                } else if start <= es {
844                    if end < ee {
845                        // New write starts before existing write, but doesn't subsume it. Trim existing write.
846                        v.trim((end - es) as usize);
847                        break;
848                    }
849                    // New write subsumes existing write entirely, remove existing write.
850                    remove.push(ee);
851                } else if end < ee {
852                    // New write starts in middle of existing write, ends before end of existing write,
853                    // put start of existing write in insert list, trim existing write.
854                    insert.push((es, v.data.clone(), v.off, (start - es) as usize));
855                    v.trim((end - es) as usize);
856                    break;
857                } else {
858                    // New write starts in middle of existing write, ends after existing write,
859                    // put start of existing write in insert list, remove existing write.
860                    insert.push((es, v.take(), v.off, (start - es) as usize));
861                    remove.push(ee);
862                }
863            }
864            for end in remove {
865                self.map.remove(&end);
866            }
867            for (start, data, off, len) in insert {
868                self.map
869                    .insert(start + len as u64, DataSlice { data, off, len });
870            }
871            self.map
872                .insert(start + len as u64, DataSlice { data, off, len });
873        }
874    }
875
876    #[cfg(feature = "pstd")]
877    /// Write to storage, existing writes which overlap with new write need to be trimmed or removed.
878    pub fn write(&mut self, start: u64, data: Data, off: usize, len: usize) {
879        if len != 0 {
880            let end = start + len as u64;
881            let mut c = self
882                .map
883                .lower_bound_mut(std::ops::Bound::Excluded(&start))
884                .with_mutable_key();
885            while let Some((eend, v)) = c.next() {
886                let ee = *eend;
887                let es = ee - v.len as u64; // Existing write Start.
888                if es >= end {
889                    // Existing write starts after end of new write, nothing to do.
890                    c.prev();
891                    break;
892                } else if start <= es {
893                    if end < ee {
894                        // New write starts before existing write, but doesn't subsume it. Trim existing write.
895                        v.trim((end - es) as usize);
896                        c.prev();
897                        break;
898                    }
899                    // New write subsumes existing write entirely, remove existing write.
900                    c.remove_prev();
901                } else if end < ee {
902                    // New write starts in middle of existing write, ends before end of existing write,
903                    // trim existing write, insert start of existing write.
904                    let (data, off, len) = (v.data.clone(), v.off, (start - es) as usize);
905                    v.trim((end - es) as usize);
906                    c.prev();
907                    c.insert_before_unchecked(es + len as u64, DataSlice { data, off, len });
908                    break;
909                } else {
910                    // New write starts in middle of existing write, ends after existing write,
911                    // Trim existing write ( modifies key, but this is ok as ordering is not affected ).
912                    v.len = (start - es) as usize;
913                    *eend = es + v.len as u64;
914                }
915            }
916            // Insert the new write.
917            c.insert_after_unchecked(start + len as u64, DataSlice { data, off, len });
918        }
919    }
920
921    /// Read from storage, taking map of existing writes into account. Unwritten ranges are read from underlying storage.
922    pub fn read(&self, start: u64, data: &mut [u8], u: &dyn BasicStorage) {
923        let len = data.len();
924        if len != 0 {
925            let mut done = 0;
926            for (&end, v) in self.map.range(start + 1..) {
927                let es = end - v.len as u64; // Existing write Start.
928                let doff = start + done as u64;
929                if es > doff {
930                    // Read from underlying storage.
931                    let a = min(len - done, (es - doff) as usize);
932                    u.read(doff, &mut data[done..done + a]);
933                    done += a;
934                    if done == len {
935                        return;
936                    }
937                }
938                // Use existing write.
939                let skip = (start + done as u64 - es) as usize;
940                let a = min(len - done, v.len - skip);
941                data[done..done + a].copy_from_slice(v.part(skip, a));
942                done += a;
943                if done == len {
944                    return;
945                }
946            }
947            u.read(start + done as u64, &mut data[done..]);
948        }
949    }
950}
951
952/// Basis for [crate::AtomicFile] ( non-async alternative ). Provides two-phase commit and buffering of writes.
953pub struct BasicAtomicFile {
954    /// The main underlying storage.
955    stg: WriteBuffer,
956    /// Temporary storage for updates during commit.
957    upd: WriteBuffer,
958    /// Map of writes.
959    map: WMap,
960    /// List of writes.
961    list: GVec<(u64, DataSlice)>,
962    size: u64,
963    stop: bool,
964}
965
966impl BasicAtomicFile {
967    /// stg is the main underlying storage, upd is temporary storage for updates during commit.
968    pub fn new(stg: Box<dyn BasicStorage>, upd: Box<dyn BasicStorage>, lim: &Limits) -> Box<Self> {
969        let size = stg.size();
970        let mut result = Box::new(Self {
971            stg: WriteBuffer::new(stg, lim.swbuf),
972            upd: WriteBuffer::new(upd, lim.uwbuf),
973            map: WMap::default(),
974            list: GVec::new(),
975            size,
976            stop: false,
977        });
978        result.init();
979        result
980    }
981
982    /// Apply outstanding updates.
983    fn init(&mut self) {
984        let end = self.upd.stg.read_u64(0);
985        let size = self.upd.stg.read_u64(8);
986        if end == 0 {
987            return;
988        }
989        assert!(end == self.upd.stg.size());
990        let mut pos = 16;
991        while pos < end {
992            let start = self.upd.stg.read_u64(pos);
993            pos += 8;
994            let len = self.upd.stg.read_u64(pos);
995            pos += 8;
996            let mut buf: GVec<u8> = gvec![0; len as usize];
997            self.upd.stg.read(pos, &mut buf);
998            pos += len;
999            self.stg.write(start, &buf);
1000        }
1001        self.stg.commit(size);
1002        self.upd.commit(0);
1003    }
1004
1005    /// Perform the specified phase ( 1 or 2 ) of a two-phase commit.
1006    pub fn commit_phase(&mut self, size: u64, phase: u8) {
1007        if self.map.is_empty() && self.list.is_empty() {
1008            return;
1009        }
1010        if phase == 1 {
1011            self.list = self.map.convert_to_vec();
1012
1013            // Write the updates to upd.
1014            // First set the end position to zero.
1015            self.upd.write_u64(0, 0);
1016            self.upd.write_u64(8, size);
1017            self.upd.commit(16); // Not clear if this is necessary.
1018
1019            // Write the update records.
1020            let mut stg_written = false;
1021            let mut pos: u64 = 16;
1022            for (start, v) in self.list.iter() {
1023                let (start, len, data) = (*start, v.len as u64, v.all());
1024                if start >= self.size {
1025                    // Writes beyond current stg size can be written directly.
1026                    stg_written = true;
1027                    self.stg.write(start, data);
1028                } else {
1029                    self.upd.write_u64(pos, start);
1030                    pos += 8;
1031                    self.upd.write_u64(pos, len);
1032                    pos += 8;
1033                    self.upd.write(pos, data);
1034                    pos += len;
1035                }
1036            }
1037            if stg_written {
1038                self.stg.commit(size);
1039            }
1040            self.upd.commit(pos); // Not clear if this is necessary.
1041
1042            // Set the end position.
1043            self.upd.write_u64(0, pos);
1044            self.upd.write_u64(8, size);
1045            self.upd.commit(pos);
1046        } else {
1047            for (start, v) in self.list.iter() {
1048                if *start < self.size {
1049                    // Writes beyond current stg size have already been written.
1050                    self.stg.write(*start, v.all());
1051                }
1052            }
1053            self.list = GVec::new();
1054            self.stg.commit(size);
1055            self.upd.commit(0);
1056        }
1057    }
1058}
1059
1060impl BasicStorage for BasicAtomicFile {
1061    fn commit(&mut self, size: u64) {
1062        if self.stop { return; }
1063        self.commit_phase(size, 1);
1064        self.commit_phase(size, 2);
1065        self.size = size;
1066    }
1067
1068    fn size(&self) -> u64 {
1069        self.size
1070    }
1071
1072    fn read(&self, start: u64, data: &mut [u8]) {
1073        self.map.read(start, data, &*self.stg.stg);
1074    }
1075
1076    fn write_data(&mut self, start: u64, data: Data, off: usize, len: usize) {
1077        self.map.write(start, data, off, len);
1078    }
1079
1080    fn write(&mut self, start: u64, data: &[u8]) {
1081        let len = data.len();
1082        let d = Arc::new(PVec::from(data));
1083        self.write_data(start, d, 0, len);
1084    }
1085
1086    fn shutdown(&mut self)
1087    {
1088        self.stop = true;
1089    }
1090}
1091
1092/// Optimized implementation of [Storage] ( unix only ).
1093#[cfg(target_family = "unix")]
1094pub struct UnixFileStorage {
1095    size: Arc<Mutex<u64>>,
1096    f: fs::File,
1097}
1098#[cfg(target_family = "unix")]
1099impl UnixFileStorage {
1100    /// Construct from filename.
1101    pub fn new(filename: &str) -> Box<Self> {
1102        let mut f = OpenOptions::new()
1103            .read(true)
1104            .write(true)
1105            .create(true)
1106            .truncate(false)
1107            .open(filename)
1108            .unwrap();
1109        let size = f.seek(SeekFrom::End(0)).unwrap();
1110        let size = Arc::new(Mutex::new(size));
1111        Box::new(Self { size, f })
1112    }
1113}
1114
1115#[cfg(target_family = "unix")]
1116impl Storage for UnixFileStorage {
1117    fn clone(&self) -> Box<dyn Storage> {
1118        Box::new(Self {
1119            size: self.size.clone(),
1120            f: self.f.try_clone().unwrap(),
1121        })
1122    }
1123}
1124
1125#[cfg(target_family = "unix")]
1126use std::os::unix::fs::FileExt;
1127
1128#[cfg(target_family = "unix")]
1129impl BasicStorage for UnixFileStorage {
1130    fn read(&self, start: u64, data: &mut [u8]) {
1131        let _ = self.f.read_at(data, start);
1132    }
1133
1134    fn write(&mut self, start: u64, data: &[u8]) {
1135        let _ = self.f.write_at(data, start);
1136    }
1137
1138    fn size(&self) -> u64 {
1139        *self.size.lock().unwrap()
1140    }
1141
1142    fn commit(&mut self, size: u64) {
1143        *self.size.lock().unwrap() = size;
1144        self.f.set_len(size).unwrap();
1145        self.f.sync_all().unwrap();
1146    }
1147}
1148
1149/// Optimized implementation of [Storage] ( windows only ).
1150#[cfg(target_family = "windows")]
1151pub struct WindowsFileStorage {
1152    size: Arc<Mutex<u64>>,
1153    f: fs::File,
1154}
1155#[cfg(target_family = "windows")]
1156impl WindowsFileStorage {
1157    /// Construct from filename.
1158    pub fn new(filename: &str) -> Box<Self> {
1159        let mut f = OpenOptions::new()
1160            .read(true)
1161            .write(true)
1162            .create(true)
1163            .truncate(false)
1164            .open(filename)
1165            .unwrap();
1166        let size = f.seek(SeekFrom::End(0)).unwrap();
1167        let size = Arc::new(Mutex::new(size));
1168        Box::new(Self { size, f })
1169    }
1170}
1171
1172#[cfg(target_family = "windows")]
1173impl Storage for WindowsFileStorage {
1174    fn clone(&self) -> Box<dyn Storage> {
1175        Box::new(Self {
1176            size: self.size.clone(),
1177            f: self.f.try_clone().unwrap(),
1178        })
1179    }
1180}
1181
1182#[cfg(target_family = "windows")]
1183use std::os::windows::fs::FileExt;
1184
1185#[cfg(target_family = "windows")]
1186impl BasicStorage for WindowsFileStorage {
1187    fn read(&self, start: u64, data: &mut [u8]) {
1188        let _ = self.f.seek_read(data, start);
1189    }
1190
1191    fn write(&mut self, start: u64, data: &[u8]) {
1192        let _ = self.f.seek_write(data, start);
1193    }
1194
1195    fn size(&self) -> u64 {
1196        *self.size.lock().unwrap()
1197    }
1198
1199    fn commit(&mut self, size: u64) {
1200        *self.size.lock().unwrap() = size;
1201        self.f.set_len(size).unwrap();
1202        self.f.sync_all().unwrap();
1203    }
1204}
1205
1206/// Optimised Storage ( varies according to platform ).
1207#[cfg(target_family = "windows")]
1208pub type MultiFileStorage = WindowsFileStorage;
1209
1210/// Optimised Storage ( varies according to platform ).
1211#[cfg(target_family = "unix")]
1212pub type MultiFileStorage = UnixFileStorage;
1213
1214/// Optimised Storage ( varies according to platform ).
1215#[cfg(not(any(target_family = "unix", target_family = "windows")))]
1216pub type MultiFileStorage = AnyFileStorage;
1217
1218/// Fast Storage for upd file ( varies according to platform ).
1219#[cfg(any(target_family = "windows", target_family = "unix"))]
1220pub type FastFileStorage = MultiFileStorage;
1221
1222/// Fast Storage for upd file ( varies according to platform ).
1223#[cfg(not(any(target_family = "windows", target_family = "unix")))]
1224pub type FastFileStorage = UpdFileStorage;
1225
1226#[cfg(test)]
1227/// Get amount of testing from environment variable TA.
1228fn test_amount() -> usize {
1229    str::parse(&std::env::var("TA").unwrap_or("1".to_string())).unwrap()
1230}
1231
1232#[test]
1233fn test_atomic_file() {
1234    use rand::Rng;
1235    /* Idea of test is to check AtomicFile and MemFile behave the same */
1236
1237    let ta = test_amount();
1238    println!(" Test amount={}", ta);
1239
1240    let mut rng = rand::thread_rng();
1241
1242    for _ in 0..100 {
1243        let mut s1 = AtomicFile::new(MemFile::new(), MemFile::new());
1244        // let mut s1 = BasicAtomicFile::new(MemFile::new(), MemFile::new(), &Limits::default() );
1245        let mut s2 = MemFile::new();
1246
1247        for _ in 0..1000 * ta {
1248            let off: usize = rng.r#gen::<usize>() % 100;
1249            let mut len = 1 + rng.r#gen::<usize>() % 20;
1250            let w: bool = rng.r#gen();
1251            if w {
1252                let mut bytes = Vec::new();
1253                while len > 0 {
1254                    len -= 1;
1255                    let b: u8 = rng.r#gen::<u8>();
1256                    bytes.push(b);
1257                }
1258                s1.write(off as u64, &bytes);
1259                s2.write(off as u64, &bytes);
1260            } else {
1261                let mut b2 = vec![0; len];
1262                let mut b3 = vec![0; len];
1263                s1.read(off as u64, &mut b2);
1264                s2.read(off as u64, &mut b3);
1265                assert!(b2 == b3);
1266            }
1267            if rng.r#gen::<usize>() % 50 == 0 {
1268                s1.commit(200);
1269                s2.commit(200);
1270            }
1271        }
1272    }
1273}