From a4e44e8f316072be592307d0e2a6692ad12b9b08 Mon Sep 17 00:00:00 2001 From: Daniel Samson Date: Tue, 21 Jul 2026 00:09:48 +0100 Subject: [PATCH] =?UTF-8?q?threads(M11):=20RwLock,=20WaitGroup,=20and=20ho?= =?UTF-8?q?st-testable=20sync=20=E2=80=94=20Phase=202=20done?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit runtime.Thread.RwLock (reader-preferring, lock/tryLock/unlock + lockShared/tryLockShared/unlockShared) and WaitGroup (start/finish/wait), both on the existing Mutex/Condition. A compile-time Futex seam gated on builtin.os.tag: the futex syscalls on danos, a spin+yield mock off-target (Zig 0.16 has no std.Thread.Futex; wake is a no-op since the state machines re-check). thread.zig is wired into zig build test, so Mutex/RwLock/WaitGroup run as host unit tests with real std.Thread threads (test blocks compile only under test, so std.Thread there is fine on freestanding). thread-rwlock QEMU case: 2 writers set both halves of a value under the exclusive lock while 3 readers check they match under the shared lock; zero half-write observations across ~150k reads. Marks Phase 2 (M7-M11) built. threading.md/threading-plan.md status updated. Gate: host zig build test covers the sync primitives; thread-rwlock PASS (3x); full Done gate 26/26 (whole thread-* suite + guardrail); build clean. --- build.zig | 16 ++ docs/threading-plan.md | 34 +-- docs/threading.md | 13 +- library/runtime/thread.zig | 222 +++++++++++++++++++- system/kernel/tests.zig | 44 ++++ system/services/thread-test/thread-test.zig | 61 ++++++ test/qemu_test.py | 8 + 7 files changed, 376 insertions(+), 22 deletions(-) diff --git a/build.zig b/build.zig index b6e6dea..ae18855 100644 --- a/build.zig +++ b/build.zig @@ -926,6 +926,22 @@ pub fn build(b: *std.Build) void { }); test_step.dependOn(&b.addRunArtifact(time_tests).step); + // runtime.Thread's lock/condvar state machines (Mutex/Condition/RwLock/WaitGroup). Its + // Futex seam falls back to std.Thread.Futex off the danos target, so the tests exercise + // them with real host threads (docs/threading-plan.md M11). Like time.zig it pulls in + // system.zig (syscall wrappers), which needs the `abi` module. + const thread_tests = b.addTest(.{ + .root_module = b.createModule(.{ + .root_source_file = b.path("library/runtime/thread.zig"), + .target = target, + .optimize = optimize, + .imports = &.{ + .{ .name = "abi", .module = abi_module }, + }, + }), + }); + test_step.dependOn(&b.addRunArtifact(thread_tests).step); + // Convenience: `zig build gen-xkeyboard-config` regenerates the layout tables from the // vendored data (offline). `fetch` (the network step) stays a manual script run. const gen_xkb = b.addSystemCommand(&.{ "python3", "tools/make-xkeyboard-config.py", "generate" }); diff --git a/docs/threading-plan.md b/docs/threading-plan.md index 0147db9..16b3599 100644 --- a/docs/threading-plan.md +++ b/docs/threading-plan.md @@ -279,8 +279,11 @@ green. cross-core parallelism, futex, and `Mutex`/`Condition`/`Semaphore`, all over a private thread ABI behind the runtime. -**Phase 2 (M7–M11): planned below** — hardening the deferred parts so threads are safe -for real workloads and reclaimed like everything else danos owns. +**Phase 2 (M7–M11): built.** Thread-safe allocation (M7), a task reaper that reclaims dead +tasks' kernel stacks (M8), endpoint-free `thread_join` (M9), per-thread `fs.base` (M10), +and `RwLock`/`WaitGroup` + host-testable sync (M11). Two things stay deferred by design +(no consumer): the Zig `threadlocal` *compiler* layer (M10) and detached-thread user-stack +reclaim (M9) — both noted in place. --- @@ -430,19 +433,24 @@ restore touches every context switch); `zig build`/`zig build test` clean. **Gate:** `thread-tls` passes; full `thread-*` suite + guardrail green. -### M11 — `RwLock`, `WaitGroup`, and host-testable sync +### M11 — `RwLock`, `WaitGroup`, and host-testable sync ✅ -- [ ] `runtime.Thread.RwLock` and `WaitGroup` on the existing `Futex`/`Mutex`/ - `Condition`. -- [ ] A compile-time `Futex` seam: syscalls on the danos target, a host-backed impl under - `zig build test`, so the `Mutex`/`Condition`/`RwLock` state machines run as host - unit tests (fast iteration, no QEMU). -- [ ] `-Dtest-case=thread-rwlock` (`smp: 4`): many readers + writers over an `RwLock` keep - an invariant (a reader never observes a half-written value); host tests cover the - lock transitions. +- [x] `runtime.Thread.RwLock` (reader-preferring: `>0` readers / `-1` writer / `0` free, + with `lock`/`tryLock`/`unlock` + `lockShared`/`tryLockShared`/`unlockShared`) and + `WaitGroup` (`start`/`finish`/`wait`), both on the existing `Mutex`/`Condition`. +- [x] A compile-time `Futex` seam gated on `builtin.os.tag == .freestanding`: the futex + syscalls on danos, a spin+yield mock off-target (Zig 0.16 has no `std.Thread.Futex`; + `wake` is a no-op since the state machines re-check). `thread.zig` is wired into + `zig build test`, so `Mutex`/`RwLock`/`WaitGroup` run as **host unit tests** with real + `std.Thread` threads (`test` blocks only compile under test). +- [x] `-Dtest-case=thread-rwlock` (`smp: 4`): 2 writers set both halves of a value under + the exclusive lock while 3 readers check the halves match under the shared lock — + zero half-write observations across ~150k reads. Host tests cover the Mutex, + RwLock, and WaitGroup state machines. -**Gate:** host `zig build test` covers the sync primitives; `thread-rwlock` passes; -guardrail green. +**Gate (met):** `zig build test` covers the sync primitives (host threads); `thread-rwlock` +passes (3×); full Done gate **26/26** (whole `thread-*` suite + guardrail); `zig build` +clean. --- diff --git a/docs/threading.md b/docs/threading.md index abcd350..72bfd70 100644 --- a/docs/threading.md +++ b/docs/threading.md @@ -2,12 +2,13 @@ A note on danos **threads** — several tasks sharing one address space — provided by a `runtime.Thread` type that mirrors the shape of Zig's `std.Thread` while keeping every -kernel entry behind the [runtime](../library/runtime). **Built** (M1–M6, see -[threading-plan.md](threading-plan.md)): `spawn`/`join`/`detach`, cross-core -parallelism, a futex (`futex_wait`/`futex_wake`), and a futex-backed -`Mutex`/`Condition`/`Semaphore`, plus `getCurrentId`/`currentCore`. Deferred by design -(no consumer yet): per-thread `threadlocal` TLS, `RwLock`/`WaitGroup`, and migrating -`join` to a futex completion word — see the plan's M5/M6 notes. The analysis is against +kernel entry behind the [runtime](../library/runtime). **Built** (M1–M11, see +[threading-plan.md](threading-plan.md)): `spawn`/`join`/`detach`, cross-core parallelism, +a futex, `Mutex`/`Condition`/`Semaphore`/`RwLock`/`WaitGroup`, `getCurrentId`/`currentCore`, +per-thread `fs.base` TLS, thread-safe allocation, and a task reaper that reclaims dead +tasks' kernel stacks. Deferred by design (no consumer yet): the Zig `threadlocal` +*compiler* layer (per-thread `fs.base` is in place, so it's runtime+linker work on top) and +detached-thread user-stack reclaim — see the plan's M9/M10 notes. The analysis is against **Zig 0.16** (the pinned toolchain); `std.Thread`'s internals move between releases, so treat upstream shapes as "0.16.x." diff --git a/library/runtime/thread.zig b/library/runtime/thread.zig index a4ec328..d50cbfd 100644 --- a/library/runtime/thread.zig +++ b/library/runtime/thread.zig @@ -11,10 +11,17 @@ //! A binary must be built multi-threaded (`addThreadedUserBinary`) before it may spawn. const std = @import("std"); +const builtin = @import("builtin"); const abi = @import("abi"); const sc = @import("system-call.zig"); const system = @import("system.zig"); +/// True in a real danos binary; false when this module is compiled for host unit tests. +/// The `Futex` seam and the test blocks below branch on it so the lock/condvar state +/// machines can be exercised on the host against `std.Thread.Futex` (docs/threading-plan.md +/// M11), while the danos build uses the futex syscalls. +const on_danos = builtin.os.tag == .freestanding; + /// A thread stack, if the caller does not override it. 64 KiB of mmap'd, zeroed pages. pub const default_stack_size: usize = 64 * 1024; @@ -121,17 +128,37 @@ pub const Thread = struct { /// the value already differs (safe against spurious returns, as in std): the /// caller re-checks its condition in a loop. pub fn wait(ptr: *const std.atomic.Value(u32), expect: u32) void { - _ = futexWait(@intFromPtr(ptr), expect, 0); + if (comptime on_danos) { + _ = futexWait(@intFromPtr(ptr), expect, 0); + } else { + // Host unit-test mock: spin+yield until the value changes (`wake` is a + // no-op — the callers re-check their condition in a loop anyway). Correct, + // if busy; fine for the state-machine tests. + while (ptr.load(.acquire) == expect) std.Thread.yield() catch {}; + } } /// As `wait`, but returns `error.Timeout` if `timeout_ns` elapses first. pub fn timedWait(ptr: *const std.atomic.Value(u32), expect: u32, timeout_ns: u64) error{Timeout}!void { - if (futexWait(@intFromPtr(ptr), expect, timeout_ns) == abi.futex_timed_out) return error.Timeout; + if (comptime on_danos) { + if (futexWait(@intFromPtr(ptr), expect, timeout_ns) == abi.futex_timed_out) return error.Timeout; + } else { + var spins: u64 = 0; + const limit = timeout_ns / 1000 + 1; + while (ptr.load(.acquire) == expect) : (spins += 1) { + if (spins >= limit) return error.Timeout; + std.Thread.yield() catch {}; + } + } } /// Wake up to `max_waiters` threads blocked on `ptr`. pub fn wake(ptr: *const std.atomic.Value(u32), max_waiters: u32) void { - _ = futexWake(@intFromPtr(ptr), max_waiters); + if (comptime on_danos) { + _ = futexWake(@intFromPtr(ptr), max_waiters); + } else { + // host mock: spin-waiters re-check their condition, so no wake is needed. + } } }; @@ -232,6 +259,96 @@ pub const Thread = struct { s.cond.signal(); } }; + + /// A reader/writer lock, `std.Thread.RwLock`-shaped: many concurrent readers OR one + /// exclusive writer. Reader-preferring (a steady stream of readers can delay a writer), + /// built on `Mutex` + `Condition` over a signed state: `>0` = that many readers hold + /// it, `-1` = a writer holds it, `0` = free. + pub const RwLock = struct { + mutex: Mutex = .{}, + cond: Condition = .{}, + state: i64 = 0, + + /// Acquire shared (read) access, blocking while a writer holds the lock. + pub fn lockShared(rw: *RwLock) void { + rw.mutex.lock(); + defer rw.mutex.unlock(); + while (rw.state < 0) rw.cond.wait(&rw.mutex); + rw.state += 1; + } + + /// Try to acquire shared access without blocking. + pub fn tryLockShared(rw: *RwLock) bool { + rw.mutex.lock(); + defer rw.mutex.unlock(); + if (rw.state < 0) return false; + rw.state += 1; + return true; + } + + /// Release shared access; wake a waiting writer once the last reader leaves. + pub fn unlockShared(rw: *RwLock) void { + rw.mutex.lock(); + defer rw.mutex.unlock(); + rw.state -= 1; + if (rw.state == 0) rw.cond.broadcast(); + } + + /// Acquire exclusive (write) access, blocking until no readers or writer remain. + pub fn lock(rw: *RwLock) void { + rw.mutex.lock(); + defer rw.mutex.unlock(); + while (rw.state != 0) rw.cond.wait(&rw.mutex); + rw.state = -1; + } + + /// Try to acquire exclusive access without blocking. + pub fn tryLock(rw: *RwLock) bool { + rw.mutex.lock(); + defer rw.mutex.unlock(); + if (rw.state != 0) return false; + rw.state = -1; + return true; + } + + /// Release exclusive access; wake all waiters (they re-check their condition). + pub fn unlock(rw: *RwLock) void { + rw.mutex.lock(); + defer rw.mutex.unlock(); + rw.state = 0; + rw.cond.broadcast(); + } + }; + + /// A `std.Thread.WaitGroup`-shaped counter: `start` before spawning work, `finish` as + /// each unit completes, `wait` blocks until the count returns to zero. + pub const WaitGroup = struct { + mutex: Mutex = .{}, + cond: Condition = .{}, + counter: usize = 0, + + /// Register one pending unit of work. + pub fn start(wg: *WaitGroup) void { + wg.mutex.lock(); + defer wg.mutex.unlock(); + wg.counter += 1; + } + + /// Mark one unit done; wake waiters if that was the last. + pub fn finish(wg: *WaitGroup) void { + wg.mutex.lock(); + defer wg.mutex.unlock(); + wg.counter -= 1; + if (wg.counter == 0) wg.cond.broadcast(); + } + + /// Block until every started unit has finished. + pub fn wait(wg: *WaitGroup) void { + wg.mutex.lock(); + defer wg.mutex.unlock(); + while (wg.counter != 0) wg.cond.wait(&wg.mutex); + } + }; }; /// thread_spawn(entry, stack_top, arg, exit_endpoint) -> tid, or a wrapped error. @@ -266,3 +383,102 @@ fn futexWait(addr: usize, expect: u32, timeout_ns: u64) usize { fn futexWake(addr: usize, count: u32) usize { return sc.systemCall2(.futex_wake, addr, count); } + +// --- host unit tests (docs/threading-plan.md M11) --------------------------- +// +// These run under `zig build test` on the host: the `Futex` seam above uses +// `std.Thread.Futex` off-danos, so the lock/condvar state machines can be exercised by +// real host threads. They are never compiled into a danos binary (test blocks only build +// under test), so their `std.Thread` use is fine even though `std.Thread` is unavailable +// on the freestanding target. + +test "Mutex serialises concurrent increments across host threads" { + var m: Thread.Mutex = .{}; + var counter: u64 = 0; + const workers = 8; + const per = 20_000; + const Ctx = struct { + m: *Thread.Mutex, + c: *u64, + fn run(ctx: @This()) void { + var i: usize = 0; + while (i < per) : (i += 1) { + ctx.m.lock(); + ctx.c.* += 1; + ctx.m.unlock(); + } + } + }; + var handles: [workers]std.Thread = undefined; + for (&handles) |*h| h.* = try std.Thread.spawn(.{}, Ctx.run, .{Ctx{ .m = &m, .c = &counter }}); + for (handles) |h| h.join(); + try std.testing.expectEqual(@as(u64, workers * per), counter); +} + +test "RwLock never lets a reader observe a half-written pair" { + var rw: Thread.RwLock = .{}; + var a: u64 = 0; + var b: u64 = 0; // invariant while a lock is held: a == b + var stop = std.atomic.Value(bool).init(false); + var ok = std.atomic.Value(bool).init(true); + + const Writer = struct { + rw: *Thread.RwLock, + a: *u64, + b: *u64, + stop: *std.atomic.Value(bool), + fn run(w: @This()) void { + var v: u64 = 1; + while (!w.stop.load(.acquire)) : (v +%= 1) { + w.rw.lock(); + w.a.* = v; // update both halves under the exclusive lock... + w.b.* = v; + w.rw.unlock(); + } + } + }; + const Reader = struct { + rw: *Thread.RwLock, + a: *u64, + b: *u64, + ok: *std.atomic.Value(bool), + fn run(r: @This()) void { + var i: usize = 0; + while (i < 200_000) : (i += 1) { + r.rw.lockShared(); + if (r.a.* != r.b.*) r.ok.store(false, .release); // ...so a reader must never see them differ + r.rw.unlockShared(); + } + } + }; + + var writers: [2]std.Thread = undefined; + for (&writers) |*w| w.* = try std.Thread.spawn(.{}, Writer.run, .{Writer{ .rw = &rw, .a = &a, .b = &b, .stop = &stop }}); + var readers: [4]std.Thread = undefined; + for (&readers) |*rd| rd.* = try std.Thread.spawn(.{}, Reader.run, .{Reader{ .rw = &rw, .a = &a, .b = &b, .ok = &ok }}); + for (readers) |rd| rd.join(); + stop.store(true, .release); + for (writers) |w| w.join(); + try std.testing.expect(ok.load(.acquire)); +} + +test "WaitGroup blocks until every started unit finishes" { + var wg: Thread.WaitGroup = .{}; + var done = std.atomic.Value(u32).init(0); + const n = 6; + const Ctx = struct { + wg: *Thread.WaitGroup, + done: *std.atomic.Value(u32), + fn run(c: @This()) void { + _ = c.done.fetchAdd(1, .monotonic); + c.wg.finish(); + } + }; + var i: usize = 0; + while (i < n) : (i += 1) wg.start(); + var handles: [n]std.Thread = undefined; + for (&handles) |*h| h.* = try std.Thread.spawn(.{}, Ctx.run, .{Ctx{ .wg = &wg, .done = &done }}); + wg.wait(); // must not return until all n finished + try std.testing.expectEqual(@as(u32, n), done.load(.acquire)); + for (handles) |h| h.join(); +} diff --git a/system/kernel/tests.zig b/system/kernel/tests.zig index cd2f72d..9116ad5 100644 --- a/system/kernel/tests.zig +++ b/system/kernel/tests.zig @@ -157,6 +157,8 @@ pub fn run(case: []const u8, boot_information: *const BootInformation) void { taskReapTest(boot_information); } else if (eql(case, "thread-tls")) { threadTlsTest(boot_information); + } else if (eql(case, "thread-rwlock")) { + threadRwlockTest(boot_information); } else if (eql(case, "args")) { argsTest(boot_information); } else if (eql(case, "init")) { @@ -1785,6 +1787,48 @@ fn threadTlsTest(boot_information: *const BootInformation) void { result(); } +/// RwLock (docs/threading-plan.md M11): `thread-test` in rwlock mode runs writers that set +/// two halves of a value under the exclusive lock and readers that check the halves match +/// under the shared lock. If the reader/writer lock were wrong, a reader would observe a +/// half-written value; zero violations across many reads → the lock holds. +fn threadRwlockTest(boot_information: *const BootInformation) void { + log("DANOS-TEST-BEGIN: thread-rwlock\n", .{}); + if (boot_information.initial_ramdisk_len == 0) { + check("bootloader handed over an initial_ramdisk", false); + result(); + return; + } + const image = @as([*]const u8, @ptrFromInt(boot_handoff.physicalToVirtual(boot_information.initial_ramdisk_base)))[0..boot_information.initial_ramdisk_len]; + const rd = initial_ramdisk.Reader.init(image) orelse { + check("initial_ramdisk image is valid", false); + result(); + return; + }; + + var started = false; + var i: u32 = 0; + while (i < rd.count) : (i += 1) { + const item = rd.entry(i) orelse continue; + if (!eql(item.name, "thread-test")) continue; + started = if (process.spawnProcess(item.blob, 4, &.{ "thread-test", "rwlock" })) true else |_| false; + break; + } + check("thread-test (rwlock mode) spawned", started); + + const ok_marker = "thread-rwlock: ok"; + const fail_marker = "thread-rwlock: FAIL"; + scheduler.setPriority(1); + const deadline = architecture.millis() + 20000; + while (architecture.millis() < deadline) { + if (bufferHas(ok_marker) or bufferHas(fail_marker)) break; + scheduler.yield(); + } + scheduler.setPriority(4); + + check("readers/writers over an RwLock never observed a half-written value", bufferHas(ok_marker) and !bufferHas(fail_marker)); + result(); +} + /// The task reaper (docs/threading-plan.md M8): a dead task's kernel stack used to be /// leaked ("no reaper yet"). Spawn and kill many ring-3 processes and confirm the total /// kernel-stack bytes return to baseline — every stack reclaimed, no leak. (Threads exit diff --git a/system/services/thread-test/thread-test.zig b/system/services/thread-test/thread-test.zig index 8d857c0..8587dd7 100644 --- a/system/services/thread-test/thread-test.zig +++ b/system/services/thread-test/thread-test.zig @@ -410,6 +410,65 @@ fn runTlsMode() void { } } +// --- M11: rwlock mode (readers/writers over an RwLock) ---------------------- + +const RwLock = runtime.Thread.RwLock; + +var rwlock = RwLock{}; +var rw_a: u64 = 0; +var rw_b: u64 = 0; // invariant while any lock is held: rw_a == rw_b +var rw_stop = std.atomic.Value(u32).init(0); +var rw_violations = std.atomic.Value(u32).init(0); +var rw_reads = std.atomic.Value(u64).init(0); + +fn rwWriter() void { + var v: u64 = 1; + while (rw_stop.load(.acquire) == 0) : (v +%= 1) { + rwlock.lock(); // exclusive: no reader may observe the gap between the two writes + rw_a = v; + rw_b = v; + rwlock.unlock(); + } +} + +fn rwReader() void { + const reads: u64 = 50_000; + var i: u64 = 0; + while (i < reads) : (i += 1) { + rwlock.lockShared(); + if (rw_a != rw_b) _ = rw_violations.fetchAdd(1, .monotonic); // saw a half-write! + rwlock.unlockShared(); + } + _ = rw_reads.fetchAdd(reads, .monotonic); +} + +fn runRwlockMode() void { + write("thread-rwlock: starting\n"); + var writers: [2]runtime.Thread = undefined; + var readers: [3]runtime.Thread = undefined; + for (&writers) |*w| { + w.* = runtime.Thread.spawn(.{}, rwWriter, .{}) catch { + write("thread-rwlock: FAIL spawn\n"); + return; + }; + } + for (&readers) |*r| { + r.* = runtime.Thread.spawn(.{}, rwReader, .{}) catch { + write("thread-rwlock: FAIL spawn\n"); + return; + }; + } + for (readers) |r| r.join(); + rw_stop.store(1, .release); // readers done → stop the writers + for (writers) |w| w.join(); + + if (rw_violations.load(.acquire) == 0 and rw_reads.load(.acquire) > 0) { + write("thread-rwlock: ok\n"); // the M11 verdict marker + } else { + write("thread-rwlock: FAIL reader observed a half-written value\n"); + } +} + pub fn main(init: runtime.process.Init) void { const mode = init.arguments.get(1) orelse "spawn"; if (std.mem.eql(u8, mode, "join")) { @@ -424,6 +483,8 @@ pub fn main(init: runtime.process.Init) void { runAllocMode(); } else if (std.mem.eql(u8, mode, "tls")) { runTlsMode(); + } else if (std.mem.eql(u8, mode, "rwlock")) { + runRwlockMode(); } else { runSpawnMode(); } diff --git a/test/qemu_test.py b/test/qemu_test.py index d613fc4..90d7ad4 100644 --- a/test/qemu_test.py +++ b/test/qemu_test.py @@ -363,6 +363,14 @@ CASES = [ "timeout": 60, "expect": r"DANOS-TEST-RESULT: PASS", "fail": r"DANOS-TEST-RESULT: FAIL"}, + + # docs/threading-plan.md M11: RwLock — readers/writers across cores; a reader never + # observes a half-written value (writers hold it exclusively). + {"name": "thread-rwlock", + "smp": 4, + "timeout": 60, + "expect": r"DANOS-TEST-RESULT: PASS", + "fail": r"DANOS-TEST-RESULT: FAIL"}, # Process arguments: argv arrives on the SysV entry stack (argv[0] = the spawned # name, argv[1..] = the system_spawn argument blob) and echoes back intact. {"name": "args",