summaryrefslogtreecommitdiff
path: root/src/io/p4.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-08-25 12:40:53 -0300
committerGabriel Schneider <[email protected]>2026-08-25 12:46:51 -0300
commitf5f8068fac59b4f16046c2022c2fc7c7e447ef4c (patch)
tree2731a3ed4e51cae09e184e25778eded5fc37d1f5 /src/io/p4.zig
downloadesp32p4-f5f8068fac59b4f16046c2022c2fc7c7e447ef4c.tar.gz
esp32p4-f5f8068fac59b4f16046c2022c2fc7c7e447ef4c.zip
zig-p4: pure-Zig ESP32-P4 toolchain
build.zig generates the linker script and drives Zig's own LLD; tools/image.zig turns the ELF into a flashable image and tools/{rom,serial}.zig speak the mask ROM loader over the UART. No CMake, ninja, idf.py, esptool, or external linker. src/soc.zig is a comptime register model over ESP-IDF's own *_reg.h headers; src/hal/ adds peripheral sequences; src/io/ implements std.Io for the chip; src/oracle/ diffs this HAL against ESP-IDF's on the die.
Diffstat (limited to 'src/io/p4.zig')
-rw-r--r--src/io/p4.zig1475
1 files changed, 1475 insertions, 0 deletions
diff --git a/src/io/p4.zig b/src/io/p4.zig
new file mode 100644
index 0000000..ef82e02
--- /dev/null
+++ b/src/io/p4.zig
@@ -0,0 +1,1475 @@
+//! `std.Io` for the ESP32-P4: a cooperative scheduler over caller-provided static stacks.
+//!
+//! This is the "FreeRTOS for Zig" layer, except that it is not a framework and there is nothing to
+//! port to: `std.Io` is an interface, and filling in enough of its vtable buys the whole ecosystem
+//! that is written against it. Implementing `futexWait`/`futexWake` and `async`/`await` here is
+//! what makes `std.Io.Mutex`, `std.Io.Condition`, `std.Io.Semaphore`, `std.Io.RwLock`,
+//! `std.Io.Queue(T)` and `std.Io.Future(T)` work on this chip - none of which appear in this file,
+//! because they are built on those primitives in std and need nothing from us
+//! (`/usr/lib/zig/std/Io.zig:1587` Mutex, `:1653` Condition, `:1874` Queue).
+//!
+//! ## What is real and what is not
+//!
+//! `std.Io.VTable` has 109 entries. Sixteen are implemented here; the other 93 come from
+//! `unimplemented.zig`, where each one panics with its own name. That is the deliberate shape of
+//! this file: a partial implementation whose gaps announce themselves, rather than a plausible
+//! stub that returns zero. Implemented:
+//!
+//! now clockResolution sleep - the timebase, from hal.systimer
+//! async concurrent await cancel - cooperative tasks on static stacks
+//! futexWait futexWaitUncancelable futexWake
+//! checkCancel recancel swapCancelProtection
+//! crashHandler random randomSecure - randomSecure reports EntropyUnavailable; see `random`
+//!
+//! Not implemented, and therefore a named panic: every `dir*`, `file*`, `net*`, `process*` and
+//! `child*` entry, `operate`, `batchAwait*`/`batchCancel`, the four `group*` entries,
+//! `lockStderr`/`tryLockStderr`/`unlockStderr`, and `progressParentFile`. So `std.Io.Group`,
+//! `std.Io.Select`, and any Reader or Writer that reaches the OS are out; the concurrency and
+//! synchronisation types listed above are in.
+//!
+//! ## The model
+//!
+//! One hart, no preemption, no timer interrupt. A task runs until it calls something that blocks -
+//! `sleep`, `futexWait`, `await`, `cancel`, or `yield` - and the scheduler is entered from inside
+//! that call. There is no scheduler task and no scheduler stack: the context that called
+//! `Runtime.init` is itself a task (the "main" slot, id 0), so a switch is always task-to-task and
+//! the first blocking call in `main` is what starts everything else. `io.async` only *assigns* a
+//! slot and marks it ready; the body first executes when the spawning context blocks or yields.
+//! That is within the contract - `async` promises a unit of concurrency, not that the work started
+//! (`Io.zig:54-58`) - and it is why a `main` that spawns seven tasks and then returns runs none of
+//! them.
+//!
+//! ## Why a futex reduces to "park this task" here, and where that stops being true
+//!
+//! A futex exists to close a race: between "I checked the value" and "I went to sleep", a waker can
+//! change the value and wake nothing, and the waiter sleeps forever. Linux closes it by doing the
+//! compare and the enqueue inside the kernel, atomically.
+//!
+//! On this chip the whole window is closed by the machine. One hart and no preemption means no
+//! other *task* can run between the compare and the enqueue - there is no instant at which control
+//! could pass to a waker - so a plain compare followed by a plain enqueue is already atomic with
+//! respect to every task. The only thing that can intervene is an **interrupt handler**, and
+//! ESP-Hosted does post from one (`_h_post_semaphore_from_isr`), so the compare-and-enqueue runs
+//! inside `hal.intr.mask()` - two CSR instructions - and `futexWake` runs inside the same critical
+//! section. That makes the pair atomic against handlers too, and `futexWake` safe to call from a
+//! CLIC handler.
+//!
+//! This stops being sound the moment either assumption goes:
+//!
+//! * **A second core.** The P4 has two HP cores and an LP core. A task on core 1 could observe
+//! the value, be preempted by nothing at all, and still race a waker on core 0, because
+//! masking interrupts on one hart says nothing about the other. A cross-core version needs a
+//! real wait queue per address plus an inter-core interrupt, not this.
+//! * **Preemption.** Add a timer interrupt that switches tasks and the "no other task can run"
+//! argument is gone; the mask covers it, but only because the mask is what the preemption
+//! would arrive through. Any preemption source that is not maskable here breaks it.
+//!
+//! Two further honest limits: waiters are found by scanning the task array rather than by a
+//! per-address wait queue (correct, O(tasks), and tasks are counted in single digits), and the scan
+//! is in slot order, so `futexWake(ptr, 1)` picks the lowest-numbered waiter rather than the
+//! longest-waiting one. Under a contended `std.Io.Mutex` that is a fairness wart, not a
+//! correctness one - the run queue itself is FIFO - but a heavily contended lock could starve a
+//! high-numbered slot.
+//!
+//! ## Static footprint
+//!
+//! Nothing here allocates, and nothing here is `comptime`-sized by a count. Measured on a real
+//! riscv32 build of this file (`footprint` below asserts both columns):
+//!
+//! @sizeOf(Task) 80 bytes per slot, in .bss - 12 of them are the suspended machine state
+//! @sizeOf(Runtime) 152 bytes including the embedded main task
+//! stack caller's whatever the slot declares
+//!
+//! On top of a slot's stack, and only while a slot is running an `io.async`/`io.concurrent` body,
+//! sit two copies carved off the top of that same stack: the argument tuple and the result. Both
+//! are a few tens of bytes for a typical call, and `claim` refuses a slot whose stack cannot hold
+//! them and still leave `min_stack_bytes`, in which case `async` runs the work eagerly instead.
+//!
+//! So a 7-task ESP-Hosted configuration with 5 KiB stacks is 35 KiB of stack + 560 bytes of slots
+//! + 152 bytes of runtime = 35.7 KiB, against the ~128 KiB of L2MEM.
+//!
+//! Code, from `nm` on a riscv32 `-OReleaseSmall` build of this file against the real `hal`, `soc`,
+//! `regs` and `mmio` modules: 2,402 bytes of out-of-line `.text` in 21 symbols for the scheduler
+//! and the sixteen implemented entries (`reschedule`, `claim`, `joinChild` and `futexBlock` are
+//! private and get inlined into their callers, so a caller pays a little more), plus 1,860 bytes
+//! of `.text.unlikely` and 3,722 bytes of `.rodata` strings for the 93 panicking stubs. The stub
+//! strings are the single largest line item here, and they are the price of a gap that announces
+//! itself: they only survive if the app's root declares a panic handler that prints the message,
+//! which `src/main.zig` does.
+//!
+//! One number worth having when sizing a stack: a task suspended inside the switch is holding
+//! about 160 bytes of its own stack for the spills the switch's clobber list forces, on top of
+//! whatever its own call chain is using (read off the `addi sp, sp, 160` frame of `reschedule` in
+//! the riscv32 asm). `stackUsed` reports the real high-water mark of any slot at any time, from
+//! the paint `init` lays down, and `dump` prints it for every task.
+//!
+//! ## Integrating this
+//!
+//! Module imports: `hal` (systimer, and `intr.mask` for the critical section), `soc` (`rom.print`
+//! and `cycles`), and `regs` + `mmio` (the one RNG register). All four already exist in build.zig.
+//!
+//! One requirement that is std's rather than this file's: the app's **root module** must declare
+//!
+//! pub const std_options: std.Options = .{ .page_size_min = 4096 };
+//!
+//! Five of the vtable's entries name `Io.File.MemoryMap`, whose `memory` field is
+//! `[]align(std.heap.page_size_min) u8` (`/usr/lib/zig/std/Io/File/MemoryMap.zig:18`), and
+//! riscv32-freestanding has no default for that (`/usr/lib/zig/std/heap.zig:48`). Defining a
+//! function whose return type is `File.MemoryMap.CreateError!File.MemoryMap` is what forces the
+//! struct to be laid out, so no stub body can dodge it, and any `std.Io` implementation on this
+//! target hits the same wall. It has to be the *root* module, because `std.options` is read as
+//! `@import("root").std_options` (`/usr/lib/zig/std/std.zig:112`) - putting the declaration in
+//! this file has no effect, which is measured rather than assumed. The value is arbitrary:
+//! nothing in this image pages, and it is only ever read for that one field's alignment.
+//!
+//! ## What is missing that a caller will notice
+//!
+//! `idle` spins instead of sleeping, because there is no timer interrupt to wake a `wfi` (see
+//! `chip.zig`). There is no priority: the run queue is FIFO. And nothing here bounds how long a
+//! task may hold the CPU, so one task that never blocks starves the rest - which is the deal with
+//! cooperative scheduling, and the reason `sleep` yields even when the deadline has already passed.
+
+const std = @import("std");
+const builtin = @import("builtin");
+const Io = std.Io;
+const Alignment = std.mem.Alignment;
+const assert = std.debug.assert;
+
+const context = @import("context.zig");
+const unimplemented = @import("unimplemented.zig");
+
+/// The chip, or the host that tests the chip's scheduler. Both sides implement the same eight
+/// declarations; see `chip.zig` for what each one costs on the die.
+const machine = if (builtin.os.tag == .freestanding) @import("chip.zig") else @import("host.zig");
+
+const Context = context.Context;
+const Switch = context.Switch;
+
+/// Written into every byte of a task's stack at `init`, so that `stackUsed` can report a real
+/// high-water mark and `checkStack` can catch an overflow at the next switch instead of at the next
+/// mystery. 0xC5 rather than 0 because a zeroed stack is indistinguishable from `.bss`.
+const paint: u8 = 0xC5;
+
+/// The least a slot's stack may be for `async` to use it, over and above the argument and result
+/// copies. Below this the slot is skipped and the work runs eagerly on the caller's stack, which is
+/// always legal (`Io.zig:54-56`) and better than a stack overflow.
+pub const min_stack_bytes: usize = 512;
+
+/// The runtime that owns the currently running task. A cooperative scheduler needs exactly one
+/// global: a task's entry function is jumped to with no arguments - the machine has nowhere to put
+/// them, see `context.zig` - so it has to find itself somehow, and this is the pointer it finds
+/// itself through. FreeRTOS spells the same thing `pxCurrentTCB`.
+///
+/// One `Runtime` per hart, therefore. Installing a second one while the first has live tasks is a
+/// programming error, and `init` asserts against it.
+var installed: ?*Runtime = null;
+
+pub const State = enum(u8) {
+ /// The slot holds no task. `async` may claim it.
+ free,
+ /// Runnable, and on the run queue.
+ ready,
+ /// Has the CPU.
+ running,
+ /// Waiting for a futex, a deadline, or a child task.
+ blocked,
+ /// Ran to completion; the result is in `result`, waiting for `await` or `cancel` to collect it.
+ done,
+};
+
+/// One unit of concurrency: the stack a task runs on, plus what the scheduler needs to know about
+/// it. The caller owns `stack` and nothing else here.
+pub const Task = struct {
+ /// The task's stack. Caller-provided and caller-sized, `.bss` or a static array; the scheduler
+ /// aligns the top down to the ABI's 16 bytes itself, so any byte slice will do.
+ ///
+ /// The main task's stack is empty: it runs on whatever stack `_start` established, and the
+ /// scheduler never needs to know where that is.
+ stack: []u8,
+
+ // -------- scheduler-private from here down --------
+
+ /// `sp`, `fp` and the resume address. Meaningful only while suspended.
+ ctx: Context = undefined,
+ state: State = .free,
+ /// Run queue link. Single-linked FIFO; the queue is only ever walked forwards.
+ next: ?*Task = null,
+ /// The address this task is parked on, if it is parked on a futex.
+ futex: ?*const u32 = null,
+ /// Machine ticks at which the scheduler must make this task ready again.
+ deadline: ?u64 = null,
+ /// The child this task is inside `await`/`cancel` for.
+ join: ?*Task = null,
+ /// Who to make ready when this task finishes.
+ awaiter: ?*Task = null,
+ /// Whether a cancel request should end the current block. False inside an uncancelable wait.
+ interruptible: bool = false,
+ /// A cancel request has been made and not yet delivered. See `Future.cancel`'s doc comment
+ /// (`Io.zig:1181-1190`): only the *next* cancelation point returns `error.Canceled`.
+ cancel_pending: bool = false,
+ /// A cancelation point in this task consumed a request. `await` needs this to decide whether a
+ /// cancelation it forwarded to a child was actually delivered.
+ cancel_acked: bool = false,
+ protection: Io.CancelProtection = .unblocked,
+ /// The type-erased body, its argument copy and its result storage, as `async` was handed them.
+ start: *const fn (context: *const anyopaque, result: *anyopaque) void = undefined,
+ arg: *const anyopaque = undefined,
+ result: []u8 = &.{},
+ /// Slot number, for diagnostics. 0 is the main task.
+ id: u8 = 0,
+};
+
+pub const Options = struct {
+ /// Nanoseconds to add to the systimer to answer `Clock.real`. Zero - the default - means
+ /// `Clock.real` reads as nanoseconds since the counter started, i.e. since boot: this board has
+ /// no RTC that survives a reset and no NTP. Set it once real time is known from outside (an
+ /// SNTP exchange, an HTTP `Date` header) and `Clock.real` becomes real.
+ real_epoch_offset_ns: i64 = 0,
+ /// How long the scheduler may find nothing runnable *and* no deadline pending before it
+ /// declares deadlock and panics with a dump of every task.
+ ///
+ /// That state is a deadlock unless an interrupt handler is about to wake somebody, which is a
+ /// legitimate thing to be waiting for - hence a timeout rather than an immediate panic. Five
+ /// seconds is long enough for any SDIO transaction on this board and short enough that a lost
+ /// wakeup is reported rather than looking like a hang.
+ deadlock_timeout_ms: u32 = 5_000,
+};
+
+pub const Runtime = struct {
+ /// Caller-provided slots, one per concurrent task. Index 0 of `tasks()` is `main`, not this.
+ slots: []Task,
+ /// The context that called `init`. Never started, only ever resumed.
+ main: Task,
+ current: *Task,
+ ready_head: ?*Task,
+ ready_tail: ?*Task,
+ real_epoch_offset_ns: i64,
+ deadlock_ticks: u64,
+ rng: Rng,
+
+ /// Take over as this hart's `std.Io`. `slots` must outlive the runtime, and so must `rt`
+ /// itself: the runtime points at its own `main` field.
+ ///
+ /// Brings the timebase up (`hal.systimer.init`, which deliberately does not reprogram the
+ /// clock source - see `chip.zig`). Touches nothing else on the chip: no interrupt is enabled,
+ /// no pad is configured, and mstatus.MIE is left exactly as the caller had it.
+ pub fn init(rt: *Runtime, slots: []Task, opts: Options) void {
+ comptime {
+ if (!context.supported) @compileError(
+ "io/p4.zig: no context switch for this architecture; see context.zig",
+ );
+ }
+ // `Task.id` is a u8 and 0 is the main task, so 254 slots is the ceiling. Anything close to
+ // it would have run out of L2MEM for stacks long before.
+ assert(slots.len < 255);
+ if (installed) |old| {
+ if (old != rt) {
+ var it = old.tasks();
+ while (it.next()) |t| assert(t.state == .free or t == &old.main);
+ }
+ }
+ machine.init();
+ rt.* = .{
+ .slots = slots,
+ .main = .{ .stack = &.{}, .state = .running, .id = 0 },
+ .current = undefined,
+ .ready_head = null,
+ .ready_tail = null,
+ .real_epoch_offset_ns = opts.real_epoch_offset_ns,
+ .deadlock_ticks = @as(u64, opts.deadlock_timeout_ms) * machine.ticks_hz / 1000,
+ .rng = .{},
+ };
+ rt.current = &rt.main;
+ for (slots, 0..) |*slot, i| {
+ const stack = slot.stack;
+ slot.* = .{ .stack = stack, .id = @intCast(i + 1) };
+ @memset(stack, paint);
+ }
+ installed = rt;
+ }
+
+ pub fn io(rt: *Runtime) Io {
+ return .{ .userdata = rt, .vtable = &vtable };
+ }
+
+ /// Put the current task at the back of the run queue and let everything else have a turn.
+ ///
+ /// Not a `std.Io` entry - there is none - but a cooperative scheduler needs a yield, and code
+ /// that only has an `Io` can get the same effect from `io.sleep(.zero, .awake)`, which is
+ /// documented to yield exactly once even for a deadline that has already passed.
+ pub fn yield(rt: *Runtime) void {
+ const me = rt.current;
+ {
+ const guard = machine.mask();
+ defer guard.release();
+ me.state = .ready;
+ rt.push(me);
+ }
+ rt.reschedule();
+ }
+
+ /// Bytes of `t`'s stack that have ever been written, from the paint laid down by `init`.
+ /// Zero for the main task, whose stack this scheduler does not own.
+ pub fn stackUsed(_: *const Runtime, t: *const Task) usize {
+ var i: usize = 0;
+ while (i < t.stack.len and t.stack[i] == paint) i += 1;
+ return t.stack.len - i;
+ }
+
+ pub fn stackFree(rt: *const Runtime, t: *const Task) usize {
+ return t.stack.len - rt.stackUsed(t);
+ }
+
+ /// Every task, main first.
+ pub fn tasks(rt: *Runtime) Iterator {
+ return .{ .rt = rt };
+ }
+
+ pub const Iterator = struct {
+ rt: *Runtime,
+ i: usize = 0,
+
+ pub fn next(it: *Iterator) ?*Task {
+ const n = it.i;
+ it.i += 1;
+ if (n == 0) return &it.rt.main;
+ if (n - 1 < it.rt.slots.len) return &it.rt.slots[n - 1];
+ return null;
+ }
+ };
+
+ /// One line per task, on the ROM console. Allocates nothing and takes no lock, so it is safe
+ /// from the crash path.
+ pub fn dump(rt: *Runtime) void {
+ machine.print("MARK IO_P4 current=%u slots=%u tick=%u\r\n", .{
+ @as(u32, rt.current.id),
+ @as(u32, @intCast(rt.slots.len)),
+ @as(u32, @truncate(machine.ticks())),
+ });
+ var it = rt.tasks();
+ while (it.next()) |t| {
+ machine.print(
+ "MARK IO_P4_TASK id=%u state=%u futex=0x%08x deadline=%u join=%u cancel=%u stack=%u/%u\r\n",
+ .{
+ @as(u32, t.id),
+ @as(u32, @intFromEnum(t.state)),
+ @as(u32, if (t.futex) |p| @truncate(@intFromPtr(p)) else 0),
+ @as(u32, @truncate(t.deadline orelse 0)),
+ @as(u32, if (t.join) |j| j.id else 255),
+ @as(u32, @intFromBool(t.cancel_pending)),
+ @as(u32, @intCast(rt.stackUsed(t))),
+ @as(u32, @intCast(t.stack.len)),
+ },
+ );
+ }
+ }
+
+ // ---------------------------------------------------------------- the run queue
+
+ /// Append to the run queue. Caller holds the interrupt mask.
+ fn push(rt: *Runtime, t: *Task) void {
+ t.next = null;
+ if (rt.ready_tail) |tail| tail.next = t else rt.ready_head = t;
+ rt.ready_tail = t;
+ }
+
+ /// Take the front of the run queue. Caller holds the interrupt mask.
+ fn pop(rt: *Runtime) ?*Task {
+ const t = rt.ready_head orelse return null;
+ rt.ready_head = t.next;
+ if (rt.ready_head == null) rt.ready_tail = null;
+ t.next = null;
+ return t;
+ }
+
+ /// Make a blocked task runnable. Caller holds the interrupt mask.
+ ///
+ /// Leaves `futex` and `deadline` set: the waiting task clears its own, after it resumes, and
+ /// the `.blocked` check here is what stops a second waker counting the same waiter twice.
+ fn wake(rt: *Runtime, t: *Task) void {
+ assert(t.state == .blocked);
+ t.state = .ready;
+ rt.push(t);
+ }
+
+ fn free(rt: *Runtime) ?*Task {
+ for (rt.slots) |*t| if (t.state == .free) return t;
+ return null;
+ }
+
+ /// Wake every task whose deadline has arrived. Caller holds the interrupt mask.
+ fn expire(rt: *Runtime, now_ticks: u64) void {
+ var it = rt.tasks();
+ while (it.next()) |t| {
+ if (t.state != .blocked) continue;
+ const d = t.deadline orelse continue;
+ if (now_ticks >= d) rt.wake(t);
+ }
+ }
+
+ /// The soonest deadline anybody is waiting for. Caller holds the interrupt mask.
+ fn earliest(rt: *Runtime) ?u64 {
+ var best: ?u64 = null;
+ var it = rt.tasks();
+ while (it.next()) |t| {
+ if (t.state != .blocked) continue;
+ const d = t.deadline orelse continue;
+ if (best == null or d < best.?) best = d;
+ }
+ return best;
+ }
+
+ // ---------------------------------------------------------------- the switch
+
+ /// Give up the CPU. The caller must already have put itself into the state it wants to be found
+ /// in: `.ready` and on the queue (a yield), or `.blocked` with a wake condition set.
+ ///
+ /// Returns when something switches back to this task. `rt.current` is set by whoever switches,
+ /// before the switch, so on the way back it is already correct.
+ fn reschedule(rt: *Runtime) void {
+ const me = rt.current;
+ assert(me.state != .running); // would be a task that gave up the CPU and stayed runnable
+ var idle_since: ?u64 = null;
+ while (true) {
+ const guard = machine.mask();
+ rt.expire(machine.ticks());
+ const next = rt.pop();
+ if (next) |n| {
+ if (n == me) {
+ // An interrupt handler re-queued us before we got here, or we were the only
+ // runnable task. Either way there is nothing to switch to.
+ me.state = .running;
+ guard.release();
+ return;
+ }
+ rt.current = n;
+ n.state = .running;
+ guard.release();
+ rt.checkStack(me);
+ var s: Switch = .{ .old = &me.ctx, .new = &n.ctx };
+ context.contextSwitch(&s);
+ return;
+ }
+ const deadline = rt.earliest();
+ guard.release();
+ rt.waitIdle(deadline, &idle_since);
+ }
+ }
+
+ /// Like `reschedule`, for a task that is never coming back. Its `Context` is still needed as
+ /// somewhere for the switch to write, and is dead the instant the switch completes.
+ fn retire(rt: *Runtime, me: *Task) noreturn {
+ var idle_since: ?u64 = null;
+ while (true) {
+ const guard = machine.mask();
+ rt.expire(machine.ticks());
+ const next = rt.pop();
+ if (next) |n| {
+ assert(n != me); // a finished task cannot be runnable
+ rt.current = n;
+ n.state = .running;
+ guard.release();
+ var s: Switch = .{ .old = &me.ctx, .new = &n.ctx };
+ context.contextSwitch(&s);
+ @panic("io.p4: a finished task was resumed");
+ }
+ const deadline = rt.earliest();
+ guard.release();
+ rt.waitIdle(deadline, &idle_since);
+ }
+ }
+
+ /// Nothing is runnable. With a deadline that is an ordinary sleep. Without one, the only thing
+ /// that can save this hart is an interrupt handler calling `futexWake` - a legitimate thing to
+ /// be waiting for, on a chip whose SDIO slave signals on a GPIO - so give it a bounded chance
+ /// and then say what happened, loudly, rather than spinning silently forever.
+ fn waitIdle(rt: *Runtime, deadline: ?u64, idle_since: *?u64) void {
+ if (deadline == null) {
+ const now_ticks = machine.ticks();
+ if (idle_since.*) |since| {
+ if (now_ticks -% since > rt.deadlock_ticks) rt.deadlock();
+ } else {
+ idle_since.* = now_ticks;
+ }
+ } else {
+ idle_since.* = null;
+ }
+ machine.idle(deadline);
+ }
+
+ fn checkStack(rt: *Runtime, t: *Task) void {
+ if (t.stack.len == 0) return;
+ if (t.stack[0] == paint) return;
+ machine.print("MARK IO_P4_STACK_OVERFLOW id=%u size=%u\r\n", .{
+ @as(u32, t.id), @as(u32, @intCast(t.stack.len)),
+ });
+ rt.dump();
+ @panic("io.p4: task stack overflow");
+ }
+
+ fn deadlock(rt: *Runtime) noreturn {
+ machine.print("MARK IO_P4_DEADLOCK no runnable task and no deadline\r\n", .{});
+ rt.dump();
+ @panic("io.p4: deadlock - every task is blocked and nothing is scheduled to wake them");
+ }
+
+ // ---------------------------------------------------------------- time
+
+ fn clockOffset(rt: *const Runtime, clock: Io.Clock) i96 {
+ return switch (clock) {
+ .real => rt.real_epoch_offset_ns,
+ // The systimer's own epoch. `awake` and `boot` are the same thing on a chip that does
+ // not suspend, and there is no per-task CPU accounting, so the two `cpu_*` clocks
+ // answer with elapsed time; `clockResolution` reports 0 for them to say so.
+ .awake, .boot, .cpu_process, .cpu_thread => 0,
+ };
+ }
+
+ /// Absolute machine tick a `Timeout` expires at, or null for `.none`.
+ fn deadlineOf(rt: *const Runtime, timeout: Io.Timeout) ?u64 {
+ return switch (timeout) {
+ .none => null,
+ .duration => |d| machine.ticks() +| ticksFromNs(d.raw.nanoseconds),
+ .deadline => |ts| ticksFromNs(ts.raw.nanoseconds - rt.clockOffset(ts.clock)),
+ };
+ }
+
+ // ---------------------------------------------------------------- tasks
+
+ /// Reserve a slot for `start`, with its argument tuple and result storage carved off the top of
+ /// the slot's own stack. Returns null if there is no free slot, or the slot's stack is too
+ /// small to hold the copies and still be a stack - in both cases the caller's fallback is to
+ /// run the work eagerly.
+ fn claim(
+ rt: *Runtime,
+ result_len: usize,
+ result_alignment: Alignment,
+ arg: []const u8,
+ arg_alignment: Alignment,
+ start: *const fn (context: *const anyopaque, result: *anyopaque) void,
+ ) ?*Task {
+ const guard = machine.mask();
+ defer guard.release();
+
+ const t = rt.free() orelse return null;
+
+ const base = @intFromPtr(t.stack.ptr);
+ const result_at = result_alignment.backward(base + t.stack.len - result_len);
+ const arg_at = arg_alignment.backward(result_at - arg.len);
+ if (arg_at < base or arg_at - base < min_stack_bytes + context.initial_stack_overhead)
+ return null;
+
+ const arg_ptr: [*]u8 = @ptrFromInt(arg_at);
+ @memcpy(arg_ptr[0..arg.len], arg);
+ const result_ptr: [*]u8 = @ptrFromInt(result_at);
+
+ t.state = .ready;
+ t.next = null;
+ t.futex = null;
+ t.deadline = null;
+ t.join = null;
+ t.awaiter = null;
+ t.interruptible = false;
+ t.cancel_pending = false;
+ t.cancel_acked = false;
+ t.protection = .unblocked;
+ t.start = start;
+ t.arg = @ptrCast(arg_ptr);
+ t.result = result_ptr[0..result_len];
+ t.ctx = context.initial(arg_at, taskEntry);
+ rt.push(t);
+ return t;
+ }
+
+ /// Wait for `child` to finish, collect its result, and give the slot back.
+ ///
+ /// `request` is `cancel` rather than `await`. The subtle case is the other one: a cancel
+ /// request that lands on *us* while we are waiting. `await` has no error to return it through,
+ /// so the request is forwarded to the child instead, and afterwards re-armed on us if the child
+ /// never actually observed it - which is what `Threaded.await` does with `recancelInner`
+ /// (`/usr/lib/zig/std/Io/Threaded.zig:2465-2473`).
+ fn joinChild(rt: *Runtime, child: *Task, result: []u8, request: bool) void {
+ const me = rt.current;
+ assert(child != me);
+ assert(child.result.len == result.len);
+
+ if (request) rt.requestCancel(child);
+ var inherited = false;
+
+ while (true) {
+ if (!request and !inherited and consumeCancel(me)) {
+ inherited = true;
+ rt.requestCancel(child);
+ }
+ const guard = machine.mask();
+ if (child.state == .done) {
+ guard.release();
+ break;
+ }
+ child.awaiter = me;
+ me.join = child;
+ me.interruptible = !(request or inherited);
+ me.state = .blocked;
+ guard.release();
+ rt.reschedule();
+ me.join = null;
+ }
+
+ if (result.len != 0) @memcpy(result, child.result);
+ const acked = child.cancel_acked;
+ child.awaiter = null;
+ child.state = .free;
+ if (inherited and !acked) me.cancel_pending = true;
+ }
+
+ /// Arm a cancel request on `t`, and end its current wait if that wait is cancelable.
+ fn requestCancel(rt: *Runtime, t: *Task) void {
+ const guard = machine.mask();
+ defer guard.release();
+ switch (t.state) {
+ .free, .done => return,
+ .ready, .running, .blocked => {},
+ }
+ t.cancel_pending = true;
+ if (t.state == .blocked and t.interruptible) rt.wake(t);
+ }
+
+ /// Park until somebody wakes `ptr`, the deadline arrives, or - if `interruptible` - a cancel
+ /// request lands. Returns without parking if the value already differs from `expected`, which
+ /// is a futex's `EAGAIN` and is what keeps a waker that got there first from being lost.
+ fn futexBlock(
+ rt: *Runtime,
+ me: *Task,
+ ptr: *const u32,
+ expected: u32,
+ deadline: ?u64,
+ interruptible: bool,
+ ) void {
+ const guard = machine.mask();
+ if (@atomicLoad(u32, ptr, .acquire) != expected) {
+ guard.release();
+ return;
+ }
+ if (deadline) |d| if (machine.ticks() >= d) {
+ guard.release();
+ return;
+ };
+ me.futex = ptr;
+ me.deadline = deadline;
+ me.interruptible = interruptible;
+ me.state = .blocked;
+ guard.release();
+ rt.reschedule();
+ me.futex = null;
+ me.deadline = null;
+ }
+};
+
+/// Where a task begins. Jumped to, not called: there is no return address and no argument, so the
+/// task finds itself through `installed` (see that declaration for why).
+fn taskEntry() callconv(.c) noreturn {
+ const rt = installed.?;
+ const me = rt.current;
+ me.start(me.arg, @ptrCast(me.result.ptr));
+
+ {
+ const guard = machine.mask();
+ defer guard.release();
+ me.state = .done;
+ if (me.awaiter) |a| if (a.state == .blocked and a.join == me) rt.wake(a);
+ }
+ rt.checkStack(me);
+ rt.retire(me);
+}
+
+/// True if `t` has an undelivered cancel request that is not blocked by cancel protection, in which
+/// case the request is consumed: only the *next* cancelation point reports it (`Io.zig:1185-1188`).
+fn consumeCancel(t: *Task) bool {
+ if (t.protection == .blocked) return false;
+ if (!t.cancel_pending) return false;
+ t.cancel_pending = false;
+ t.cancel_acked = true;
+ return true;
+}
+
+fn cast(userdata: ?*anyopaque) *Runtime {
+ return @ptrCast(@alignCast(userdata));
+}
+
+// -------------------------------------------------------------------- tick arithmetic
+
+/// Machine ticks to nanoseconds, exactly. 16 MHz is 62.5 ns per tick, so the conversion is
+/// `* 125 / 2` and not a multiply by 62 - which would drift by 0.8%, i.e. 43 seconds a day.
+fn nsFromTicks(t: u64) i96 {
+ return @divTrunc(@as(i96, t) * 125, 2);
+}
+
+/// Nanoseconds to machine ticks, rounded **up**, so a deadline is never reached early. Saturates
+/// rather than wrapping or trapping: `Duration.max` is `maxInt(i96)` nanoseconds, which is a
+/// request never to wake, and the largest tick count this returns is 36,000 years at 16 MHz.
+fn ticksFromNs(ns: i96) u64 {
+ if (ns <= 0) return 0;
+ const limit: i96 = @divFloor(@as(i96, std.math.maxInt(u64)) * 125, 2);
+ if (ns >= limit) return std.math.maxInt(u64);
+ return @intCast(@divFloor(ns * 2 + 124, 125));
+}
+
+// -------------------------------------------------------------------- entropy
+
+/// The random source behind `random`.
+///
+/// `WDEV_RND_REG` is real on this die - it is `LP_SYSTEM_REG_RNG_DATA_REG`, and `esp_random` reads
+/// exactly it (`chip.zig`) - but the *entropy* behind it is not established for this image, and
+/// this file will not claim otherwise. ESP-IDF's own documentation is precise about the conditions
+/// (`docs/en/api-reference/system/random.rst:37-46`): the hardware produces true random numbers
+/// while the RF subsystem is enabled, or while `bootloader_random_enable` has the SAR ADC entropy
+/// source on, or while the IDF second-stage bootloader is running - "if none of the above
+/// conditions are true, the output of the RNG should be considered as pseudo-random only". The P4
+/// has no RF at all, this image is not started by the IDF bootloader, and nothing here calls
+/// `bootloader_random_enable`. There is a hardware secondary source that is always mixed in - an
+/// asynchronous oscillator sampled for metastability (`random.rst:98-101`) - which is why the
+/// register is worth reading at all, but "always mixed in" is not the same as "measured on this
+/// board", and I cannot measure this board.
+///
+/// So: **`random` is not cryptographic**, and says so. It is a PCG32 (`std.Random.Pcg`) seeded from
+/// four hardware words mixed with the systimer and the cycle counter, re-stirred with a fresh
+/// hardware word whenever the pacing interval has elapsed. `randomSecure` reports
+/// `error.EntropyUnavailable` rather than pretending. Closing that gap is a small, separate job:
+/// port `bootloader_random_enable`, then measure.
+const Rng = struct {
+ prng: std.Random.Pcg = undefined,
+ seeded: bool = false,
+ last_sample: u64 = 0,
+
+ /// ESP-IDF paces `esp_random` at one *byte* per 14 microseconds on this part -
+ /// `APB_CYCLE_WAIT_NUM` is `CONFIG_ESP_DEFAULT_CPU_FREQ_MHZ * 14` APB cycles per byte
+ /// (`hw_random.c:41-44`), which at the P4's 1:1 CPU:APB ratio is the ~75 kHz byte rate the
+ /// comment there says the RNG was tested at. Draining it faster returns the generator's state
+ /// rather than new entropy. Expressed in systimer ticks because that is the clock this file
+ /// trusts, and because the constant in IDF assumes a 360 MHz CPU that this board is not
+ /// running.
+ const sample_ticks: u64 = machine.ticks_hz * 14 / 1_000_000;
+
+ fn fill(r: *Rng, buffer: []u8) void {
+ if (!r.seeded) r.seed();
+ r.stir();
+ r.prng.fill(buffer);
+ }
+
+ /// Four paced hardware words, which costs ~56 microseconds of spinning, once, on the first call
+ /// to `random`. Deliberately not done in `init`: a runtime that never asks for randomness
+ /// should not pay for it, and `init` runs before anything can tolerate a delay.
+ fn seed(r: *Rng) void {
+ var s: u64 = machine.noise() ^ (machine.ticks() << 17);
+ for (0..4) |_| {
+ s = s *% 0x9E37_79B9_7F4A_7C15 ^ r.sample();
+ }
+ r.prng = .init(s);
+ r.seeded = true;
+ }
+
+ /// Mix in a fresh hardware word if one is due. Never waits: the caller asked for random bytes,
+ /// not for a delay, and the generator is sound without it.
+ fn stir(r: *Rng) void {
+ const now_ticks = machine.ticks();
+ if (now_ticks -% r.last_sample < sample_ticks) return;
+ r.last_sample = now_ticks;
+ r.prng.s ^= machine.entropyWord();
+ }
+
+ /// One hardware word, waiting out the pacing interval first.
+ fn sample(r: *Rng) u32 {
+ const target = r.last_sample +% sample_ticks;
+ while (machine.ticks() -% r.last_sample < sample_ticks) machine.idle(target);
+ r.last_sample = machine.ticks();
+ return machine.entropyWord();
+ }
+};
+
+// -------------------------------------------------------------------- the vtable
+
+/// Start from "nothing is implemented" and overwrite what is. The list of assignments below is the
+/// authoritative answer to "what works": anything not named here panics with its own name.
+const vtable: Io.VTable = build: {
+ var v = unimplemented.vtable;
+
+ v.now = now;
+ v.clockResolution = clockResolution;
+ v.sleep = sleep;
+
+ v.async = asyncTask;
+ v.concurrent = concurrent;
+ v.await = awaitTask;
+ v.cancel = cancelTask;
+
+ v.futexWait = futexWait;
+ v.futexWaitUncancelable = futexWaitUncancelable;
+ v.futexWake = futexWake;
+
+ v.checkCancel = checkCancel;
+ v.recancel = recancel;
+ v.swapCancelProtection = swapCancelProtection;
+
+ v.crashHandler = crashHandler;
+ v.random = random;
+ v.randomSecure = randomSecure;
+
+ break :build v;
+};
+
+fn now(userdata: ?*anyopaque, clock: Io.Clock) Io.Timestamp {
+ const rt = cast(userdata);
+ return .{ .nanoseconds = nsFromTicks(machine.ticks()) + rt.clockOffset(clock) };
+}
+
+fn clockResolution(_: ?*anyopaque, clock: Io.Clock) Io.Clock.ResolutionError!Io.Duration {
+ return switch (clock) {
+ // 1e9/16e6 = 62.5 ns, truncated: `Duration` counts whole nanoseconds and the tick is not a
+ // whole number of them. The conversion in `nsFromTicks` keeps the half.
+ .real, .awake, .boot => .fromNanoseconds(@divTrunc(std.time.ns_per_s, @as(i96, machine.ticks_hz))),
+ // Zero means "unsupported" (`Io.zig:787-790`). There is one process and no per-task CPU
+ // accounting, so `now` answers these with elapsed time and this says not to believe it.
+ .cpu_process, .cpu_thread => .zero,
+ };
+}
+
+fn sleep(userdata: ?*anyopaque, timeout: Io.Timeout) Io.Cancelable!void {
+ // `.none` means no timeout at all, i.e. nothing to wait for. Matches `Threaded.sleep`
+ // (`Threaded.zig:11577`), and notably is *not* a cancelation point.
+ if (timeout == .none) return;
+
+ const rt = cast(userdata);
+ const me = rt.current;
+ if (consumeCancel(me)) return error.Canceled;
+ const deadline = rt.deadlineOf(timeout).?;
+
+ while (true) {
+ {
+ const guard = machine.mask();
+ defer guard.release();
+ if (machine.ticks() >= deadline) {
+ // The deadline has already passed, but a `sleep` that returns without ever leaving
+ // the CPU makes `io.sleep(.zero, .awake)` useless as a yield - and it is the only
+ // yield `std.Io` exposes. So go round the run queue exactly once.
+ me.deadline = null;
+ me.interruptible = true;
+ me.state = .ready;
+ rt.push(me);
+ } else {
+ me.deadline = deadline;
+ me.interruptible = true;
+ me.state = .blocked;
+ }
+ }
+ rt.reschedule();
+ me.deadline = null;
+ if (consumeCancel(me)) return error.Canceled;
+ if (machine.ticks() >= deadline) return;
+ }
+}
+
+fn asyncTask(
+ userdata: ?*anyopaque,
+ result: []u8,
+ result_alignment: Alignment,
+ arg: []const u8,
+ arg_alignment: Alignment,
+ start: *const fn (context: *const anyopaque, result: *anyopaque) void,
+) ?*Io.AnyFuture {
+ const rt = cast(userdata);
+ const t = rt.claim(result.len, result_alignment, arg, arg_alignment, start) orelse {
+ // No slot: run it here and now. `await` will be a no-op (`Io.zig:54-56`).
+ start(arg.ptr, result.ptr);
+ return null;
+ };
+ return @ptrCast(t);
+}
+
+fn concurrent(
+ userdata: ?*anyopaque,
+ result_len: usize,
+ result_alignment: Alignment,
+ arg: []const u8,
+ arg_alignment: Alignment,
+ start: *const fn (context: *const anyopaque, result: *anyopaque) void,
+) Io.ConcurrentError!*Io.AnyFuture {
+ const rt = cast(userdata);
+ // `concurrent` promises the caller can block on something the task will unblock. A cooperative
+ // task can do that - it runs whenever the caller blocks - so the only failure is running out of
+ // slots. Unlike `async` there is no eager fallback: running the body inline is exactly the
+ // guarantee `concurrent` exists to rule out.
+ const t = rt.claim(result_len, result_alignment, arg, arg_alignment, start) orelse
+ return error.ConcurrencyUnavailable;
+ return @ptrCast(t);
+}
+
+fn awaitTask(userdata: ?*anyopaque, any_future: *Io.AnyFuture, result: []u8, _: Alignment) void {
+ const rt = cast(userdata);
+ rt.joinChild(@ptrCast(@alignCast(any_future)), result, false);
+}
+
+fn cancelTask(userdata: ?*anyopaque, any_future: *Io.AnyFuture, result: []u8, _: Alignment) void {
+ const rt = cast(userdata);
+ rt.joinChild(@ptrCast(@alignCast(any_future)), result, true);
+}
+
+fn futexWait(userdata: ?*anyopaque, ptr: *const u32, expected: u32, timeout: Io.Timeout) Io.Cancelable!void {
+ const rt = cast(userdata);
+ const me = rt.current;
+ // Cancelation is checked before the value, which is the order `Threaded` gets from calling
+ // `Syscall.start()` before the futex syscall (`Threaded.zig:986`): an already-canceled task
+ // does not acquire a lock on its way out.
+ if (consumeCancel(me)) return error.Canceled;
+ rt.futexBlock(me, ptr, expected, rt.deadlineOf(timeout), true);
+ if (consumeCancel(me)) return error.Canceled;
+}
+
+fn futexWaitUncancelable(userdata: ?*anyopaque, ptr: *const u32, expected: u32) void {
+ const rt = cast(userdata);
+ rt.futexBlock(rt.current, ptr, expected, null, false);
+}
+
+fn futexWake(userdata: ?*anyopaque, ptr: *const u32, max_waiters: u32) void {
+ const rt = cast(userdata);
+ const guard = machine.mask();
+ defer guard.release();
+ var woken: u32 = 0;
+ var it = rt.tasks();
+ while (it.next()) |t| {
+ if (woken >= max_waiters) break;
+ if (t.state != .blocked) continue;
+ if (t.futex != ptr) continue;
+ rt.wake(t);
+ woken += 1;
+ }
+}
+
+fn checkCancel(userdata: ?*anyopaque) Io.Cancelable!void {
+ if (consumeCancel(cast(userdata).current)) return error.Canceled;
+}
+
+fn recancel(userdata: ?*anyopaque) void {
+ const me = cast(userdata).current;
+ assert(!me.cancel_pending); // `recancel` with a request already pending
+ me.cancel_pending = true;
+}
+
+fn swapCancelProtection(userdata: ?*anyopaque, new: Io.CancelProtection) Io.CancelProtection {
+ const me = cast(userdata).current;
+ const old = me.protection;
+ me.protection = new;
+ return old;
+}
+
+/// Called from `std.debug`'s panic path (`/usr/lib/zig/std/debug.zig:536`) and from the segfault
+/// handler (`:1641`), both of which go on to print the panic themselves.
+///
+/// So this prints and **returns**; it does not park. Parking here would swallow the panic message,
+/// which is the one thing worth having. No allocation, no lock, no scheduling: `dump` is straight
+/// `ets_printf`, and marking the crashing task cancel-protected keeps any cleanup that runs after
+/// this from being interrupted by a cancelation - the same thing `Threaded.crashHandler` does
+/// (`Threaded.zig:2066-2072`).
+fn crashHandler(userdata: ?*anyopaque) void {
+ const rt = cast(userdata);
+ rt.current.cancel_pending = false;
+ rt.current.protection = .blocked;
+ machine.print("MARK IO_P4_CRASH\r\n", .{});
+ rt.dump();
+}
+
+fn random(userdata: ?*anyopaque, buffer: []u8) void {
+ cast(userdata).rng.fill(buffer);
+}
+
+/// See `Rng`: the hardware register is real, its entropy on this image is not established, and
+/// guessing is worse than saying so.
+fn randomSecure(_: ?*anyopaque, _: []u8) Io.RandomSecureError!void {
+ return error.EntropyUnavailable;
+}
+
+// -------------------------------------------------------------------- static declaration helper
+
+/// A whole runtime as one `.bss` object: `task_count` stacks of `stack_bytes` each, the slots that
+/// describe them, and the `Runtime`.
+///
+/// ```
+/// var pool: p4.Static(7, 5 * 1024) = .{};
+/// const io = pool.init(.{}).io();
+/// ```
+pub fn Static(comptime task_count: usize, comptime stack_bytes: usize) type {
+ comptime {
+ if (stack_bytes < min_stack_bytes) @compileError("stack_bytes is smaller than min_stack_bytes");
+ if (stack_bytes % context.stack_align != 0) @compileError(
+ "stack_bytes must be a multiple of the ABI's stack alignment, so that every slot's " ++
+ "stack top is aligned without waste",
+ );
+ }
+ return struct {
+ stacks: [task_count][stack_bytes]u8 align(context.stack_align) = undefined,
+ slots: [task_count]Task = undefined,
+ runtime: Runtime = undefined,
+
+ /// Total `.bss` this declaration costs.
+ pub const size = @sizeOf(@This());
+
+ pub fn init(self: *@This(), opts: Options) *Runtime {
+ for (&self.slots, &self.stacks) |*slot, *stack| slot.* = .{ .stack = stack };
+ self.runtime.init(&self.slots, opts);
+ return &self.runtime;
+ }
+ };
+}
+
+// -------------------------------------------------------------------- tests
+//
+// These run on the host, against `host.zig`'s virtual clock and `std.Io.fiber`'s x86-64 context
+// switch, which is the same scheduler with a different two files underneath it. What they cannot
+// test is the riscv32 asm in `context.zig` - that only runs on the die.
+
+const testing = std.testing;
+
+/// Deliberately file-scope: 8 KiB stacks in a test function's frame is not what a stack is for.
+/// Unreferenced outside tests, so it is not analysed - let alone emitted - in a firmware build.
+var test_pool: Static(4, 8 * 1024) = .{};
+
+fn testRuntime() *Runtime {
+ return test_pool.init(.{});
+}
+
+const Trace = struct {
+ buf: [64]u8 = undefined,
+ len: usize = 0,
+
+ fn put(t: *Trace, c: u8) void {
+ t.buf[t.len] = c;
+ t.len += 1;
+ }
+ fn seen(t: *const Trace) []const u8 {
+ return t.buf[0..t.len];
+ }
+};
+
+/// The numbers quoted in this file's header. Not a tuning knob - a regression alarm: `Task` growing
+/// silently is how a 35 KiB budget becomes 40. The chip's figures are the 32-bit column, and they
+/// were read off a real riscv32 build; the host runs the 64-bit column so that this test is not
+/// vacuous where it can actually run.
+pub const footprint = struct {
+ pub const task_bytes: usize = if (@sizeOf(usize) == 4) 80 else 128;
+ pub const runtime_bytes: usize = if (@sizeOf(usize) == 4) 152 else 216;
+ pub const context_bytes: usize = 3 * @sizeOf(usize);
+};
+
+test footprint {
+ try testing.expectEqual(footprint.task_bytes, @sizeOf(Task));
+ try testing.expectEqual(footprint.runtime_bytes, @sizeOf(Runtime));
+ try testing.expectEqual(footprint.context_bytes, @sizeOf(Context));
+}
+
+fn tick(io: Io, trace: *Trace, mark: u8, rounds: usize) void {
+ for (0..rounds) |_| {
+ trace.put(mark);
+ io.sleep(.zero, .awake) catch return;
+ }
+}
+
+test "three tasks round-robin through the scheduler" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ var a = io.async(tick, .{ io, &trace, 'a', 3 });
+ var b = io.async(tick, .{ io, &trace, 'b', 3 });
+ var c = io.async(tick, .{ io, &trace, 'c', 3 });
+ a.await(io);
+ b.await(io);
+ c.await(io);
+
+ // Spawn order is queue order, and each task yields after every mark, so the interleaving is
+ // exact. `main` awaits `a` first and so is not in the rotation.
+ try testing.expectEqualStrings("abcabcabc", trace.seen());
+}
+
+fn locker(io: Io, m: *Io.Mutex, trace: *Trace, mark: u8) void {
+ m.lock(io) catch return;
+ defer m.unlock(io);
+ trace.put(mark);
+ // Hold the lock across a yield, so the other task must actually block on the futex rather than
+ // finding it free.
+ io.sleep(.zero, .awake) catch {};
+ trace.put(std.ascii.toUpper(mark));
+}
+
+test "std.Io.Mutex serialises two tasks through the futex" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+ var m: Io.Mutex = .init;
+
+ var a = io.async(locker, .{ io, &m, &trace, 'a' });
+ var b = io.async(locker, .{ io, &m, &trace, 'b' });
+ a.await(io);
+ b.await(io);
+
+ // Interleaved would be "abAB"; serialised is each task's pair adjacent.
+ try testing.expectEqualStrings("aAbB", trace.seen());
+ try testing.expect(m.tryLock());
+}
+
+fn producer(io: Io, q: *Io.Queue(u32), n: u32) void {
+ var i: u32 = 0;
+ while (i < n) : (i += 1) q.putOne(io, i) catch return;
+ q.close(io);
+}
+
+fn consumer(io: Io, q: *Io.Queue(u32), sum: *u32) void {
+ while (true) {
+ const v = q.getOne(io) catch return;
+ sum.* += v;
+ }
+}
+
+test "std.Io.Queue passes items between two tasks" {
+ const rt = testRuntime();
+ const io = rt.io();
+ // Capacity 2 against 10 items, so both directions block and both directions get woken.
+ var storage: [2]u32 = undefined;
+ var q: Io.Queue(u32) = .init(&storage);
+ var sum: u32 = 0;
+
+ var p = io.async(producer, .{ io, &q, 10 });
+ var c = io.async(consumer, .{ io, &q, &sum });
+ p.await(io);
+ c.await(io);
+
+ try testing.expectEqual(@as(u32, 45), sum);
+}
+
+fn napper(io: Io, trace: *Trace, mark: u8, ms: u64) void {
+ io.sleep(.fromMilliseconds(@intCast(ms)), .awake) catch return;
+ trace.put(mark);
+}
+
+test "sleep orders three tasks by deadline, not by spawn order" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ const t0 = Io.Timestamp.now(io, .awake);
+ var a = io.async(napper, .{ io, &trace, 'a', @as(u64, 30) });
+ var b = io.async(napper, .{ io, &trace, 'b', @as(u64, 10) });
+ var c = io.async(napper, .{ io, &trace, 'c', @as(u64, 20) });
+ a.await(io);
+ b.await(io);
+ c.await(io);
+
+ try testing.expectEqualStrings("bca", trace.seen());
+ // And the clock really moved, by at least the longest sleep.
+ const elapsed = t0.durationTo(Io.Timestamp.now(io, .awake));
+ try testing.expect(elapsed.toMilliseconds() >= 30);
+}
+
+fn napAndCatch(io: Io, trace: *Trace) void {
+ trace.put('s');
+ io.sleep(.fromMilliseconds(1000), .awake) catch |err| switch (err) {
+ error.Canceled => {
+ trace.put('x');
+ return;
+ },
+ };
+ trace.put('e'); // must not be reached
+}
+
+test "a canceled task does not run to completion" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ var f = io.async(napAndCatch, .{ io, &trace });
+ // The body has not run yet - `async` only assigned a slot - so let it reach its sleep first.
+ rt.yield();
+ try testing.expectEqualStrings("s", trace.seen());
+
+ f.cancel(io);
+ try testing.expectEqualStrings("sx", trace.seen());
+}
+
+test "cancel before the body ever runs still runs it, and its first cancelation point reports" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ var f = io.async(napAndCatch, .{ io, &trace });
+ f.cancel(io);
+ // Same as `Threaded`: the function always runs; cancelation is delivered at the first
+ // cancelation point inside it.
+ try testing.expectEqualStrings("sx", trace.seen());
+}
+
+fn protected(io: Io, trace: *Trace) void {
+ const old = io.swapCancelProtection(.blocked);
+ io.sleep(.fromMilliseconds(5), .awake) catch unreachable; // protection is on: cannot fail
+ trace.put('p');
+ _ = io.swapCancelProtection(old);
+ io.checkCancel() catch {
+ trace.put('x');
+ return;
+ };
+ trace.put('e');
+}
+
+test "cancel protection defers delivery to the next unprotected point" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ var f = io.async(protected, .{ io, &trace });
+ f.cancel(io);
+ try testing.expectEqualStrings("px", trace.seen());
+}
+
+fn waiter(io: Io, word: *std.atomic.Value(u32), trace: *Trace) void {
+ while (word.load(.acquire) == 0) {
+ io.futexWait(u32, &word.raw, 0) catch return;
+ }
+ trace.put('w');
+}
+
+fn poster(io: Io, word: *std.atomic.Value(u32), trace: *Trace) void {
+ trace.put('p');
+ word.store(1, .release);
+ io.futexWake(u32, &word.raw, 1);
+}
+
+test "futexWait parks and futexWake releases exactly one waiter" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+ var word: std.atomic.Value(u32) = .init(0);
+
+ var w = io.async(waiter, .{ io, &word, &trace });
+ var p = io.async(poster, .{ io, &word, &trace });
+ w.await(io);
+ p.await(io);
+
+ try testing.expectEqualStrings("pw", trace.seen());
+}
+
+test "futexWait returns immediately when the value already differs" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var word: std.atomic.Value(u32) = .init(7);
+ // Would deadlock if it parked: nobody is going to wake it.
+ try io.futexWait(u32, &word.raw, 0);
+}
+
+test "futexWaitTimeout unblocks on its deadline with nothing else runnable" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var word: std.atomic.Value(u32) = .init(0);
+
+ // Nobody is going to wake this, so the deadline has to - which also drives the scheduler's
+ // "nothing runnable, one deadline pending" idle path from the main task, the same path the
+ // deadlock watchdog must *not* fire on.
+ const t0 = Io.Timestamp.now(io, .awake);
+ try io.futexWaitTimeout(u32, &word.raw, 0, .{ .duration = .{
+ .raw = .fromMilliseconds(25),
+ .clock = .awake,
+ } });
+ const elapsed = t0.durationTo(Io.Timestamp.now(io, .awake));
+ try testing.expect(elapsed.toMilliseconds() >= 25);
+ try testing.expectEqual(@as(u32, 0), word.load(.monotonic));
+}
+
+fn add(a: u32, b: u32) u32 {
+ return a + b;
+}
+
+test "async falls back to running eagerly when every slot is taken" {
+ const rt = testRuntime();
+ const io = rt.io();
+
+ // Fill all four slots with tasks that will not finish until woken.
+ var word: std.atomic.Value(u32) = .init(0);
+ var trace: Trace = .{};
+ var held: [4]Io.Future(void) = undefined;
+ for (&held) |*f| f.* = io.async(waiter, .{ io, &word, &trace });
+
+ var eager = io.async(add, .{ 20, 22 });
+ try testing.expectEqual(@as(?*Io.AnyFuture, null), eager.any_future);
+ try testing.expectEqual(@as(u32, 42), eager.await(io));
+
+ word.store(1, .release);
+ io.futexWake(u32, &word.raw, 4);
+ for (&held) |*f| f.await(io);
+ try testing.expectEqualStrings("wwww", trace.seen());
+}
+
+test "concurrent reports ConcurrencyUnavailable instead of running inline" {
+ const rt = testRuntime();
+ const io = rt.io();
+
+ var word: std.atomic.Value(u32) = .init(0);
+ var trace: Trace = .{};
+ var held: [4]Io.Future(void) = undefined;
+ for (&held) |*f| f.* = io.async(waiter, .{ io, &word, &trace });
+
+ try testing.expectError(error.ConcurrencyUnavailable, io.concurrent(add, .{ 1, 2 }));
+
+ word.store(1, .release);
+ io.futexWake(u32, &word.raw, 4);
+ for (&held) |*f| f.await(io);
+}
+
+test "now is monotonic and converts ticks exactly" {
+ const rt = testRuntime();
+ const io = rt.io();
+
+ const a = Io.Timestamp.now(io, .awake);
+ try io.sleep(.fromMilliseconds(7), .awake);
+ const b = Io.Timestamp.now(io, .awake);
+ try testing.expect(b.nanoseconds >= a.nanoseconds);
+ try testing.expect(a.durationTo(b).toMilliseconds() >= 7);
+
+ // 62.5 ns a tick, kept exact: one tick is 62 ns and two are 125, not 124.
+ try testing.expectEqual(@as(i96, 62), nsFromTicks(1));
+ try testing.expectEqual(@as(i96, 125), nsFromTicks(2));
+ try testing.expectEqual(@as(i96, 1_000_000_000), nsFromTicks(machine.ticks_hz));
+ // And back, rounding up so a deadline is never early.
+ try testing.expectEqual(@as(u64, 1), ticksFromNs(1));
+ try testing.expectEqual(@as(u64, 1), ticksFromNs(62));
+ try testing.expectEqual(@as(u64, 2), ticksFromNs(63));
+ try testing.expectEqual(@as(u64, machine.ticks_hz), ticksFromNs(1_000_000_000));
+ try testing.expectEqual(@as(u64, std.math.maxInt(u64)), ticksFromNs(Io.Duration.max.nanoseconds));
+
+ // `real` is the same counter plus an offset the caller sets; `cpu_*` report resolution 0 to say
+ // they are not really implemented.
+ try testing.expectEqual(Io.Duration{ .nanoseconds = 62 }, try io.vtable.clockResolution(io.userdata, .awake));
+ try testing.expectEqual(Io.Duration.zero, try io.vtable.clockResolution(io.userdata, .cpu_thread));
+ rt.real_epoch_offset_ns = 1_700_000_000 * std.time.ns_per_s;
+ try testing.expect(Io.Timestamp.now(io, .real).toSeconds() > 1_600_000_000);
+ rt.real_epoch_offset_ns = 0;
+}
+
+test "sleep(.none) is not a wait and not a cancelation point" {
+ const rt = testRuntime();
+ const io = rt.io();
+ rt.current.cancel_pending = true;
+ try io.vtable.sleep(io.userdata, .none);
+ try testing.expect(rt.current.cancel_pending);
+ rt.current.cancel_pending = false;
+}
+
+test "random fills, is not all zero, and does not repeat itself" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var a: [32]u8 = @splat(0);
+ var b: [32]u8 = @splat(0);
+ io.random(&a);
+ io.random(&b);
+ try testing.expect(!std.mem.allEqual(u8, &a, 0));
+ try testing.expect(!std.mem.eql(u8, &a, &b));
+ // The one thing this port will not pretend about.
+ try testing.expectError(error.EntropyUnavailable, io.randomSecure(&a));
+}
+
+test "a task's stack high-water mark is measurable" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var f = io.async(add, .{ 1, 2 });
+ _ = f.await(io);
+ // Slot 1 ran `add` through the trampoline; something was written, and nowhere near 8 KiB.
+ const used = rt.stackUsed(&rt.slots[0]);
+ try testing.expect(used > 0);
+ try testing.expect(used < 8 * 1024);
+ try testing.expectEqual(@as(usize, 0), rt.stackUsed(&rt.main));
+}
+
+fn childAcks(io: Io, trace: *Trace) void {
+ io.sleep(.fromMilliseconds(1000), .awake) catch {
+ trace.put('c');
+ return;
+ };
+ trace.put('C');
+}
+
+fn childNeverAcks(io: Io, trace: *Trace) void {
+ // Blocks, so the parent really has to wait for it, but never observes the cancelation - which
+ // is the case where the parent has to take its own request back.
+ const old = io.swapCancelProtection(.blocked);
+ io.sleep(.fromMilliseconds(5), .awake) catch unreachable;
+ _ = io.swapCancelProtection(old);
+ trace.put('C');
+}
+
+fn parentAwaits(io: Io, trace: *Trace, acks: bool) void {
+ var f = if (acks) io.async(childAcks, .{ io, trace }) else io.async(childNeverAcks, .{ io, trace });
+ f.await(io);
+ trace.put('r');
+ io.checkCancel() catch {
+ trace.put('x');
+ return;
+ };
+ trace.put('e');
+}
+
+test "a cancel that lands during await is forwarded to the child" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ var f = io.async(parentAwaits, .{ io, &trace, true });
+ f.cancel(io);
+
+ // 'c': the child's sleep reported `error.Canceled`, so the child consumed the request.
+ // 'r': the parent's `await` returned. 'e': and the parent's own next cancelation point is
+ // *clean*, because the cancelation was delivered into the child rather than to the parent.
+ try testing.expectEqualStrings("cre", trace.seen());
+}
+
+test "a cancel the child never acknowledges is re-armed on the parent" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+
+ var f = io.async(parentAwaits, .{ io, &trace, false });
+ f.cancel(io);
+
+ // 'C': the child ran to completion under cancel protection. 'r': await returned. 'x': and the
+ // request the parent gave away is back, so the parent's next cancelation point reports it -
+ // `Threaded.await` does the same thing through `recancelInner` (`Threaded.zig:2470`).
+ try testing.expectEqualStrings("Crx", trace.seen());
+}
+
+fn condWaiter(io: Io, m: *Io.Mutex, c: *Io.Condition, flag: *bool, trace: *Trace) void {
+ m.lock(io) catch return;
+ defer m.unlock(io);
+ while (!flag.*) c.wait(io, m) catch return;
+ trace.put('w');
+}
+
+fn condSignaler(io: Io, m: *Io.Mutex, c: *Io.Condition, flag: *bool, trace: *Trace) void {
+ m.lock(io) catch return;
+ flag.* = true;
+ m.unlock(io);
+ trace.put('s');
+ c.signal(io);
+}
+
+test "std.Io.Condition signals across tasks" {
+ const rt = testRuntime();
+ const io = rt.io();
+ var trace: Trace = .{};
+ var m: Io.Mutex = .init;
+ var c: Io.Condition = .init;
+ var flag = false;
+
+ // `Condition` is the most futex-dependent thing in std: an epoch word, a packed
+ // waiters/signals state, and a mutex handed back and forth across the wait.
+ var w = io.async(condWaiter, .{ io, &m, &c, &flag, &trace });
+ var s = io.async(condSignaler, .{ io, &m, &c, &flag, &trace });
+ w.await(io);
+ s.await(io);
+
+ try testing.expectEqualStrings("sw", trace.seen());
+ try testing.expect(m.tryLock());
+}