threads(M3): join, detach, and cross-core parallelism
thread_spawn takes a 4th arg, an exit-endpoint handle: spawnThreadSupervised resolves and refcounts it under the spawn lock (like spawnProcessSupervised), so a thread's death posts a child-exit notification carrying its tid. runtime Thread.join blocks in replyWait on that (private) endpoint for its tid, then munmaps the stack; detach relinquishes the join (stack reclaimed at process exit, for now). New current_core=39 syscall + Thread.currentCore() lets a worker observe which core it ran on. The closure now lives at the top of the thread's own (private) stack instead of the heap, so spawn/join never touch the not-yet-thread-safe runtime heap. thread-test gains a join mode: 4 workers x 100k atomic increments, joined, with counter == N*K and >1 core stamped (real parallelism), plus a detached worker. Gate thread-join PASS (4x, non-flaky); 17 guardrail/M1/M2 cases green; build + host tests clean.
This commit is contained in:
+2
-1
@@ -63,8 +63,9 @@ pub const SystemCall = enum(u64) {
|
||||
shm_create = 34, // shm_create(len) -> vaddr (rax), handle (rdx): a shareable, zeroed, cacheable RAM region mapped into this AS; the handle is a capability passed to another process as an ipc_call send_cap (docs/display-v2.md)
|
||||
shm_map = 35, // shm_map(cap) -> vaddr: map the shared region named by a received capability into this AS (the same physical pages the creator sees)
|
||||
shm_physical = 36, // shm_physical(cap) -> paddr: the guest-physical base of a shared region held by capability, so a driver can program it into a device (e.g. virtio-gpu attach_backing); the pages are contiguous (docs/display-v2.md)
|
||||
thread_spawn = 37, // thread_spawn(entry, stack_top, arg) -> tid: start a task that shares the caller's address space at `entry` on `stack_top`, with `arg` in rdi (docs/threading.md)
|
||||
thread_spawn = 37, // thread_spawn(entry, stack_top, arg, exit_endpoint) -> tid: start a task sharing the caller's address space at `entry` on `stack_top`, `arg` in rdi; exit_endpoint (a handle, or no_cap) is notified when it ends — how join waits (docs/threading.md)
|
||||
thread_exit = 38, // thread_exit(): end the calling thread, dropping one reference to its address space (destroyed on the last)
|
||||
current_core = 39, // current_core() -> index: the dense 0-based index of the core the caller is running on (for parallelism/affinity introspection)
|
||||
_,
|
||||
};
|
||||
|
||||
|
||||
@@ -228,6 +228,7 @@ fn system_call(state: *architecture.CpuState) void {
|
||||
.shm_map => systemShmMap(state),
|
||||
.shm_physical => systemShmPhysical(state),
|
||||
.thread_spawn => systemThreadSpawn(state),
|
||||
.current_core => systemCurrentCore(state),
|
||||
.thread_exit => {
|
||||
// A thread ends like a process exit(0), but only this task: its
|
||||
// resources are released and its address-space reference dropped (the
|
||||
@@ -661,14 +662,36 @@ fn systemThreadSpawn(state: *architecture.CpuState) void {
|
||||
const entry = architecture.systemCallArg(state, 0);
|
||||
const stack_top = architecture.systemCallArg(state, 1);
|
||||
const arg = architecture.systemCallArg(state, 2);
|
||||
const exit_handle = architecture.systemCallArg(state, 3);
|
||||
const t = scheduler.current();
|
||||
if (t.aspace == 0) return fail(state); // kernel tasks own no address space to share
|
||||
if (entry == 0 or entry >= user_half_end) return fail(state);
|
||||
if (stack_top == 0 or stack_top > user_half_end) return fail(state);
|
||||
const tid = scheduler.spawnThread(t.aspace, entry, stack_top, arg, t.priority, t.id) orelse return fail(state);
|
||||
// The endpoint the thread notifies on exit (how join waits), or none.
|
||||
const exit_endpoint: ?*ipc.Endpoint = if (exit_handle == abi.no_cap)
|
||||
null
|
||||
else
|
||||
ipc.resolveHandle(t, exit_handle) orelse return failErr(state, ipc.EBADF);
|
||||
const tid = spawnThreadSupervised(t.aspace, entry, stack_top, arg, t.priority, t.id, exit_endpoint) orelse return fail(state);
|
||||
architecture.setSystemCallResult(state, tid);
|
||||
}
|
||||
|
||||
/// Spawn a thread sharing `aspace`, taking the exit-endpoint reference under the **same**
|
||||
/// lock as the spawn (as `spawnProcessSupervised` does), so the thread cannot die before
|
||||
/// its reference exists. Returns the new thread id, or null on resource exhaustion.
|
||||
fn spawnThreadSupervised(aspace: u64, entry: u64, stack_top: u64, arg: u64, priority: scheduler.Priority, supervisor: u32, exit_endpoint: ?*ipc.Endpoint) ?u32 {
|
||||
const flags = sync.enter();
|
||||
defer sync.leave(flags);
|
||||
const tid = scheduler.spawnUserLocked(aspace, entry, stack_top, arg, priority, "thread", supervisor, if (exit_endpoint) |e| @ptrCast(e) else null) orelse return null;
|
||||
if (exit_endpoint) |endpoint| endpoint.refcount += 1; // the thread holds it birth-to-death
|
||||
return tid;
|
||||
}
|
||||
|
||||
/// current_core() -> index: the dense 0-based index of the core the caller runs on.
|
||||
fn systemCurrentCore(state: *architecture.CpuState) void {
|
||||
architecture.setSystemCallResult(state, scheduler.currentCpuIndex());
|
||||
}
|
||||
|
||||
/// process_enumerate(buffer, maximum) -> total: snapshot the task table into the
|
||||
/// caller's buffer (up to `maximum` `abi.ProcessDescriptor` entries), returning
|
||||
/// the total live-task count — the exact shape of `device_enumerate`, so a `ps`
|
||||
|
||||
@@ -424,17 +424,6 @@ pub fn spawnUserLocked(aspace: u64, entry: u64, user_sp: u64, user_arg: u64, pri
|
||||
return t.id;
|
||||
}
|
||||
|
||||
/// Spawn a **thread**: a user task that shares an *existing* address space `aspace`
|
||||
/// (docs/threading.md), starting at `entry` on `user_sp` with `arg` delivered in its
|
||||
/// rdi. Takes a reference to `aspace` (destroyed only when the last thread on it
|
||||
/// exits). Acquires the kernel lock itself. `supervisor` is the spawning process.
|
||||
/// Returns the new thread's id, or null if the task table is full / out of memory.
|
||||
pub fn spawnThread(aspace: u64, entry: u64, user_sp: u64, arg: u64, priority: Priority, supervisor: u32) ?u32 {
|
||||
const flags = sync.enter();
|
||||
defer sync.leave(flags);
|
||||
return spawnUserLocked(aspace, entry, user_sp, arg, priority, "thread", supervisor, null);
|
||||
}
|
||||
|
||||
/// The first thing a fresh user task runs (in ring 0, via task_trampoline). It
|
||||
/// drops to ring 3 at the task's recorded entry/stack. Reading them from the
|
||||
/// Task avoids smuggling values through callee-saved registers across the
|
||||
|
||||
@@ -143,6 +143,8 @@ pub fn run(case: []const u8, boot_information: *const BootInformation) void {
|
||||
aspaceRefcountTest(boot_information);
|
||||
} else if (eql(case, "thread-spawn")) {
|
||||
threadSpawnTest(boot_information);
|
||||
} else if (eql(case, "thread-join")) {
|
||||
threadJoinTest(boot_information);
|
||||
} else if (eql(case, "args")) {
|
||||
argsTest(boot_information);
|
||||
} else if (eql(case, "init")) {
|
||||
@@ -1504,6 +1506,51 @@ fn threadSpawnTest(boot_information: *const BootInformation) void {
|
||||
result();
|
||||
}
|
||||
|
||||
/// Thread join + parallelism (docs/threading-plan.md M3): `thread-test` in join mode
|
||||
/// spawns N workers that each do K atomic increments on a shared counter and stamp the
|
||||
/// core they ran on; it `join`s all N and asserts the total is exactly N*K (every worker
|
||||
/// ran, join waited for each) and that >1 core was used (genuine parallelism), then a
|
||||
/// detached worker proves `detach`. Its single verdict marker is the case result.
|
||||
fn threadJoinTest(boot_information: *const BootInformation) void {
|
||||
log("DANOS-TEST-BEGIN: thread-join\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;
|
||||
};
|
||||
|
||||
// Spawn thread-test in join mode (argv selects the mode).
|
||||
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", "join" })) true else |_| false;
|
||||
break;
|
||||
}
|
||||
check("thread-test (join mode) spawned", started);
|
||||
|
||||
const ok_marker = "thread-test: join ok";
|
||||
const fail_marker = "thread-test: FAIL";
|
||||
scheduler.setPriority(1);
|
||||
const deadline = architecture.millis() + 15000;
|
||||
while (architecture.millis() < deadline) {
|
||||
if (bufferHas(ok_marker) or bufferHas(fail_marker)) break;
|
||||
scheduler.yield();
|
||||
}
|
||||
scheduler.setPriority(4);
|
||||
|
||||
check("N worker threads joined; counter exact (N*K) and >1 core used", bufferHas(ok_marker));
|
||||
check("no thread failure reported", !bufferHas(fail_marker));
|
||||
result();
|
||||
}
|
||||
|
||||
/// The full PID-1 path: the bootloader read /system/services/init off the boot volume and
|
||||
/// handed it over; load it as a user ELF and spawn it as a real ring-3 process
|
||||
/// — the same call the normal boot path makes — then confirm it beats. init
|
||||
|
||||
@@ -1,47 +1,123 @@
|
||||
//! thread-test — the first multi-threaded danos binary (docs/threading-plan.md M2).
|
||||
//! thread-test — danos's multi-threaded exerciser (docs/threading-plan.md M2, M3).
|
||||
//!
|
||||
//! Proves `runtime.Thread.spawn` starts a task in the **same address space**: the main
|
||||
//! thread spawns a worker, the worker writes a shared global and signals `done`, and the
|
||||
//! main thread — polling that shared memory — observes the write. Seeing the write proves
|
||||
//! the two tasks share one address space (a separate process could not touch this memory).
|
||||
//! The `thread-test: child ran in shared aspace ok` line is the case's marker.
|
||||
//! Two modes, chosen by argv[1] (default "spawn"):
|
||||
//! spawn — M2: one worker writes a shared global; the main thread observes it, proving
|
||||
//! `runtime.Thread.spawn` started a task in the **same** address space.
|
||||
//! join — M3: N workers each do K atomic increments on a shared counter and stamp the
|
||||
//! core they ran on; the main thread `join`s all N and checks the total is
|
||||
//! exactly N*K (every worker ran, join waited) and that >1 core was used
|
||||
//! (genuine parallelism). Then a detached worker proves `detach` runs and
|
||||
//! needs no join.
|
||||
//!
|
||||
//! Built multi-threaded (`addThreadedUserBinary`), so the poll below is a real atomic
|
||||
//! load the compiler must re-read — under a single-threaded build it could be hoisted.
|
||||
//! Built multi-threaded (`addThreadedUserBinary`) so atomics/shared reads are real.
|
||||
|
||||
const std = @import("std");
|
||||
const runtime = @import("runtime");
|
||||
|
||||
/// Written by the worker thread, read by main — the shared-address-space evidence.
|
||||
var shared_value: u32 = 0;
|
||||
/// Release/acquire handshake: publishes the `shared_value` write to the reader.
|
||||
var done = std.atomic.Value(u32).init(0);
|
||||
fn write(comptime s: []const u8) void {
|
||||
_ = runtime.system.write(s);
|
||||
}
|
||||
|
||||
// --- M2: spawn mode ---------------------------------------------------------
|
||||
|
||||
var shared_value: u32 = 0;
|
||||
var spawn_done = std.atomic.Value(u32).init(0);
|
||||
const sentinel: u32 = 0xA5A5;
|
||||
|
||||
fn worker() void {
|
||||
shared_value = sentinel; // a plain write to a global we share with main
|
||||
done.store(1, .release); // ...published by this release store
|
||||
fn spawnWorker() void {
|
||||
shared_value = sentinel;
|
||||
spawn_done.store(1, .release);
|
||||
}
|
||||
|
||||
pub fn main() void {
|
||||
_ = runtime.system.write("thread-test: starting\n");
|
||||
|
||||
_ = runtime.Thread.spawn(.{}, worker, .{}) catch {
|
||||
_ = runtime.system.write("thread-test: FAIL spawn refused\n");
|
||||
fn runSpawnMode() void {
|
||||
write("thread-test: starting\n");
|
||||
_ = runtime.Thread.spawn(.{}, spawnWorker, .{}) catch {
|
||||
write("thread-test: FAIL spawn refused\n");
|
||||
return;
|
||||
};
|
||||
|
||||
// Bounded wait for the worker to run and publish. yield() keeps the core useful;
|
||||
// the acquire load pairs with the worker's release store.
|
||||
var spins: usize = 0;
|
||||
while (done.load(.acquire) == 0 and spins < 50_000_000) : (spins += 1) {
|
||||
while (spawn_done.load(.acquire) == 0 and spins < 50_000_000) : (spins += 1) {
|
||||
runtime.system.yield();
|
||||
}
|
||||
|
||||
if (done.load(.acquire) == 1 and shared_value == sentinel) {
|
||||
_ = runtime.system.write("thread-test: child ran in shared aspace ok\n");
|
||||
if (spawn_done.load(.acquire) == 1 and shared_value == sentinel) {
|
||||
write("thread-test: child ran in shared aspace ok\n");
|
||||
} else {
|
||||
_ = runtime.system.write("thread-test: FAIL worker did not update shared memory\n");
|
||||
write("thread-test: FAIL worker did not update shared memory\n");
|
||||
}
|
||||
}
|
||||
|
||||
// --- M3: join mode ----------------------------------------------------------
|
||||
|
||||
const worker_count: u32 = 4;
|
||||
const iterations: u64 = 100_000;
|
||||
|
||||
var counter = std.atomic.Value(u64).init(0);
|
||||
var cores_seen = std.atomic.Value(u32).init(0);
|
||||
|
||||
fn joinWorker() void {
|
||||
var i: u64 = 0;
|
||||
while (i < iterations) : (i += 1) {
|
||||
_ = counter.fetchAdd(1, .monotonic);
|
||||
if (i % 1000 == 0) stampCore(); // periodic: catches cross-core migration too
|
||||
}
|
||||
stampCore();
|
||||
}
|
||||
|
||||
fn stampCore() void {
|
||||
const core = runtime.Thread.currentCore();
|
||||
if (core < 32) _ = cores_seen.fetchOr(@as(u32, 1) << @intCast(core), .monotonic);
|
||||
}
|
||||
|
||||
var detach_done = std.atomic.Value(u32).init(0);
|
||||
|
||||
fn detachWorker() void {
|
||||
detach_done.store(1, .release);
|
||||
}
|
||||
|
||||
fn runJoinMode() void {
|
||||
write("thread-test: join mode starting\n");
|
||||
|
||||
var threads: [worker_count]runtime.Thread = undefined;
|
||||
var spawned: u32 = 0;
|
||||
while (spawned < worker_count) : (spawned += 1) {
|
||||
threads[spawned] = runtime.Thread.spawn(.{}, joinWorker, .{}) catch break;
|
||||
}
|
||||
if (spawned != worker_count) {
|
||||
write("thread-test: FAIL could not spawn all workers\n");
|
||||
return;
|
||||
}
|
||||
for (threads[0..spawned]) |t| t.join();
|
||||
|
||||
const total = counter.load(.acquire);
|
||||
const cores = @popCount(cores_seen.load(.acquire));
|
||||
if (total != worker_count * iterations) {
|
||||
write("thread-test: FAIL counter mismatch (a worker was lost or join did not wait)\n");
|
||||
return;
|
||||
}
|
||||
if (cores <= 1) {
|
||||
write("thread-test: FAIL workers never ran on more than one core\n");
|
||||
return;
|
||||
}
|
||||
|
||||
// detach: the worker runs and we never join it.
|
||||
const dt = runtime.Thread.spawn(.{}, detachWorker, .{}) catch {
|
||||
write("thread-test: FAIL detach spawn refused\n");
|
||||
return;
|
||||
};
|
||||
dt.detach();
|
||||
var spins: usize = 0;
|
||||
while (detach_done.load(.acquire) == 0 and spins < 50_000_000) : (spins += 1) {
|
||||
runtime.system.yield();
|
||||
}
|
||||
if (detach_done.load(.acquire) != 1) {
|
||||
write("thread-test: FAIL detached worker did not run\n");
|
||||
return;
|
||||
}
|
||||
|
||||
write("thread-test: join ok\n"); // the M3 verdict marker
|
||||
}
|
||||
|
||||
pub fn main(init: runtime.process.Init) void {
|
||||
const mode = init.arguments.get(1) orelse "spawn";
|
||||
if (std.mem.eql(u8, mode, "join")) runJoinMode() else runSpawnMode();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user