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