summaryrefslogtreecommitdiff
path: root/src/net/hosted_os.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/net/hosted_os.zig')
-rw-r--r--src/net/hosted_os.zig890
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));
+}