Compare commits
7
Commits
d26515706e
...
def34e71fc
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
def34e71fc | ||
|
|
1b33f48acd | ||
|
|
b3a8147bd7 | ||
|
|
0730e77530 | ||
|
|
73df864fd2 | ||
|
|
11e363896f | ||
|
|
6e8b02d771 |
@@ -69,6 +69,36 @@ fn addUserBinary(
|
||||
acpi_ids_module: *std.Build.Module,
|
||||
name: []const u8,
|
||||
root: []const u8,
|
||||
) *std.Build.Step.Compile {
|
||||
return addUserBinaryImpl(b, target, runtime_module, mmio_module, xkeyboard_config_module, acpi_ids_module, name, root, false);
|
||||
}
|
||||
|
||||
/// As `addUserBinary`, but built multi-threaded (`single_threaded = false`) so real
|
||||
/// atomics/TLS work — required before a binary may call `runtime.Thread.spawn`
|
||||
/// (docs/threading.md). Threads are a deliberate per-binary opt-in.
|
||||
fn addThreadedUserBinary(
|
||||
b: *std.Build,
|
||||
target: std.Build.ResolvedTarget,
|
||||
runtime_module: *std.Build.Module,
|
||||
mmio_module: *std.Build.Module,
|
||||
xkeyboard_config_module: *std.Build.Module,
|
||||
acpi_ids_module: *std.Build.Module,
|
||||
name: []const u8,
|
||||
root: []const u8,
|
||||
) *std.Build.Step.Compile {
|
||||
return addUserBinaryImpl(b, target, runtime_module, mmio_module, xkeyboard_config_module, acpi_ids_module, name, root, true);
|
||||
}
|
||||
|
||||
fn addUserBinaryImpl(
|
||||
b: *std.Build,
|
||||
target: std.Build.ResolvedTarget,
|
||||
runtime_module: *std.Build.Module,
|
||||
mmio_module: *std.Build.Module,
|
||||
xkeyboard_config_module: *std.Build.Module,
|
||||
acpi_ids_module: *std.Build.Module,
|
||||
name: []const u8,
|
||||
root: []const u8,
|
||||
threaded: bool,
|
||||
) *std.Build.Step.Compile {
|
||||
// Settings (target, optimize, code model, ...) live on the root module only;
|
||||
// the program and runtime modules leave theirs null and inherit them.
|
||||
@@ -93,7 +123,7 @@ fn addUserBinary(
|
||||
.target = target,
|
||||
.optimize = .ReleaseSmall,
|
||||
.code_model = .large,
|
||||
.single_threaded = true,
|
||||
.single_threaded = !threaded, // a threaded binary needs real atomics/TLS
|
||||
.sanitize_c = .off,
|
||||
.stack_check = false,
|
||||
.stack_protector = false,
|
||||
@@ -548,6 +578,9 @@ pub fn build(b: *std.Build) void {
|
||||
const args_echo_exe = addUserBinary(b, kernel_target, runtime_module, mmio_module, xkeyboard_config_module, acpi_ids_module, "args-echo", "system/services/args-echo/args-echo.zig");
|
||||
const process_test_exe = addUserBinary(b, kernel_target, runtime_module, mmio_module, xkeyboard_config_module, acpi_ids_module, "process-test", "system/services/process-test/process-test.zig");
|
||||
const log_flush_exe = addUserBinary(b, kernel_target, runtime_module, mmio_module, xkeyboard_config_module, acpi_ids_module, "log-flush", "system/services/log-flush/log-flush.zig");
|
||||
// The first multi-threaded binary: exercises runtime.Thread over the thread ABI
|
||||
// (docs/threading.md). Built threaded so its shared-memory poll is real.
|
||||
const thread_test_exe = addThreadedUserBinary(b, kernel_target, runtime_module, mmio_module, xkeyboard_config_module, acpi_ids_module, "thread-test", "system/services/thread-test/thread-test.zig");
|
||||
|
||||
// Pack the user binaries into the initial_ramdisk image with the host-side Python tool
|
||||
// (the container format is trivial, and Python sidesteps std API churn). Args:
|
||||
@@ -591,6 +624,8 @@ pub fn build(b: *std.Build) void {
|
||||
mk_run.addFileArg(pci_bus_exe.getEmittedBin());
|
||||
mk_run.addArg("crash-test");
|
||||
mk_run.addFileArg(crash_test_exe.getEmittedBin());
|
||||
mk_run.addArg("thread-test");
|
||||
mk_run.addFileArg(thread_test_exe.getEmittedBin());
|
||||
mk_run.addArg("device-list");
|
||||
mk_run.addFileArg(device_list_exe.getEmittedBin());
|
||||
mk_run.addArg("discovery");
|
||||
|
||||
@@ -104,6 +104,13 @@ Start with the north star:
|
||||
port to **one seam** (`std.os.danos`), so we build `runtime.os` (→ that seam) plus a
|
||||
thin `runtime.fs`, retire the `posix` shim, and follow a phased path to
|
||||
`zig build-exe hello.zig` running on danos — **not** Linux-ABI emulation.
|
||||
- **[threading.md](threading.md) — threads, the std-shaped way.** **Built** (M1–M6):
|
||||
`runtime.Thread` mirrors `std.Thread`'s API (spawn/join/detach, Mutex/Condition/
|
||||
Semaphore) over a **private** thread ABI — several tasks sharing one address space via
|
||||
a `thread_spawn` syscall, futex-backed blocking, aspace refcounting. Why it's the
|
||||
native type and not literal `std.Thread` (the [private ABI](syscall.md)), and why
|
||||
threads stay a narrow opt-in against the [resilience](resilience.md) default. Build
|
||||
plan + gates: [threading-plan.md](threading-plan.md).
|
||||
- **[vdso.md](vdso.md) — the vDSO, the public system-call boundary.** A design note
|
||||
(not built yet) on keeping `abi.zig` genuinely private: a kernel-supplied, C-ABI
|
||||
entry blob mapped into every process as the *only* way into the kernel — so the
|
||||
|
||||
@@ -0,0 +1,290 @@
|
||||
# Threading — build plan (`runtime.Thread` over a private thread ABI)
|
||||
|
||||
The ordered, checkpointable build-out for [threading.md](threading.md). Each milestone
|
||||
lands on its own and ends in a **verifiable gate** — shaped for a `/loop` run, like
|
||||
[display-v2-plan.md](display-v2-plan.md). Read threading.md first for the *why*.
|
||||
|
||||
## Locked decisions (do not relitigate)
|
||||
|
||||
- **`runtime.Thread` mirrors `std.Thread`'s API; the implementation is danos-native.**
|
||||
Not literal `std.Thread` — that would break the [private ABI](syscall.md).
|
||||
- **Threads are a narrow, per-binary opt-in.** Default concurrency stays process + IPC
|
||||
([resilience.md](resilience.md)); only a service that asks is built
|
||||
`single_threaded = false`.
|
||||
- **Blocking is futex-backed, never spin-backed** — waiters park in the kernel so an
|
||||
idle core still halts ([halting.md](halting.md)).
|
||||
- **New syscalls are private**: extend [abi.zig](../system/abi.zig) `SystemCall` after
|
||||
`shm_physical = 36` (`thread_spawn = 37`, `thread_exit = 38`, `current_core = 39`,
|
||||
`futex_wait = 40`, `futex_wake = 41`) + a `library/runtime` wrapper; user code never names a number.
|
||||
- **Restart granularity stays the process** — a faulting thread kills its process; the
|
||||
supervisor restarts the process, which respawns its threads.
|
||||
|
||||
## Conventions
|
||||
|
||||
Follow [coding-standards.md](coding-standards.md): spell out non-acronym abbreviations,
|
||||
kebab-case file names, no `Co-Authored-By` trailers. New user binaries go through
|
||||
`addUserBinary` (with the new `threaded` flag where a binary spawns threads) and get
|
||||
packed into the initial-ramdisk; new syscalls extend [abi.zig](../system/abi.zig)
|
||||
`SystemCall` + a `library/runtime` wrapper; test services live beside the code they
|
||||
exercise and register a `ServiceId` if they must be looked up.
|
||||
|
||||
## How to verify along the way
|
||||
|
||||
**Every gate is serial-checkable — no screenshots** (this plan runs unattended). A
|
||||
thread proves it ran by writing to **shared memory** the parent reads back, and proves
|
||||
parallelism by stamping the **core index** it ran on (like the `smp`/`affinity` cases).
|
||||
|
||||
- `zig build test` — host unit tests (closure packing, mutex state machine, futex
|
||||
wrapper encodings).
|
||||
- `python3 test/qemu_test.py <case>` — boots the kernel in QEMU; asserts on serial
|
||||
markers. Thread cases set `smp: true` (real parallelism) and bump `mem` (they boot
|
||||
the process/scheduler stack); each milestone **adds its case to `CASES`** so its gate
|
||||
is runnable.
|
||||
- **Guardrail every milestone:** the concurrency-sensitive existing cases stay green —
|
||||
`smoke`, `sched`, `priority`, `smp`, `affinity`, `process`, `process-kill`,
|
||||
`supervision`, `fault-recovery`, `vfs-client-death`, `ipc`/`ipc-cap`,
|
||||
`display-service`. A threading change that regresses those is rejected.
|
||||
|
||||
## Unattended execution (the loop contract)
|
||||
|
||||
This plan runs to completion **without human input**. Every design choice is already
|
||||
fixed in *Locked decisions*; the checkboxes are the only state. A loop iteration must:
|
||||
|
||||
1. **Resume** at the first milestone that still has an unchecked `- [ ]`. (All earlier
|
||||
milestones are done — do not revisit them.)
|
||||
2. **Work on a branch.** On the first iteration, branch off `main` (e.g. `threading`);
|
||||
never commit threading work to `main`. All work stays local — **do not push**.
|
||||
3. **Implement** every unchecked item in that milestone, including adding its
|
||||
`-Dtest-case` to `CASES` in [test/qemu_test.py](../test/qemu_test.py) (with
|
||||
`smp: true` / a `mem` bump where noted) so the gate is runnable.
|
||||
4. **Run the gate**: `python3 test/qemu_test.py <case>`, then the full **guardrail
|
||||
set**, then `zig build` (clean) and `zig build test` (green).
|
||||
5. **Decide, do not ask:**
|
||||
- **Green** = the milestone's case prints its stated marker(s) and reports `PASS`,
|
||||
the whole guardrail set passes, `zig build` is clean, and host tests are green.
|
||||
→ tick this milestone's boxes **and** its `**Gate:**`-referenced case, `git commit`
|
||||
(`threads(M<n>): <summary>`, no `Co-Authored-By` trailer per
|
||||
[coding-standards.md](coding-standards.md)), and continue to the next milestone in
|
||||
the same iteration if budget remains; otherwise let the loop re-fire.
|
||||
- **Red** = anything above fails. Diagnose from the captured serial log
|
||||
(`zig-out/qemu-test/<case>-failed-serial.log`) and fix in place, then re-run — up to
|
||||
**3 fix attempts** for that gate. A concurrency case that fails then passes on a
|
||||
bare re-run is **flaky, not green**: re-run it **twice more** and treat green only
|
||||
if it passes all; otherwise fix the race (a real threading bug), don't paper over
|
||||
it.
|
||||
6. **A genuinely ambiguous fork is not a stop.** Pick the option most consistent with
|
||||
[threading.md](threading.md)'s *Locked decisions*, note the choice in the commit
|
||||
message, and continue. Do not pause for confirmation on in-scope, reversible work —
|
||||
this plan is that authorization.
|
||||
|
||||
**The only stop conditions:**
|
||||
|
||||
- **Done** — every milestone box is checked (M1–M6), `zig build` clean, whole
|
||||
`thread-*` suite + guardrail green. Update threading.md's status line to "built" (that
|
||||
is M6's own task) and stop.
|
||||
- **Blocked** — a gate is still red after 3 fix attempts, or a step needs something
|
||||
outside the repo (a toolchain change, new hardware, a decision no locked decision
|
||||
covers). Append `> **BLOCKED (M<n>):** <what failed, what was tried, the serial
|
||||
marker missing>` under that milestone, commit the WIP on the branch, and stop. Do not
|
||||
thrash further and do not silently skip the milestone.
|
||||
|
||||
Nothing else warrants stopping — not "should I proceed?", not "is this right?". The
|
||||
checkboxes + git history are the resumable record; the next iteration picks up from the
|
||||
first unchecked box.
|
||||
|
||||
---
|
||||
|
||||
## M1 — Address-space refcount (kernel foundation, no API, no behaviour change) ✅
|
||||
|
||||
The one invariant change threads require, landed and proven **before** anything shares
|
||||
an address space. Today aspace is 1:1 with a task and teardown destroys it on any user
|
||||
task's exit; make destruction happen on the **last** exit.
|
||||
|
||||
- [x] A refcount keyed by the address-space root, held in `scheduler.zig`
|
||||
(`aspace_refs`): `retainAspace` takes a reference in `spawnUserLocked` (on the
|
||||
success path, after the slot + stack are secured), all under the big kernel lock.
|
||||
- [x] Both task-teardown paths ([scheduler.zig](../system/kernel/scheduler.zig):
|
||||
`exitUserLocked` and `destroyTaskLocked`) call `releaseAspace`, which decrements
|
||||
and only `destroyAddressSpace`s at **zero**; an unretained space (hand-built test
|
||||
spaces) is destroyed directly, preserving prior behaviour.
|
||||
- [x] `-Dtest-case=aspace-refcount`: spawn and reap several ring-3 processes in sequence
|
||||
and assert (via test-observable `liveAspaceCount`/`aspaceDestroyCount`) that the
|
||||
live-space count returns to **baseline** and destructions advance by exactly that
|
||||
many — each space destroyed exactly once, no leak, no double-free. (Refcount
|
||||
observables, not raw frame counts, since kernel stacks are still leaked on exit.)
|
||||
|
||||
**Gate (met):** `python3 test/qemu_test.py aspace-refcount` passes
|
||||
(`aspace-refcount: spaces released to baseline ok` → `DANOS-TEST-RESULT: PASS`), and the
|
||||
full guardrail set passes unchanged — 13/13 (`smoke`, `sched`, `priority`, `smp`,
|
||||
`affinity`, `process`, `process-kill`, `supervision`, `fault-recovery`,
|
||||
`vfs-client-death`, `ipc`, `ipc-cap`, `display-service`); default `zig build` clean,
|
||||
`zig build test` green. The reframing is invisible until an aspace is actually shared.
|
||||
|
||||
## M2 — `thread_spawn` + `thread_exit`: a thread runs in the shared address space ✅
|
||||
|
||||
Spawn only — no join yet. Prove a second task executes in the **caller's** address
|
||||
space and exits cleanly.
|
||||
|
||||
- [x] [abi.zig](../system/abi.zig): `thread_spawn = 37`, `thread_exit = 38`. Handlers in
|
||||
process.zig; `thread_spawn` calls `scheduler.spawnThread` (shares the caller's
|
||||
aspace, `retainAspace`); `thread_exit` ends the task like a process `exit(0)`
|
||||
(`terminateCurrent` → `releaseAspace`). The closure pointer is delivered in the new
|
||||
thread's **rdi** via a new `jump_to_user_arg` asm path (`t.user_arg`, 0 for a
|
||||
process) — no naked runtime asm.
|
||||
- [x] `library/runtime/thread.zig` (barrel-exported as `runtime.Thread`): `spawn` maps a
|
||||
stack (`mmap`), heap-allocates the `{args}` closure, and calls
|
||||
`thread_spawn(&Closure.entry, stack_top, closure)`; `Closure.entry` (a plain C-ABI
|
||||
Zig fn, closure in rdi) runs the function and calls `thread_exit`. Stack top is
|
||||
16-aligned-minus-8 for the C entry.
|
||||
- [x] A `threaded` flag on the user-binary recipe (`addThreadedUserBinary` →
|
||||
`single_threaded = false`); `thread-test` is the first opt-in binary.
|
||||
- [x] `-Dtest-case=thread-spawn`: `thread-test` spawns a worker that writes a sentinel to
|
||||
a **shared** global and release-stores `done`; the main thread acquire-polls `done`
|
||||
and asserts the shared global holds the sentinel — proof the worker ran in the same
|
||||
address space.
|
||||
|
||||
**Gate (met):** `python3 test/qemu_test.py thread-spawn` passes
|
||||
(`thread-test: child ran in shared aspace ok` → `DANOS-TEST-RESULT: PASS`); guardrail set
|
||||
16/16 green (incl. `args`/`init`/`process`, which exercise the new `jump_to_user_arg`
|
||||
process path with arg 0) plus `aspace-refcount`; `zig build` clean, `zig build test`
|
||||
green.
|
||||
|
||||
> **Note (deferred to M3+):** the mmap arena is per-*task* (`heap_next`), so two threads
|
||||
> in one aspace that both `mmap` would collide. Fine for M2 (only the parent maps, for the
|
||||
> child's stack); make the arena per-aspace and the runtime heap thread-safe alongside the
|
||||
> `Mutex` work (M5).
|
||||
|
||||
## M3 — `join` + `detach` + real parallelism ✅
|
||||
|
||||
- [x] `join` over the existing exit-notification path
|
||||
([process-lifecycle.md](process-lifecycle.md)): `thread_spawn` gained a 4th arg, an
|
||||
`exit_endpoint` handle (resolved + refcounted like `spawnProcessSupervised`, via
|
||||
`spawnThreadSupervised`); `join` blocks in `ipc_reply_wait` on that endpoint until
|
||||
the child-exit notice for its `tid`, then `munmap`s the stack. `detach` relinquishes
|
||||
the join right (its stack is reclaimed at process exit — kernel-reaper reclaim for
|
||||
detached threads is deferred; see note).
|
||||
- [x] `runtime.Thread.join` / `detach`, plus `Thread.currentCore()` (a new `current_core`
|
||||
= 39 syscall) for the parallelism proof. `getCurrentId` deferred to M6 (TLS), where
|
||||
a lighter self-id fits. The closure now rides the **thread's own stack** (not the
|
||||
heap) — private per thread, so spawn/join touch no shared heap.
|
||||
- [x] `-Dtest-case=thread-join` (`smp: 4`): `thread-test` join mode spawns N=4 workers
|
||||
that each do K=100k `@atomicRmw`-increments on a shared counter and stamp the core
|
||||
they ran on; the main thread joins all N and asserts `counter == N*K` **and**
|
||||
`@popCount(cores_seen) > 1` (genuine cross-core parallelism), then a detached worker
|
||||
proves `detach` runs without a join.
|
||||
|
||||
**Gate (met):** `python3 test/qemu_test.py thread-join` passes (`thread-test: join ok` →
|
||||
`DANOS-TEST-RESULT: PASS`), robust across 4 runs; guardrail 17/17 green (incl. `smp`,
|
||||
`affinity`, `process-kill`, and `args`/`init`/`process` on the exit-endpoint spawn path)
|
||||
plus `aspace-refcount`/`thread-spawn`; `zig build` clean, `zig build test` green.
|
||||
|
||||
> **Note (deferred):** a detached thread's stack is freed only at process exit (not by the
|
||||
> reaper on thread exit) — kernel user-stack tracking + reclaim is a later refinement. And
|
||||
> the runtime heap is still not thread-safe: threads that both allocate concurrently would
|
||||
> race (the thread *machinery* avoids the heap, but worker code sharing an allocator does
|
||||
> not). Both fold into the M5 `Mutex`/allocator work.
|
||||
|
||||
## M4 — Futex: the one blocking primitive ✅
|
||||
|
||||
- [x] [abi.zig](../system/abi.zig): `futex_wait = 40`, `futex_wake = 41`. A waiter is a
|
||||
`.blocked` task tagged with `Task.futex_addr` (no queue linkage);
|
||||
`futex_wait(addr, expected, timeout_ns)` reads the user word under the big lock,
|
||||
parks iff `*addr == expected`, and returns on wake or timeout; `futex_wake(addr,
|
||||
count)` scans the task table and readies up to `count` matching waiters (same
|
||||
address space). No spinning — a parked waiter leaves its core free to `hlt`. A
|
||||
timed wait also sets `wake_at`, so the timer's `wakeExpired` wakes it; `futex_addr`
|
||||
staying non-zero (only `futex_wake` clears it) is how the waiter tells timeout from
|
||||
a real wake.
|
||||
- [x] `runtime.Thread.Futex` (`wait` / `timedWait` / `wake`) over the syscall wrappers.
|
||||
- [x] `-Dtest-case=thread-futex` (`smp: 4`): a waiter thread prints `waiting` and
|
||||
`futex_wait`s on a word; the main thread publishes it, prints `waking`, and
|
||||
`futex_wake`s; the waiter prints `woke`. Then a `timedWait` on an unwoken word
|
||||
reports `error.Timeout`.
|
||||
|
||||
**Gate (met):** `python3 test/qemu_test.py thread-futex` passes, robust across 3 runs —
|
||||
the case's **ordered** regex asserts `waiting → waking → woke → PASS` on the serial
|
||||
stream (the handoff proof), and `thread-futex: timeout ok` confirms the timeout.
|
||||
Guardrail 18/18 green (incl. `sleep`/`event`/`ipc` blocking paths) + `aspace-refcount`,
|
||||
`thread-spawn`, `thread-join`; `zig build` clean, `zig build test` green.
|
||||
|
||||
> **Note:** the kernel test checks only the freshest verdict marker via `bufferHas` (the
|
||||
> in-memory log ring buffer evicts older lines); ordering is asserted against the full
|
||||
> serial stream by the qemu regex instead.
|
||||
|
||||
## M5 — `Mutex` + `Condition` + `Semaphore` ✅
|
||||
|
||||
- [x] `runtime.Thread.Mutex` (three-state futex mutex: CAS fast path, `futex_wait`/`wake`
|
||||
slow path), `Condition` (`wait`/`timedWait`/`signal`/`broadcast`, a futex sequence
|
||||
counter), `Semaphore` (permits over `Mutex`+`Condition`) — the same state machines
|
||||
`std.Thread` uses, ported onto our `Futex`.
|
||||
- [x] `-Dtest-case=thread-mutex` (`smp: 4`): a bounded producer/consumer — 2 producers +
|
||||
2 consumers over one `Mutex` and two `Condition`s move N=2000 unique items through
|
||||
an 8-slot ring; the consumed checksum and tally match exactly (no lost/duplicated
|
||||
item, no overrun) under real cross-core contention. The small ring forces producers
|
||||
to block on full and consumers on empty, exercising `Condition.wait`.
|
||||
|
||||
**Gate (met):** `python3 test/qemu_test.py thread-mutex` passes (`thread-mutex: ok` →
|
||||
`DANOS-TEST-RESULT: PASS`), robust across 3 runs; guardrail 17/17 green (incl.
|
||||
`sleep`/`event`/`ipc`) + all M1–M4 thread cases; `zig build` clean, `zig build test`
|
||||
green.
|
||||
|
||||
> **Deferred (with rationale):**
|
||||
> - **`join` → futex completion word** — the exit-endpoint join (M3) is correct and
|
||||
> tested. A futex-completion join needs the *kernel* to clear+wake a word after the
|
||||
> thread is fully off its stack (a CLONE_CHILD_CLEARTID-style mechanism); doing it in
|
||||
> the thread's own trampoline would let `join` `munmap` the stack while the thread still
|
||||
> runs on it (use-after-free). Left on the exit-endpoint path; the kernel clear-on-exit
|
||||
> is a later, separate refinement.
|
||||
> - **Host unit tests for the state machines** — `Mutex`/`Condition` bottom out in the
|
||||
> `futex_*` syscalls, unavailable on the host without a mockable `Futex` seam. The QEMU
|
||||
> `thread-mutex` gate exercises them under real concurrency instead; a host-side mock is
|
||||
> future work.
|
||||
|
||||
## M6 — `getCurrentId`, docs, and CI wiring ✅
|
||||
|
||||
- [x] `getCurrentId` via a small `thread_self = 42` syscall (`runtime.Thread.getCurrentId`
|
||||
returns the kernel task id). **Per-thread `threadlocal` TLS is deferred** — no
|
||||
consumer needs it, and it would require context-switching `fs.base` per task (real
|
||||
kernel + per-switch cost) for an unused feature; threaded binaries have run fine
|
||||
without it through M2–M5. threading.md's TLS reasoning already scoped it as
|
||||
deferred-unless-needed. When a consumer appears, the shape is: `thread_spawn`
|
||||
allocates a per-thread TLS block, sets `fs.base`, and the context switch saves/
|
||||
restores it.
|
||||
- [x] `RwLock` / `WaitGroup` deferred (no consumer yet); they slot onto the same
|
||||
`Futex`/`Mutex`/`Condition` when wanted.
|
||||
- [x] All `thread-*` cases wired into [test/qemu_test.py](../test/qemu_test.py)
|
||||
(`thread-spawn`/`-join`/`-futex`/`-mutex`/`-id`); threading.md + docs/README.md
|
||||
status updated to **built**; the worked example is threading.md's win-condition.
|
||||
- [x] `-Dtest-case=thread-id` (`smp: 4`): two workers read `getCurrentId`; the main
|
||||
thread confirms all three ids are non-zero and distinct — each thread has its own
|
||||
kernel identity. (Renamed from `thread-tls`, which implied `threadlocal`.)
|
||||
|
||||
**Gate (met):** `python3 test/qemu_test.py thread-id` passes; the whole `thread-*` suite
|
||||
(`thread-spawn`/`-join`/`-futex`/`-mutex`/`-id`) plus the full guardrail set pass; default
|
||||
`zig build` clean, `zig build test` green.
|
||||
|
||||
---
|
||||
|
||||
## Status: built
|
||||
|
||||
M1–M6 complete. danos has `runtime.Thread` — `spawn`/`join`/`detach`, cross-core
|
||||
parallelism, futex, and `Mutex`/`Condition`/`Semaphore`, all over a private thread ABI
|
||||
behind the runtime. Deferred (with rationale, no consumer yet): `threadlocal` TLS,
|
||||
`RwLock`/`WaitGroup`, kernel clear-on-exit for a futex-completion `join`, a per-aspace
|
||||
mmap arena / thread-safe runtime heap, and host-side unit tests via a mockable `Futex`.
|
||||
|
||||
---
|
||||
|
||||
## Deferred (explicitly not in this plan)
|
||||
|
||||
- **Cross-process shared-memory futex** — the `(aspace, vaddr)` key can become a
|
||||
physical-address key so two processes share a futex through an [shm](display-v2.md)
|
||||
region. Not needed for intra-process threads.
|
||||
- **Per-thread priorities / affinity distinct from the process** — threads inherit the
|
||||
process priority ([scheduling.md](scheduling.md)); revisit only if it earns its keep.
|
||||
- **Per-thread signal delivery** — signals stay process-scoped
|
||||
([process-lifecycle.md](process-lifecycle.md)).
|
||||
- **A `pthread`/POSIX surface** — the API is `std.Thread`-shaped Zig, nothing more.
|
||||
- **A real `std.Thread` backend** — arrives with self-hosting
|
||||
([zig-self-hosting.md](zig-self-hosting.md)); it sits on these same primitives, so it
|
||||
swaps the impl under `runtime.Thread`, not the call sites.
|
||||
@@ -0,0 +1,316 @@
|
||||
# Threading: `runtime.Thread`, a std-shaped API over a private thread ABI
|
||||
|
||||
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
|
||||
**Zig 0.16** (the pinned toolchain); `std.Thread`'s internals move between releases, so
|
||||
treat upstream shapes as "0.16.x."
|
||||
|
||||
## The win condition
|
||||
|
||||
A danos service can write
|
||||
|
||||
```zig
|
||||
const t = try runtime.Thread.spawn(.{}, worker, .{ctx});
|
||||
// ... do other work concurrently ...
|
||||
t.join();
|
||||
```
|
||||
|
||||
and get real parallelism across cores — with `runtime.Thread.Mutex`,
|
||||
`runtime.Thread.Condition`, and `runtime.Thread.Semaphore` available for
|
||||
coordination — **without any code path reaching the kernel except through the
|
||||
runtime**. The call sites read exactly like `std.Thread`, so the day danos becomes a
|
||||
real Zig target (see [self-hosting](#the-self-hosting-endgame)) we swap the
|
||||
implementation underneath, not the API above.
|
||||
|
||||
## Locked decisions (do not relitigate)
|
||||
|
||||
- **We build `runtime.Thread`, not literal `std.Thread`.** It mirrors std's *API and
|
||||
features*; the implementation underneath is danos-native. See
|
||||
[Why not literal std.Thread](#why-not-literal-stdthread).
|
||||
- **Threads are a narrow, opt-in capability — not the default concurrency tool.** The
|
||||
default for resilience stays **process + IPC** ([resilience.md](resilience.md),
|
||||
[ipc.md](ipc.md)). See [Where threads fit](#where-threads-fit-the-resilience-tension).
|
||||
- **Blocking synchronization is futex-backed, never spin-backed.** Waiters sleep in
|
||||
the kernel so an idle core still halts ([halting.md](halting.md)).
|
||||
- **Per-binary opt-in to multi-threaded codegen.** Only a service that asks for
|
||||
threads is built `single_threaded = false`; the rest stay lean and single-threaded.
|
||||
- **The thread ABI is private.** New syscalls extend [abi.zig](../system/abi.zig)
|
||||
`SystemCall` and are reached only through `library/runtime` wrappers, exactly like
|
||||
every other danos syscall ([syscall.md](syscall.md)) — numbers stay renumberable.
|
||||
|
||||
## Why not literal `std.Thread`
|
||||
|
||||
danos's ABI invariant is that the **runtime is the sole holder of the syscall ABI**,
|
||||
and that ABI is private and renumberable ([syscall.md](syscall.md) — "unstable
|
||||
private ABI"). That is a security and evolvability asset: no compiled binary can
|
||||
hardcode a syscall number, and the kernel can renumber freely because only the
|
||||
runtime — rebuilt in lockstep — knows the mapping.
|
||||
|
||||
`std.Thread` is incompatible with that invariant on two counts:
|
||||
|
||||
1. **It selects its backend from `builtin.os.tag`, and issues syscalls directly.**
|
||||
danos targets `.os_tag = .freestanding` ([build.zig](../build.zig)), for which
|
||||
`std.Thread` resolves to an unsupported stub that `@compileError`s. Adding a real
|
||||
backend would either bake danos syscall numbers into std (breaking ABI privacy and
|
||||
renumbering) or fork std to route back through the runtime — a permanent rebase
|
||||
cost that buys nothing the native type doesn't.
|
||||
2. **Our user binaries are built `single_threaded = true`** ([build.zig](../build.zig)
|
||||
`addUserBinary`), which compiles threading out entirely and makes atomics and TLS
|
||||
single-threaded. Threads need this flipped per binary regardless.
|
||||
|
||||
So we take the *shape* of `std.Thread`, not the *type*. The cost of replicating the
|
||||
surface (spawn/join/Mutex/Condition) is small; the cost of the std type is the ABI
|
||||
invariant.
|
||||
|
||||
## Where threads fit: the resilience tension
|
||||
|
||||
Threads are in genuine tension with a resilience-first microkernel, and it is worth
|
||||
being explicit so we do not reach for them by reflex.
|
||||
|
||||
The reason danos pays for a microkernel is **fault isolation**
|
||||
([resilience.md](resilience.md)): a component corrupts its own address space, faults,
|
||||
and is **restarted** without touching anyone else — because the boundary *is* the
|
||||
address space. Threads deliberately remove that boundary *within* a process:
|
||||
|
||||
- Threads share one address space, so one thread's stray write corrupts them all —
|
||||
there is no isolation **between** threads.
|
||||
- Threads share fate: a fault in any thread, or a "kill the process" decision, takes
|
||||
down **all** of them. Restartability lives at the process level, not the thread
|
||||
level.
|
||||
- Shared mutable state reintroduces data races — the failure class the
|
||||
isolate-and-message model was chosen to avoid.
|
||||
|
||||
**Therefore:** the default answer to "make X concurrent" stays *another process over
|
||||
IPC* (isolated, independently restartable) or a single event loop with several
|
||||
message sources. Reach for a thread only inside **one** service that needs genuine
|
||||
**shared-memory, low-latency parallelism** and can accept intra-service fate-sharing —
|
||||
e.g. a compositor splitting tile compositing across cores, where per-tile IPC would be
|
||||
too chatty. "Input on one thread, display on another" is *not* that case; it wants two
|
||||
processes. The isolation boundary stays at process granularity.
|
||||
|
||||
## The API surface (mirrors `std.Thread`)
|
||||
|
||||
Lives in `library/runtime/thread.zig`, re-exported as `runtime.Thread`.
|
||||
|
||||
```zig
|
||||
pub const Thread = struct {
|
||||
pub const Id = u32; // the kernel task id
|
||||
pub const SpawnConfig = struct {
|
||||
stack_size: usize = default_stack_size,
|
||||
allocator: ?std.mem.Allocator = null, // for the closure + stack bookkeeping
|
||||
};
|
||||
pub const SpawnError = error{ OutOfMemory, ThreadQuotaExceeded, SystemResources };
|
||||
|
||||
pub fn spawn(config: SpawnConfig, comptime function: anytype, args: anytype) SpawnError!Thread;
|
||||
pub fn join(self: Thread) void; // block until the thread ends, reclaim its stack
|
||||
pub fn detach(self: Thread) void; // give up the right to join; kernel reclaims on exit
|
||||
pub fn getCurrentId() Id;
|
||||
pub fn yield() void; // -> existing `yield` syscall
|
||||
|
||||
pub const Mutex = struct { pub fn lock(*Mutex) void; pub fn tryLock(*Mutex) bool; pub fn unlock(*Mutex) void; };
|
||||
pub const Condition = struct { pub fn wait(*Condition, *Mutex) void; pub fn timedWait(*Condition, *Mutex, u64) error{Timeout}!void; pub fn signal(*Condition) void; pub fn broadcast(*Condition) void; };
|
||||
pub const Semaphore = struct { pub fn wait(*Semaphore) void; pub fn post(*Semaphore) void; };
|
||||
pub const Futex = struct { pub fn wait(*const atomic.Value(u32), u32) void; pub fn timedWait(...) error{Timeout}!void; pub fn wake(*const atomic.Value(u32), u32) void; };
|
||||
// RwLock / ResetEvent / WaitGroup follow the same pattern, added as needed.
|
||||
};
|
||||
```
|
||||
|
||||
Deviations from `std.Thread`, called out honestly:
|
||||
|
||||
- **The thread function's return value is discarded** (as `std.Thread.join` returns
|
||||
`void`). Return data through shared state or a `Semaphore`/`Condition`, not the
|
||||
return.
|
||||
- `getCpuCount()` maps to the existing SMP core count ([smp.md](smp.md)); a service
|
||||
rarely needs it.
|
||||
|
||||
## Kernel primitives (new private syscalls)
|
||||
|
||||
Four new entries extend [abi.zig](../system/abi.zig) `SystemCall` after
|
||||
`shm_physical = 36`, each with a `library/runtime` wrapper:
|
||||
|
||||
| Syscall | Signature | Purpose |
|
||||
|---|---|---|
|
||||
| `thread_spawn` | `(entry, stack_top, arg) -> tid` | create a task sharing the **caller's** address space |
|
||||
| `thread_exit` | `(stack_base, stack_len)` | end the calling thread; hand back its stack range for reclaim |
|
||||
| `futex_wait` | `(addr, expected, timeout_ns) -> status` | block if `*addr == expected`, until woken or timeout |
|
||||
| `futex_wake` | `(addr, count) -> woken` | wake up to `count` waiters on `addr` |
|
||||
|
||||
Plus one invariant change with no new syscall: **address-space reference counting**.
|
||||
|
||||
## Mechanics
|
||||
|
||||
### Address-space reference counting
|
||||
|
||||
Today an address space is 1:1 with a task: `spawnUserLocked` records `aspace` on the
|
||||
Task, and teardown does `destroyAddressSpace(t.aspace)` when **any** user task exits
|
||||
([scheduler.zig](../system/kernel/scheduler.zig)). With threads, several tasks share
|
||||
one `aspace`, so the first to exit would rip the address space out from under its
|
||||
siblings.
|
||||
|
||||
Fix: a small refcount keyed by the address-space root (`createAddressSpace` in
|
||||
[process.zig](../system/kernel/process.zig) sets it to 1). `thread_spawn` increments
|
||||
it; task teardown decrements and only calls `destroyAddressSpace` at **zero**. All of
|
||||
this is already under the big kernel lock, so no new locking. This is the one piece
|
||||
that must land and be proven before anything shares an address space.
|
||||
|
||||
### `thread_spawn` and the trampoline
|
||||
|
||||
The scheduler already accepts an arbitrary `aspace` and does **not** smuggle values
|
||||
through registers — `startUserTask` reads the entry/stack from the Task and
|
||||
`jumpToUser`s ([scheduler.zig](../system/kernel/scheduler.zig)). That makes the thread
|
||||
path clean:
|
||||
|
||||
1. The runtime's `spawn` `mmap`s a stack (syscall `4`), heap-allocates a closure —
|
||||
`{ fn_ptr, args_tuple, completion }`, the std "Instance" pattern — and writes the
|
||||
closure pointer to the **top word of the new stack**.
|
||||
2. It calls `thread_spawn(entry = &threadTrampoline, stack_top, arg = closure_ptr)`.
|
||||
The kernel calls the same `spawnUserLocked` path with the **caller's aspace**
|
||||
(refcount++), `entry`, and `user_sp = stack_top`.
|
||||
3. `threadTrampoline` (a small runtime shim) reads the closure off its stack, calls
|
||||
the user function, then calls `thread_exit`. No new register ABI — the closure
|
||||
pointer rides the stack the runtime set up, mirroring how `startUserTask` avoids
|
||||
register smuggling.
|
||||
|
||||
Unlike a process start, there is **no** System V argc/argv/auxv block
|
||||
([sysv.md](sysv.md)) — a thread stack carries only the closure pointer.
|
||||
|
||||
### Lifetime: exit, join, detach, stack reclaim
|
||||
|
||||
- **`thread_exit`** marks the task dead and hands the kernel the thread's user-stack
|
||||
range. The kernel reaps the task on the scheduler (already running on a *kernel*
|
||||
stack, so it can safely unmap the user stack), decrements the aspace refcount, and
|
||||
frees the task slot.
|
||||
- **`join` — Stage 1** reuses the existing exit-notification machinery
|
||||
([process-lifecycle.md](process-lifecycle.md)): `spawn` passes a per-thread
|
||||
`exit_endpoint`, and `join` blocks in `ipc_reply_wait` until the child-exit
|
||||
notification for that `tid` arrives, then `munmap`s the stack. No futex needed to
|
||||
land spawn/join.
|
||||
- **`join` — Stage 2 refinement** migrates to the std shape: a `completion` word in
|
||||
the closure that `thread_exit`'s trampoline `futex_wake`s and `join` `futex_wait`s
|
||||
on — dropping the per-thread endpoint. Kept as a refinement so Stage 1 ships first.
|
||||
- **`detach`** relinquishes the join right; the kernel reclaims the stack and slot on
|
||||
`thread_exit` (a detached thread's stack range is unmapped by the reaper, since no
|
||||
joiner will).
|
||||
|
||||
### Futex, and the sync primitives on top
|
||||
|
||||
`futex_wait`/`futex_wake` are the one blocking primitive; `Mutex`, `Condition`, and
|
||||
`Semaphore` are ordinary user-space state machines over an `atomic.Value(u32)` that
|
||||
call the futex wrappers on the slow path — the same construction `std.Thread` uses,
|
||||
so the algorithms port directly.
|
||||
|
||||
Keying: threads share an address space, so a **virtual address within that aspace**
|
||||
identifies a futex uniquely; the kernel keys its wait queue by `(aspace_root, vaddr)`.
|
||||
Keying by the **physical** address instead (translate `vaddr -> paddr` on entry) is a
|
||||
deliberate forward door: it lets two *processes* share a futex through an
|
||||
[shm](display-v2.md) region later, without changing the API. We start with the
|
||||
private-per-aspace key and note the physical-key upgrade.
|
||||
|
||||
No spinning: a contended lock parks the task in the kernel and the core is free to run
|
||||
other work or `hlt` ([halting.md](halting.md)). This is why futex is a locked
|
||||
decision, not a "maybe later."
|
||||
|
||||
### TLS and `getCurrentId`
|
||||
|
||||
danos sets up no `fs.base` TLS today (fine under `single_threaded`). Two scoped needs:
|
||||
|
||||
- **`getCurrentId`** returns the kernel task id — either a trivial syscall or, better,
|
||||
a value the runtime stashes in a per-thread control block.
|
||||
- **`threadlocal` variables** need a real per-thread TLS block and `fs.base` set per
|
||||
thread. `thread_spawn` sets `fs.base` to a runtime-allocated per-thread block; full
|
||||
`threadlocal` support is Stage 3, only if a consumer needs it. Nothing in the core
|
||||
spawn/join/mutex path requires `threadlocal`.
|
||||
|
||||
### Build: multi-threaded codegen, opt-in
|
||||
|
||||
`addUserBinary` gains a `threaded: bool = false` parameter; when set it builds that
|
||||
binary `single_threaded = false` so atomics and (later) TLS are real. Threads and
|
||||
atomics are unsound in a `single_threaded` image, so a binary must opt in **before**
|
||||
it may call `runtime.Thread.spawn`. Everyone else stays single-threaded and lean.
|
||||
|
||||
## Interaction with the rest of the kernel
|
||||
|
||||
- **Scheduler / SMP** ([scheduling.md](scheduling.md), [smp.md](smp.md)): a thread is
|
||||
just another `Task` with an `aspace` shared with its siblings; the existing
|
||||
per-core ready queues, priorities, and affinity apply unchanged. Threads of one
|
||||
process can run on different cores simultaneously — that is the point.
|
||||
- **Halting** ([halting.md](halting.md)): futex-parked waiters keep the "idle core
|
||||
halts" property intact under lock contention — no busy-wait.
|
||||
- **Lifecycle** ([process-lifecycle.md](process-lifecycle.md)): killing a process
|
||||
must kill *all* its threads and only then drop the last aspace ref. The kill path
|
||||
already targets a process; it fans out to every task on that aspace.
|
||||
- **Resilience** ([resilience.md](resilience.md)): a faulting thread kills its whole
|
||||
process (shared fate). The supervisor restarts the **process**, which respawns its
|
||||
threads from a known-good state — restart granularity stays the process.
|
||||
|
||||
## Build-out plan (staged, each gate serial-checkable)
|
||||
|
||||
The ordered, `/loop`-runnable milestones live in
|
||||
**[threading-plan.md](threading-plan.md)** (shaped like
|
||||
[display-v2-plan.md](display-v2-plan.md)): every milestone lands on its own and ends in
|
||||
a verifiable gate (`python3 test/qemu_test.py <case>`, asserting serial markers;
|
||||
`zig build test` for host unit tests). The stages below are the shape it expands.
|
||||
|
||||
- **Stage 0 — address-space refcount.** Refcount on the aspace root; teardown destroys
|
||||
at zero. No API yet; nothing shares an aspace, so refcount is 1 everywhere.
|
||||
*Gate:* the full QEMU suite stays green (no regression) — proves the reframing is
|
||||
invisible until used.
|
||||
- **Stage 1 — spawn / join / detach.** `thread_spawn` + `thread_exit`, the trampoline,
|
||||
stacks via `mmap`, join over the exit-endpoint, the `threaded` build flag.
|
||||
*Gate:* `-Dtest-case=thread-spawn` — a threaded test service spawns N threads that
|
||||
each `@atomicRmw`-increment a shared counter, the parent joins all N, and asserts
|
||||
the total is exactly N × iterations. Runs `smp` (multi-core) to prove real
|
||||
parallelism.
|
||||
- **Stage 2 — blocking synchronization.** `futex_wait`/`futex_wake` + `Futex`,
|
||||
`Mutex`, `Condition`, `Semaphore`; optionally migrate join to a futex completion
|
||||
word. *Gate:* `-Dtest-case=thread-mutex` — a bounded producer/consumer over a
|
||||
`Mutex` + `Condition` moves K items with no lost wakeups and no busy-wait (assert
|
||||
the consumer blocked, e.g. via a low idle tick count).
|
||||
- **Stage 3 — polish.** Per-thread TLS / `fs.base` and `threadlocal` (only if a
|
||||
consumer needs it), `RwLock`/`WaitGroup` as demanded, and this doc's cases wired
|
||||
into [test/qemu_test.py](../test/qemu_test.py).
|
||||
|
||||
## Conventions
|
||||
|
||||
Follow [coding-standards.md](coding-standards.md): spell out non-acronym
|
||||
abbreviations, kebab-case file names, no `Co-Authored-By` trailers. New syscalls
|
||||
extend [abi.zig](../system/abi.zig) `SystemCall` + a `library/runtime` wrapper
|
||||
([syscall.md](syscall.md)). `runtime.Thread` is a first-class runtime module, the same
|
||||
way `runtime.process` ([process-lifecycle.md](process-lifecycle.md)) and `runtime.ipc`
|
||||
are — user code never names a syscall.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- **No preemptive user-space signals delivered to a specific thread.** Signals stay
|
||||
process-scoped ([process-lifecycle.md](process-lifecycle.md)).
|
||||
- **No thread priorities distinct from the process.** Threads inherit the process
|
||||
priority; per-thread priority is a later question if it ever earns its keep.
|
||||
- **No cross-process shared-memory futex yet** — the physical-address key leaves the
|
||||
door open, but the first cut is private-per-aspace.
|
||||
- **No `pthread`/POSIX surface.** The API is `std.Thread`-shaped Zig, nothing more.
|
||||
|
||||
## The self-hosting endgame
|
||||
|
||||
When danos becomes a real Zig target and we (eventually) add a danos backend to std
|
||||
([zig-self-hosting.md](zig-self-hosting.md)), `std.Thread` can sit *on top of* these
|
||||
same kernel primitives — the danos `std.Thread.Impl` would call the very
|
||||
`thread_spawn`/`futex_*` wrappers `runtime.Thread` already uses. Because
|
||||
`runtime.Thread` was built API-compatible from day one, that transition swaps the
|
||||
implementation, not a single call site. Designing to the std shape now is what makes
|
||||
the later self-hosting lift cheap.
|
||||
|
||||
## Further reading
|
||||
|
||||
- [scheduling.md](scheduling.md), [smp.md](smp.md) — the task model these threads join.
|
||||
- [resilience.md](resilience.md), [vision.md](vision.md) — why isolation is the default
|
||||
and threads are the exception.
|
||||
- [syscall.md](syscall.md), [ipc.md](ipc.md) — the private ABI and the messaging model
|
||||
threads sit beside.
|
||||
- [halting.md](halting.md) — the idle/halt property futex-backed blocking preserves.
|
||||
- [zig-self-hosting.md](zig-self-hosting.md) — the target this bends toward.
|
||||
@@ -70,6 +70,10 @@ pub const panic = start.panic;
|
||||
/// Process entry types: the `Init` handed to `main`, and its `Arguments`.
|
||||
pub const process = @import("process.zig");
|
||||
|
||||
/// Threads: `runtime.Thread`, std.Thread-shaped, over the private thread ABI
|
||||
/// (docs/threading.md). A binary must be built multi-threaded to spawn.
|
||||
pub const Thread = @import("thread.zig").Thread;
|
||||
|
||||
/// The service harness: one replyWait loop folding requests, signals, and
|
||||
/// notifications into callbacks (docs/process-lifecycle.md).
|
||||
pub const service = @import("service.zig");
|
||||
|
||||
@@ -0,0 +1,260 @@
|
||||
//! `runtime.Thread` — threads for danos, shaped like Zig's `std.Thread` but built on
|
||||
//! danos's private thread ABI (docs/threading.md). Several tasks share one address
|
||||
//! space; `spawn` starts one, the kernel delivers the closure pointer in the new
|
||||
//! thread's rdi, a plain Zig trampoline runs the user function and calls `thread_exit`,
|
||||
//! and `join` blocks on the thread's exit notification. See docs/threading.md for why
|
||||
//! this mirrors `std.Thread`'s API rather than being the literal type.
|
||||
//!
|
||||
//! The closure (the function's captured args) lives at the **top of the thread's own
|
||||
//! stack**, not the heap — each thread's stack is private, so there is no shared-heap
|
||||
//! concurrency in the spawn/join machinery (the runtime heap is not yet thread-safe).
|
||||
//! A binary must be built multi-threaded (`addThreadedUserBinary`) before it may spawn.
|
||||
|
||||
const std = @import("std");
|
||||
const abi = @import("abi");
|
||||
const sc = @import("system-call.zig");
|
||||
const system = @import("system.zig");
|
||||
const ipc = @import("ipc.zig");
|
||||
|
||||
/// 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;
|
||||
|
||||
pub const Thread = struct {
|
||||
/// The kernel task id of the spawned thread.
|
||||
tid: u32,
|
||||
/// The endpoint the kernel notifies when this thread ends — what `join` blocks on.
|
||||
exit_endpoint: ipc.Handle,
|
||||
/// The mmap'd stack, reclaimed by `join` (or at process exit after `detach`).
|
||||
stack_base: usize,
|
||||
stack_size: usize,
|
||||
|
||||
pub const Id = u32;
|
||||
|
||||
pub const SpawnConfig = struct {
|
||||
/// Bytes of stack, rounded up to whole pages by the kernel's mmap.
|
||||
stack_size: usize = default_stack_size,
|
||||
};
|
||||
|
||||
pub const SpawnError = error{
|
||||
/// The kernel refused the thread, the stack mmap failed, or no endpoint was free.
|
||||
SystemResources,
|
||||
};
|
||||
|
||||
/// Start `function(args...)` on a new thread sharing this address space. Mirrors
|
||||
/// `std.Thread.spawn`. The thread's return value is discarded (as in `std.Thread`);
|
||||
/// return data through shared state.
|
||||
pub fn spawn(config: SpawnConfig, comptime function: anytype, args: anytype) SpawnError!Thread {
|
||||
const Args = @TypeOf(args);
|
||||
const Closure = struct {
|
||||
args: Args,
|
||||
/// Entered directly by the kernel with `self` in rdi (C ABI). Runs the user
|
||||
/// function, then ends the thread — never returns.
|
||||
fn entry(self_addr: usize) callconv(.c) noreturn {
|
||||
const self: *@This() = @ptrFromInt(self_addr);
|
||||
@call(.auto, function, self.args);
|
||||
exitThread();
|
||||
}
|
||||
};
|
||||
|
||||
// The endpoint the kernel posts this thread's exit notification to.
|
||||
const endpoint = ipc.createIpcEndpoint() orelse return error.SystemResources;
|
||||
|
||||
const base = system.mmap(config.stack_size, system.PROT_READ | system.PROT_WRITE);
|
||||
if (system.mmapFailed(base)) return error.SystemResources;
|
||||
|
||||
// Lay the closure at the very top of the thread's own stack, then start the
|
||||
// thread's rsp just below it (16-aligned minus 8, the alignment a `call` leaves
|
||||
// for a C-ABI entry) so the growing stack never overwrites the args.
|
||||
var closure_addr = (base + config.stack_size) - @sizeOf(Closure);
|
||||
closure_addr &= ~@as(usize, @alignOf(Closure) - 1); // align the closure down
|
||||
const closure: *Closure = @ptrFromInt(closure_addr);
|
||||
closure.* = .{ .args = args };
|
||||
|
||||
var stack_top = closure_addr & ~@as(usize, 15); // 16-align below the closure
|
||||
stack_top -= 8; // ...then rsp % 16 == 8 at the C entry
|
||||
|
||||
const tid = threadSpawn(@intFromPtr(&Closure.entry), stack_top, closure_addr, endpoint);
|
||||
if (threadSpawnFailed(tid)) {
|
||||
_ = system.munmap(base, config.stack_size);
|
||||
return error.SystemResources;
|
||||
}
|
||||
return .{ .tid = @intCast(tid), .exit_endpoint = endpoint, .stack_base = base, .stack_size = config.stack_size };
|
||||
}
|
||||
|
||||
/// Block until this thread finishes, then reclaim its stack. Mirrors
|
||||
/// `std.Thread.join`. The exit endpoint is private to this thread, so the first
|
||||
/// child-exit notification on it is this thread's.
|
||||
pub fn join(self: Thread) void {
|
||||
var receive: [0]u8 = undefined;
|
||||
while (true) {
|
||||
const got = ipc.replyWait(self.exit_endpoint, &.{}, &receive, null);
|
||||
if (got.isChildExit() and got.childProcessId() == self.tid) break;
|
||||
}
|
||||
_ = system.munmap(self.stack_base, self.stack_size);
|
||||
}
|
||||
|
||||
/// Relinquish the right to join: never wait for or reclaim this thread. Its stack is
|
||||
/// reclaimed at process exit (docs/threading-plan.md M3 — kernel-reaper stack reclaim
|
||||
/// for detached threads is a later refinement). Mirrors `std.Thread.detach`.
|
||||
pub fn detach(self: Thread) void {
|
||||
_ = self;
|
||||
}
|
||||
|
||||
/// The calling thread's id (its kernel task id). Mirrors `std.Thread.getCurrentId`.
|
||||
pub fn getCurrentId() Id {
|
||||
return @intCast(sc.systemCall0(.thread_self));
|
||||
}
|
||||
|
||||
/// The dense 0-based index of the core the calling thread is running on. A danos
|
||||
/// extension beyond `std.Thread`, used to observe genuine cross-core parallelism.
|
||||
pub fn currentCore() Id {
|
||||
return @intCast(sc.systemCall0(.current_core));
|
||||
}
|
||||
|
||||
/// `std.Thread.Futex`-shaped block/wake on a `u32` atomic — the primitive the
|
||||
/// blocking `Mutex`/`Condition`/`Semaphore` are built on. Waiters park in the
|
||||
/// kernel (no busy-wait), so an idle core still halts (docs/halting.md).
|
||||
pub const Futex = struct {
|
||||
/// Block while `ptr.* == expect`. Returns when woken by `wake`, or promptly if
|
||||
/// 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);
|
||||
}
|
||||
|
||||
/// 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;
|
||||
}
|
||||
|
||||
/// 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);
|
||||
}
|
||||
};
|
||||
|
||||
/// A mutual-exclusion lock, `std.Thread.Mutex`-shaped. The classic three-state
|
||||
/// futex mutex (unlocked / locked / contended): the fast path is a single CAS, and
|
||||
/// only a contended lock ever enters the kernel.
|
||||
pub const Mutex = struct {
|
||||
state: std.atomic.Value(u32) = std.atomic.Value(u32).init(unlocked),
|
||||
|
||||
const unlocked: u32 = 0;
|
||||
const locked: u32 = 1;
|
||||
const contended: u32 = 2;
|
||||
|
||||
/// Try to take the lock without blocking; returns whether it was acquired.
|
||||
pub fn tryLock(m: *Mutex) bool {
|
||||
return m.state.cmpxchgStrong(unlocked, locked, .acquire, .monotonic) == null;
|
||||
}
|
||||
|
||||
/// Acquire the lock, blocking in the kernel while it is contended.
|
||||
pub fn lock(m: *Mutex) void {
|
||||
if (m.state.cmpxchgStrong(unlocked, locked, .acquire, .monotonic) != null) m.lockSlow();
|
||||
}
|
||||
|
||||
fn lockSlow(m: *Mutex) void {
|
||||
@branchHint(.cold);
|
||||
// Mark the lock contended and take it as soon as it falls unlocked; park on
|
||||
// the futex while it stays contended. Marking contended may cause a spurious
|
||||
// wake on unlock (harmless), never a missed one.
|
||||
while (m.state.swap(contended, .acquire) != unlocked) {
|
||||
Futex.wait(&m.state, contended);
|
||||
}
|
||||
}
|
||||
|
||||
/// Release the lock; wake one waiter if the lock was contended.
|
||||
pub fn unlock(m: *Mutex) void {
|
||||
if (m.state.swap(unlocked, .release) == contended) Futex.wake(&m.state, 1);
|
||||
}
|
||||
};
|
||||
|
||||
/// A condition variable, `std.Thread.Condition`-shaped. Spurious wakeups are
|
||||
/// allowed — always wait in a predicate loop with the mutex held. Built on a futex
|
||||
/// sequence counter: a waiter samples the seq, drops the mutex, and parks until the
|
||||
/// seq changes (a signal that races the unlock bumps the seq, so it is not missed).
|
||||
pub const Condition = struct {
|
||||
seq: std.atomic.Value(u32) = std.atomic.Value(u32).init(0),
|
||||
|
||||
/// Atomically release `mutex` and block until signalled, then re-acquire it.
|
||||
pub fn wait(c: *Condition, mutex: *Mutex) void {
|
||||
const seq = c.seq.load(.acquire);
|
||||
mutex.unlock();
|
||||
Futex.wait(&c.seq, seq);
|
||||
mutex.lock();
|
||||
}
|
||||
|
||||
/// As `wait`, but returns `error.Timeout` if `timeout_ns` elapses first. The
|
||||
/// mutex is re-acquired either way.
|
||||
pub fn timedWait(c: *Condition, mutex: *Mutex, timeout_ns: u64) error{Timeout}!void {
|
||||
const seq = c.seq.load(.acquire);
|
||||
mutex.unlock();
|
||||
const timed_out = if (Futex.timedWait(&c.seq, seq, timeout_ns)) |_| false else |_| true;
|
||||
mutex.lock();
|
||||
if (timed_out) return error.Timeout;
|
||||
}
|
||||
|
||||
/// Wake one waiter.
|
||||
pub fn signal(c: *Condition) void {
|
||||
_ = c.seq.fetchAdd(1, .release);
|
||||
Futex.wake(&c.seq, 1);
|
||||
}
|
||||
|
||||
/// Wake all waiters.
|
||||
pub fn broadcast(c: *Condition) void {
|
||||
_ = c.seq.fetchAdd(1, .release);
|
||||
Futex.wake(&c.seq, std.math.maxInt(u32));
|
||||
}
|
||||
};
|
||||
|
||||
/// A counting semaphore, `std.Thread.Semaphore`-shaped: a permit count guarded by a
|
||||
/// `Mutex` + `Condition`.
|
||||
pub const Semaphore = struct {
|
||||
mutex: Mutex = .{},
|
||||
cond: Condition = .{},
|
||||
permits: usize = 0,
|
||||
|
||||
/// Take a permit, blocking until one is available.
|
||||
pub fn wait(s: *Semaphore) void {
|
||||
s.mutex.lock();
|
||||
defer s.mutex.unlock();
|
||||
while (s.permits == 0) s.cond.wait(&s.mutex);
|
||||
s.permits -= 1;
|
||||
}
|
||||
|
||||
/// Return a permit and wake a waiter.
|
||||
pub fn post(s: *Semaphore) void {
|
||||
s.mutex.lock();
|
||||
defer s.mutex.unlock();
|
||||
s.permits += 1;
|
||||
s.cond.signal();
|
||||
}
|
||||
};
|
||||
};
|
||||
|
||||
/// thread_spawn(entry, stack_top, arg, exit_endpoint) -> tid, or a wrapped error.
|
||||
fn threadSpawn(entry: usize, stack_top: usize, arg: usize, exit_endpoint: ipc.Handle) usize {
|
||||
return sc.systemCall4(.thread_spawn, entry, stack_top, arg, exit_endpoint);
|
||||
}
|
||||
|
||||
/// The kernel returns a real (small) task id on success and a wrapped `-1` on failure;
|
||||
/// no valid task id ever exceeds a u32.
|
||||
inline fn threadSpawnFailed(ret: usize) bool {
|
||||
return ret > std.math.maxInt(u32);
|
||||
}
|
||||
|
||||
/// End the calling thread. Never returns.
|
||||
fn exitThread() noreturn {
|
||||
_ = sc.systemCall0(.thread_exit);
|
||||
unreachable;
|
||||
}
|
||||
|
||||
/// futex_wait(addr, expect, timeout_ns) -> status (abi.futex_*).
|
||||
fn futexWait(addr: usize, expect: u32, timeout_ns: u64) usize {
|
||||
return sc.systemCall3(.futex_wait, addr, expect, timeout_ns);
|
||||
}
|
||||
|
||||
/// futex_wake(addr, count) -> number woken.
|
||||
fn futexWake(addr: usize, count: u32) usize {
|
||||
return sc.systemCall2(.futex_wake, addr, count);
|
||||
}
|
||||
@@ -63,9 +63,20 @@ 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, 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)
|
||||
futex_wait = 40, // futex_wait(addr, expected, timeout_ns) -> status: if *addr == expected, block until woken or the timeout; returns futex_woken/mismatch/timed_out (docs/threading.md)
|
||||
futex_wake = 41, // futex_wake(addr, count) -> woken: wake up to `count` tasks blocked in futex_wait on `addr` in this address space
|
||||
thread_self = 42, // thread_self() -> tid: the calling thread's kernel task id (runtime.Thread.getCurrentId)
|
||||
_,
|
||||
};
|
||||
|
||||
/// `futex_wait` return codes (in rax).
|
||||
pub const futex_woken: u64 = 0; // woken by a futex_wake
|
||||
pub const futex_mismatch: u64 = 1; // *addr != expected on entry; the caller did not block
|
||||
pub const futex_timed_out: u64 = 2; // the timeout elapsed before a wake
|
||||
|
||||
/// How a process ended — recorded by the kernel at death, queried by the
|
||||
/// supervisor with `process_exit_reason`, and the input to its restart decision
|
||||
/// (docs/process-lifecycle.md): a clean exit meant to stop, a fault wants a
|
||||
|
||||
@@ -687,6 +687,15 @@ pub fn jumpToUser(entry: u64, stack_top: u64) noreturn {
|
||||
jump_to_user(entry, stack_top);
|
||||
}
|
||||
|
||||
/// As `jumpToUser`, but delivers `arg0` in the user's `rdi` — how a fresh thread
|
||||
/// receives its closure pointer (docs/threading.md). A normal process is dropped
|
||||
/// with `arg0 = 0`, which its `_start` ignores (it reads argv off the stack).
|
||||
extern fn jump_to_user_arg(rip: u64, rsp: u64, arg0: u64) callconv(.c) noreturn;
|
||||
|
||||
pub fn jumpToUserArg(entry: u64, stack_top: u64, arg0: u64) noreturn {
|
||||
jump_to_user_arg(entry, stack_top, arg0);
|
||||
}
|
||||
|
||||
/// Route CPU exceptions to `handler`, which receives the trap frame and does not
|
||||
/// return. Until set, faults just halt the core.
|
||||
pub fn setFaultHandler(handler: *const fn (*const CpuState) noreturn) void {
|
||||
|
||||
@@ -123,6 +123,22 @@ jump_to_user:
|
||||
swapgs # user GS base (isr_common/syscall swap back on entry)
|
||||
iretq
|
||||
|
||||
# jump_to_user_arg(rdi = user rip, rsi = user rsp, rdx = user rdi/arg0): as
|
||||
# jump_to_user, but delivers arg0 in the user's rdi — how a fresh **thread**
|
||||
# receives its closure pointer (docs/threading.md). rdi carries the rip only until
|
||||
# it is pushed into the iretq frame, after which we overwrite it with the arg.
|
||||
.global jump_to_user_arg
|
||||
jump_to_user_arg:
|
||||
cli
|
||||
push $0x1B # user SS (0x18 | RPL 3)
|
||||
push %rsi # user RSP
|
||||
push $0x202 # RFLAGS: IF | reserved-1
|
||||
push $0x23 # user CS (0x20 | RPL 3)
|
||||
push %rdi # user RIP (consumes rdi)
|
||||
mov %rdx, %rdi # user rdi = arg0 (the thread's closure pointer)
|
||||
swapgs # user GS base (isr_common/syscall swap back on entry)
|
||||
iretq
|
||||
|
||||
# --- ring 3 entry/exit ------------------------------------------------------
|
||||
|
||||
# enter_user(rdi = user rip, rsi = user rsp, rdx = &TSS.rsp0)
|
||||
|
||||
+106
-1
@@ -227,6 +227,20 @@ fn system_call(state: *architecture.CpuState) void {
|
||||
.shm_create => systemShmCreate(state),
|
||||
.shm_map => systemShmMap(state),
|
||||
.shm_physical => systemShmPhysical(state),
|
||||
.thread_spawn => systemThreadSpawn(state),
|
||||
.current_core => systemCurrentCore(state),
|
||||
.thread_self => systemThreadSelf(state),
|
||||
.futex_wait => systemFutexWait(state),
|
||||
.futex_wake => systemFutexWake(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
|
||||
// space survives while sibling threads hold it). docs/threading.md.
|
||||
if (scheduler.currentIsUserProcess()) {
|
||||
scheduler.current().exit_reason = .exited;
|
||||
terminateCurrent();
|
||||
} else architecture.userExit();
|
||||
},
|
||||
_ => fail(state),
|
||||
}
|
||||
}
|
||||
@@ -642,6 +656,97 @@ fn systemSpawn(state: *architecture.CpuState) void {
|
||||
fail(state); // no bundled binary by that name
|
||||
}
|
||||
|
||||
/// thread_spawn(entry, stack_top, arg) -> tid: start a task that shares the **caller's**
|
||||
/// address space (docs/threading.md). The runtime supplies `entry` (its thread
|
||||
/// trampoline), a stack it mmap'd, and the closure pointer, which the kernel delivers in
|
||||
/// the new thread's rdi. The entry and stack must lie in the user half; the new thread is
|
||||
/// supervised by the caller and inherits its priority. Only a user process may spawn.
|
||||
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);
|
||||
// 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());
|
||||
}
|
||||
|
||||
/// thread_self() -> tid: the calling thread's kernel task id.
|
||||
fn systemThreadSelf(state: *architecture.CpuState) void {
|
||||
architecture.setSystemCallResult(state, scheduler.currentId());
|
||||
}
|
||||
|
||||
/// futex_wait(addr, expected, timeout_ns) -> status (docs/threading.md): if the 4-byte
|
||||
/// user word at `addr` still equals `expected`, block until a futex_wake on `addr` or
|
||||
/// (if timeout_ns > 0) the deadline. The compare and the block are one critical section,
|
||||
/// so a concurrent futex_wake cannot slip between them. Returns futex_woken / mismatch /
|
||||
/// timed_out.
|
||||
fn systemFutexWait(state: *architecture.CpuState) void {
|
||||
const addr = architecture.systemCallArg(state, 0);
|
||||
const expected: u32 = @truncate(architecture.systemCallArg(state, 1));
|
||||
const timeout_ns = architecture.systemCallArg(state, 2);
|
||||
const t = scheduler.current();
|
||||
if (t.aspace == 0) return fail(state);
|
||||
if (addr == 0 or (addr & 3) != 0 or addr + 4 > user_half_end) return fail(state);
|
||||
|
||||
const flags = sync.enter();
|
||||
var word_bytes: [4]u8 = undefined;
|
||||
if (!ipc.copyFromUser(t.aspace, addr, &word_bytes)) {
|
||||
sync.leave(flags);
|
||||
return fail(state);
|
||||
}
|
||||
if (std.mem.readInt(u32, &word_bytes, .little) != expected) {
|
||||
sync.leave(flags);
|
||||
architecture.setSystemCallResult(state, abi.futex_mismatch);
|
||||
return;
|
||||
}
|
||||
const timeout_ms = if (timeout_ns == 0) 0 else (timeout_ns + 999_999) / 1_000_000;
|
||||
const result = scheduler.futexWaitLocked(addr, timeout_ms);
|
||||
sync.leave(flags);
|
||||
architecture.setSystemCallResult(state, switch (result) {
|
||||
.woken => abi.futex_woken,
|
||||
.timed_out => abi.futex_timed_out,
|
||||
});
|
||||
}
|
||||
|
||||
/// futex_wake(addr, count) -> woken: wake up to `count` tasks blocked in futex_wait on
|
||||
/// `addr` in the caller's address space.
|
||||
fn systemFutexWake(state: *architecture.CpuState) void {
|
||||
const addr = architecture.systemCallArg(state, 0);
|
||||
const count: u32 = @truncate(architecture.systemCallArg(state, 1));
|
||||
const t = scheduler.current();
|
||||
if (t.aspace == 0) return fail(state);
|
||||
if (addr == 0 or (addr & 3) != 0 or addr + 4 > user_half_end) return fail(state);
|
||||
const flags = sync.enter();
|
||||
const woken = scheduler.futexWakeLocked(t.aspace, addr, count);
|
||||
sync.leave(flags);
|
||||
architecture.setSystemCallResult(state, woken);
|
||||
}
|
||||
|
||||
/// 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`
|
||||
@@ -1433,7 +1538,7 @@ pub fn spawnProcessSupervised(image: []const u8, priority: u3, argv: []const []c
|
||||
architecture.mapUserPageInto(aspace, page_virtual, stack_frame, true, false); // RW + NX
|
||||
}
|
||||
|
||||
const child = scheduler.spawnUserLocked(aspace, parsed.entry, user_sp, priority, argv[0], supervisor, if (exit_endpoint) |endpoint| @ptrCast(endpoint) else null) orelse
|
||||
const child = scheduler.spawnUserLocked(aspace, parsed.entry, user_sp, 0, priority, argv[0], supervisor, if (exit_endpoint) |endpoint| @ptrCast(endpoint) else null) orelse
|
||||
return error.OutOfMemory;
|
||||
// The child holds a reference to its exit endpoint from birth to death. Taken
|
||||
// only now, after nothing can fail; the lock is still held, so the child
|
||||
|
||||
+120
-4
@@ -78,6 +78,11 @@ pub const Task = struct {
|
||||
aspace: u64 = 0,
|
||||
user_ip: u64 = 0, // user-mode entry point (user task only)
|
||||
user_sp: u64 = 0, // user-mode stack pointer (user task only)
|
||||
user_arg: u64 = 0, // value delivered in the user's rdi at first entry: 0 for a
|
||||
// process (its _start ignores it), the closure pointer for a thread (docs/threading.md)
|
||||
// The user address this task is blocked on in futex_wait (0 = not futex-waiting).
|
||||
// Cleared to 0 by futexWakeLocked as the "woken, not timed out" signal (docs/threading.md).
|
||||
futex_addr: u64 = 0,
|
||||
// Next free virtual address in this task's mmap grant arena (0 = uninitialised;
|
||||
// process.zig lazily seeds it to the arena base on the first mmap). Bumped up
|
||||
// as the user heap grows; user task only.
|
||||
@@ -136,6 +141,67 @@ pub const ipc_maximum_handles = 16;
|
||||
pub const HandleObject = struct { kind: u8, ptr: *anyopaque };
|
||||
|
||||
var tasks = [_]Task{.{}} ** maximum_tasks;
|
||||
|
||||
/// Address-space reference counts: one live entry per address space, counting the
|
||||
/// tasks that share it. An address space is 1:1 with a process today; threads
|
||||
/// (docs/threading.md) will push a count above 1, and `destroyAddressSpace` must run
|
||||
/// only when the **last** task on an address space exits. All access is under the big
|
||||
/// kernel lock. There can be no more live address spaces than tasks, so the table is
|
||||
/// sized to the task pool and never overflows in practice.
|
||||
const AspaceRef = struct { root: u64 = 0, count: u32 = 0 };
|
||||
var aspace_refs = [_]AspaceRef{.{}} ** maximum_tasks;
|
||||
var aspace_destroy_count: u64 = 0;
|
||||
|
||||
/// Take a reference to address space `root` (0 = a kernel task, which owns none).
|
||||
/// Returns false only if the ref table is full — bounded by `maximum_tasks`, so in
|
||||
/// practice it never is. Caller holds the kernel lock.
|
||||
fn retainAspace(root: u64) bool {
|
||||
if (root == 0) return true;
|
||||
var free: ?*AspaceRef = null;
|
||||
for (&aspace_refs) |*entry| {
|
||||
if (entry.count != 0 and entry.root == root) {
|
||||
entry.count += 1;
|
||||
return true;
|
||||
}
|
||||
if (entry.count == 0 and free == null) free = entry;
|
||||
}
|
||||
const slot = free orelse return false;
|
||||
slot.* = .{ .root = root, .count = 1 };
|
||||
return true;
|
||||
}
|
||||
|
||||
/// Drop a reference to `root`; destroy the address space when the **last** one drops.
|
||||
/// A `root` with no entry — never retained, e.g. a hand-built test space — is
|
||||
/// destroyed directly, preserving the pre-refcount behaviour. Caller holds the lock.
|
||||
fn releaseAspace(root: u64) void {
|
||||
if (root == 0) return;
|
||||
for (&aspace_refs) |*entry| {
|
||||
if (entry.count == 0 or entry.root != root) continue;
|
||||
entry.count -= 1;
|
||||
if (entry.count == 0) {
|
||||
entry.root = 0;
|
||||
architecture.destroyAddressSpace(root);
|
||||
aspace_destroy_count += 1;
|
||||
}
|
||||
return;
|
||||
}
|
||||
architecture.destroyAddressSpace(root);
|
||||
aspace_destroy_count += 1;
|
||||
}
|
||||
|
||||
/// Test-observable: how many address spaces are live (entries with a nonzero count).
|
||||
pub fn liveAspaceCount() u32 {
|
||||
var live: u32 = 0;
|
||||
for (&aspace_refs) |*entry| {
|
||||
if (entry.count != 0) live += 1;
|
||||
}
|
||||
return live;
|
||||
}
|
||||
|
||||
/// Test-observable: total address-space destructions since boot.
|
||||
pub fn aspaceDestroyCount() u64 {
|
||||
return aspace_destroy_count;
|
||||
}
|
||||
var next_id: u32 = 1;
|
||||
|
||||
/// Per-CPU scheduler state: the task each core is running, its own idle task, and a
|
||||
@@ -327,9 +393,15 @@ pub fn spawnOn(entry: *const fn () void, priority: Priority, cpu: u32) bool {
|
||||
/// out of memory.
|
||||
/// **Caller must hold the kernel lock** (the loader that builds `aspace` holds it
|
||||
/// across the whole spawn, so the address space and the task appear atomically).
|
||||
pub fn spawnUserLocked(aspace: u64, entry: u64, user_sp: u64, priority: Priority, task_name: []const u8, supervisor: u32, exit_endpoint: ?*anyopaque) ?u32 {
|
||||
pub fn spawnUserLocked(aspace: u64, entry: u64, user_sp: u64, user_arg: u64, priority: Priority, task_name: []const u8, supervisor: u32, exit_endpoint: ?*anyopaque) ?u32 {
|
||||
const t = freeSlot() orelse return null;
|
||||
const stack = heap.allocator().alloc(u8, stack_size) catch return null;
|
||||
// Take this task's reference to the address space before we commit the slot, so a
|
||||
// failure here leaves nothing to unwind (the caller still owns the raw `aspace`).
|
||||
if (!retainAspace(aspace)) {
|
||||
heap.allocator().free(stack);
|
||||
return null;
|
||||
}
|
||||
t.* = .{
|
||||
.id = next_id,
|
||||
.state = .ready,
|
||||
@@ -338,6 +410,7 @@ pub fn spawnUserLocked(aspace: u64, entry: u64, user_sp: u64, priority: Priority
|
||||
.aspace = aspace,
|
||||
.user_ip = entry,
|
||||
.user_sp = user_sp,
|
||||
.user_arg = user_arg,
|
||||
.supervisor = supervisor,
|
||||
.exit_endpoint = exit_endpoint,
|
||||
};
|
||||
@@ -363,7 +436,7 @@ fn startUserTask() void {
|
||||
// No serial chatter here: this runs on every spawn, unserialized against
|
||||
// user-space writes, and its output used to shear concurrent log lines in
|
||||
// half — the largest source of corrupted markers in the QEMU scenarios.
|
||||
architecture.jumpToUser(t.user_ip, t.user_sp); // noreturn
|
||||
architecture.jumpToUserArg(t.user_ip, t.user_sp, t.user_arg); // noreturn (arg0 = 0 for a process)
|
||||
}
|
||||
|
||||
/// The unlocked task-creation primitive. Caller must hold the kernel lock (or be the
|
||||
@@ -446,6 +519,49 @@ pub fn sleep(ms: u64) void {
|
||||
sync.leave(flags);
|
||||
}
|
||||
|
||||
// --- futex: block/wake on a user address (docs/threading.md) ----------------
|
||||
//
|
||||
// A futex waiter is not linked into any queue — it is simply a `.blocked` task
|
||||
// tagged with the address it waits on (`futex_addr`). Waking scans the task table
|
||||
// (bounded) for matching waiters. A timed wait also sets `wake_at`, so the timer's
|
||||
// `wakeExpired` can wake it; `futex_addr` stays non-zero in that case, which is how
|
||||
// the waiter tells a timeout from a real wake.
|
||||
|
||||
pub const FutexResult = enum { woken, timed_out };
|
||||
|
||||
/// Block the current task on futex `addr` until woken, or (if `timeout_ms > 0`) the
|
||||
/// deadline. **Precondition:** the big kernel lock is held and the caller has already
|
||||
/// checked, under this same lock, that the futex word equals the expected value — so
|
||||
/// no wake can be missed. Returns with the lock still held.
|
||||
pub fn futexWaitLocked(addr: u64, timeout_ms: u64) FutexResult {
|
||||
const t = current();
|
||||
t.futex_addr = addr;
|
||||
t.wake_at = if (timeout_ms > 0) architecture.millis() + timeout_ms else 0;
|
||||
t.state = .blocked;
|
||||
schedule(); // woken by futexWakeLocked (clears futex_addr) or wakeExpired (timeout)
|
||||
const woken = t.futex_addr == 0;
|
||||
t.futex_addr = 0;
|
||||
t.wake_at = 0;
|
||||
return if (woken) .woken else .timed_out;
|
||||
}
|
||||
|
||||
/// Wake up to `count` tasks blocked in `futex_wait` on `addr` in address space
|
||||
/// `aspace`. Precondition: the big kernel lock is held. Returns how many woke.
|
||||
pub fn futexWakeLocked(aspace: u64, addr: u64, count: u32) u32 {
|
||||
var woken: u32 = 0;
|
||||
for (&tasks) |*t| {
|
||||
if (woken >= count) break;
|
||||
if (t.state == .blocked and t.aspace == aspace and t.futex_addr == addr) {
|
||||
t.futex_addr = 0; // the "woken, not timed out" signal to futexWaitLocked
|
||||
t.wake_at = 0;
|
||||
t.state = .ready;
|
||||
enqueue(t);
|
||||
woken += 1;
|
||||
}
|
||||
}
|
||||
return woken;
|
||||
}
|
||||
|
||||
// --- event-based blocking -------------------------------------------------
|
||||
//
|
||||
// A WaitQueue is a set of tasks blocked waiting for something (a resource, a
|
||||
@@ -720,7 +836,7 @@ pub fn exitUserLocked() noreturn {
|
||||
const kroot = architecture.kernelPageTable();
|
||||
architecture.loadPageTable(kroot); // off the process tables before freeing them
|
||||
pc.loaded_aspace = kroot;
|
||||
architecture.destroyAddressSpace(as);
|
||||
releaseAspace(as); // destroys only when this was the last task on the space
|
||||
}
|
||||
dying.state = .free;
|
||||
dying.aspace = 0;
|
||||
@@ -741,7 +857,7 @@ pub fn exitUserLocked() noreturn {
|
||||
/// task isn't running). The kernel stack is leaked, as in `exitUser` (no reaper
|
||||
/// yet). Precondition: the big kernel lock is held.
|
||||
pub fn destroyTaskLocked(t: *Task) void {
|
||||
if (t.aspace != 0) architecture.destroyAddressSpace(t.aspace);
|
||||
if (t.aspace != 0) releaseAspace(t.aspace); // destroys only on the last reference
|
||||
t.aspace = 0;
|
||||
t.kill_pending = false;
|
||||
t.in_system_call = false;
|
||||
|
||||
+261
-1
@@ -139,6 +139,18 @@ pub fn run(case: []const u8, boot_information: *const BootInformation) void {
|
||||
userPfTest();
|
||||
} else if (eql(case, "fault-recovery")) {
|
||||
faultRecoveryTest(boot_information);
|
||||
} else if (eql(case, "aspace-refcount")) {
|
||||
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, "thread-futex")) {
|
||||
threadFutexTest(boot_information);
|
||||
} else if (eql(case, "thread-mutex")) {
|
||||
threadMutexTest(boot_information);
|
||||
} else if (eql(case, "thread-id")) {
|
||||
threadIdTest(boot_information);
|
||||
} else if (eql(case, "args")) {
|
||||
argsTest(boot_information);
|
||||
} else if (eql(case, "init")) {
|
||||
@@ -1374,7 +1386,7 @@ fn spawnFaultingProcess() ?u32 {
|
||||
architecture.mapUserPageInto(aspace, process.stack_base_virtual, stack_frame, true, false); // RW + NX
|
||||
|
||||
// Supervised by the calling test task, so exitReasonOf can read the verdict.
|
||||
const id = scheduler.spawnUserLocked(aspace, process.code_virtual, process.stack_base_virtual + abi.page_size, 4, "fault-probe", scheduler.currentId(), null) orelse {
|
||||
const id = scheduler.spawnUserLocked(aspace, process.code_virtual, process.stack_base_virtual + abi.page_size, 0, 4, "fault-probe", scheduler.currentId(), null) orelse {
|
||||
architecture.destroyAddressSpace(aspace);
|
||||
return null;
|
||||
};
|
||||
@@ -1429,6 +1441,254 @@ fn faultRecoveryTest(boot_information: *const BootInformation) void {
|
||||
result();
|
||||
}
|
||||
|
||||
/// Address-space refcount (docs/threading-plan.md M1): every process holds exactly one
|
||||
/// reference to its address space, released when it dies, so `destroyAddressSpace` runs
|
||||
/// exactly once per space — no leak, no double-free. Spawn and kill several ring-3
|
||||
/// processes (the faulting probe, reaped by the kernel) and confirm the count of live
|
||||
/// address spaces returns to baseline while destructions advance by exactly that many.
|
||||
/// This is the foundation threads (shared address spaces) build on: the refactor must be
|
||||
/// invisible while every space still has exactly one task.
|
||||
fn aspaceRefcountTest(boot_information: *const BootInformation) void {
|
||||
_ = boot_information;
|
||||
log("DANOS-TEST-BEGIN: aspace-refcount\n", .{});
|
||||
const base_live = scheduler.liveAspaceCount();
|
||||
const base_destroyed = scheduler.aspaceDestroyCount();
|
||||
const rounds: u32 = 5;
|
||||
var killed: u32 = 0;
|
||||
var round: u32 = 0;
|
||||
while (round < rounds) : (round += 1) {
|
||||
process.fault_kill_count = 0;
|
||||
const probe = spawnFaultingProcess() orelse break;
|
||||
_ = probe;
|
||||
// Let the probe fault on its first instruction and be reaped.
|
||||
scheduler.setPriority(1);
|
||||
const deadline = architecture.millis() + 5000;
|
||||
while (process.fault_kill_count < 1 and architecture.millis() < deadline) scheduler.yield();
|
||||
scheduler.setPriority(4);
|
||||
if (process.fault_kill_count >= 1) killed += 1;
|
||||
}
|
||||
check("all probes spawned and were killed", killed == rounds);
|
||||
check("live address-space count returned to baseline", scheduler.liveAspaceCount() == base_live);
|
||||
check("each address space destroyed exactly once", scheduler.aspaceDestroyCount() == base_destroyed + rounds);
|
||||
if (killed == rounds and scheduler.liveAspaceCount() == base_live and
|
||||
scheduler.aspaceDestroyCount() == base_destroyed + rounds)
|
||||
log("aspace-refcount: spaces released to baseline ok\n", .{});
|
||||
result();
|
||||
}
|
||||
|
||||
/// Thread spawn (docs/threading-plan.md M2): the `thread-test` service spawns a worker
|
||||
/// thread that writes a shared global; the main thread, polling that memory, observes the
|
||||
/// write — proving `runtime.Thread.spawn` started a task in the **same** address space
|
||||
/// (a separate process could not touch it). The service's own marker is the verdict.
|
||||
fn threadSpawnTest(boot_information: *const BootInformation) void {
|
||||
log("DANOS-TEST-BEGIN: thread-spawn\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;
|
||||
};
|
||||
|
||||
check("thread-test spawned", spawnNamed(rd, "thread-test"));
|
||||
|
||||
// Wait for the service's verdict marker (it polls shared memory the worker wrote).
|
||||
const ok_marker = "thread-test: child ran in shared aspace ok";
|
||||
const fail_marker = "thread-test: FAIL";
|
||||
scheduler.setPriority(1);
|
||||
const deadline = architecture.millis() + 12000;
|
||||
while (architecture.millis() < deadline) {
|
||||
if (bufferHas(ok_marker) or bufferHas(fail_marker)) break;
|
||||
scheduler.yield();
|
||||
}
|
||||
scheduler.setPriority(4);
|
||||
|
||||
check("a worker thread ran in the shared address space (shared write observed)", bufferHas(ok_marker));
|
||||
check("the thread path reported no failure", !bufferHas(fail_marker));
|
||||
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();
|
||||
}
|
||||
|
||||
/// Futex (docs/threading-plan.md M4): `thread-test` in futex mode has a waiter thread
|
||||
/// block in `futex_wait` on a word; the main thread publishes the word and `futex_wake`s
|
||||
/// it. The serial order `waiting → waking → woke` shows the kernel handoff (the waiter
|
||||
/// parked and was woken, not spun), and a `timedWait` on an unwoken word reports a
|
||||
/// timeout. The verdict marker is emitted only after both hold.
|
||||
fn threadFutexTest(boot_information: *const BootInformation) void {
|
||||
log("DANOS-TEST-BEGIN: thread-futex\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", "futex" })) true else |_| false;
|
||||
break;
|
||||
}
|
||||
check("thread-test (futex mode) spawned", started);
|
||||
|
||||
const ok_marker = "thread-futex: ok";
|
||||
const fail_marker = "thread-futex: 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);
|
||||
|
||||
// Only the freshest marker is checked here — the kernel's log ring buffer may have
|
||||
// evicted the earlier ones by now. The waiting/waking/woke ordering (the handoff
|
||||
// proof) is asserted against the full serial stream by the qemu case's regex; the
|
||||
// verdict marker is emitted by thread-test only after the wake AND the timeout hold.
|
||||
check("futex handoff + timeout completed (verdict reached, no failure)", bufferHas(ok_marker) and !bufferHas(fail_marker));
|
||||
result();
|
||||
}
|
||||
|
||||
/// Mutex + Condition (docs/threading-plan.md M5): `thread-test` in mutex mode runs a
|
||||
/// bounded producer/consumer — P producers and C consumers over one `Mutex` and two
|
||||
/// `Condition`s move N unique items through a small ring. Every item is produced once;
|
||||
/// if the lock and condition variables are correct under real cross-core contention,
|
||||
/// the consumed checksum and tally match exactly (no lost or duplicated item, no
|
||||
/// overrun). The verdict marker is emitted only when both match.
|
||||
fn threadMutexTest(boot_information: *const BootInformation) void {
|
||||
log("DANOS-TEST-BEGIN: thread-mutex\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", "mutex" })) true else |_| false;
|
||||
break;
|
||||
}
|
||||
check("thread-test (mutex mode) spawned", started);
|
||||
|
||||
const ok_marker = "thread-mutex: ok";
|
||||
const fail_marker = "thread-mutex: 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("producer/consumer over Mutex+Condition moved every item exactly once", bufferHas(ok_marker) and !bufferHas(fail_marker));
|
||||
result();
|
||||
}
|
||||
|
||||
/// Thread identity (docs/threading-plan.md M6): `thread-test` in id mode spawns two
|
||||
/// workers that each read `runtime.Thread.getCurrentId`; the main thread confirms all
|
||||
/// three ids are non-zero and distinct — proof each thread has its own kernel identity.
|
||||
fn threadIdTest(boot_information: *const BootInformation) void {
|
||||
log("DANOS-TEST-BEGIN: thread-id\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", "id" })) true else |_| false;
|
||||
break;
|
||||
}
|
||||
check("thread-test (id mode) spawned", started);
|
||||
|
||||
const ok_marker = "thread-id: ok";
|
||||
const fail_marker = "thread-id: FAIL";
|
||||
scheduler.setPriority(1);
|
||||
const deadline = architecture.millis() + 12000;
|
||||
while (architecture.millis() < deadline) {
|
||||
if (bufferHas(ok_marker) or bufferHas(fail_marker)) break;
|
||||
scheduler.yield();
|
||||
}
|
||||
scheduler.setPriority(4);
|
||||
|
||||
check("each thread has a distinct, non-zero getCurrentId", bufferHas(ok_marker) and !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
|
||||
|
||||
@@ -0,0 +1,310 @@
|
||||
//! thread-test — danos's multi-threaded exerciser (docs/threading-plan.md M2, M3).
|
||||
//!
|
||||
//! 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 atomics/shared reads are real.
|
||||
|
||||
const std = @import("std");
|
||||
const runtime = @import("runtime");
|
||||
|
||||
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 spawnWorker() void {
|
||||
shared_value = sentinel;
|
||||
spawn_done.store(1, .release);
|
||||
}
|
||||
|
||||
fn runSpawnMode() void {
|
||||
write("thread-test: starting\n");
|
||||
_ = runtime.Thread.spawn(.{}, spawnWorker, .{}) catch {
|
||||
write("thread-test: FAIL spawn refused\n");
|
||||
return;
|
||||
};
|
||||
var spins: usize = 0;
|
||||
while (spawn_done.load(.acquire) == 0 and spins < 50_000_000) : (spins += 1) {
|
||||
runtime.system.yield();
|
||||
}
|
||||
if (spawn_done.load(.acquire) == 1 and shared_value == sentinel) {
|
||||
write("thread-test: child ran in shared aspace ok\n");
|
||||
} else {
|
||||
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
|
||||
}
|
||||
|
||||
// --- M4: futex mode ---------------------------------------------------------
|
||||
|
||||
const Futex = runtime.Thread.Futex;
|
||||
|
||||
var futex_word = std.atomic.Value(u32).init(0);
|
||||
var waiter_parked = std.atomic.Value(u32).init(0);
|
||||
|
||||
fn futexWaiter() void {
|
||||
write("thread-futex: waiting\n");
|
||||
waiter_parked.store(1, .release);
|
||||
// Block while the word is still 0; the waker sets it to 1 and wakes us.
|
||||
while (futex_word.load(.acquire) == 0) {
|
||||
Futex.wait(&futex_word, 0);
|
||||
}
|
||||
write("thread-futex: woke\n");
|
||||
}
|
||||
|
||||
fn runFutexMode() void {
|
||||
write("thread-futex: starting\n");
|
||||
|
||||
const waiter = runtime.Thread.spawn(.{}, futexWaiter, .{}) catch {
|
||||
write("thread-futex: FAIL spawn refused\n");
|
||||
return;
|
||||
};
|
||||
// Let the waiter reach its wait, then give it a beat to actually park in-kernel.
|
||||
var spins: usize = 0;
|
||||
while (waiter_parked.load(.acquire) == 0 and spins < 50_000_000) : (spins += 1) {
|
||||
runtime.system.yield();
|
||||
}
|
||||
runtime.system.sleep(50);
|
||||
|
||||
// The handshake: publish the value, then wake the parked waiter.
|
||||
futex_word.store(1, .release);
|
||||
write("thread-futex: waking\n");
|
||||
Futex.wake(&futex_word, 1);
|
||||
|
||||
waiter.join(); // returns once the waiter woke and printed "woke"
|
||||
|
||||
// Timeout: nobody ever wakes this word, so timedWait must report a timeout.
|
||||
var lonely = std.atomic.Value(u32).init(0);
|
||||
if (Futex.timedWait(&lonely, 0, 100_000_000)) |_| {
|
||||
write("thread-futex: FAIL timedWait did not time out\n");
|
||||
return;
|
||||
} else |_| {}
|
||||
write("thread-futex: timeout ok\n");
|
||||
|
||||
write("thread-futex: ok\n"); // the M4 verdict marker
|
||||
}
|
||||
|
||||
// --- M5: mutex mode (bounded producer/consumer over Mutex + Condition) ------
|
||||
|
||||
const Mutex = runtime.Thread.Mutex;
|
||||
const Condition = runtime.Thread.Condition;
|
||||
|
||||
const producers: u32 = 2;
|
||||
const consumers: u32 = 2;
|
||||
const per_producer: u32 = 1000;
|
||||
const per_consumer: u32 = 1000; // producers*per_producer == consumers*per_consumer (balanced)
|
||||
const total_items: u32 = producers * per_producer;
|
||||
const ring_cap: usize = 8; // small, so producers block on full and consumers on empty
|
||||
|
||||
var ring: [ring_cap]u32 = undefined;
|
||||
var ring_count: usize = 0;
|
||||
var ring_head: usize = 0;
|
||||
var ring_tail: usize = 0;
|
||||
|
||||
var pc_mutex = Mutex{};
|
||||
var not_full = Condition{};
|
||||
var not_empty = Condition{};
|
||||
|
||||
// Verified outside the lock: the checksum and tally of everything consumed.
|
||||
var consumed_sum = std.atomic.Value(u64).init(0);
|
||||
var consumed_count = std.atomic.Value(u32).init(0);
|
||||
|
||||
fn producer(base: u32) void {
|
||||
var i: u32 = 0;
|
||||
while (i < per_producer) : (i += 1) {
|
||||
const item = base + i;
|
||||
pc_mutex.lock();
|
||||
while (ring_count == ring_cap) not_full.wait(&pc_mutex);
|
||||
ring[ring_tail] = item;
|
||||
ring_tail = (ring_tail + 1) % ring_cap;
|
||||
ring_count += 1;
|
||||
pc_mutex.unlock();
|
||||
not_empty.signal();
|
||||
}
|
||||
}
|
||||
|
||||
fn consumer() void {
|
||||
var i: u32 = 0;
|
||||
while (i < per_consumer) : (i += 1) {
|
||||
pc_mutex.lock();
|
||||
while (ring_count == 0) not_empty.wait(&pc_mutex);
|
||||
const item = ring[ring_head];
|
||||
ring_head = (ring_head + 1) % ring_cap;
|
||||
ring_count -= 1;
|
||||
pc_mutex.unlock();
|
||||
not_full.signal();
|
||||
_ = consumed_sum.fetchAdd(item, .monotonic);
|
||||
_ = consumed_count.fetchAdd(1, .monotonic);
|
||||
}
|
||||
}
|
||||
|
||||
fn runMutexMode() void {
|
||||
write("thread-mutex: starting\n");
|
||||
|
||||
var threads: [producers + consumers]runtime.Thread = undefined;
|
||||
var n: usize = 0;
|
||||
var p: u32 = 0;
|
||||
while (p < producers) : (p += 1) {
|
||||
threads[n] = runtime.Thread.spawn(.{}, producer, .{p * per_producer}) catch {
|
||||
write("thread-mutex: FAIL producer spawn\n");
|
||||
return;
|
||||
};
|
||||
n += 1;
|
||||
}
|
||||
var c: u32 = 0;
|
||||
while (c < consumers) : (c += 1) {
|
||||
threads[n] = runtime.Thread.spawn(.{}, consumer, .{}) catch {
|
||||
write("thread-mutex: FAIL consumer spawn\n");
|
||||
return;
|
||||
};
|
||||
n += 1;
|
||||
}
|
||||
for (threads[0..n]) |t| t.join();
|
||||
|
||||
// Every item 0..total_items-1 was produced exactly once; if the mutex/condition are
|
||||
// correct, each was consumed exactly once, so the checksum matches.
|
||||
const expected_sum: u64 = @as(u64, total_items) * (total_items - 1) / 2;
|
||||
if (consumed_count.load(.acquire) != total_items) {
|
||||
write("thread-mutex: FAIL wrong number of items consumed\n");
|
||||
return;
|
||||
}
|
||||
if (consumed_sum.load(.acquire) != expected_sum) {
|
||||
write("thread-mutex: FAIL checksum mismatch (item lost or duplicated)\n");
|
||||
return;
|
||||
}
|
||||
write("thread-mutex: ok\n"); // the M5 verdict marker
|
||||
}
|
||||
|
||||
// --- M6: id mode (getCurrentId identity) ------------------------------------
|
||||
|
||||
var worker_ids: [2]std.atomic.Value(u32) = .{ std.atomic.Value(u32).init(0), std.atomic.Value(u32).init(0) };
|
||||
|
||||
fn idWorker(slot: usize) void {
|
||||
worker_ids[slot].store(runtime.Thread.getCurrentId(), .release);
|
||||
}
|
||||
|
||||
fn runIdMode() void {
|
||||
write("thread-id: starting\n");
|
||||
const main_id = runtime.Thread.getCurrentId();
|
||||
|
||||
const t0 = runtime.Thread.spawn(.{}, idWorker, .{@as(usize, 0)}) catch {
|
||||
write("thread-id: FAIL spawn\n");
|
||||
return;
|
||||
};
|
||||
const t1 = runtime.Thread.spawn(.{}, idWorker, .{@as(usize, 1)}) catch {
|
||||
write("thread-id: FAIL spawn\n");
|
||||
return;
|
||||
};
|
||||
t0.join();
|
||||
t1.join();
|
||||
|
||||
const id0 = worker_ids[0].load(.acquire);
|
||||
const id1 = worker_ids[1].load(.acquire);
|
||||
// Each thread has its own kernel task id: all three distinct and non-zero.
|
||||
if (main_id == 0 or id0 == 0 or id1 == 0) {
|
||||
write("thread-id: FAIL a thread reported id 0\n");
|
||||
return;
|
||||
}
|
||||
if (id0 == id1 or id0 == main_id or id1 == main_id) {
|
||||
write("thread-id: FAIL thread ids collided\n");
|
||||
return;
|
||||
}
|
||||
write("thread-id: ok\n"); // the M6 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 if (std.mem.eql(u8, mode, "futex")) {
|
||||
runFutexMode();
|
||||
} else if (std.mem.eql(u8, mode, "mutex")) {
|
||||
runMutexMode();
|
||||
} else if (std.mem.eql(u8, mode, "id")) {
|
||||
runIdMode();
|
||||
} else {
|
||||
runSpawnMode();
|
||||
}
|
||||
}
|
||||
@@ -293,6 +293,52 @@ CASES = [
|
||||
"timeout": 60,
|
||||
"expect": r"DANOS-TEST-RESULT: PASS",
|
||||
"fail": r"DANOS-TEST-RESULT: FAIL"},
|
||||
|
||||
# docs/threading-plan.md M1: address-space refcount — spaces destroyed exactly
|
||||
# once per process, no leak/double-free (the foundation shared-aspace threads need).
|
||||
{"name": "aspace-refcount",
|
||||
"timeout": 60,
|
||||
"expect": r"DANOS-TEST-RESULT: PASS",
|
||||
"fail": r"DANOS-TEST-RESULT: FAIL"},
|
||||
|
||||
# docs/threading-plan.md M2: runtime.Thread.spawn — a worker thread runs in the
|
||||
# caller's address space (a shared-memory write, observed by the main thread).
|
||||
{"name": "thread-spawn",
|
||||
"timeout": 60,
|
||||
"expect": r"DANOS-TEST-RESULT: PASS",
|
||||
"fail": r"DANOS-TEST-RESULT: FAIL"},
|
||||
|
||||
# docs/threading-plan.md M3: join + parallelism — N workers each do K atomic
|
||||
# increments (total exactly N*K after join) and run on >1 core; plus detach.
|
||||
{"name": "thread-join",
|
||||
"smp": 4,
|
||||
"timeout": 60,
|
||||
"expect": r"DANOS-TEST-RESULT: PASS",
|
||||
"fail": r"DANOS-TEST-RESULT: FAIL"},
|
||||
|
||||
# docs/threading-plan.md M4: futex — a thread parks in futex_wait and is woken by
|
||||
# futex_wake (serial order waiting/waking/woke), and timedWait reports a timeout.
|
||||
{"name": "thread-futex",
|
||||
"smp": 4,
|
||||
"timeout": 60,
|
||||
"expect": r"thread-futex: waiting[\s\S]*thread-futex: waking[\s\S]*thread-futex: woke[\s\S]*DANOS-TEST-RESULT: PASS",
|
||||
"fail": r"DANOS-TEST-RESULT: FAIL"},
|
||||
|
||||
# docs/threading-plan.md M5: Mutex + Condition — a bounded producer/consumer moves
|
||||
# N unique items across cores; the consumed checksum matches exactly (no loss).
|
||||
{"name": "thread-mutex",
|
||||
"smp": 4,
|
||||
"timeout": 60,
|
||||
"expect": r"DANOS-TEST-RESULT: PASS",
|
||||
"fail": r"DANOS-TEST-RESULT: FAIL"},
|
||||
|
||||
# docs/threading-plan.md M6: thread identity — getCurrentId is distinct and non-zero
|
||||
# for the main thread and two workers.
|
||||
{"name": "thread-id",
|
||||
"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",
|
||||
|
||||
Reference in New Issue
Block a user