diff options
| author | Gabriel Schneider <[email protected]> | 2026-08-25 12:40:53 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-08-25 12:46:51 -0300 |
| commit | f5f8068fac59b4f16046c2022c2fc7c7e447ef4c (patch) | |
| tree | 2731a3ed4e51cae09e184e25778eded5fc37d1f5 /src/net/hosted_os.zig | |
| download | esp32p4-f5f8068fac59b4f16046c2022c2fc7c7e447ef4c.tar.gz esp32p4-f5f8068fac59b4f16046c2022c2fc7c7e447ef4c.zip | |
zig-p4: pure-Zig ESP32-P4 toolchain
build.zig generates the linker script and drives Zig's own LLD; tools/image.zig
turns the ELF into a flashable image and tools/{rom,serial}.zig speak the mask
ROM loader over the UART. No CMake, ninja, idf.py, esptool, or external linker.
src/soc.zig is a comptime register model over ESP-IDF's own *_reg.h headers;
src/hal/ adds peripheral sequences; src/io/ implements std.Io for the chip;
src/oracle/ diffs this HAL against ESP-IDF's on the die.
Diffstat (limited to 'src/net/hosted_os.zig')
| -rw-r--r-- | src/net/hosted_os.zig | 890 |
1 files changed, 890 insertions, 0 deletions
diff --git a/src/net/hosted_os.zig b/src/net/hosted_os.zig new file mode 100644 index 0000000..d844177 --- /dev/null +++ b/src/net/hosted_os.zig @@ -0,0 +1,890 @@ +//! ESP-Hosted's OS objects - mutex, counting semaphore, fixed-capacity queue, thread, software +//! timer - expressed in `std.Io`, with FreeRTOS's exact observable behaviour. +//! +//! Nothing here reimplements a synchronisation primitive. `std.Io.Mutex`, `std.Io.Semaphore` and +//! `std.Io.TypeErasedQueue` do the blocking; this file supplies only the three things ESP-Hosted +//! needs that they do not have: +//! +//! 1. **The timeout dialect.** `_h_lock_mutex`, `_h_get_semaphore` and `_h_dequeue_item` all take +//! an `int`, where 0 means "do not block", a negative value means "block forever", and a +//! positive value means a bounded wait. The unit of that positive value is *not* the same in +//! all three - see `Wait`. +//! 2. **The return codes.** `RET_OK`/`RET_FAIL`/`RET_INVALID`/`RET_FAIL_TIMEOUT` from +//! `port_esp_hosted_host_os.h:86-91`, which the C caller branches on. +//! 3. **The initial state.** A FreeRTOS semaphore created by +//! `hosted_create_semaphore` (`port_esp_hosted_host_os.c:523-547`) is given *once* before it +//! is returned, so it starts with one permit, and callers rely on that: `sdio_drv.c:1504`, +//! `:1508` and `:1540` each take that permit back immediately after creating the semaphore. A +//! semaphore that started at zero would leave every count in the transport off by one. +//! +//! This file is deliberately free of hardware and of the C ABI, so it runs on the host under +//! `std.Io.Threaded` and the tests below are real tests. + +const std = @import("std"); +const assert = std.debug.assert; +const Io = std.Io; +const Allocator = std.mem.Allocator; + +/// `port_esp_hosted_host_os.h:86-91`. +pub const ret = struct { + pub const ok: c_int = 0; + pub const fail: c_int = -1; + pub const invalid: c_int = -2; + pub const fail_mem: c_int = -3; + pub const fail4: c_int = -4; + pub const fail_timeout: c_int = -5; +}; + +/// The clock everything here measures against. `.awake` is `std.Io`'s monotonic clock; on this +/// board it is `hal.systimer`'s fixed 16 MHz, which does not move when the CPU clock does. +pub const clock: Io.Clock = .awake; + +/// How often a bounded wait re-checks. +/// +/// `std.Io.Semaphore` and `std.Io.TypeErasedQueue` have no timed acquire, and neither does +/// `std.Io.Mutex`; only `futexWaitTimeout` does, and reaching for it would mean rebuilding those +/// three primitives instead of using them. So a *bounded* wait polls, and an unbounded one - which +/// is what every hot path in ESP-Hosted actually uses - blocks properly with no polling at all. +/// +/// The cost is bounded and small: a bounded wait is used in exactly one place in the tree, +/// `rpc_core.c:844`, the synchronous-RPC response wait, whose timeout is measured in seconds. One +/// millisecond of added latency on a request that is allowed to take five seconds is not worth a +/// hand-rolled futex semaphore. +pub const poll_interval_ms: u32 = 1; + +/// The three shapes an ESP-Hosted timeout argument can take. +pub const Wait = union(enum) { + /// `0` - try, do not block. + immediate, + /// Negative, i.e. `HOSTED_BLOCKING` (-1) or `HOSTED_BLOCK_MAX` (`portMAX_DELAY`, which reaches + /// an `int` parameter as -1). + forever, + /// A bounded wait, in milliseconds. + bounded_ms: u32, + + /// The dialect used by `_h_lock_mutex` and `_h_get_semaphore`: a positive value is + /// milliseconds (`port_esp_hosted_host_os.c:452`, `:573`). + pub fn fromMillis(timeout: c_int) Wait { + if (timeout == 0) return .immediate; + if (timeout < 0) return .forever; + return .{ .bounded_ms = @intCast(timeout) }; + } + + /// The dialect used by `_h_dequeue_item`: a positive value is *seconds*, because the + /// implementation converts it with `SEC_TO_MILLISEC` before `pdMS_TO_TICKS` + /// (`port_esp_hosted_host_os.c:336`). The asymmetry with `fromMillis` is not a mistake in this + /// file; it is a mistake in ESP-Hosted that this file has to reproduce. No caller in the tree + /// passes a positive value to a queue, so nothing depends on it today. + pub fn fromQueueTimeout(timeout: c_int) Wait { + if (timeout == 0) return .immediate; + if (timeout < 0) return .forever; + return .{ .bounded_ms = @as(u32, @intCast(timeout)) *| 1000 }; + } +}; + +/// Milliseconds on `clock`, which is what `_h_get_time_ms` returns. +/// +/// The narrowing to `u64` before the division is not cosmetic. `Io.Timestamp.nanoseconds` is `i96`, +/// and `@divFloor` on an `i96` compiles to a call to compiler_rt's `__divti3` - a 128-bit software +/// division, on every call, on a 90 MHz in-order core. Narrowing first turns that into +/// `__udivdi3`, a 64-bit one. Both were read out of the object's undefined-symbol list rather than +/// guessed; ReleaseSmall declines to strength-reduce even a constant 64-bit divisor, so the +/// libcall stays, but it is now half the width. The range given up is imaginary: 2^64 nanoseconds +/// is 584 years, and this clock starts at boot. +pub fn nowMs(io: Io) u64 { + const ns = clock.now(io).nanoseconds; + if (ns <= 0) return 0; + const ns64: u64 = @intCast(ns); + return ns64 / std.time.ns_per_ms; +} + +/// Sleep one poll interval. Reports cancelation so bounded waits abandon promptly rather than +/// spinning out the full timeout after the task has been asked to stop. +fn pollSleep(io: Io) error{Canceled}!void { + return io.sleep(.fromMilliseconds(poll_interval_ms), clock); +} + +// --------------------------------------------------------------------------------------- Mutex + +/// FreeRTOS gives ESP-Hosted a *recursive-capable* mutex handle but ESP-Hosted never recurses on +/// one: every use is a bracketed `SDIO_LOCK`/`SDIO_UNLOCK` or equivalent, and every one of the ten +/// call sites in the tree passes `HOSTED_BLOCK_MAX`. So a plain `std.Io.Mutex` is the whole +/// requirement. +pub const Mutex = struct { + inner: Io.Mutex = .init, + + pub fn lock(m: *Mutex, io: Io, w: Wait) c_int { + switch (w) { + .immediate => return if (m.inner.tryLock()) ret.ok else ret.fail, + .forever => { + m.inner.lockUncancelable(io); + return ret.ok; + }, + .bounded_ms => |ms| { + const deadline = nowMs(io) + ms; + while (true) { + if (m.inner.tryLock()) return ret.ok; + if (nowMs(io) >= deadline) return ret.fail; + pollSleep(io) catch return ret.fail; + } + }, + } + } + + pub fn unlock(m: *Mutex, io: Io) c_int { + m.inner.unlock(io); + return ret.ok; + } +}; + +// ----------------------------------------------------------------------------------- Semaphore + +/// A counting semaphore with FreeRTOS's cap and FreeRTOS's initial count. +/// +/// The blocking path is `std.Io.Semaphore` untouched. What is added around it: +/// +/// * a **maximum count**, because `xSemaphoreCreateCounting(maxCount, 0)` refuses a give past +/// `maxCount` and `std.Io.Semaphore` has no ceiling. `sdio_drv.c:1502` sizes +/// `sem_to_slave_queue` at `tx_queue_size * MAX_PRIORITY_QUEUES` precisely so that the +/// semaphore saturates when the queues do. +/// * a **non-blocking take**, which `std.Io.Semaphore` does not expose. It is the tail of +/// `Semaphore.wait` (`std/Io/Semaphore.zig:18-24`) with the `cond.wait` loop removed, using +/// the same public fields, so it takes and releases the same mutex in the same order. +/// * an **ISR-deferred post**; see `postFromIsr`. +pub const Semaphore = struct { + inner: Io.Semaphore, + max: u32, + /// Posts an interrupt handler could not deliver directly. Folded in by the next task-side + /// operation on this semaphore. + isr_posts: std.atomic.Value(u32) = .init(0), + + /// `maxCount` as ESP-Hosted passes it: `<= 1` means a binary semaphore. + /// + /// Starts with one permit, matching `port_esp_hosted_host_os.c:544` - see the file header. + pub fn init(max_count: u32) Semaphore { + return .{ .inner = .{ .permits = 1 }, .max = @max(max_count, 1) }; + } + + pub fn post(s: *Semaphore, io: Io) c_int { + // Fold in anything an interrupt deferred, so every task-side entry point closes that + // window and not just the waiting ones. Two instructions when nothing is pending. + s.drainIsrPosts(io); + return if (s.add(io, 1) == 1) ret.ok else ret.fail; + } + + /// `_h_post_semaphore_from_isr`, and the one entry in the whole table whose FreeRTOS meaning + /// does not survive the move to a cooperative scheduler intact. + /// + /// FreeRTOS has `xSemaphoreGiveFromISR`, which manipulates the semaphore inside a port-level + /// critical section and then asks for a context switch on return from the interrupt. Neither + /// half exists here. `std.Io.Semaphore.post` takes the semaphore's own `Io.Mutex`, and an + /// interrupt that blocked on a mutex held by the task it interrupted would deadlock the core - + /// there is no other task to run and no preemption to run it. + /// + /// What is safe on this runtime, confirmed with the runtime's author: `io.futexWake` runs + /// inside a critical section that clears `mstatus.MIE`, touches only the run queue, and never + /// takes a task-held lock. `Io.Mutex.tryLock` is a single compare-exchange. So: + /// + /// * if the mutex is free, the post happens inline and completely. On a single core with + /// interrupts already masked, no task can observe the intermediate state. + /// * if the mutex is held, the interrupted task is *running* and holds it - `Io.Condition` + /// releases the mutex before it blocks (`std/Io.zig:1689`), so nobody ever sleeps holding + /// it. The post is recorded in `isr_posts` and folded in by that task's next operation on + /// this semaphore, which is a few instructions away. + /// + /// The residual hole: if the interrupt lands in that few-instruction window *and* the only + /// other participant is already blocked in `wait`, the deferred post sits until someone else + /// touches the semaphore. `drainIsrPosts` exists so an application can close it from an + /// interrupt epilogue. On the SDIO transport this is moot: `_h_post_semaphore_from_isr` has + /// exactly two callers in the tree, `spi_drv.c:181` and `:190`, plus `spi_hd_drv.c:174`, and + /// none of them is compiled for SDIO. + /// + /// Returns `RET_OK` if the post was delivered or deferred, never fails: an interrupt has + /// nowhere to report a failure to. + pub fn postFromIsr(s: *Semaphore, io: Io) c_int { + if (s.inner.mutex.tryLock()) { + defer s.inner.mutex.unlock(io); + if (s.inner.permits < s.max) { + s.inner.permits += 1; + s.inner.cond.signal(io); + } + return ret.ok; + } + _ = s.isr_posts.fetchAdd(1, .release); + return ret.ok; + } + + /// Fold any interrupt-deferred posts into the semaphore. Safe and cheap to call from a task at + /// any time; a no-op when nothing is pending. + pub fn drainIsrPosts(s: *Semaphore, io: Io) void { + const pending = s.isr_posts.swap(0, .acquire); + if (pending != 0) _ = s.add(io, pending); + } + + /// Add `n` permits, saturating at `max`. Returns how many were actually added. + fn add(s: *Semaphore, io: Io, n: u32) u32 { + s.inner.mutex.lockUncancelable(io); + defer s.inner.mutex.unlock(io); + const room = s.max -| @as(u32, @intCast(s.inner.permits)); + const added = @min(room, n); + if (added == 0) return 0; + s.inner.permits += added; + // One signal per permit: `Io.Condition.signal` releases exactly one waiter. + for (0..added) |_| s.inner.cond.signal(io); + return added; + } + + /// `_h_get_semaphore`. Returns 0 on success and `RET_FAIL_TIMEOUT` otherwise, which is what + /// `port_esp_hosted_host_os.c:577-579` returns and what `rpc_core.c:844` tests. + pub fn wait(s: *Semaphore, io: Io, w: Wait) c_int { + switch (w) { + .immediate => return if (s.tryTake(io)) ret.ok else ret.fail_timeout, + .forever => { + s.drainIsrPosts(io); + // Cancelation is reported as RET_FAIL_TIMEOUT, which is the only failure code + // `hosted_get_semaphore` ever returns (port_esp_hosted_host_os.c:579) and + // therefore the only one callers test for. + s.inner.wait(io) catch return ret.fail_timeout; + return ret.ok; + }, + .bounded_ms => |ms| { + const deadline = nowMs(io) + ms; + while (true) { + if (s.tryTake(io)) return ret.ok; + if (nowMs(io) >= deadline) return ret.fail_timeout; + pollSleep(io) catch return ret.fail; + } + }, + } + } + + /// Take a permit if one is available. The body is `Semaphore.wait` + /// (`std/Io/Semaphore.zig:18-24`) minus its `cond.wait` loop. + pub fn tryTake(s: *Semaphore, io: Io) bool { + s.drainIsrPosts(io); + s.inner.mutex.lockUncancelable(io); + defer s.inner.mutex.unlock(io); + if (s.inner.permits == 0) return false; + s.inner.permits -= 1; + if (s.inner.permits > 0) s.inner.cond.signal(io); + return true; + } + + pub fn count(s: *Semaphore, io: Io) u32 { + s.inner.mutex.lockUncancelable(io); + defer s.inner.mutex.unlock(io); + return @intCast(s.inner.permits); + } +}; + +// --------------------------------------------------------------------------------------- Queue + +/// A fixed-capacity queue of runtime-sized items. +/// +/// `_h_create_queue(qnum_elem, qitem_size)` fixes the element size at *run* time, so +/// `std.Io.Queue(Elem)` - which needs the type at compile time - cannot be used, but +/// `std.Io.TypeErasedQueue` can: it is a byte ring with `min`-byte put and get, which is exactly a +/// queue of fixed-size records once every operation moves `item_size` bytes. +/// +/// That the ring only ever moves whole items is what makes the non-blocking forms exact. The +/// buffer is `count * item_size` bytes and every transfer is `item_size`, so the occupied length is +/// always a multiple of `item_size`; a `min = 0` put therefore either fits the whole item or moves +/// nothing at all, and can never leave half a record in the ring. +pub const Queue = struct { + inner: Io.TypeErasedQueue, + item_size: u32, + /// Owned; freed by `destroy`. + buffer: []u8, + + pub fn create(gpa: Allocator, count: u32, item_size: u32) ?*Queue { + assert(item_size > 0); + const q = gpa.create(Queue) catch return null; + const buf = gpa.alloc(u8, @as(usize, count) * item_size) catch { + gpa.destroy(q); + return null; + }; + q.* = .{ .inner = .init(buf), .item_size = item_size, .buffer = buf }; + return q; + } + + pub fn destroy(q: *Queue, io: Io, gpa: Allocator) void { + q.inner.close(io); + gpa.free(q.buffer); + gpa.destroy(q); + } + + /// `_h_queue_item`. `item` points at one `item_size` record, which is copied into the queue - + /// FreeRTOS's `xQueueSendToBack` copies too, which is why every caller passes `&handle` rather + /// than a heap pointer. + pub fn send(q: *Queue, io: Io, item: [*]const u8, w: Wait) c_int { + const n = q.item_size; + const slice = item[0..n]; + switch (w) { + .immediate => { + const put = q.inner.put(io, slice, 0) catch return ret.fail; + return if (put == n) ret.ok else ret.fail; + }, + .forever => { + // Uncancelable, deliberately. A cancelable blocking put can be interrupted + // *after* it has copied part of a record into the ring, and since the ring's + // occupied length is what makes the non-blocking forms exact, a half record + // desynchronises every subsequent transfer. FreeRTOS's portMAX_DELAY does not + // return early either. The cost is that a canceled task blocked here stays + // blocked - which it would anyway: `Future.cancel` signals only the *next* + // cancelation point, and ESP-Hosted's task bodies loop straight back into the + // queue. See `Thread.cancel`. + const put = q.inner.putUncancelable(io, slice, n) catch return ret.fail; + return if (put == n) ret.ok else ret.fail; + }, + .bounded_ms => |ms| { + const deadline = nowMs(io) + ms; + while (true) { + const put = q.inner.put(io, slice, 0) catch return ret.fail; + if (put == n) return ret.ok; + assert(put == 0); // a partial record would corrupt the ring + if (nowMs(io) >= deadline) return ret.fail; + pollSleep(io) catch return ret.fail; + } + }, + } + } + + /// `_h_dequeue_item`. Returns 0 on success, `RET_FAIL` on timeout - note the asymmetry with + /// `Semaphore.wait`, which returns `RET_FAIL_TIMEOUT`; `port_esp_hosted_host_os.c:342` really + /// does return the plain failure code here. + pub fn receive(q: *Queue, io: Io, out: [*]u8, w: Wait) c_int { + const n = q.item_size; + const slice = out[0..n]; + switch (w) { + .immediate => { + const got = q.inner.get(io, slice, 0) catch return ret.fail; + return if (got == n) ret.ok else ret.fail; + }, + .forever => { + // Uncancelable for the same reason as `send`. + const got = q.inner.getUncancelable(io, slice, n) catch return ret.fail; + return if (got == n) ret.ok else ret.fail; + }, + .bounded_ms => |ms| { + const deadline = nowMs(io) + ms; + while (true) { + const got = q.inner.get(io, slice, 0) catch return ret.fail; + if (got == n) return ret.ok; + assert(got == 0); + if (nowMs(io) >= deadline) return ret.fail; + pollSleep(io) catch return ret.fail; + } + }, + } + } + + /// `_h_queue_msg_waiting` = `uxQueueMessagesWaiting`, which counts *buffered* items only and + /// not producers blocked with an item in hand. + pub fn waiting(q: *Queue, io: Io) c_int { + q.inner.mutex.lockUncancelable(io); + defer q.inner.mutex.unlock(io); + return @intCast(q.inner.len / q.item_size); + } + + /// `_h_reset_queue` = `xQueueReset`: discard everything buffered. Blocked producers and + /// consumers are left alone, which is also what FreeRTOS does for waiting *receivers*; it + /// differs in that FreeRTOS re-evaluates blocked senders. No caller in the tree uses this. + pub fn reset(q: *Queue, io: Io) c_int { + q.inner.mutex.lockUncancelable(io); + defer q.inner.mutex.unlock(io); + q.inner.start = 0; + q.inner.len = 0; + return ret.ok; + } +}; + +// -------------------------------------------------------------------------------------- Thread + +/// ESP-Hosted's task entry point: `void (*)(void const *)`, called once and never expected to +/// return (`port_esp_hosted_host_os.c:163`, and every body in the tree is a `while (1)` loop). +pub const StartRoutine = *const fn (?*const anyopaque) callconv(.c) void; + +pub const Thread = struct { + future: Io.Future(void), + name: [*:0]const u8, + + fn trampoline(start: StartRoutine, arg: ?*const anyopaque) void { + start(arg); + } + + /// `io.concurrent`, not `io.async`, and the difference is the whole point of the entry. + /// + /// `xTaskCreate` returns a task that exists and will run whatever its creator does next. + /// `io.async` promises less: the implementation is allowed to run the body inline before + /// returning, which for an ESP-Hosted task body - an unconditional `while (1)` - would never + /// return and would deadlock initialisation on the spot. `io.concurrent` forbids exactly that + /// (`std/Io.zig:2358-2364`) and reports `error.ConcurrencyUnavailable` when no unit of + /// concurrency is free. + /// + /// Turning that error into NULL is right: `_h_thread_create` is documented to return NULL on + /// failure and its callers check (`rpc_core.c:582`, `sdio_drv.c:1543`). A task pool one slot + /// too small then produces a legible "thread creation failed" instead of a hang. + pub fn create(io: Io, gpa: Allocator, name: [*:0]const u8, start: StartRoutine, arg: ?*const anyopaque) ?*Thread { + const t = gpa.create(Thread) catch return null; + t.* = .{ + .future = io.concurrent(trampoline, .{ start, arg }) catch { + gpa.destroy(t); + return null; + }, + .name = name, + }; + return t; + } + + /// `_h_thread_cancel` maps to `Future.cancel`, and this is the second place FreeRTOS's model + /// does not fit. + /// + /// `vTaskDelete` destroys a task from outside, wherever it happens to be. `std.Io`'s cancel is + /// cooperative: it asks, then *waits for the task body to return*. ESP-Hosted's task bodies + /// never return - `sdio_read_task`, `rpc_rx_thread` and the rest are unconditional loops - so + /// this call completes only if the body happens to exit, and otherwise blocks. + /// + /// That is survivable because of where it is called from: `cancel_rpc_threads` + /// (`rpc_core.c`) and the transport teardown paths, both of which run only when the host is + /// about to restart the slave. It is not survivable as a routine operation, and if a teardown + /// path becomes routine the fix is a `killTask` on the runtime that reclaims the slot without + /// unwinding, not a change here: there is no way to unwind a C frame from Zig. + pub fn cancel(t: *Thread, io: Io, gpa: Allocator) c_int { + t.future.cancel(io); + gpa.destroy(t); + return ret.ok; + } +}; + +// ------------------------------------------------------------------------------- software timers + +pub const TimerHandler = *const fn (?*anyopaque) callconv(.c) void; + +pub const TimerKind = enum(c_int) { + /// `H_TIMER_TYPE_ONESHOT`, port_esp_hosted_host_os.h:39. + oneshot = 0, + /// `H_TIMER_TYPE_PERIODIC`. + periodic = 1, +}; + +pub const Timer = struct { + handler: TimerHandler = undefined, + arg: ?*anyopaque = null, + /// Absolute deadline on `clock`, in milliseconds. + deadline_ms: u64 = 0, + /// 0 for a one-shot. + period_ms: u32 = 0, + in_use: bool = false, +}; + +/// One task servicing every software timer, rather than one task per timer. +/// +/// ESP-IDF backs `_h_timer_start` with `esp_timer`, which has its own dedicated task. Doing the +/// obvious thing here - `io.async` per timer - would cost one whole task slot and one whole static +/// stack per timer, and ESP-Hosted starts up to three concurrently: the slave-unresponsive timer +/// (`transport_drv.c:188`), the per-request asynchronous RPC timeout (`rpc_core.c:215`), and the +/// power-save timer. At the stack sizes this runtime needs that is 15 KB to run three sleeps. +/// +/// So: one task, an array of slots, and a futex word that a `start` or `stop` bumps to make the +/// service task recompute its next deadline. Static footprint is `@sizeOf(Timer)` (24 bytes on +/// rv32) per slot plus one task stack. +/// +/// Handlers run on the service task, not in an interrupt, so they may block. `init_timeout_cb` +/// (`transport_drv.c`) calls `_h_restart_host`, which never returns, and that is fine here. +pub fn TimerService(comptime slot_count: usize) type { + return struct { + const Self = @This(); + + slots: [slot_count]Timer = @splat(.{}), + /// Bumped whenever a slot is armed or disarmed; the service task waits on it. + epoch: std.atomic.Value(u32) = .init(0), + /// Guards `slots`. A plain `Io.Mutex`: every critical section here is a few dozen + /// instructions and never blocks. + mutex: Io.Mutex = .init, + task: ?*Thread = null, + stopping: bool = false, + + pub fn start(self: *Self, io: Io, gpa: Allocator) bool { + if (self.task != null) return true; + const t = gpa.create(Thread) catch return false; + // Concurrent for the same reason as `Thread.create`: the service loop never returns, + // so an implementation permitted to run it inline would never return from `start`. + t.* = .{ + .future = io.concurrent(service, .{ self, io }) catch { + gpa.destroy(t); + return false; + }, + .name = "hosted_timers", + }; + self.task = t; + return true; + } + + /// Arm a slot. Returns its index, or null when every slot is in use. + pub fn arm(self: *Self, io: Io, ms: u32, kind: TimerKind, handler: TimerHandler, arg: ?*anyopaque) ?usize { + self.mutex.lockUncancelable(io); + const idx = blk: { + for (&self.slots, 0..) |*s, i| if (!s.in_use) break :blk i; + self.mutex.unlock(io); + return null; + }; + self.slots[idx] = .{ + .handler = handler, + .arg = arg, + .deadline_ms = nowMs(io) + ms, + .period_ms = if (kind == .periodic) ms else 0, + .in_use = true, + }; + self.mutex.unlock(io); + self.kick(io); + return idx; + } + + pub fn disarm(self: *Self, io: Io, idx: usize) c_int { + if (idx >= slot_count) return ret.invalid; + self.mutex.lockUncancelable(io); + const was = self.slots[idx].in_use; + self.slots[idx].in_use = false; + self.mutex.unlock(io); + self.kick(io); + return if (was) ret.ok else ret.fail; + } + + fn kick(self: *Self, io: Io) void { + _ = self.epoch.fetchAdd(1, .release); + io.futexWake(u32, &self.epoch.raw, 1); + } + + fn service(self: *Self, io: Io) void { + while (!self.stopping) { + const seen = self.epoch.load(.acquire); + const now = nowMs(io); + + // Fire everything due, collecting the handlers first so none of them runs while + // the slot table is locked: a handler may arm or disarm a timer. + var due: [slot_count]struct { h: TimerHandler, a: ?*anyopaque } = undefined; + var due_len: usize = 0; + var next_deadline: ?u64 = null; + + self.mutex.lockUncancelable(io); + for (&self.slots) |*s| { + if (!s.in_use) continue; + if (s.deadline_ms <= now) { + due[due_len] = .{ .h = s.handler, .a = s.arg }; + due_len += 1; + if (s.period_ms == 0) { + s.in_use = false; + } else { + s.deadline_ms = now + s.period_ms; + } + } + if (s.in_use) { + if (next_deadline == null or s.deadline_ms < next_deadline.?) + next_deadline = s.deadline_ms; + } + } + self.mutex.unlock(io); + + for (due[0..due_len]) |d| d.h(d.a); + if (due_len != 0) continue; + + if (next_deadline) |dl| { + const remaining = dl -| nowMs(io); + io.futexWaitTimeout(u32, &self.epoch.raw, seen, .{ + .duration = .{ .clock = clock, .raw = .fromMilliseconds(@intCast(remaining)) }, + }) catch return; + } else { + io.futexWait(u32, &self.epoch.raw, seen) catch return; + } + } + } + + pub fn stop(self: *Self, io: Io, gpa: Allocator) void { + const t = self.task orelse return; + self.stopping = true; + self.kick(io); + t.future.cancel(io); + gpa.destroy(t); + self.task = null; + } + }; +} + +// ---------------------------------------------------------------------------------------- tests + +const testing = std.testing; + +fn hostIo() struct { threaded: *Io.Threaded, io: Io } { + const t = testing.allocator.create(Io.Threaded) catch unreachable; + t.* = .init(testing.allocator, .{}); + return .{ .threaded = t, .io = t.io() }; +} + +test "Semaphore starts with one permit, as FreeRTOS's create+give does" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + + // sdio_drv.c:1502-1504 creates a counting semaphore and immediately takes the permit that + // hosted_create_semaphore left behind. If the count started at zero this take would fail and + // every subsequent count would be one too high. + var s = Semaphore.init(60); + try testing.expectEqual(@as(u32, 1), s.count(io)); + try testing.expectEqual(ret.ok, s.wait(io, .immediate)); + try testing.expectEqual(@as(u32, 0), s.count(io)); + + // Empty: a non-blocking take reports RET_FAIL_TIMEOUT, which is the code rpc_core.c:844 tests. + try testing.expectEqual(ret.fail_timeout, s.wait(io, .immediate)); +} + +test "Semaphore counts, saturates at max, and times out" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + + var s = Semaphore.init(3); + // Starts at 1; two more posts reach the cap. + try testing.expectEqual(ret.ok, s.post(io)); + try testing.expectEqual(ret.ok, s.post(io)); + try testing.expectEqual(@as(u32, 3), s.count(io)); + // FreeRTOS's xSemaphoreGive returns pdFALSE past maxCount, and so does this. + try testing.expectEqual(ret.fail, s.post(io)); + try testing.expectEqual(@as(u32, 3), s.count(io)); + + for (0..3) |_| try testing.expectEqual(ret.ok, s.wait(io, .immediate)); + + // A bounded wait on an empty semaphore returns RET_FAIL_TIMEOUT, and takes at least as long as + // it was asked to. + const before = nowMs(io); + try testing.expectEqual(ret.fail_timeout, s.wait(io, .{ .bounded_ms = 25 })); + try testing.expect(nowMs(io) - before >= 25); +} + +test "Semaphore: a blocked waiter is released by a post from another task" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + + var s = Semaphore.init(4); + try testing.expectEqual(ret.ok, s.wait(io, .immediate)); // drain the initial permit + + const Worker = struct { + fn run(sem: *Semaphore, i: Io) c_int { + return sem.wait(i, .forever); + } + }; + var f = io.async(Worker.run, .{ &s, io }); + // Give the waiter time to actually block, then release it. + try io.sleep(.fromMilliseconds(20), clock); + try testing.expectEqual(ret.ok, s.post(io)); + try testing.expectEqual(ret.ok, f.await(io)); + try testing.expectEqual(@as(u32, 0), s.count(io)); +} + +test "Semaphore: an interrupt-deferred post is folded in by the next task-side operation" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + + var s = Semaphore.init(4); + try testing.expectEqual(ret.ok, s.wait(io, .immediate)); + + // Simulate the contended case: hold the semaphore's mutex, so postFromIsr cannot deliver + // inline and must defer. This is the exact window described on `postFromIsr`. + s.inner.mutex.lockUncancelable(io); + try testing.expectEqual(ret.ok, s.postFromIsr(io)); + try testing.expectEqual(@as(u32, 1), s.isr_posts.load(.acquire)); + s.inner.mutex.unlock(io); + + // The next task-side touch delivers it. + try testing.expectEqual(ret.ok, s.wait(io, .immediate)); + try testing.expectEqual(@as(u32, 0), s.isr_posts.load(.acquire)); + + // Uncontended, it lands directly. + try testing.expectEqual(ret.ok, s.postFromIsr(io)); + try testing.expectEqual(@as(u32, 0), s.isr_posts.load(.acquire)); + try testing.expectEqual(@as(u32, 1), s.count(io)); +} + +test "Queue: fixed-capacity records, non-blocking edges, and message count" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + const gpa = testing.allocator; + + // 24 bytes is sizeof(interface_buffer_handle_t) on rv32, which is what every transport queue + // in ESP-Hosted carries. + const item_size = 24; + const q = Queue.create(gpa, 4, item_size).?; + defer q.destroy(io, gpa); + + try testing.expectEqual(@as(c_int, 0), q.waiting(io)); + // Empty, non-blocking: RET_FAIL, and note it is RET_FAIL and not RET_FAIL_TIMEOUT. + var out: [item_size]u8 = undefined; + try testing.expectEqual(ret.fail, q.receive(io, &out, .immediate)); + + var item: [item_size]u8 = undefined; + for (0..4) |i| { + @memset(&item, @intCast(i)); + try testing.expectEqual(ret.ok, q.send(io, &item, .immediate)); + try testing.expectEqual(@as(c_int, @intCast(i + 1)), q.waiting(io)); + } + // Full: a non-blocking send fails and leaves no partial record behind. + @memset(&item, 0xFF); + try testing.expectEqual(ret.fail, q.send(io, &item, .immediate)); + try testing.expectEqual(@as(c_int, 4), q.waiting(io)); + + // FIFO order, whole records. + for (0..4) |i| { + try testing.expectEqual(ret.ok, q.receive(io, &out, .forever)); + try testing.expect(std.mem.allEqual(u8, &out, @intCast(i))); + } + try testing.expectEqual(@as(c_int, 0), q.waiting(io)); + + // A bounded receive on an empty queue waits and then fails. + const before = nowMs(io); + try testing.expectEqual(ret.fail, q.receive(io, &out, .{ .bounded_ms = 25 })); + try testing.expect(nowMs(io) - before >= 25); +} + +test "Queue: blocking receive is woken by a producer, and reset discards" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + const gpa = testing.allocator; + + const q = Queue.create(gpa, 2, 4).?; + defer q.destroy(io, gpa); + + const Consumer = struct { + fn run(queue: *Queue, i: Io) u32 { + var buf: [4]u8 = undefined; + if (queue.receive(i, &buf, .forever) != ret.ok) return 0xDEAD; + return std.mem.readInt(u32, &buf, .little); + } + }; + var f = io.async(Consumer.run, .{ q, io }); + try io.sleep(.fromMilliseconds(20), clock); + + var word: [4]u8 = undefined; + std.mem.writeInt(u32, &word, 0xC0FFEE, .little); + try testing.expectEqual(ret.ok, q.send(io, &word, .forever)); + try testing.expectEqual(@as(u32, 0xC0FFEE), f.await(io)); + + // reset drops buffered records. + try testing.expectEqual(ret.ok, q.send(io, &word, .immediate)); + try testing.expectEqual(ret.ok, q.send(io, &word, .immediate)); + try testing.expectEqual(@as(c_int, 2), q.waiting(io)); + _ = q.reset(io); + try testing.expectEqual(@as(c_int, 0), q.waiting(io)); +} + +test "Mutex: the three timeout dialects" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + + var m: Mutex = .{}; + try testing.expectEqual(ret.ok, m.lock(io, .forever)); + // Held: a non-blocking lock fails rather than deadlocking. + try testing.expectEqual(ret.fail, m.lock(io, .immediate)); + const before = nowMs(io); + try testing.expectEqual(ret.fail, m.lock(io, .{ .bounded_ms = 25 })); + try testing.expect(nowMs(io) - before >= 25); + try testing.expectEqual(ret.ok, m.unlock(io)); + try testing.expectEqual(ret.ok, m.lock(io, .immediate)); + try testing.expectEqual(ret.ok, m.unlock(io)); +} + +test "Wait: the two timeout dialects ESP-Hosted uses" { + // _h_lock_mutex and _h_get_semaphore: positive means milliseconds. + try testing.expectEqual(Wait.immediate, Wait.fromMillis(0)); + try testing.expectEqual(Wait.forever, Wait.fromMillis(-1)); + // HOSTED_BLOCK_MAX is portMAX_DELAY, 0xFFFFFFFF, which reaches an `int` parameter as -1. + try testing.expectEqual(Wait.forever, Wait.fromMillis(@bitCast(@as(u32, 0xFFFF_FFFF)))); + try testing.expectEqual(Wait{ .bounded_ms = 5000 }, Wait.fromMillis(5000)); + + // _h_dequeue_item: positive means seconds. port_esp_hosted_host_os.c:336. + try testing.expectEqual(Wait{ .bounded_ms = 5000 }, Wait.fromQueueTimeout(5)); + try testing.expectEqual(Wait.forever, Wait.fromQueueTimeout(-1)); +} + +test "TimerService: one-shot fires once, periodic repeats, stop cancels" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + const gpa = testing.allocator; + + const Counter = struct { + var oneshot: u32 = 0; + var periodic: u32 = 0; + fn bumpOneshot(_: ?*anyopaque) callconv(.c) void { + oneshot += 1; + } + fn bumpPeriodic(_: ?*anyopaque) callconv(.c) void { + periodic += 1; + } + }; + Counter.oneshot = 0; + Counter.periodic = 0; + + var svc: TimerService(4) = .{}; + try testing.expect(svc.start(io, gpa)); + defer svc.stop(io, gpa); + + _ = svc.arm(io, 10, .oneshot, Counter.bumpOneshot, null).?; + const p = svc.arm(io, 10, .periodic, Counter.bumpPeriodic, null).?; + + try io.sleep(.fromMilliseconds(120), clock); + try testing.expectEqual(@as(u32, 1), Counter.oneshot); + try testing.expect(Counter.periodic >= 3); + + // Disarming stops it; the count must not move afterwards. + try testing.expectEqual(ret.ok, svc.disarm(io, p)); + const frozen = Counter.periodic; + try io.sleep(.fromMilliseconds(60), clock); + try testing.expectEqual(frozen, Counter.periodic); + // Disarming an already-disarmed slot reports failure, as esp_timer_stop does. + try testing.expectEqual(ret.fail, svc.disarm(io, p)); +} + +test "TimerService: slot exhaustion is reported, not fatal" { + var h = hostIo(); + defer { + h.threaded.deinit(); + testing.allocator.destroy(h.threaded); + } + const io = h.io; + + const Nop = struct { + fn f(_: ?*anyopaque) callconv(.c) void {} + }; + var svc: TimerService(2) = .{}; + _ = svc.arm(io, 10_000, .oneshot, Nop.f, null).?; + _ = svc.arm(io, 10_000, .oneshot, Nop.f, null).?; + try testing.expectEqual(@as(?usize, null), svc.arm(io, 10_000, .oneshot, Nop.f, null)); +} |
