//! `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()); }