threads(M11): RwLock, WaitGroup, and host-testable sync — Phase 2 done

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.
This commit is contained in:
2026-07-21 00:09:48 +01:00
parent c7e9b5a4f6
commit a4e44e8f31
7 changed files with 376 additions and 22 deletions
+219 -3
View File
@@ -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();
}