Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 12 additions & 2 deletions nodedb-wal/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -63,8 +63,6 @@ getrandom_03 = { workspace = true }
getrandom_02 = { workspace = true }

[dev-dependencies]
tokio = { workspace = true }
fluxbench = { workspace = true }
tracing-subscriber = { workspace = true }
# Named again here so the diagnostics integration tests can act as the host —
# call `init` and read reports back. Same published version as the optional
Expand All @@ -73,6 +71,18 @@ faultbox = { workspace = true }
tempfile = { workspace = true }
rand = { workspace = true }

# Native-only dev-dependencies. `tokio = { workspace = true }` inherits
# `features = ["full"]` and `fluxbench` declares its own normal
# `tokio = { features = ["full"] }`; tokio rejects fs/io-std/net/process/
# rt-multi-thread/signal on wasm, and Cargo unifies features across the whole
# `cargo test` invocation, so either entry breaks every wasm32 test target in
# the build. Neither is reachable from a wasm test: this crate has no tokio call
# sites, and `fluxbench` is used only by benches/wal_throughput.rs, which
# `cargo test` does not build. Native test and bench builds keep both.
[target.'cfg(not(target_arch = "wasm32"))'.dev-dependencies]
tokio = { workspace = true }
fluxbench = { workspace = true }

[target.'cfg(target_os = "linux")'.dev-dependencies]
io-uring = { workspace = true }

Expand Down
10 changes: 10 additions & 0 deletions nodedb-wal/src/double_write/raw_io.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,16 @@ pub(crate) fn full_capacity_slice(buf: &AlignedBuf) -> &[u8] {
}

/// `pwrite`-retry helper that handles short writes.
///
/// Reports [`WalError::Unsupported`] on wasm32 rather than falling back to a
/// seek-and-write. A fallback would mirror the slot successfully and the writer
/// would then report `DwbProtection::Active`, but `recover_record` is a stub on
/// that target (slot recovery needs `pread`), so nothing could ever read the
/// slot back. The protection would be claimed and unusable. Failing here keeps
/// the writer's own degradation signal honest: the append still succeeds
/// (`DoubleWriteBuffer` failures are deliberately not fatal to the WAL), while
/// `dwb_protection()` reports `Degraded(WriteFailed)` and every record is
/// counted as unprotected.
pub(crate) fn pwrite_all(file: &File, data: &[u8], offset: u64) -> Result<()> {
#[cfg(not(target_arch = "wasm32"))]
{
Expand Down
30 changes: 27 additions & 3 deletions nodedb-wal/src/segment/atomic_io.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,16 +48,36 @@ use crate::error::{Result, WalError};
/// before the directory entry is persisted causes the file to "disappear"
/// on reboot. Calling fsync on the directory fd ensures the metadata
/// (filename, inode pointer) is on stable storage.
///
/// The exception is `wasm32-wasip1`: wasi preview1 has no directory fsync, so on
/// that target this is a documented no-op and the durability guarantee above
/// does not hold there. The alternative was failing every caller that renames
/// and fsyncs on that target, which turns a weaker guarantee into an unusable
/// one.
pub fn fsync_directory(dir: &Path) -> Result<()> {
// Crash injection: the directory entry never reaches stable storage.
// Every caller must treat this as a durability failure, not a warning.
nodedb_types::fail_point_err!("wal::fsync_directory", |detail: String| WalError::Io(
std::io::Error::other(format!("failpoint wal::fsync_directory: {detail}"))
));

let dir_file = fs::File::open(dir).map_err(WalError::Io)?;
dir_file.sync_all().map_err(WalError::Io)?;
Ok(())
#[cfg(target_arch = "wasm32")]
{
// `File::sync_all` on a directory is not implemented by wasi preview1
// and there is no weaker syscall with the same guarantee. The rename
// itself still succeeds; what is lost on this target is the assurance
// that the rename survives a host crash. The failpoint above still
// fires here, so crash injection keeps working.
let _ = dir;
Ok(())
}

#[cfg(not(target_arch = "wasm32"))]
{
let dir_file = fs::File::open(dir).map_err(WalError::Io)?;
dir_file.sync_all().map_err(WalError::Io)?;
Ok(())
}
}

fn invalid_input(detail: String) -> WalError {
Expand Down Expand Up @@ -103,6 +123,10 @@ fn tmp_name(name: &str) -> String {
/// 4. `fsync_directory(dir)` — forces the directory entry durable so the new
/// name survives power loss.
///
/// Step 4 is a no-op on `wasm32-wasip1`, which has no directory fsync: there the
/// rename is atomic with respect to the process but its durability across a host
/// crash is not assured. Steps 1-3 are unchanged on every target.
///
/// `name` must be one plain path component; anything else returns
/// `InvalidInput` before a byte is written. Both paths are built here from
/// `dir`, so the rename stays inside one filesystem directory and the single
Expand Down
87 changes: 72 additions & 15 deletions nodedb-wal/src/writer/flush.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,20 @@
// SPDX-License-Identifier: Apache-2.0

use crate::error::Result;
// Only the unix `pwrite` path constructs an error value directly; elsewhere
// failures propagate as `Result` from calls that build their own.
#[cfg(unix)]
// The unix `pwrite` arm and the wasm seek+write arm both construct an error
// value directly; elsewhere failures propagate as `Result` from calls that
// build their own.
#[cfg(any(unix, target_arch = "wasm32"))]
use crate::error::WalError;

use super::core::WalWriter;

// The write arms below cover unix and wasm32. A target in neither family would
// compile no arm at all and reinstate the silent success this file's wasi arm
// exists to prevent, so fail the build instead of writing nothing.
#[cfg(not(any(unix, target_arch = "wasm32")))]
compile_error!("nodedb-wal has no WAL write arm for this target family");

impl WalWriter {
/// Flush the aligned buffer to the file.
///
Expand Down Expand Up @@ -42,8 +49,12 @@ impl WalWriter {
self.buffer.as_slice()
};

// Use pwrite to write at the exact offset, retrying on short writes.
#[cfg(unix)]
// Write at the exact offset. The unix path uses `pwrite` and retries
// short writes. `cfg(unix)` is false on `wasm32-wasip1` — its
// `target_family` is `wasm`, not `unix` — so wasi needs its own arm:
// without one this function advanced `file_offset`, cleared the buffer
// and reported the flush with nothing written at all.
#[cfg(all(unix, not(target_arch = "wasm32")))]
{
use std::os::unix::io::AsRawFd;
let fd = self.file.as_raw_fd();
Expand All @@ -59,7 +70,8 @@ impl WalWriter {
)
};
if written < 0 {
return Err(write_error(
return Err(classify_write_error(
std::io::Error::last_os_error(),
"WAL segment append",
write_offset,
remaining.len() as u64,
Expand All @@ -79,6 +91,43 @@ impl WalWriter {
}
}

// This arm covers every wasm32 target, including
// `wasm32-unknown-unknown`, where libc defines no `pwrite` at all, and
// std's wasi positional API (`std::os::wasi::fs`) is still unstable.
// The segment is only ever appended to, so seeking to `file_offset` and
// writing is equivalent to a positional write. An error returns before
// the shared bookkeeping below, so a failed flush leaves the buffer and
// `file_offset` untouched for a retry — the same property the `pwrite`
// arm has.
#[cfg(target_arch = "wasm32")]
{
use std::io::{Seek as _, SeekFrom, Write as _};
self.file
.seek(SeekFrom::Start(self.file_offset))
.map_err(WalError::Io)?;
// Crash injection: this arm's write fails with a full device. It
// sits inside the arm, so a test that arms it proves the arm is what
// ran — and it is classified exactly as the real failure below is,
// by the same call, which is the only way to reach that classifier
// on a runtime whose writes cannot be made to fail on demand.
nodedb_types::fail_point_err!("wal::wasm_flush_write", |_detail: String| {
classify_write_error(
std::io::Error::new(std::io::ErrorKind::StorageFull, "device full"),
"WAL segment append",
self.file_offset,
data.len() as u64,
)
});
self.file.write_all(data).map_err(|err| {
classify_write_error(
err,
"WAL segment append",
self.file_offset,
data.len() as u64,
)
})?;
}

self.file_offset += data.len() as u64;
self.buffer.clear();

Expand All @@ -90,22 +139,30 @@ impl WalWriter {
}
}

/// Classify the current `errno` from a failed WAL write.
/// Classify a failed WAL write.
///
/// A full device is called out separately from generic I/O failure: it is not
/// transient, retrying cannot succeed, and the caller must stop acknowledging
/// writes rather than treat it as a passing error. `offset` and `pending` say
/// where the batch stalled and how much of it never reached the file, which is
/// what a report needs to describe the write that could not complete.
///
/// Gated to match its only call site: the `pwrite` loop is unix-only, so on
/// other targets (wasm32) this would be dead code and a `-D warnings` build
/// would reject it.
#[cfg(unix)]
fn write_error(context: &'static str, offset: u64, pending: u64) -> WalError {
let err = std::io::Error::last_os_error();
#[cfg(unix)]
if err.raw_os_error() == Some(libc::ENOSPC) {
/// Takes the error rather than reading `errno` again: the wasm arm is handed a
/// real `io::Error` by `write_all`, and re-deriving it from the thread's errno
/// would classify whatever happened to be there last.
///
/// Keyed on `ErrorKind::StorageFull` rather than on `libc::ENOSPC`: std maps the
/// full-device errno to that kind on every target that has a filesystem, and
/// `libc` defines no constants at all for `wasm32-unknown-unknown`, which this
/// crate is also compiled for.
#[cfg(any(unix, target_arch = "wasm32"))]
fn classify_write_error(
err: std::io::Error,
context: &'static str,
offset: u64,
pending: u64,
) -> WalError {
if err.kind() == std::io::ErrorKind::StorageFull {
let out_of_space = WalError::OutOfSpace { context };
crate::diag::out_of_space(&out_of_space, context, offset, pending);
return out_of_space;
Expand Down
138 changes: 138 additions & 0 deletions nodedb-wal/tests/wasi_append.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,138 @@
// SPDX-License-Identifier: Apache-2.0

//! The WAL's write path on `wasm32-wasip1`, the target the wasm test job runs.
//!
//! `WalWriter::flush_buffer` wrote the batch with a `#[cfg(unix)]` arm and
//! nothing else. `cfg(unix)` is false on this target — `target_family` is `wasm`
//! — so the write was not compiled at all, while the statements after it still
//! ran: the offset advanced, the buffer was cleared and the flush reported
//! success. The test below is the difference between "the append landed" and
//! "the writer said it landed", and it is decided by the file, not by the
//! return value.
//!
//! `fsync_directory` had the same shape of problem: it opened the directory and
//! called `sync_all`, which wasi preview1 cannot do, so every caller that
//! fsyncs a directory after a rename failed on this target.
//!
//! Notes for whoever runs this:
//!
//! - Run it with `--nocapture`. A panic aborts the process on this target, so
//! without it a failed assertion arrives as a wasm trap and the message that
//! says which byte count disagreed is lost.
//! - `wasmtime --dir=.` preopens the working directory (Cargo runs the test
//! binary from the package root), and `tempfile`'s default base is
//! `std::env::temp_dir()`, which is `unimplemented!()` in std for wasi. A
//! directory under the preopen is the only place a wasm test can write.
//! - A panicking test skips its own cleanup, so an aborted run leaves the
//! tempdir behind next to the package.
//! - The failure-path test needs the crate's own failpoints, so it runs with
//! `--features failpoints`. The other two run either way.
//!
//! wasm-only on purpose: the native append path is covered by `tests/wal_suite`.

#![cfg(target_arch = "wasm32")]

use nodedb_wal::WalWriter;
use nodedb_wal::segment::atomic_io::{atomic_write_fsync, fsync_directory};
use std::fs;

fn tempdir() -> tempfile::TempDir {
tempfile::Builder::new()
.tempdir_in(".")
.expect("the wasm test runner preopens the working directory")
}

const PAYLOAD: &[u8] = b"wasip1 append probe";

#[test]
fn an_appended_record_reaches_the_file() {
let dir = tempdir();
let path = dir.path().join("append.wal");

let mut writer = WalWriter::open_without_direct_io(&path).expect("open the segment");
writer.append(1, 0, 0, 0, PAYLOAD).expect("append");
writer.sync().expect("sync");

let counted = writer.file_offset();
let written = fs::read(&path).expect("read the segment back");

assert_eq!(
written.len() as u64,
counted,
"sync() reported success with {counted} bytes written, but the segment holds {} bytes",
written.len()
);
assert!(
written.windows(PAYLOAD.len()).any(|w| w == PAYLOAD),
"the segment does not contain the appended payload"
);
}

#[test]
fn a_checkpoint_can_fsync_its_directory() {
let dir = tempdir();

fsync_directory(dir.path()).expect("wasi preview1 has no directory fsync, not a failure");
atomic_write_fsync(dir.path(), "payload.ckpt", b"hello wasi")
.expect("atomic_write_fsync renames then fsyncs the directory");
assert_eq!(
fs::read(dir.path().join("payload.ckpt")).expect("read back"),
b"hello wasi",
"the checkpoint bytes did not survive the round trip"
);
}

/// A failed write is reported as a failure and advances nothing, so the batch is
/// retried byte-for-byte at the same offset.
///
/// The injection is `wal::wasm_flush_write`, which sits *inside* the wasi write
/// arm: with the arm removed the injection cannot fire and this test fails, so
/// it is a chokepoint for the arm rather than for the wrapper above it. The
/// error it produces goes through the same `classify_write_error` call the real
/// `write_all` failure does, which is the only way to reach that classifier on a
/// runtime whose writes cannot be made to fail on demand. It carries the error
/// kind a full device produces, so what the assertion pins is the classifier's
/// rule — not the runtime's errno translation, which std owns and which a real
/// full device would exercise.
///
/// Needs the feature, and is compiled out without it rather than weakened:
/// without `failpoints` the injection expands to nothing and these assertions
/// would describe a successful flush.
///
/// ```text
/// CARGO_TARGET_WASM32_WASIP1_RUNNER="wasmtime --dir=." \
/// cargo test -p nodedb-wal --target wasm32-wasip1 --features failpoints \
/// --test wasi_append -- --nocapture
/// ```
#[cfg(feature = "failpoints")]
#[test]
fn a_failed_flush_reports_failure_and_advances_nothing() {
use nodedb_wal::error::WalError;

let dir = tempdir();
let path = dir.path().join("failed.wal");

let mut writer = WalWriter::open_without_direct_io(&path).expect("open the segment");
writer.append(1, 0, 0, 0, PAYLOAD).expect("append");
let before = writer.file_offset();

let guard = nodedb_types::fail_point::FailGuard::fail("wal::wasm_flush_write", "device full");
let result = writer.sync();
drop(guard);

let err = result.expect_err("an armed write failpoint must not be reported as a flush");
assert!(
matches!(err, WalError::OutOfSpace { .. }),
"a full device must be classified as OutOfSpace on this target too, got {err:?}"
);
assert_eq!(
writer.file_offset(),
before,
"a failed flush advanced the offset, so a retry would skip the batch"
);
assert_eq!(
fs::metadata(&path).map(|m| m.len()).unwrap_or(0),
0,
"a failed flush left bytes in the segment"
);
}
Loading