//! non-blocking socket helpers (ip4/ip6, libc via std.c). //! shared by the fiber backend, the scheduler, and the demos. const std = @import("std"); const builtin = @import("builtin"); const c = std.c; pub const wouldBlock = switch (builtin.os.tag) { // EAGAIN == EWOULDBLOCK on linux and darwin else => c.E.AGAIN, }; fn sockaddrIn(octets: [4]u8, port: u16) c.sockaddr.in { return .{ .family = c.AF.INET, .port = std.mem.nativeToBig(u16, port), .addr = @bitCast(octets), }; } pub fn setNonblock(fd: c.fd_t) !void { const flags = c.fcntl(fd, c.F.GETFL, @as(c_int, 0)); if (flags < 0) return error.FcntlFailed; const O_NONBLOCK: c_int = @bitCast(c.O{ .NONBLOCK = true }); if (c.fcntl(fd, c.F.SETFL, flags | O_NONBLOCK) < 0) return error.FcntlFailed; } /// blocking connect (loopback: instantaneous), then flips non-blocking. /// good enough for steady-state bench setup; storms use tcpConnectIp4Start. pub fn tcpConnectIp4(octets: [4]u8, port: u16) !c.fd_t { const fd = c.socket(c.AF.INET, c.SOCK.STREAM, 0); if (fd < 0) return error.SocketFailed; errdefer _ = c.close(fd); const sa = sockaddrIn(octets, port); if (c.connect(fd, @ptrCast(&sa), @sizeOf(c.sockaddr.in)) != 0) return error.ConnectFailed; try setNonblock(fd); return fd; } /// begin a non-blocking connect. register with write interest; on the /// writable event call connectResult to learn the outcome. pub fn tcpConnectIp4Start(octets: [4]u8, port: u16) !c.fd_t { const fd = c.socket(c.AF.INET, c.SOCK.STREAM, 0); if (fd < 0) return error.SocketFailed; errdefer _ = c.close(fd); try setNonblock(fd); const sa = sockaddrIn(octets, port); if (c.connect(fd, @ptrCast(&sa), @sizeOf(c.sockaddr.in)) != 0) { const e = c.errno(@as(isize, -1)); if (e != .INPROGRESS) return error.ConnectFailed; } return fd; } /// begin a non-blocking ip6 connect. same contract as tcpConnectIp4Start. pub fn tcpConnectIp6Start(bytes: [16]u8, port: u16, flow: u32, scope: u32) !c.fd_t { const fd = c.socket(c.AF.INET6, c.SOCK.STREAM, 0); if (fd < 0) return error.SocketFailed; errdefer _ = c.close(fd); try setNonblock(fd); var sa = std.mem.zeroes(c.sockaddr.in6); sa.family = c.AF.INET6; sa.port = std.mem.nativeToBig(u16, port); sa.flowinfo = flow; sa.addr = bytes; sa.scope_id = scope; if (c.connect(fd, @ptrCast(&sa), @sizeOf(c.sockaddr.in6)) != 0) { const e = c.errno(@as(isize, -1)); if (e != .INPROGRESS) return error.ConnectFailed; } return fd; } /// after writability on an in-progress connect: did it succeed? pub fn connectResult(fd: c.fd_t) !void { var so_err: c_int = 0; var len: c.socklen_t = @sizeOf(c_int); if (c.getsockopt(fd, c.SOL.SOCKET, c.SO.ERROR, &so_err, &len) != 0) return error.GetSockOptFailed; if (so_err != 0) return error.ConnectFailed; } pub fn tcpListenIp4(octets: [4]u8, port: u16, backlog: u31) !struct { fd: c.fd_t, port: u16 } { const fd = c.socket(c.AF.INET, c.SOCK.STREAM, 0); if (fd < 0) return error.SocketFailed; errdefer _ = c.close(fd); const one: c_int = 1; _ = c.setsockopt(fd, c.SOL.SOCKET, c.SO.REUSEADDR, &one, @sizeOf(c_int)); const sa = sockaddrIn(octets, port); if (c.bind(fd, @ptrCast(&sa), @sizeOf(c.sockaddr.in)) != 0) return error.BindFailed; if (c.listen(fd, backlog) != 0) return error.ListenFailed; var bound: c.sockaddr.in = undefined; var len: c.socklen_t = @sizeOf(c.sockaddr.in); if (c.getsockname(fd, @ptrCast(&bound), &len) != 0) return error.GetSockNameFailed; try setNonblock(fd); return .{ .fd = fd, .port = std.mem.bigToNative(u16, bound.port) }; } /// accept one pending connection, non-blocking. null when drained. pub fn acceptNonblock(listen_fd: c.fd_t) !?c.fd_t { const fd = c.accept(listen_fd, null, null); if (fd < 0) { if (c.errno(fd) == wouldBlock) return null; return error.AcceptFailed; } try setNonblock(fd); return fd; } /// read until EAGAIN; returns total bytes read. 0 with eof=false means EAGAIN /// hit immediately (spurious wakeup). pub fn drainRead(fd: c.fd_t, buf: []u8, total: *u64) enum { open, eof, err } { while (true) { const n = c.read(fd, buf.ptr, buf.len); if (n > 0) { total.* += @intCast(n); continue; } if (n == 0) return .eof; return if (c.errno(n) == wouldBlock) .open else .err; } } /// PUMP semantics: writes `block` REPEATEDLY until the socket buffer is full /// (EAGAIN). this is a throughput firehose, not a send-once helper — a reader /// that keeps draining means this never returns. for one-shot sends use /// writeOnce. pub fn drainWrite(fd: c.fd_t, block: []const u8, total: *u64) enum { open, err } { while (true) { const n = c.write(fd, block.ptr, block.len); if (n > 0) { total.* += @intCast(n); continue; } return if (c.errno(n) == wouldBlock) .open else .err; } } /// single best-effort write (may be short on a full buffer). pub fn writeOnce(fd: c.fd_t, data: []const u8) !usize { const n = c.write(fd, data.ptr, data.len); if (n < 0) { if (c.errno(n) == wouldBlock) return 0; return error.WriteFailed; } return @intCast(n); }