diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-14 14:10:28 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-14 14:20:25 -0300 |
| commit | 5f24c4a2284a0af84fb9d116f7a39f6a58e42ae9 (patch) | |
| tree | 8679763439492361fe99ea01f5190e5204fa590c /src/transport.zig | |
| download | cloud9-5f24c4a2284a0af84fb9d116f7a39f6a58e42ae9.tar.gz cloud9-5f24c4a2284a0af84fb9d116f7a39f6a58e42ae9.zip | |
Implement base 9P2000 sessions, shared transports, and conformance probes
Diffstat (limited to 'src/transport.zig')
| -rw-r--r-- | src/transport.zig | 207 |
1 files changed, 207 insertions, 0 deletions
diff --git a/src/transport.zig b/src/transport.zig new file mode 100644 index 0000000..d7619be --- /dev/null +++ b/src/transport.zig @@ -0,0 +1,207 @@ +//! 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; +} |
