1use std::io::{self, ErrorKind};
2
3#[cfg(any(unix, windows))]
4use std::{
5 fs::File,
6 path::Path,
7 sync::{Mutex, MutexGuard},
8};
9
10#[cfg(unix)]
11use std::os::unix::fs::FileExt;
12#[cfg(windows)]
13use std::os::windows::fs::FileExt;
14
15#[cfg(any(unix, windows))]
16struct FileIo {
17 #[cfg(unix)]
18 file: File,
19 #[cfg(windows)]
20 file: std::sync::Arc<File>,
21 /// OneNote's coordination bytes, unlocked when dropped.
22 #[cfg(windows)]
23 locks: Vec<file_guard::FileGuard<std::sync::Arc<File>>>,
24 _process: MutexGuard<'static, ()>,
25 unlock_on_drop: bool,
26}
27
28#[cfg(any(unix, windows))]
29impl FileIo {
30 fn open(path: impl AsRef<Path>, write: bool) -> io::Result<Self> {
31 static PROCESS: Mutex<()> = Mutex::new(());
32 let process = PROCESS
33 .lock()
34 .map_err(|_| io::Error::other("A file operation panicked in this process"))?;
35 let mut options = File::options();
36 options.read(true).write(write);
37 #[cfg(target_os = "macos")]
38 {
39 use std::os::unix::fs::OpenOptionsExt;
40 // SMB can lose exclusion when separate opens race with flock. On smbfs these are
41 // share modes, taken as OneNote takes them: a writer's shared lock denies only
42 // other writers, and a reader takes none, as `stable` sees past a commit.
43 let smb = nix::sys::statfs::statfs(path.as_ref())
44 .is_ok_and(|fs| fs.filesystem_type_name() == "smbfs");
45 let lock = match (write, smb) {
46 (true, false) => nix::libc::O_EXLOCK,
47 (false, true) => 0,
48 _ => nix::libc::O_SHLOCK,
49 };
50 options.custom_flags(lock | nix::libc::O_NONBLOCK);
51 }
52 #[cfg(windows)]
53 {
54 use std::os::windows::fs::OpenOptionsExt;
55 // As OneNote opens a section: a reader shares it with everyone, a writer denies
56 // other writers (FILE_SHARE_READ, FILE_SHARE_WRITE, FILE_SHARE_DELETE).
57 options.share_mode(if write { 0x1 | 0x4 } else { 0x1 | 0x2 | 0x4 });
58 }
59 #[cfg(windows)]
60 let file = std::sync::Arc::new(options.open(path).map_err(|error| {
61 // ERROR_SHARING_VIOLATION: another writer has it open.
62 match error.raw_os_error() {
63 Some(32) => ErrorKind::WouldBlock.into(),
64 _ => error,
65 }
66 })?);
67 #[cfg(unix)]
68 let file = options.open(path.as_ref())?;
69 #[cfg(all(unix, not(target_os = "macos")))]
70 {
71 use std::os::unix::fs::MetadataExt;
72 file.try_lock()?;
73 // Locked after it was opened, the file may since have been superseded.
74 let (open, named) = (file.metadata()?, std::fs::metadata(path)?);
75 if (open.dev(), open.ino()) != (named.dev(), named.ino()) {
76 return Err(ErrorKind::ResourceBusy.into());
77 }
78 }
79 Ok(Self {
80 #[cfg(windows)]
81 locks: onenote_locks(&file, write)?,
82 file,
83 _process: process,
84 unlock_on_drop: true,
85 })
86 }
87
88 #[cfg(unix)]
89 fn release(&mut self) -> io::Result<()> {
90 // A failed unlock may have reached the server; Drop must not repeat it.
91 self.unlock_on_drop = false;
92 self.file.unlock()
93 }
94
95 #[cfg(windows)]
96 fn release(&mut self) -> io::Result<()> {
97 self.unlock_on_drop = false;
98 self.locks.clear();
99 Ok(())
100 }
101
102 fn finish(mut self, result: Result<(), CommitError>) -> Result<(), CommitError> {
103 let released = self.release();
104 result?;
105 released.map_err(|error| CommitError {
106 state: CommitState::Committed,
107 error,
108 })
109 }
110}
111
112#[cfg(unix)]
113impl Drop for FileIo {
114 fn drop(&mut self) {
115 if self.unlock_on_drop {
116 let _ = self.file.unlock();
117 }
118 }
119}
120
121/// Windows takes OneNote 2010's own locks on the file, one byte each past any data, so
122/// neither app's locks bar the other's reads: the reader byte shared, and to write, the
123/// writer byte exclusively, as `notebook::smb` takes them on a share.
124#[cfg(windows)]
125fn onenote_locks(
126 file: &std::sync::Arc<File>,
127 write: bool,
128) -> io::Result<Vec<file_guard::FileGuard<std::sync::Arc<File>>>> {
129 use file_guard::Lock;
130 let reader = file_guard::try_lock(file.clone(), Lock::Shared, 0xffff_fffb, 1)?;
131 let mut locks = vec![reader];
132 if write {
133 locks.push(file_guard::try_lock(
134 file.clone(),
135 Lock::Exclusive,
136 0xffff_fffd,
137 1,
138 )?);
139 }
140 Ok(locks)
141}
142
143/// Places a file in its notebook the way OneNote does on adoption: the header's
144/// `guidAncestor` becomes the parent table of contents' file identity and `crcName` the CRC
145/// of `name` (a section's file name, a group's folder name). OneNote re-identifies a file
146/// whose header disagrees with its location, which orphans its TOC entry. The caller holds
147/// OneNote-compatible exclusion on `io`.
148pub fn place(io: &mut impl CommitIo, ancestor: [u8; 16], name: &str) -> io::Result<()> {
149 let mut header = [0; 1024];
150 if io.read_at(0, &mut header)? != header.len() {
151 return Err(io::Error::from(ErrorKind::UnexpectedEof));
152 }
153 crate::Header::parse(&header)
154 .map_err(|error| io::Error::new(ErrorKind::InvalidData, error.message))?;
155 let placement = crate::create::placement(ancestor, name);
156 if io.write_at(128, &placement)? != placement.len() {
157 return Err(io::Error::from(ErrorKind::WriteZero));
158 }
159 io.flush()
160}
161
162/// `place` under the conservative filesystem adapter's whole-file exclusion.
163#[cfg(any(unix, windows))]
164pub fn place_file(path: impl AsRef<Path>, ancestor: [u8; 16], name: &str) -> io::Result<()> {
165 let mut io = FileIo::open(path, true)?;
166 let result = place(&mut io, ancestor, name);
167 let released = io.release();
168 result?;
169 released
170}
171
172/// Reads a snapshot, excluding writers as their commits exclude it, except on an SMB mount,
173/// where it takes no lock, as OneNote's readers take none (macOS shares it with other readers).
174/// A read that meets a commit in progress is read again, then refused as `WouldBlock`.
175#[cfg(any(unix, windows))]
176pub fn read_file(path: impl AsRef<Path>) -> io::Result<Vec<u8>> {
177 read_file_limited(path, usize::MAX)
178}
179
180/// `read_file`, rejecting a snapshot larger than the byte limit.
181/// A size failure returns `FileTooLarge` without a partial snapshot.
182#[cfg(any(unix, windows))]
183pub fn read_file_limited(path: impl AsRef<Path>, limit: usize) -> io::Result<Vec<u8>> {
184 let mut io = FileIo::open(path, false)?;
185 let result = stable(|offset, output| io.read_at(offset, output), limit);
186 let released = io.release();
187 let bytes = result?;
188 released?;
189 Ok(bytes)
190}
191
192/// How many times a read that meets a commit in progress is made before it is refused.
193#[cfg(any(unix, windows))]
194const TRIES: usize = 3;
195
196/// The file through `read`, read again where a commit tore it: its header changed while it was
197/// read, as commits write the header last, or it ends short of the header's length.
198#[cfg(any(unix, windows))]
199fn stable(
200 mut read: impl FnMut(u64, &mut [u8]) -> io::Result<usize>,
201 limit: usize,
202) -> io::Result<Vec<u8>> {
203 let mut block = vec![0; 1 << 16];
204 for _ in 0..TRIES {
205 let mut bytes = Vec::new();
206 loop {
207 let size = (limit.saturating_add(1) - bytes.len()).min(block.len());
208 match read(bytes.len() as u64, &mut block[..size]) {
209 Ok(0) => break,
210 Ok(count) if count <= size => bytes.extend_from_slice(&block[..count]),
211 Ok(_) => return Err(ErrorKind::InvalidData.into()),
212 Err(error) if error.kind() == ErrorKind::Interrupted => {}
213 Err(error) => return Err(error),
214 }
215 if bytes.len() > limit {
216 return Err(ErrorKind::FileTooLarge.into());
217 }
218 }
219 let mut header = vec![0; bytes.len().min(1024)];
220 match crate::snapshot::read_exact(&mut read, 0, &mut header) {
221 Ok(()) => {}
222 Err(error) if error.kind() == ErrorKind::UnexpectedEof => continue,
223 Err(error) => return Err(error),
224 }
225 let whole = crate::Header::parse(&bytes)
226 .map_or(true, |parsed| parsed.expected_length <= bytes.len() as u64);
227 if header == bytes[..header.len()] && whole {
228 return Ok(bytes);
229 }
230 }
231 Err(ErrorKind::WouldBlock.into())
232}
233
234#[cfg(any(unix, windows))]
235impl CommitIo for FileIo {
236 fn read_at(&mut self, offset: u64, bytes: &mut [u8]) -> io::Result<usize> {
237 #[cfg(unix)]
238 return self.file.read_at(bytes, offset);
239 #[cfg(windows)]
240 return self.file.seek_read(bytes, offset);
241 }
242
243 fn write_at(&mut self, offset: u64, bytes: &[u8]) -> io::Result<usize> {
244 #[cfg(unix)]
245 return self.file.write_at(bytes, offset);
246 #[cfg(windows)]
247 return self.file.seek_write(bytes, offset);
248 }
249
250 fn flush(&mut self) -> io::Result<()> {
251 crate::flush::flush(&self.file)
252 }
253}
254
255/// `confirm` under the same whole-file exclusion `commit_file` uses.
256#[cfg(any(unix, windows))]
257pub fn confirm_file(path: impl AsRef<Path>, base: &Stamp) -> Result<(), CommitError> {
258 let mut io = FileIo::open(path, true).map_err(|error| CommitError {
259 state: CommitState::NotCommitted,
260 error,
261 })?;
262 let result = confirm(&mut io, base);
263 io.finish(result)
264}
265
266/// Puts the file at `with` in the place of the file at `path`, provided `path` still has
267/// `base`'s stamp, under the exclusion `commit_file` takes: the whole-image write OneNote 2010
268/// makes when it writes a section anew.
269#[cfg(any(unix, windows))]
270pub fn supersede_file(
271 path: impl AsRef<Path>,
272 base: &Stamp,
273 with: impl AsRef<Path>,
274) -> Result<(), CommitError> {
275 let (path, with) = (path.as_ref(), with.as_ref());
276 let failed = |state| move |error| CommitError { state, error };
277 // Windows renames nothing over an open file, so there the old file goes aside first and is
278 // deleted once released, as OneNote's maintenance does.
279 let mut aside = path.as_os_str().to_owned();
280 if cfg!(windows) {
281 aside.push(".old");
282 }
283 let aside = std::path::PathBuf::from(aside);
284 let mut io = FileIo::open(path, true).map_err(failed(CommitState::NotCommitted))?;
285 let result = base
286 .check(&mut io)
287 .and_then(|()| std::fs::rename(path, &aside))
288 .map_err(failed(CommitState::NotCommitted))
289 .and_then(|()| {
290 std::fs::rename(with, path).map_err(|error| CommitError {
291 state: match std::fs::rename(&aside, path) {
292 Ok(()) => CommitState::NotCommitted,
293 Err(_) => CommitState::Unknown,
294 },
295 error,
296 })
297 });
298 io.finish(result)?;
299 if aside != path {
300 std::fs::remove_file(aside).map_err(failed(CommitState::Committed))?;
301 }
302 Ok(())
303}
304
305/// Checks that the file still has `base`'s stamp and flushes it, then refreshes its header
306/// version metadata. No revision is added; reread before committing on `base` again.
307/// The caller must hold OneNote-compatible exclusion and independently establish which
308/// intents the image contains. A successful read alone is not a durable acknowledgement.
309pub fn confirm(io: &mut impl CommitIo, base: &Stamp) -> Result<(), CommitError> {
310 let mut state = CommitState::NotCommitted;
311 let result = (|| -> io::Result<()> {
312 let header = crate::Header::parse(&base.header).map_err(io::Error::other)?;
313 let generation = header
314 .generation
315 .checked_add(1)
316 .ok_or(ErrorKind::InvalidData)?;
317 let mut version = [0; 40];
318 version[..16].copy_from_slice(&crate::write::fresh_guid().map_err(io::Error::other)?);
319 version[16..24].copy_from_slice(&generation.to_le_bytes());
320 version[24..].copy_from_slice(&crate::write::fresh_guid().map_err(io::Error::other)?);
321 base.check(io)?;
322 state = CommitState::Unknown;
323 io.flush()?;
324 write_all(io, 212, &version)?;
325 io.flush()
326 })();
327 result.map_err(|error| CommitError { state, error })
328}
329
330/// The caller must hold OneNote-compatible exclusion for the entire operation.
331/// Flush must make preceding writes durable before subsequent writes can persist.
332pub trait CommitIo {
333 fn read_at(&mut self, offset: u64, bytes: &mut [u8]) -> io::Result<usize>;
334 fn write_at(&mut self, offset: u64, bytes: &[u8]) -> io::Result<usize>;
335 fn flush(&mut self) -> io::Result<()>;
336}
337
338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
339pub enum CommitState {
340 /// Publication did not occur; preparation bytes may remain. Reread before retrying.
341 NotCommitted,
342 /// Publication may have persisted. Reread and reconcile the intent before retrying.
343 Unknown,
344 /// Publication was durably acknowledged, but cleanup failed. Do not replay the edit.
345 Committed,
346}
347
348#[derive(Debug)]
349pub struct CommitError {
350 pub state: CommitState,
351 pub error: io::Error,
352}
353
354impl std::fmt::Display for CommitError {
355 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
356 write!(formatter, "{:?}: {}", self.state, self.error)
357 }
358}
359
360impl std::error::Error for CommitError {
361 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
362 Some(&self.error)
363 }
364}
365
366fn write_all(io: &mut impl CommitIo, mut offset: u64, mut bytes: &[u8]) -> io::Result<()> {
367 while !bytes.is_empty() {
368 match io.write_at(offset, bytes) {
369 Ok(0) => return Err(io::Error::from(ErrorKind::WriteZero)),
370 Ok(count) if count <= bytes.len() => {
371 offset += count as u64;
372 bytes = &bytes[count..];
373 }
374 Ok(_) => {
375 return Err(io::Error::new(
376 ErrorKind::InvalidData,
377 "Storage returned an excessive write count",
378 ));
379 }
380 Err(error) if error.kind() == ErrorKind::Interrupted => {}
381 Err(error) => return Err(error),
382 }
383 }
384 Ok(())
385}
386
387/// What a commit requires unchanged since its snapshot: the header, which every committed
388/// transaction and every placement rewrites (MS-ONESTORE 2.3.1 `guidFileVersion`), and the
389/// length appends start from. Equal stamps name the same committed image.
390#[derive(Clone, Debug, PartialEq, Eq)]
391pub struct Stamp {
392 pub header: [u8; 1024],
393 pub length: u64,
394}
395
396impl Stamp {
397 pub fn of(image: &[u8]) -> Result<Self, crate::Error> {
398 Ok(Self {
399 header: image.first_chunk().copied().ok_or(crate::Error {
400 offset: 0,
401 message: "Truncated revision-store header",
402 })?,
403 length: image.len() as u64,
404 })
405 }
406
407 /// Checks that the file `io` reads still has this stamp, reading its header and probing its
408 /// length without reading the body; `ResourceBusy` when it moved on. The caller holds
409 /// OneNote-compatible exclusion for whatever the check guards.
410 pub fn check(&self, io: &mut impl CommitIo) -> io::Result<()> {
411 let mut header = [0; 1024];
412 crate::snapshot::read_exact(
413 &mut |offset, output| io.read_at(offset, output),
414 0,
415 &mut header,
416 )?;
417 if header != self.header {
418 return Err(io::Error::new(
419 ErrorKind::ResourceBusy,
420 "The file header changed after the edit snapshot",
421 ));
422 }
423 let last = self.length.checked_sub(1).ok_or(ErrorKind::InvalidInput)?;
424 let mut tail = [0; 2];
425 let count = loop {
426 match io.read_at(last, &mut tail) {
427 Err(error) if error.kind() == ErrorKind::Interrupted => {}
428 result => break result?,
429 }
430 };
431 if count != 1 {
432 return Err(io::Error::new(
433 ErrorKind::ResourceBusy,
434 "The file length changed after the edit snapshot",
435 ));
436 }
437 Ok(())
438 }
439}
440
441/// One change to a revision store as the bytes committing it writes: `append` at the base
442/// length, `patches` inside the base's data area (list tails and the transaction log), and
443/// the header that publishes them.
444#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
445#[serde(into = "Wire", try_from = "Wire")]
446pub struct Transaction {
447 pub(crate) base: Stamp,
448 pub(crate) append: Vec<u8>,
449 pub(crate) patches: Vec<(u64, Vec<u8>)>,
450 pub(crate) header: [u8; 1024],
451}
452
453/// A transaction as serialized: headers as byte strings, since serde stops at 32-byte arrays.
454#[derive(serde::Serialize, serde::Deserialize)]
455struct Wire {
456 base: Vec<u8>,
457 length: u64,
458 append: Vec<u8>,
459 patches: Vec<(u64, Vec<u8>)>,
460 header: Vec<u8>,
461}
462
463impl From<Transaction> for Wire {
464 fn from(transaction: Transaction) -> Self {
465 Self {
466 base: transaction.base.header.to_vec(),
467 length: transaction.base.length,
468 append: transaction.append,
469 patches: transaction.patches,
470 header: transaction.header.to_vec(),
471 }
472 }
473}
474
475impl TryFrom<Wire> for Transaction {
476 type Error = &'static str;
477
478 fn try_from(wire: Wire) -> Result<Self, Self::Error> {
479 let header = |bytes: Vec<u8>| {
480 <[u8; 1024]>::try_from(bytes).map_err(|_| "A transaction header is 1024 bytes")
481 };
482 Ok(Self {
483 base: Stamp {
484 header: header(wire.base)?,
485 length: wire.length,
486 },
487 append: wire.append,
488 patches: wire.patches,
489 header: header(wire.header)?,
490 })
491 }
492}
493
494impl Transaction {
495 /// Combines a following transaction into one guarded publication, retaining every revision.
496 pub fn extend(&mut self, next: Transaction) -> Result<(), crate::Error> {
497 let length = self.base.length + self.append.len() as u64;
498 if next.base.header != self.header || next.base.length != length {
499 return Err(crate::Error {
500 offset: 0,
501 message: "Transactions are not consecutive",
502 });
503 }
504 if next.patches.iter().any(|(offset, bytes)| {
505 *offset < 1024 || offset.saturating_add(bytes.len() as u64) > length
506 }) {
507 return Err(crate::Error {
508 offset: 0,
509 message: "A patch outside the base's data",
510 });
511 }
512 for (offset, bytes) in next.patches {
513 let earlier = (self.base.length.saturating_sub(offset) as usize).min(bytes.len());
514 if earlier != 0 {
515 self.patches.push((offset, bytes[..earlier].to_vec()));
516 }
517 if earlier < bytes.len() {
518 let at = (offset + earlier as u64 - self.base.length) as usize;
519 self.append[at..at + bytes.len() - earlier].copy_from_slice(&bytes[earlier..]);
520 }
521 }
522 self.append.extend_from_slice(&next.append);
523 self.header = next.header;
524 Ok(())
525 }
526
527 /// The image this transaction applies to.
528 pub fn base(&self) -> &Stamp {
529 &self.base
530 }
531
532 /// The transaction as bytes `from_bytes` reads back, to carry it to another machine to
533 /// commit: the base header and length, the new header, the appended bytes' length and
534 /// the bytes, then each patch's offset, length and bytes. Integers are little-endian.
535 pub fn to_bytes(&self) -> Vec<u8> {
536 let mut bytes = Vec::with_capacity(2064 + self.append.len());
537 bytes.extend_from_slice(&self.base.header);
538 bytes.extend_from_slice(&self.base.length.to_le_bytes());
539 bytes.extend_from_slice(&self.header);
540 bytes.extend_from_slice(&(self.append.len() as u64).to_le_bytes());
541 bytes.extend_from_slice(&self.append);
542 for (offset, patch) in &self.patches {
543 bytes.extend_from_slice(&offset.to_le_bytes());
544 bytes.extend_from_slice(&(patch.len() as u32).to_le_bytes());
545 bytes.extend_from_slice(patch);
546 }
547 bytes
548 }
549
550 /// Reads `to_bytes`'s form, refusing a patch outside the base's data area, which no
551 /// transaction writes.
552 pub fn from_bytes(bytes: &[u8]) -> Result<Self, crate::Error> {
553 let malformed = |message| crate::Error { offset: 0, message };
554 fn take<'a>(rest: &mut &'a [u8], length: usize) -> Result<&'a [u8], crate::Error> {
555 let (taken, after) = rest.split_at_checked(length).ok_or(crate::Error {
556 offset: 0,
557 message: "A truncated transaction",
558 })?;
559 *rest = after;
560 Ok(taken)
561 }
562 let rest = &mut &bytes[..];
563 let base_header: [u8; 1024] = take(rest, 1024)?.try_into().expect("1024 bytes");
564 let length = u64::from_le_bytes(take(rest, 8)?.try_into().expect("8 bytes"));
565 let header: [u8; 1024] = take(rest, 1024)?.try_into().expect("1024 bytes");
566 let appended = u64::from_le_bytes(take(rest, 8)?.try_into().expect("8 bytes"));
567 let appended = usize::try_from(appended).map_err(|_| malformed("A huge append"))?;
568 let append = take(rest, appended)?.to_vec();
569 let mut patches = Vec::new();
570 while !rest.is_empty() {
571 let offset = u64::from_le_bytes(take(rest, 8)?.try_into().expect("8 bytes"));
572 let size = u32::from_le_bytes(take(rest, 4)?.try_into().expect("4 bytes"));
573 let patch = take(rest, size as usize)?;
574 if offset < 1024 || offset.saturating_add(u64::from(size)) > length {
575 return Err(malformed("A patch outside the base's data"));
576 }
577 patches.push((offset, patch.to_vec()));
578 }
579 Ok(Self {
580 base: Stamp {
581 header: base_header,
582 length,
583 },
584 append,
585 patches,
586 header,
587 })
588 }
589
590 /// The bytes a commit writes, by offset, in the order `apply` writes them; the header,
591 /// last, covers bytes 0..1024.
592 pub fn writes(&self) -> impl Iterator<Item = (u64, &[u8])> {
593 std::iter::once((self.base.length, self.append.as_slice()))
594 .chain(
595 self.patches
596 .iter()
597 .map(|(offset, bytes)| (*offset, bytes.as_slice())),
598 )
599 .chain(std::iter::once((0, self.header.as_slice())))
600 }
601
602 /// Writes this transaction into its base image, as a successful commit leaves the file.
603 pub fn apply(&self, image: &mut Vec<u8>) -> Result<(), crate::Error> {
604 if Stamp::of(image)? != self.base {
605 return Err(crate::Error {
606 offset: 0,
607 message: "The image is not this transaction's base",
608 });
609 }
610 image.extend_from_slice(&self.append);
611 for (offset, bytes) in &self.patches {
612 let offset = *offset as usize;
613 image[offset..offset + bytes.len()].copy_from_slice(bytes);
614 }
615 image[..1024].copy_from_slice(&self.header);
616 Ok(())
617 }
618
619 /// Publishes under caller-held OneNote-compatible exclusion, provided the file still has
620 /// the base stamp; otherwise returns ResourceBusy without writing. Appended data and
621 /// patches are flushed before the header, and the transaction count commits them
622 /// (MS-ONESTORE 2.3.3). Flush must make preceding writes durable before later ones.
623 pub fn commit(&self, io: &mut impl CommitIo) -> Result<(), CommitError> {
624 let mut state = CommitState::NotCommitted;
625 let result = (|| -> io::Result<()> {
626 self.base.check(io)?;
627 if self.append.is_empty() && self.patches.is_empty() && self.header == self.base.header
628 {
629 state = CommitState::Unknown;
630 return io.flush();
631 }
632 write_all(io, self.base.length, &self.append)?;
633 for (offset, bytes) in &self.patches {
634 write_all(io, *offset, bytes)?;
635 }
636 io.flush()?;
637 write_all(io, 100, &self.header[100..212])?;
638 write_all(io, 252, &self.header[252..])?;
639 io.flush()?;
640 state = CommitState::Unknown;
641 // At a counter carry the highest changed byte commits; the lower ones follow.
642 if let Some(highest) = (96..100).rfind(|at| self.base.header[*at] != self.header[*at]) {
643 write_all(io, highest as u64, &self.header[highest..highest + 1])?;
644 io.flush()?;
645 if highest > 96 {
646 write_all(io, 96, &self.header[96..highest])?;
647 io.flush()?;
648 }
649 }
650 // Native readers cache the version GUID without rechecking the transaction count.
651 write_all(io, 212, &self.header[212..252])?;
652 io.flush()
653 })();
654 result.map_err(|error| CommitError { state, error })
655 }
656
657 /// `commit` under the conservative filesystem adapter's whole-file exclusion.
658 #[cfg(any(unix, windows))]
659 pub fn commit_file(&self, path: impl AsRef<Path>) -> Result<(), CommitError> {
660 let mut io = FileIo::open(path, true).map_err(|error| CommitError {
661 state: CommitState::NotCommitted,
662 error,
663 })?;
664 let result = self.commit(&mut io);
665 io.finish(result)
666 }
667}
668
669#[cfg(all(test, any(unix, windows)))]
670mod tests {
671 use super::*;
672
673 /// Reads `bytes`, changing a header byte on each of the first `changes` rereads of it.
674 fn reader(bytes: &[u8], mut changes: usize) -> impl FnMut(u64, &mut [u8]) -> io::Result<usize> {
675 let mut bytes = bytes.to_vec();
676 let mut reads = 0;
677 move |offset, output| {
678 if offset == 0 {
679 reads += 1;
680 // Each pass reads the header twice: with the body, then to check it.
681 if reads % 2 == 0 && changes > 0 {
682 changes -= 1;
683 bytes[1000] ^= 1;
684 }
685 }
686 let rest = bytes.get(offset as usize..).unwrap_or_default();
687 let count = rest.len().min(output.len());
688 output[..count].copy_from_slice(&rest[..count]);
689 Ok(count)
690 }
691 }
692
693 #[test]
694 fn consecutive_transactions_publish_the_same_bytes_together() {
695 let original = vec![1; 2048];
696 let first = Transaction {
697 base: Stamp::of(&original).unwrap(),
698 append: vec![2; 32],
699 patches: vec![(1024, vec![3; 8])],
700 header: [4; 1024],
701 };
702 let mut separate = original.clone();
703 first.apply(&mut separate).unwrap();
704 let second = Transaction {
705 base: Stamp::of(&separate).unwrap(),
706 append: vec![5; 16],
707 patches: vec![(1028, vec![6; 8]), (2044, vec![7; 12])],
708 header: [8; 1024],
709 };
710 second.apply(&mut separate).unwrap();
711 let mut combined = first.clone();
712 assert!(combined.extend(first).is_err());
713 combined.extend(second).unwrap();
714 let combined = Transaction::from_bytes(&combined.to_bytes()).unwrap();
715 let mut together = original;
716 combined.apply(&mut together).unwrap();
717 assert_eq!(together, separate);
718 }
719
720 #[test]
721 fn transactions_read_back_from_bytes() {
722 let transaction = Transaction {
723 base: Stamp {
724 header: [1; 1024],
725 length: 4096,
726 },
727 append: vec![2; 300],
728 patches: vec![(1024, vec![3; 8]), (4000, vec![4; 96])],
729 header: [5; 1024],
730 };
731 let bytes = transaction.to_bytes();
732 assert_eq!(Transaction::from_bytes(&bytes).unwrap(), transaction);
733 assert!(Transaction::from_bytes(&bytes[..bytes.len() - 1]).is_err());
734 let outside = Transaction {
735 patches: vec![(4000, vec![4; 97])],
736 ..transaction.clone()
737 };
738 assert!(Transaction::from_bytes(&outside.to_bytes()).is_err());
739 let header = Transaction {
740 patches: vec![(1000, vec![4; 8])],
741 ..transaction
742 };
743 assert!(Transaction::from_bytes(&header.to_bytes()).is_err());
744 }
745
746 #[test]
747 fn a_read_torn_by_a_commit_is_read_again_then_refused() {
748 let section = crate::create_section("Torn.one", "Text", "Fixture").unwrap();
749 assert_eq!(stable(reader(&section, 0), section.len()).unwrap(), section);
750 let mut changed = section.clone();
751 changed[1000] ^= 1;
752 assert_eq!(stable(reader(&section, 1), section.len()).unwrap(), changed);
753 assert_eq!(
754 stable(reader(&section, TRIES), section.len())
755 .unwrap_err()
756 .kind(),
757 ErrorKind::WouldBlock
758 );
759 // Storage shortened ahead of the header that publishes its new length.
760 let short = &section[..section.len() - 1];
761 assert_eq!(
762 stable(reader(short, 0), section.len()).unwrap_err().kind(),
763 ErrorKind::WouldBlock
764 );
765 assert_eq!(
766 stable(reader(&section, 0), section.len() - 1)
767 .unwrap_err()
768 .kind(),
769 ErrorKind::FileTooLarge
770 );
771 // Files that are not revision stores read as they are.
772 assert_eq!(
773 stable(reader(b"unfinished", 0), 100).unwrap(),
774 b"unfinished"
775 );
776 }
777
778 /// OneNote's reads go on through Snowbound's reads and writes, and its writers wait for
779 /// Snowbound's, as its opens and coordination bytes meet Snowbound's.
780 #[cfg(windows)]
781 #[test]
782 fn windows_takes_onenote_s_opens_and_bytes() {
783 use file_guard::Lock;
784 use std::io::Read;
785 use std::os::windows::fs::OpenOptionsExt;
786 let folder = std::env::temp_dir().join(format!("onestore-locks-{}", std::process::id()));
787 std::fs::create_dir_all(&folder).unwrap();
788 let path = folder.join("Locks.one");
789 std::fs::write(&path, b"section").unwrap();
790 let onenote = |write: bool| {
791 File::options()
792 .read(true)
793 .write(write)
794 .share_mode(if write { 0x5 } else { 0x7 })
795 .open(&path)
796 .map(std::sync::Arc::new)
797 };
798 let busy = |result: io::Result<file_guard::FileGuard<std::sync::Arc<File>>>| {
799 result.is_err_and(|error| error.kind() == ErrorKind::WouldBlock)
800 };
801
802 let reading = FileIo::open(&path, false).unwrap();
803 let writer = onenote(true).unwrap();
804 let reader_byte = file_guard::try_lock(writer.clone(), Lock::Shared, 0xffff_fffb, 1);
805 assert!(reader_byte.is_ok(), "OneNote writes beside a reader");
806 drop((reader_byte, writer));
807 assert!(!busy(file_guard::try_lock(
808 onenote(false).unwrap(),
809 Lock::Exclusive,
810 0xffff_fffd,
811 1
812 )));
813 drop(reading);
814
815 let writing = FileIo::open(&path, true).unwrap();
816 let mut text = String::new();
817 (&*onenote(false).unwrap())
818 .read_to_string(&mut text)
819 .unwrap();
820 assert_eq!(text, "section", "OneNote reads through a commit");
821 assert!(onenote(true).is_err(), "a second writer can't open it");
822 let other = onenote(false).unwrap();
823 assert!(busy(file_guard::try_lock(
824 other,
825 Lock::Exclusive,
826 0xffff_fffd,
827 1
828 )));
829 drop(writing);
830 std::fs::remove_dir_all(&folder).unwrap();
831 }
832}