| 1 | use std::io::{self, ErrorKind}; |
| 2 | |
| 3 | #[cfg(any(unix, windows))] |
| 4 | use std::{ |
| 5 | fs::File, |
| 6 | path::Path, |
| 7 | sync::{Mutex, MutexGuard}, |
| 8 | }; |
| 9 | |
| 10 | #[cfg(unix)] |
| 11 | use std::os::unix::fs::FileExt; |
| 12 | #[cfg(windows)] |
| 13 | use std::os::windows::fs::FileExt; |
| 14 | |
| 15 | #[cfg(any(unix, windows))] |
| 16 | struct 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))] |
| 29 | impl 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)] |
| 113 | impl 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)] |
| 125 | fn 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`. |
| 148 | pub 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))] |
| 164 | pub 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))] |
| 176 | pub 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))] |
| 183 | pub 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))] |
| 194 | const 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))] |
| 199 | fn 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))] |
| 235 | impl 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))] |
| 257 | pub 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))] |
| 270 | pub 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. |
| 309 | pub 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. |
| 332 | pub 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)] |
| 339 | pub 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)] |
| 349 | pub struct CommitError { |
| 350 | pub state: CommitState, |
| 351 | pub error: io::Error, |
| 352 | } |
| 353 | |
| 354 | impl 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 | |
| 360 | impl std::error::Error for CommitError { |
| 361 | fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { |
| 362 | Some(&self.error) |
| 363 | } |
| 364 | } |
| 365 | |
| 366 | fn 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)] |
| 391 | pub struct Stamp { |
| 392 | pub header: [u8; 1024], |
| 393 | pub length: u64, |
| 394 | } |
| 395 | |
| 396 | impl 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")] |
| 446 | pub 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)] |
| 455 | struct Wire { |
| 456 | base: Vec<u8>, |
| 457 | length: u64, |
| 458 | append: Vec<u8>, |
| 459 | patches: Vec<(u64, Vec<u8>)>, |
| 460 | header: Vec<u8>, |
| 461 | } |
| 462 | |
| 463 | impl 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 | |
| 475 | impl 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 | |
| 494 | impl 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)))] |
| 670 | mod 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 | } |