//! Stream transport adapters. Namespace paths, mounting, connection limits, and //! event-loop scheduling belong to the application. POSIX descriptors are owned //! by the caller; std.Io streams retain the standard library ownership contract. const std = @import("std"); const libc = std.c; const wire = @import("wire.zig"); pub const darwin = @import("builtin").os.tag.isDarwin(); pub const sun_path_len = @typeInfo(@FieldType(libc.sockaddr.un, "path")).array.len; pub const Address = union(enum) { unix: [:0]const u8, tcp: std.Io.net.IpAddress }; pub const Error = error{ Socket, Flags, SocketOption, BadAddress, Bind, Listen, Connect, Closed, Timeout, Io }; /// Read one frame from any std.Io reader, including TCP, Unix sockets, or files. pub fn readFrame(reader: *std.Io.Reader, buffer: []u8, msize: u32) ![]const u8 { if (buffer.len < wire.header_len) return error.NoSpace; try reader.readSliceAll(buffer[0..4]); const len = wire.frameLen(buffer[0..4]).?; if (len < wire.header_len) return error.BadValue; if (len > msize or len > buffer.len) return error.Overlong; try reader.readSliceAll(buffer[4..len]); return buffer[0..len]; } /// Batching and flushing are explicit: this does not flush the writer. pub fn writeFrame(writer: *std.Io.Writer, frame: []const u8, msize: u32) !void { if (frame.len > msize) return error.Overlong; _ = try wire.decode(frame); try writer.writeAll(frame); } pub fn connect(io: std.Io, address: Address) !std.Io.net.Stream { return switch (address) { .tcp => |ip| ip.connect(io, .{ .mode = .stream }), .unix => |path| (try std.Io.net.UnixAddress.init(path)).connect(io), }; } pub fn listen(io: std.Io, address: Address, backlog: u31) !std.Io.net.Server { return switch (address) { .tcp => |ip| ip.listen(io, .{ .kernel_backlog = backlog }), .unix => |path| (try std.Io.net.UnixAddress.init(path)).listen(io, .{ .kernel_backlog = backlog }), }; } pub fn nowMs() i64 { var ts: libc.timespec = undefined; if (libc.clock_gettime(.MONOTONIC, &ts) != 0) return std.math.maxInt(i64); return @as(i64, ts.sec) * std.time.ms_per_s + @divTrunc(ts.nsec, std.time.ns_per_ms); } pub fn configure(fd: c_int, tcp: bool) Error!void { const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0)); if (flags < 0) return error.Flags; var options: libc.O = @bitCast(@as(u32, @bitCast(flags))); options.NONBLOCK = true; if (libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(options))))) != 0 or libc.fcntl(fd, libc.F.SETFD, @as(c_int, libc.FD_CLOEXEC)) != 0) return error.Flags; const on: c_int = 1; if (tcp and libc.setsockopt(fd, libc.IPPROTO.TCP, libc.TCP.NODELAY, &on, @sizeOf(c_int)) != 0) return error.SocketOption; if (comptime darwin) { if (libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.NOSIGPIPE, &on, @sizeOf(c_int)) != 0) return error.SocketOption; } } pub fn ipSockaddr(ip: std.Io.net.IpAddress, out: *libc.sockaddr.storage) libc.socklen_t { return switch (ip) { .ip4 => |a| blk: { const addr: *libc.sockaddr.in = @ptrCast(out); addr.* = .{ .port = std.mem.nativeToBig(u16, a.port), .addr = @bitCast(a.bytes) }; break :blk @sizeOf(libc.sockaddr.in); }, .ip6 => |a| blk: { const addr: *libc.sockaddr.in6 = @ptrCast(out); addr.* = .{ .port = std.mem.nativeToBig(u16, a.port), .addr = a.bytes, .flowinfo = 0, .scope_id = a.interface.index }; break :blk @sizeOf(libc.sockaddr.in6); }, }; } fn sockaddr(address: Address, out: *libc.sockaddr.storage) Error!libc.socklen_t { return switch (address) { .tcp => |ip| ipSockaddr(ip, out), .unix => |path| blk: { if (path.len >= sun_path_len or std.mem.indexOfScalar(u8, path, 0) != null) return error.BadAddress; const addr: *libc.sockaddr.un = @ptrCast(out); addr.* = .{ .path = @splat(0) }; @memcpy(addr.path[0..path.len], path); break :blk @sizeOf(libc.sockaddr.un); }, }; } pub fn wait(fd: c_int, events: i16, deadline_ms: i64) Error!void { while (true) { const left = deadline_ms -| nowMs(); if (left <= 0) return error.Timeout; var fds = [1]libc.pollfd{.{ .fd = fd, .events = events, .revents = 0 }}; const rc = libc.poll(&fds, 1, @intCast(@min(left, std.math.maxInt(c_int)))); if (rc < 0) { if (libc.errno(rc) == .INTR) continue; return error.Io; } if (rc == 0) continue; if (nowMs() >= deadline_ms) return error.Timeout; if (fds[0].revents & events != 0) return; return error.Closed; } } /// Connect a nonblocking socket with an absolute monotonic deadline. pub fn connectFd(address: Address, deadline_ms: i64) Error!c_int { var addr: libc.sockaddr.storage = undefined; const len = try sockaddr(address, &addr); const fd = libc.socket(addr.family, libc.SOCK.STREAM, 0); if (fd < 0) return error.Socket; errdefer close(fd); try configure(fd, address == .tcp); const rc = libc.connect(fd, @ptrCast(&addr), len); if (rc != 0) { switch (libc.errno(rc)) { .INPROGRESS, .ALREADY, .INTR => {}, else => return error.Connect, } try wait(fd, @intCast(libc.POLL.OUT), deadline_ms); var status: c_int = 0; var size: libc.socklen_t = @sizeOf(c_int); if (libc.getsockopt(fd, libc.SOL.SOCKET, libc.SO.ERROR, @ptrCast(&status), &size) != 0) return error.Connect; if (status != 0) return error.Connect; } return fd; } /// Does not unlink Unix paths or alter their permissions. pub fn listenFd(address: Address, backlog: u31) Error!c_int { var addr: libc.sockaddr.storage = undefined; const len = try sockaddr(address, &addr); const fd = libc.socket(addr.family, libc.SOCK.STREAM, 0); if (fd < 0) return error.Socket; errdefer close(fd); try configure(fd, false); if (address == .tcp) { const on: c_int = 1; if (libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.REUSEADDR, &on, @sizeOf(c_int)) != 0) return error.SocketOption; if (address.tcp == .ip6) { const v6only = if (darwin) 27 else std.os.linux.IPV6.V6ONLY; if (libc.setsockopt(fd, libc.IPPROTO.IPV6, v6only, &on, @sizeOf(c_int)) != 0) return error.SocketOption; } } if (libc.bind(fd, @ptrCast(&addr), len) != 0) return error.Bind; if (libc.listen(fd, backlog) != 0) return error.Listen; return fd; } pub fn acceptFd(listener: c_int, tcp: bool) Error!?c_int { const fd = libc.accept(listener, null, null); if (fd < 0) return switch (libc.errno(fd)) { .AGAIN, .INTR, .CONNABORTED => null, else => error.Socket, }; errdefer close(fd); try configure(fd, tcp); return fd; } /// null means retry after readiness; zero means EOF. Empty buffers are forbidden. pub fn read(fd: c_int, buffer: []u8) Error!?usize { std.debug.assert(buffer.len > 0); const n = libc.read(fd, buffer.ptr, buffer.len); if (n < 0) return switch (libc.errno(n)) { .INTR, .AGAIN => null, else => error.Io, }; return @intCast(n); } pub fn write(fd: c_int, bytes: []const u8) Error!?usize { std.debug.assert(bytes.len > 0); const n = libc.send(fd, bytes.ptr, bytes.len, if (darwin) 0 else libc.MSG.NOSIGNAL); if (n < 0) return switch (libc.errno(n)) { .INTR, .AGAIN => null, else => error.Io, }; if (n == 0) return error.Closed; return @intCast(n); } pub fn close(fd: c_int) void { _ = libc.close(fd); } /// Probe a Unix listener without changing namespace entries. Uncertainty is live. pub fn isListening(path: [:0]const u8) bool { var addr: libc.sockaddr.un = .{ .path = @splat(0) }; if (path.len + 1 > sun_path_len) return true; // cannot ask; assume occupied @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]); const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0); if (fd < 0) return true; defer _ = libc.close(fd); configure(fd, false) catch return true; if (libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) == 0) return true; return libc.errno(-1) != .CONNREFUSED; }