Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337//! Zio: the hybrid std.Io backend.//!//! wraps an inner Io.Threaded and delegates the entire vtable to it, then//! overrides the ops that benefit from the fiber scheduler://!//! concurrent/await/cancel -> fibers on the sched loop thread//! sleep -> sched timer heap (when called from a fiber)//! netConnectIp/Read/Write/Close -> park on poller readiness//! futexWait/futexWake -> fiber-aware futex table; this makes//! Io.Mutex, Io.Condition and Io.Queue work//! correctly for fiber callers with zero//! changes to their code//! netLookup -> spilled to a pool of plain threads zio//! owns (spill.zig) so DNS never blocks the//! loop; the caller's fiber parks until the//! lookup completes//!//! every op detects its caller domain: on-a-fiber -> scheduler path;//! anywhere else (main thread, pool threads) -> delegate to Io.Threaded.//! this is what makes the backend a drop-in `const Backend = zio.Zio;`.const std = @import("std");const builtin = @import("builtin");const c = std.c;const Io = std.Io;const sched_mod = @import("fiber/sched.zig");const spill_mod = @import("spill.zig");const engine = @import("net.zig");const Sched = sched_mod.Sched;const Fiber = sched_mod.Fiber;
const max_result_size = 512;const max_context_size = 1024;
pub const Zio = struct { allocator: std.mem.Allocator, threaded: Io.Threaded, vtable: Io.VTable, sched: Sched, loop_thread: ?std.Thread, dns_spill: spill_mod.SpillPool(DnsRequest),
/// identifies the loop thread; ops compare against Sched.tls_active. pub fn init(zio: *Zio, allocator: std.mem.Allocator, options: Io.Threaded.InitOptions) !void { zio.allocator = allocator; zio.threaded = .init(allocator, options); zio.sched = try Sched.init(allocator); try zio.sched.attachWake(); zio.loop_thread = null; zio.dns_spill = .{};
const inner = zio.threaded.io(); zio.vtable = inner.vtable.*; zio.vtable.concurrent = vtConcurrent; zio.vtable.await = vtAwait; zio.vtable.cancel = vtCancel; zio.vtable.checkCancel = vtCheckCancel; zio.vtable.sleep = vtSleep; zio.vtable.netConnectIp = vtNetConnectIp; zio.vtable.netAccept = vtNetAccept; zio.vtable.operate = vtOperate; zio.vtable.futexWait = vtFutexWait; zio.vtable.futexWaitUncancelable = vtFutexWaitUncancelable; zio.vtable.futexWake = vtFutexWake; zio.vtable.netLookup = vtNetLookup; zio.vtable.netClose = vtNetClose; zio.vtable.groupAsync = vtGroupAsync; zio.vtable.groupConcurrent = vtGroupConcurrent; zio.vtable.groupAwait = vtGroupAwait; zio.vtable.groupCancel = vtGroupCancel; }
pub fn start(zio: *Zio) !void { zio.loop_thread = try std.Thread.spawn(.{}, loopMain, .{zio}); }
pub fn deinit(zio: *Zio) void { // before stopping the loop: in-flight lookups complete and their // wakes drain through a live sched, so no parked fiber is stranded. zio.dns_spill.deinit(); if (zio.loop_thread) |t| { zio.sched.requestStop(); t.join(); zio.loop_thread = null; } zio.sched.deinit(); zio.threaded.deinit(); }
pub fn io(zio: *Zio) Io { return .{ .userdata = &zio.threaded, .vtable = &zio.vtable }; }
fn loopMain(zio: *Zio) void { zio.sched.runForever(); }
fn fromUserdata(userdata: ?*anyopaque) *Zio { const t: *Io.Threaded = @ptrCast(@alignCast(userdata.?)); return @fieldParentPtr("threaded", t); }
/// scheduler + poller + spill snapshot for diagnostics. call from a /// fiber (loop thread); exported by consumers as metrics so a wedged /// loop is visible in a scrape instead of requiring a debugger. pub const DebugStats = struct { sched: sched_mod.Stats, poller_registrations: usize, dns_queued: usize, dns_workers: u32, dns_idle: u32, /// live fibers grouped by the function they run. a leak names its /// own culprit — resolve the address with `info symbol`. starts: [8]StartCount = @splat(.{}),
pub const StartCount = struct { fn_addr: usize = 0, count: usize = 0 }; };
pub fn debugStats(zio: *Zio) DebugStats { const sp = zio.dns_spill.stats(); var out: DebugStats = .{ .sched = zio.sched.stats(), .poller_registrations = zio.sched.poller.registrationCount(), .dns_queued = sp.queued, .dns_workers = sp.workers, .dns_idle = sp.idle, }; var f = zio.sched.all_fibers; while (f) |x| : (f = x.all_next) { const tp = x.task orelse continue; const t: *Task = @ptrCast(@alignCast(tp)); const addr: usize = if (t.group_start) |g| @intFromPtr(g) else @intFromPtr(t.start); var slot: ?*DebugStats.StartCount = null; for (&out.starts) |*sc| { if (sc.fn_addr == addr) { slot = sc; break; } if (sc.fn_addr == 0 and slot == null) slot = sc; } if (slot) |sc| { sc.fn_addr = addr; sc.count += 1; } } return out; }
/// non-null when the calling code is a fiber owned by this Zio's sched. fn currentFiber(zio: *Zio) ?*Fiber { if (Sched.tls_active != &zio.sched) return null; return zio.sched.current; }};
// --- task: unit of concurrency behind *AnyFuture ---------------------------
const Task = struct { zio: *Zio, start: *const fn (context: *const anyopaque, result: *anyopaque) void, context: [max_context_size]u8 align(16), result: [max_result_size]u8 align(16), /// 0 = running, 1 = done. futex word for foreign-thread awaiters. state: std.atomic.Value(u32), cancel_requested: std.atomic.Value(bool), /// fiber awaiting this task, parked on the sched (sched thread only). awaiter: ?*Fiber, /// the fiber executing this task (sched thread only). fiber: ?*Fiber, /// two owners. futures: the completing fiber and the awaiting caller. /// group members: the completing fiber and the group's teardown. refs: std.atomic.Value(u32), /// non-null for group members: completion is accounted to the group /// instead of a future, and `group_start` runs instead of `start`. group: ?*FiberGroup, group_start: ?*const fn (context: *const anyopaque) void, /// this task's node in its group's member list (null when not in a /// group, or once reaped). gives O(1) unlink on completion. member: ?*FiberGroup.Member = null,
fn unref(task: *Task) void { if (task.refs.fetchSub(1, .acq_rel) == 1) { task.zio.allocator.destroy(task); } }};
fn taskFiberEntry(s: *Sched, arg: ?*anyopaque) void { const task: *Task = @ptrCast(@alignCast(arg.?)); task.fiber = s.current; s.current.?.task = task; if (task.group_start) |gs| gs(&task.context) else task.start(&task.context, &task.result); task.fiber = null; task.state.store(1, .release); // wake a fiber awaiter (same thread) and any foreign-thread awaiter. if (task.awaiter) |af| { task.awaiter = null; if (af.state == .parked) s.pushReady(af); } zioFutexWake(task.zio, &task.state.raw, 1); // the group ref keeps `task` alive past this unref; memberDone only // touches the group, whose teardown cannot start until it runs. const group = task.group; if (group) |fg| { // memberDone drops the group's ref; the runner's ref is dropped // after, so `task` stays valid across the call. fg.memberDone(task); task.unref(); } else { task.unref(); }}
// --- groups: fiber-native Io.Group -----------------------------------------//// std spawns group tasks through the vtable (HostName.connectMany's// per-address connect attempts are the load-bearing case). delegating them// to the inner Threaded pool put fiber-originated work on foreign threads:// blocking connects tying up a worker per dead host, and blocking fds// handed back into the fiber domain. groups are therefore implemented// uniformly on the sched — members are ordinary Tasks (so the existing// cancel/interrupt machinery applies), and `outstanding` doubles as the// await futex word for both domains.const FiberGroup = struct { zio: *Zio, /// members not yet finished. await parks on this word; the last /// memberDone wakes both domains through zioFutexWake. outstanding: std.atomic.Value(u32) = .init(0), /// spinlock guarding `members` and `canceled`. spawns can come from any /// thread; critical sections are pointer pushes. lock: std.atomic.Value(u32) = .init(0), /// every member ever spawned; holds a task ref each, freed at teardown. /// kept whole so cancel can walk it without racing member completion. members: ?*Member = null, canceled: bool = false,
/// doubly linked so a finishing member unlinks in O(1). a singly /// linked list forces an O(n) scan per completion, i.e. O(n^2) for a /// group that serves n connections over its life. const Member = struct { task: *Task, prev: ?*Member = null, next: ?*Member = null };
fn lockSpin(fg: *FiberGroup) void { while (fg.lock.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { std.atomic.spinLoopHint(); } }
fn unlock(fg: *FiberGroup) void { fg.lock.store(0, .release); }
/// a member finished: drop it from the list and release its task ref /// NOW. holding members until await/cancel teardown is unbounded for a /// long-lived group — a server that spawns one member per accepted /// connection (std's websocket server does) accumulates a Task plus a /// Member for every connection it has ever served. fn memberDone(fg: *FiberGroup, task: *Task) void { fg.lockSpin(); const found = task.member; if (found) |m| { if (m.prev) |p| p.next = m.next else fg.members = m.next; if (m.next) |n| n.prev = m.prev; task.member = null; } fg.unlock(); if (found) |m| { m.task.unref(); fg.zio.allocator.destroy(m); } if (fg.outstanding.fetchSub(1, .acq_rel) == 1) { zioFutexWake(fg.zio, &fg.outstanding.raw, std.math.maxInt(u32)); } }
/// all members finished and no more can be spawned: release everything /// and clear the token, as the Group.await/cancel contract requires. fn teardown(fg: *FiberGroup, group: *Io.Group) void { std.debug.assert(fg.outstanding.load(.acquire) == 0); var cur = fg.members; while (cur) |m| { const next = m.next; m.task.member = null; m.task.unref(); fg.zio.allocator.destroy(m); cur = next; } fg.members = null; group.token.store(null, .release); fg.zio.allocator.destroy(fg); }};
/// resolve (or lazily create) the FiberGroup behind a Group's token.fn groupResolve(zio: *Zio, group: *Io.Group) error{OutOfMemory}!*FiberGroup { if (group.token.load(.acquire)) |t| return @ptrCast(@alignCast(t)); const fg = try zio.allocator.create(FiberGroup); fg.* = .{ .zio = zio }; if (group.token.cmpxchgStrong(null, fg, .acq_rel, .acquire)) |won| { // lost the race: another thread installed its FiberGroup first. zio.allocator.destroy(fg); return @ptrCast(@alignCast(won.?)); } return fg;}
fn groupSpawnMember( zio: *Zio, group: *Io.Group, context: []const u8, start: *const fn (context: *const anyopaque) void,) error{ OutOfMemory, ConcurrencyUnavailable }!void { std.debug.assert(context.len <= max_context_size); const fg = try groupResolve(zio, group); const task = zio.allocator.create(Task) catch return error.OutOfMemory; task.* = .{ .zio = zio, .start = undefined, .context = undefined, .result = undefined, .state = .init(0), .cancel_requested = .init(false), .awaiter = null, .fiber = null, .refs = .init(2), .group = fg, .group_start = start, }; @memcpy(task.context[0..context.len], context); const member = zio.allocator.create(FiberGroup.Member) catch { zio.allocator.destroy(task); return error.OutOfMemory; }; member.* = .{ .task = task }; task.member = member;
fg.lockSpin(); if (fg.canceled) task.cancel_requested.store(true, .release); member.next = fg.members; if (fg.members) |h| h.prev = member; fg.members = member; fg.unlock(); _ = fg.outstanding.fetchAdd(1, .acq_rel);
const spawn_err = if (zio.currentFiber() != null) zio.sched.spawn(taskFiberEntry, task) else zio.sched.injectSpawn(taskFiberEntry, task); spawn_err catch { // roll back: it never ran, so reap the member and drop the runner's // ref and its completion account. fg.memberDone(task); task.unref(); return error.ConcurrencyUnavailable; };}
fn vtGroupAsync( userdata: ?*anyopaque, group: *Io.Group, context: []const u8, context_alignment: std.mem.Alignment, start: *const fn (context: *const anyopaque) void,) void { _ = context_alignment; const zio = Zio.fromUserdata(userdata); groupSpawnMember(zio, group, context, start) catch { // async permits eager inline execution — same fallback as Threaded. start(context.ptr); };}
fn vtGroupConcurrent( userdata: ?*anyopaque, group: *Io.Group, context: []const u8, context_alignment: std.mem.Alignment, start: *const fn (context: *const anyopaque) void,) Io.ConcurrentError!void { _ = context_alignment; const zio = Zio.fromUserdata(userdata); groupSpawnMember(zio, group, context, start) catch return error.ConcurrencyUnavailable;}
fn groupWaitDrained(zio: *Zio, fg: *FiberGroup, cancelable: bool) Io.Cancelable!void { if (zio.currentFiber()) |f| { while (true) { const n = fg.outstanding.load(.acquire); if (n == 0) return; zio.sched.futexWaitFiber(&fg.outstanding.raw, n, null); if (cancelable and fiberCanceled(f)) return error.Canceled; } } const inner = zio.threaded.io(); while (true) { const n = fg.outstanding.load(.acquire); if (n == 0) return; inner.vtable.futexWaitUncancelable(inner.userdata, &fg.outstanding.raw, n); }}
fn vtGroupAwait(userdata: ?*anyopaque, group: *Io.Group, token: *anyopaque) Io.Cancelable!void { const zio = Zio.fromUserdata(userdata); const fg: *FiberGroup = @ptrCast(@alignCast(token)); // on Canceled the token stays set: the caller's Group.cancel follows and // performs the actual teardown. try groupWaitDrained(zio, fg, true); fg.teardown(group);}
fn vtGroupCancel(userdata: ?*anyopaque, group: *Io.Group, token: *anyopaque) void { const zio = Zio.fromUserdata(userdata); const fg: *FiberGroup = @ptrCast(@alignCast(token)); // the walk must hold the lock: members are unlinked and freed as they // finish, so an unlocked walk can dereference a freed node. this is a // shutdown path, so the longer critical section is acceptable; members // spawned after `canceled` is set cancel themselves. fg.lockSpin(); fg.canceled = true; var cur = fg.members; while (cur) |m| : (cur = m.next) { m.task.cancel_requested.store(true, .release); if (Sched.tls_active == &zio.sched) { zio.sched.interruptTask(m.task); } else { zio.sched.injectInterrupt(m.task) catch {}; } } fg.unlock(); groupWaitDrained(zio, fg, false) catch unreachable; fg.teardown(group);}
fn vtConcurrent( userdata: ?*anyopaque, result_len: usize, result_alignment: std.mem.Alignment, context: []const u8, context_alignment: std.mem.Alignment, start: *const fn (context: *const anyopaque, result: *anyopaque) void,) Io.ConcurrentError!*Io.AnyFuture { const zio = Zio.fromUserdata(userdata); std.debug.assert(result_len <= max_result_size); std.debug.assert(context.len <= max_context_size); std.debug.assert(result_alignment.compare(.lte, .@"16")); std.debug.assert(context_alignment.compare(.lte, .@"16"));
const task = zio.allocator.create(Task) catch return error.ConcurrencyUnavailable; task.* = .{ .zio = zio, .start = start, .context = undefined, .result = undefined, .state = .init(0), .cancel_requested = .init(false), .awaiter = null, .fiber = null, .refs = .init(2), .group = null, .group_start = null, }; @memcpy(task.context[0..context.len], context);
if (zio.currentFiber() != null) { // spawning from a fiber: same thread, spawn directly. zio.sched.spawn(taskFiberEntry, task) catch { zio.allocator.destroy(task); return error.ConcurrencyUnavailable; }; } else { zio.sched.injectSpawn(taskFiberEntry, task) catch { zio.allocator.destroy(task); return error.ConcurrencyUnavailable; }; } return @ptrCast(task);}
fn awaitImpl(zio: *Zio, task: *Task, result: []u8) void { if (zio.currentFiber()) |f| { while (task.state.load(.acquire) == 0) { task.awaiter = f; zio.sched.parkCurrent(); } } else { // foreign thread: wait on the task state word through the inner // Threaded futex (real futex/ulock). completion wakes via zioFutexWake. const inner = zio.threaded.io(); while (task.state.load(.acquire) == 0) { inner.vtable.futexWaitUncancelable(inner.userdata, &task.state.raw, 0); } } @memcpy(result, task.result[0..result.len]); task.unref();}
fn vtAwait(userdata: ?*anyopaque, any_future: *Io.AnyFuture, result: []u8, result_alignment: std.mem.Alignment) void { _ = result_alignment; const zio = Zio.fromUserdata(userdata); const task: *Task = @ptrCast(@alignCast(any_future)); awaitImpl(zio, task, result);}
fn vtCancel(userdata: ?*anyopaque, any_future: *Io.AnyFuture, result: []u8, result_alignment: std.mem.Alignment) void { _ = result_alignment; const zio = Zio.fromUserdata(userdata); const task: *Task = @ptrCast(@alignCast(any_future)); task.cancel_requested.store(true, .release); // wake the task's fiber if it is parked so it can observe cancellation. zio.sched.injectInterrupt(task) catch {}; awaitImpl(zio, task, result);}
fn vtCheckCancel(userdata: ?*anyopaque) Io.Cancelable!void { const zio = Zio.fromUserdata(userdata); if (zio.currentFiber()) |f| { if (fiberCanceled(f)) return error.Canceled; return; } const inner = zio.threaded.io(); return inner.vtable.checkCancel(inner.userdata);}
fn fiberCanceled(f: *Fiber) bool { const task_ptr = f.task orelse return false; const task: *Task = @ptrCast(@alignCast(task_ptr)); return task.cancel_requested.load(.acquire);}
// --- sleep ------------------------------------------------------------------
/// nanoseconds -> whole milliseconds, rounding UP so that any non-zero/// duration is at least 1ms.////// truncating instead is a trap with teeth: the sched's clock is/// millisecond-granular, so a sub-millisecond sleep truncated to 0 yields a/// deadline of *now*, `sleepUntil`'s `while (nowMs() < deadline)` is false/// immediately, and the sleep returns WITHOUT EVER PARKING. callers that/// poll — `while (!done) io.sleep(100us)` is a common shape — then busy-spin/// the loop thread instead of yielding it. one such caller starves every/// other fiber on that thread: accept loops go deaf, io throughput collapses,/// and the process still looks alive. that was the 2026-08-04 field failure.fn nsToMsCeil(ns: i128) u64 { if (ns <= 0) return 0; const ms = @divTrunc(ns + std.time.ns_per_ms - 1, std.time.ns_per_ms); return @intCast(ms);}
fn timeoutToDeadlineMs(zio: *Zio, timeout: Io.Timeout) ?u64 { // all supported clocks are mapped onto the sched's monotonic ms clock. const inner = zio.threaded.io(); const now_ns = Io.Timestamp.now(inner, .awake).nanoseconds; return switch (timeout) { .none => null, .duration => |d| sched_mod.nowMs() + nsToMsCeil(d.raw.nanoseconds), .deadline => |d| blk: { const delta_ns = d.raw.nanoseconds - now_ns; if (delta_ns <= 0) break :blk sched_mod.nowMs(); break :blk sched_mod.nowMs() + nsToMsCeil(delta_ns); }, };}
fn vtSleep(userdata: ?*anyopaque, timeout: Io.Timeout) Io.Cancelable!void { const zio = Zio.fromUserdata(userdata); if (zio.currentFiber()) |f| { if (fiberCanceled(f)) return error.Canceled; const deadline = timeoutToDeadlineMs(zio, timeout) orelse { // sleep forever: park until canceled. while (true) { zio.sched.parkCurrent(); if (fiberCanceled(f)) return error.Canceled; } }; zio.sched.sleepUntil(deadline); if (fiberCanceled(f)) return error.Canceled; return; } const inner = zio.threaded.io(); return inner.vtable.sleep(inner.userdata, timeout);}
// --- networking --------------------------------------------------------------
fn vtNetConnectIp( userdata: ?*anyopaque, address: *const Io.net.IpAddress, options: Io.net.IpAddress.ConnectOptions,) Io.net.IpAddress.ConnectError!Io.net.Socket { const zio = Zio.fromUserdata(userdata); const f = zio.currentFiber() orelse { const inner = zio.threaded.io(); return inner.vtable.netConnectIp(inner.userdata, address, options); }; if (options.mode != .stream) { // datagram connect is address assignment, not a handshake — the inner // call cannot stall the loop. the fd, however, is about to enter the // fiber domain, whose read/write paths park on EAGAIN and therefore // require every fd they touch to be non-blocking. const inner = zio.threaded.io(); const sock = try inner.vtable.netConnectIp(inner.userdata, address, options); engine.setNonblock(sock.handle) catch return error.Unexpected; return sock; }
// both families take the non-blocking path: a blocking connect(2) from a // fiber seizes the loop thread for up to a full TCP timeout per dead // host. the ip6 fall-through to the inner Threaded connect did exactly // that (found via a loop-thread backtrace stuck in __libc_connect with // a sockaddr_in6). const fd = switch (address.*) { .ip4 => |a| engine.tcpConnectIp4Start(a.bytes, a.port) catch return error.ConnectionRefused, .ip6 => |a| engine.tcpConnectIp6Start(a.bytes, a.port, a.flow, a.interface.index) catch return error.ConnectionRefused, }; errdefer _ = c.close(fd); while (true) { if (fiberCanceled(f)) return error.Canceled; zio.sched.waitWritable(fd); if (fiberCanceled(f)) return error.Canceled; engine.connectResult(fd) catch return error.ConnectionRefused; var peer: c.sockaddr.storage = undefined; var len: c.socklen_t = @sizeOf(c.sockaddr.storage); if (c.getpeername(fd, @ptrCast(&peer), &len) == 0) break; if (c.errno(@as(isize, -1)) != .NOTCONN) return error.ConnectionRefused; } switch (address.*) { .ip4 => { var bound: c.sockaddr.in = undefined; var blen: c.socklen_t = @sizeOf(c.sockaddr.in); _ = c.getsockname(fd, @ptrCast(&bound), &blen); return .{ .handle = fd, .address = .{ .ip4 = .{ .bytes = @bitCast(bound.addr), .port = std.mem.bigToNative(u16, bound.port), } } }; }, .ip6 => { var bound: c.sockaddr.in6 = undefined; var blen: c.socklen_t = @sizeOf(c.sockaddr.in6); _ = c.getsockname(fd, @ptrCast(&bound), &blen); return .{ .handle = fd, .address = .{ .ip6 = .{ .bytes = bound.addr, .port = std.mem.bigToNative(u16, bound.port), .flow = bound.flowinfo, .interface = .{ .index = bound.scope_id }, } } }; }, }}
fn vtNetAccept( userdata: ?*anyopaque, server: Io.net.Socket.Handle, options: Io.net.Server.AcceptOptions,) Io.net.Server.AcceptError!Io.net.Socket { const zio = Zio.fromUserdata(userdata); const f = zio.currentFiber() orelse { const inner = zio.threaded.io(); return inner.vtable.netAccept(inner.userdata, server, options); }; // the listener may have been created by the (blocking) Threaded path; // flip it non-blocking so the loop never blocks in accept. idempotent. engine.setNonblock(server) catch return error.SocketNotListening; while (true) { if (fiberCanceled(f)) return error.Canceled; if (f.wake_closed) { f.wake_closed = false; return error.SocketNotListening; // listener closed under us } const maybe = engine.acceptNonblock(server) catch return error.SocketNotListening; if (maybe) |fd| { var peer: c.sockaddr.in = undefined; var plen: c.socklen_t = @sizeOf(c.sockaddr.in); _ = c.getpeername(fd, @ptrCast(&peer), &plen); return .{ .handle = fd, .address = .{ .ip4 = .{ .bytes = @bitCast(peer.addr), .port = std.mem.bigToNative(u16, peer.port), } } }; } zio.sched.waitReadable(server); }}
/// closing an fd strands anyone parked on it: the kernel drops closed fds/// from the epoll set without an event, so the fiber waits forever. release/// them first, flagged so they return an error instead of touching an fd/// whose number may already have been reused.fn vtNetClose(userdata: ?*anyopaque, sockets: []const Io.net.Socket) void { const zio = Zio.fromUserdata(userdata); for (sockets) |s| { if (Sched.tls_active == &zio.sched) { zio.sched.closeFdWaiters(s.handle); } else { zio.sched.injectCloseFd(s.handle) catch {}; } } const inner = zio.threaded.io(); inner.vtable.netClose(inner.userdata, sockets);}
fn vtOperate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Operation.Result { const zio = Zio.fromUserdata(userdata); if (zio.currentFiber()) |f| switch (operation) { // control data needs recvmsg; only plain reads take the fiber path. .net_read => |op| if (op.control.len == 0) return .{ .net_read = if (fiberRead(zio, f, op)) |n| .{ .data_len = n } else |err| switch (err) { error.Canceled => return error.Canceled, else => |e| e, } }, .net_write => |op| return .{ .net_write = fiberWrite(zio, f, op) catch |err| switch (err) { error.Canceled => return error.Canceled, else => |e| e, } }, else => {}, }; const inner = zio.threaded.io(); return inner.vtable.operate(inner.userdata, operation);}
fn fiberRead(zio: *Zio, f: *Fiber, op: Io.Operation.NetRead) (Io.Cancelable || Io.Operation.NetRead.Error)!usize { const src = op.socket_handle; const buf = for (op.data) |d| { if (d.len > 0) break d; } else return 0; // invariant: every fd enters the fiber domain non-blocking — the fiber // connect paths and acceptNonblock establish it, and fiber-native groups // keep connectMany's attempts on those paths. a blocking fd here would // seize the loop thread in read(2). while (true) { if (fiberCanceled(f)) return error.Canceled; if (f.wake_closed) { f.wake_closed = false; return error.Unexpected; // fd closed while we were parked } const n = c.read(src, buf.ptr, buf.len); if (n >= 0) return @intCast(n); if (c.errno(n) != engine.wouldBlock) return error.Unexpected; zio.sched.waitReadable(src); }}
fn fiberWrite(zio: *Zio, f: *Fiber, op: Io.Operation.NetWrite) (Io.Cancelable || Io.Operation.NetWrite.Error)!usize { // short writes are allowed by the contract: send the first non-empty // chunk, parking on EAGAIN until at least one byte lands. const chunk: []const u8 = if (op.header.len > 0) op.header else blk: { if (op.data.len == 0) return 0; for (op.data[0 .. op.data.len - 1]) |d| { if (d.len > 0) break :blk d; } const last = op.data[op.data.len - 1]; if (last.len == 0 or op.splat == 0) return 0; break :blk last; }; // same non-blocking invariant as fiberRead. while (true) { if (fiberCanceled(f)) return error.Canceled; if (f.wake_closed) { f.wake_closed = false; return error.Unexpected; // fd closed while we were parked } const n = c.write(op.socket_handle, chunk.ptr, chunk.len); if (n >= 0) return @intCast(n); if (c.errno(n) != engine.wouldBlock) return error.Unexpected; zio.sched.waitWritable(op.socket_handle); }}
// --- fiber-aware futex --------------------------------------------------------// Io.Mutex, Io.Condition and Io.Queue all park through these two ops. fibers// park on the sched futex table; foreign threads keep the inner Threaded// futex. wakes always hit both domains.
fn vtFutexWait(userdata: ?*anyopaque, ptr: *const u32, expected: u32, timeout: Io.Timeout) Io.Cancelable!void { const zio = Zio.fromUserdata(userdata); if (zio.currentFiber()) |f| { if (fiberCanceled(f)) return error.Canceled; const deadline = timeoutToDeadlineMs(zio, timeout); zio.sched.futexWaitFiber(ptr, expected, deadline); if (fiberCanceled(f)) return error.Canceled; return; } const inner = zio.threaded.io(); return inner.vtable.futexWait(inner.userdata, ptr, expected, timeout);}
fn vtFutexWaitUncancelable(userdata: ?*anyopaque, ptr: *const u32, expected: u32) void { const zio = Zio.fromUserdata(userdata); if (zio.currentFiber() != null) { zio.sched.futexWaitFiber(ptr, expected, null); return; } const inner = zio.threaded.io(); inner.vtable.futexWaitUncancelable(inner.userdata, ptr, expected);}
fn zioFutexWake(zio: *Zio, ptr: *const u32, max_waiters: u32) void { // fiber domain (thread-safe: goes through the sched's external queue // when called off-loop) ... zio.sched.futexWakeFibers(ptr, max_waiters); // ... and thread domain. const inner = zio.threaded.io(); inner.vtable.futexWake(inner.userdata, ptr, max_waiters);}
fn vtFutexWake(userdata: ?*anyopaque, ptr: *const u32, max_waiters: u32) void { zioFutexWake(Zio.fromUserdata(userdata), ptr, max_waiters);}
// --- dns spill -----------------------------------------------------------------
/// a blocking netLookup executed on the SpillPool on behalf of a parked/// fiber. lives on the calling fiber's stack frame; the frame cannot unwind/// before run() signals completion, which is what makes the borrowed queue/// (Threaded's contract: not closed until netLookup returns) and the/// intrusive link safe. see spill.zig for why it is neither detached, nor an/// inner-Threaded task, nor run on the loop thread.const DnsRequest = struct { zio: *Zio, host_name: Io.net.HostName, queue: *Io.Queue(Io.net.HostName.LookupResult), options: Io.net.HostName.LookupOptions, /// 0 = pending, 1 = done. the caller's fiber parks on this word. state: std.atomic.Value(u32) = .init(0), err: ?Io.net.HostName.LookupError = null, next: ?*DnsRequest = null,
pub fn run(req: *DnsRequest) void { const inner = req.zio.threaded.io(); inner.vtable.netLookup(inner.userdata, req.host_name, req.queue, req.options) catch |e| { req.err = e; }; req.state.store(1, .release); // wake by address through the sched's inject queue — never // dereferences the request, which may already be resuming. req.zio.sched.futexWakeFibers(&req.state.raw, std.math.maxInt(u32)); }};
fn vtNetLookup( userdata: ?*anyopaque, host_name: Io.net.HostName, queue: *Io.Queue(Io.net.HostName.LookupResult), options: Io.net.HostName.LookupOptions,) Io.net.HostName.LookupError!void { const zio = Zio.fromUserdata(userdata); const inner = zio.threaded.io(); if (zio.currentFiber() == null) { return inner.vtable.netLookup(inner.userdata, host_name, queue, options); } // fiber caller: getaddrinfo must not run on the loop thread. hand it to // the spill pool and park this fiber until it completes. var req: DnsRequest = .{ .zio = zio, .host_name = host_name, .queue = queue, .options = options, }; zio.dns_spill.submit(&req) catch |e| switch (e) { // pool gone (shutdown) or unspawnable: the lookup still has to // happen and the loop thread is all that's left. blocking here is // the degraded mode, not the design. error.ShuttingDown, error.SystemResources => return inner.vtable.netLookup(inner.userdata, host_name, queue, options), }; // park until the pool signals — uncancelably: getaddrinfo cannot be // aborted, and returning early would free the frame it writes into. while (req.state.load(.acquire) == 0) { zio.sched.futexWaitFiber(&req.state.raw, 0, null); } if (req.err) |e| return e;}
test "Zio: concurrent + await through the vtable (fiber tasks)" { var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
const Worker = struct { fn double(x: u64) u64 { return x * 2; } }; var fut = try io.concurrent(Worker.double, .{21}); const answer = fut.await(io); try std.testing.expectEqual(@as(u64, 42), answer);}
test "Zio: io.sleep on a fiber task goes through the sched timers" { var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
const Worker = struct { fn nap(io_inner: Io) u64 { io_inner.sleep(Io.Duration.fromMilliseconds(20), .awake) catch return 0; return 7; } }; const t0 = sched_mod.nowMs(); var fut = try io.concurrent(Worker.nap, .{io}); const r = fut.await(io); const elapsed = sched_mod.nowMs() - t0; try std.testing.expectEqual(@as(u64, 7), r); try std.testing.expect(elapsed >= 15);}
test "Zio: blocking-style net through Io.net on fibers (tcp echo)" { var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
// listener stays on the inner Threaded (accept isn't overridden yet) — // exactly the hybrid contract: unoverridden ops fall back to threads. const bind_addr = try Io.net.IpAddress.parse("127.0.0.1", 0); var server = try Io.net.IpAddress.listen(&bind_addr, io, .{}); defer server.deinit(io); const addr = server.socket.address;
const Server = struct { fn serve(io_inner: Io, srv: *Io.net.Server) u64 { const stream = srv.accept(io_inner) catch return 0; defer stream.close(io_inner); var rbuf: [128]u8 = undefined; var wbuf: [128]u8 = undefined; var reader = stream.reader(io_inner, &rbuf); var writer = stream.writer(io_inner, &wbuf); var total: u64 = 0; // fixed-size protocol: read exactly one 11-byte ping per round. // (readSliceShort would park waiting to fill its buffer — short // reads only happen at stream end, not at message boundaries.) var chunk: [11]u8 = undefined; while (true) { reader.interface.readSliceAll(&chunk) catch return total; writer.interface.writeAll(&chunk) catch return total; writer.interface.flush() catch return total; total += chunk.len; } } };
const Client = struct { fn ping(io_inner: Io, a: *const Io.net.IpAddress) u64 { const stream = Io.net.IpAddress.connect(a, io_inner, .{ .mode = .stream }) catch return 0; defer stream.close(io_inner); var rbuf: [128]u8 = undefined; var wbuf: [128]u8 = undefined; var reader = stream.reader(io_inner, &rbuf); var writer = stream.writer(io_inner, &wbuf); var echoed: u64 = 0; var round: usize = 0; while (round < 50) : (round += 1) { writer.interface.writeAll("fiber echo!") catch return echoed; writer.interface.flush() catch return echoed; var got: [11]u8 = undefined; reader.interface.readSliceAll(&got) catch return echoed; echoed += got.len; } return echoed; } };
var server_fut = try io.concurrent(Server.serve, .{ io, &server }); var client_fut = try io.concurrent(Client.ping, .{ io, &addr }); const echoed = client_fut.await(io); const served = server_fut.await(io); try std.testing.expectEqual(@as(u64, 550), echoed); try std.testing.expectEqual(@as(u64, 550), served);}
test "Zio: Io.Mutex contended across fiber and pool-thread domains" { var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
const Shared = struct { var mutex: Io.Mutex = .init; var counter: u64 = 0;
fn fiberSide(io_inner: Io) u64 { var i: u32 = 0; while (i < 1000) : (i += 1) { mutex.lock(io_inner) catch return 0; counter += 1; mutex.unlock(io_inner); } return 1; } fn threadSide(zio_io: Io) u64 { var i: u32 = 0; while (i < 1000) : (i += 1) { mutex.lock(zio_io) catch return 0; counter += 1; mutex.unlock(zio_io); } return 1; } }; Shared.counter = 0;
// fiber domain and OS-thread domain BOTH operate the shared mutex // through the ZIO io — the zlay thread_pool shape (workers pop with // w.io = the zio io). This is the only supported contract for a // primitive shared across domains: the zio vtable's futexWake reaches // both the sched futex table and the kernel futex, and its futexWait // routes foreign threads to the kernel. Operating the same primitive // through the INNER Threaded io instead deadlocks about half the time: // the inner unlock's wake is kernel-only, so a fiber parked in the // sched futex table never hears it — a permanently lost wake. This // test originally encoded that broken shape and hung intermittently. const inner = zio.threaded.io(); var f1 = try io.concurrent(Shared.fiberSide, .{io}); var f2 = try inner.concurrent(Shared.threadSide, .{io}); _ = f1.await(io); _ = f2.await(inner); try std.testing.expectEqual(@as(u64, 2000), Shared.counter);}
test "Zio: cancel reaches a fiber parked on io promptly" { var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
// a connected socket pair where no data ever arrives: the reader fiber // parks in netRead indefinitely; cancel must wake it with error.Canceled. const listener = try engine.tcpListenIp4(.{ 127, 0, 0, 1 }, 0, 4); defer _ = c.close(listener.fd);
const Reader = struct { fn readForever(io_inner: Io, port: u16) u64 { const addr = Io.net.IpAddress{ .ip4 = .{ .bytes = .{ 127, 0, 0, 1 }, .port = port } }; const stream = Io.net.IpAddress.connect(&addr, io_inner, .{ .mode = .stream }) catch return 99; defer stream.close(io_inner); var rbuf: [64]u8 = undefined; var reader = stream.reader(io_inner, &rbuf); var byte: [1]u8 = undefined; reader.interface.readSliceAll(&byte) catch return 1; // canceled -> ReadFailed path return 0; } };
var fut = try io.concurrent(Reader.readForever, .{ io, listener.port }); // let it connect and park (accept side never sends anything) const accepted = try engine.acceptNonblock(listener.fd) orelse blk: { var spins: u32 = 0; while (spins < 500) : (spins += 1) { const ts: c.timespec = .{ .sec = 0, .nsec = 5 * std.time.ns_per_ms }; _ = c.nanosleep(&ts, null); if (try engine.acceptNonblock(listener.fd)) |fd2| break :blk fd2; } return error.NeverConnected; }; defer _ = c.close(accepted); { const ts: c.timespec = .{ .sec = 0, .nsec = 50 * std.time.ns_per_ms }; _ = c.nanosleep(&ts, null); // give the fiber time to park in netRead }
const t0 = sched_mod.nowMs(); const r = fut.cancel(io); // must return promptly, not hang const elapsed = sched_mod.nowMs() - t0; try std.testing.expectEqual(@as(u64, 1), r); try std.testing.expect(elapsed < 500);}
test { _ = Zio;}
test "Zio: group members run as fibers; await drains; cancel interrupts" { var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
const Ctx = struct { var counter: std.atomic.Value(u32) = .init(0); var on_fiber: std.atomic.Value(u32) = .init(0);
fn member(z: *Zio, n: u32) Io.Cancelable!void { if (z.currentFiber() != null) _ = on_fiber.fetchAdd(1, .monotonic); var i: u32 = 0; while (i < n) : (i += 1) { io_g.sleep(Io.Duration.fromMicroseconds(100), .awake) catch {}; } _ = counter.fetchAdd(n, .monotonic); }
fn napMember(_: *Zio) Io.Cancelable!void { // parks long; only a cancel interrupt ends it early. io_g.sleep(Io.Duration.fromSeconds(3600), .awake) catch return error.Canceled; return error.Canceled; }
var io_g: Io = undefined;
fn driver(z: *Zio) void { var group: Io.Group = .init; group.async(io_g, member, .{ z, 1 }); group.async(io_g, member, .{ z, 2 }); group.async(io_g, member, .{ z, 3 }); group.await(io_g) catch unreachable; std.debug.assert(counter.load(.acquire) == 6); std.debug.assert(on_fiber.load(.acquire) == 3);
var cancel_group: Io.Group = .init; cancel_group.async(io_g, napMember, .{z}); cancel_group.cancel(io_g); } }; Ctx.io_g = io;
var fut = try io.concurrent(Ctx.driver, .{&zio}); fut.await(io); try std.testing.expectEqual(@as(u32, 6), Ctx.counter.load(.acquire));}
test "Zio: reader and writer fibers on one fd both wake (poller arm collision)" { // one epoll registration per fd + one-shot MOD means a writer's arm can // clobber a parked reader's interest and token, stranding the reader. var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
var fds: [2]c.fd_t = undefined; try std.testing.expect(c.socketpair(c.AF.UNIX, c.SOCK.STREAM, 0, &fds) == 0); const shared = fds[0]; const peer = fds[1]; defer _ = c.close(peer); try engine.setNonblock(shared); try engine.setNonblock(peer);
const Ctx = struct { var reader_done = std.atomic.Value(bool).init(false); var writer_done = std.atomic.Value(bool).init(false);
fn reader(io_: Io, fd: c.fd_t) void { var byte: [1]u8 = undefined; var vec = [_][]u8{byte[0..]}; _ = vtOperate(io_.userdata, .{ .net_read = .{ .socket_handle = fd, .data = &vec } }) catch {}; reader_done.store(true, .release); }
fn writer(io_: Io, fd: c.fd_t) void { // fill the send buffer until the fiber parks on writability, // then keep writing as the peer drains. // darwin AF_UNIX socketpair buffers are 8 KiB (linux ~208 KiB), // and the main thread frees at most one buffer per drain pass. // 128 KiB stays comfortably inside 50 passes on both; the // original 512 KiB target was arithmetically impossible on // darwin — the writer could never finish, and the close() at // the end then stranded its parked fiber (kqueue had no // close-while-parked handling), turning a guaranteed failure // into a suite hang. var chunk: [4096]u8 = @splat(0xaa); var total: usize = 0; while (total < 128 * 1024) { const data = [_][]const u8{chunk[0..]}; const result = vtOperate(io_.userdata, .{ .net_write = .{ .socket_handle = fd, .data = &data } }) catch break; const n = result.net_write catch break; total += n; } writer_done.store(true, .release); } };
var rf = try io.concurrent(Ctx.reader, .{ io, shared }); var wf = try io.concurrent(Ctx.writer, .{ io, shared });
// give both fibers time to park on the same fd (reader first, writer // after the send buffer fills), then drain the peer and send one byte. std.Thread.yield() catch {}; var spin: u32 = 0; while (spin < 50) : (spin += 1) { var drain: [65536]u8 = undefined; while (c.read(peer, &drain, drain.len) > 0) {} _ = c.write(peer, "x", 1); if (Ctx.reader_done.load(.acquire) and Ctx.writer_done.load(.acquire)) break; var pause: u32 = 0; while (pause < 2_000_000) : (pause += 1) std.atomic.spinLoopHint(); }
const reader_ok = Ctx.reader_done.load(.acquire); const writer_ok = Ctx.writer_done.load(.acquire); if (!reader_ok or !writer_ok) { std.debug.print("ARM COLLISION: reader_done={} writer_done={}\n", .{ reader_ok, writer_ok }); } _ = c.close(shared); rf.await(io); wf.await(io); try std.testing.expect(reader_ok); try std.testing.expect(writer_ok);}
test "Zio: closing an fd releases the fiber parked on it" { // regression: zio did not override netClose, so a close stranded any // fiber parked on that fd — the kernel drops closed fds from the epoll // set without an event. under connection churn that leaked one fiber // per connection, and a stranded ACCEPT fiber is a listener that goes // deaf while the rest of the process looks healthy. var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
var fds: [2]c.fd_t = undefined; try std.testing.expect(c.socketpair(c.AF.UNIX, c.SOCK.STREAM, 0, &fds) == 0); try engine.setNonblock(fds[0]); try engine.setNonblock(fds[1]); defer _ = c.close(fds[1]);
const Ctx = struct { var returned = std.atomic.Value(bool).init(false);
fn reader(io_: Io, fd: c.fd_t) void { var byte: [1]u8 = undefined; var vec = [_][]u8{byte[0..]}; // parks: nothing will ever be written to this socket _ = vtOperate(io_.userdata, .{ .net_read = .{ .socket_handle = fd, .data = &vec } }) catch {}; returned.store(true, .release); } };
var rf = try io.concurrent(Ctx.reader, .{ io, fds[0] }); // let the reader reach its park before closing under it io.sleep(Io.Duration.fromMilliseconds(50), .awake) catch {}; try std.testing.expect(!Ctx.returned.load(.acquire));
const sockets = [_]Io.net.Socket{.{ .handle = fds[0], .address = .{ .ip4 = .unspecified(0) } }}; io.vtable.netClose(io.userdata, &sockets);
// without the fix this never returns and the test hangs var spins: u32 = 0; while (!Ctx.returned.load(.acquire) and spins < 200) : (spins += 1) { io.sleep(Io.Duration.fromMilliseconds(10), .awake) catch {}; } rf.await(io); try std.testing.expect(Ctx.returned.load(.acquire));}
test "Zio: a sub-millisecond sleep parks instead of busy-spinning" { // regression: the sched clock is millisecond-granular and durations were // truncated, so io.sleep(100us) produced a deadline of `now` and returned // without parking. `while (!done) io.sleep(100us)` — the shape used by // every DbRequest wait in zlay — then spun the loop thread, starving the // accept fibers. that was the 2026-08-04 field failure. var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
const Ctx = struct { fn sleeper(io_: Io, out: *u64) void { const t0 = sched_mod.nowMs(); var i: u32 = 0; // 200 sub-millisecond sleeps. spinning finishes ~instantly; // actually parking takes at least ~200ms. while (i < 200) : (i += 1) { io_.sleep(Io.Duration.fromMicroseconds(100), .awake) catch {}; } out.* = sched_mod.nowMs() -| t0; } }; var elapsed: u64 = 0; var fut = try io.concurrent(Ctx.sleeper, .{ io, &elapsed }); fut.await(io);
// each 100us sleep must round up to a real 1ms park. try std.testing.expect(elapsed >= 150);}
test "Zio: cancelling a future whose task already finished is safe" { // zlay's forceReconnect (admin/hosts/reconnect) sets a shutdown flag on a // worker, releases the lock, then cancels that worker's future. A healthy // worker can observe the flag and exit in the gap, so cancel routinely // lands on an already-completed task. That must not be a use-after-free. // // It is not, by construction: Task starts at refs=2 ("the completing // fiber and the awaiting caller"), so completion drops it to 1 and the // Task outlives its fiber until the future is awaited or cancelled -- // exactly once, which is the std.Io contract. vtCancel also operates on // the Task, not the Fiber, and interruptTask only compares `f.task` // pointers against the live-fiber list rather than dereferencing a dead // one. This pins that down so a future refcount change cannot quietly // break the recovery path. var zio: Zio = undefined; try zio.init(std.heap.c_allocator, .{}); defer zio.deinit(); try zio.start(); const io = zio.io();
const Worker = struct { fn finishImmediately(done: *std.atomic.Value(bool)) u64 { done.store(true, .release); return 7; } };
var done = std.atomic.Value(bool).init(false); var fut = try io.concurrent(Worker.finishImmediately, .{&done});
// wait for the task to actually run to completion, so the cancel below // cannot be racing an in-flight fiber. var spins: u32 = 0; while (!done.load(.acquire) and spins < 2000) : (spins += 1) { const ts: c.timespec = .{ .sec = 0, .nsec = 2 * std.time.ns_per_ms }; _ = c.nanosleep(&ts, null); } try std.testing.expect(done.load(.acquire)); // give the runner its own unref (state -> done) before cancelling spins = 0; while (spins < 50) : (spins += 1) { const ts: c.timespec = .{ .sec = 0, .nsec = 2 * std.time.ns_per_ms }; _ = c.nanosleep(&ts, null); }
// the operation under test: cancel a future that has already finished. const result = fut.cancel(io); try std.testing.expectEqual(@as(u64, 7), result);}