diff options
Diffstat (limited to 'src/io/p4.zig')
| -rw-r--r-- | src/io/p4.zig | 1475 |
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()); +} |
