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