summaryrefslogtreecommitdiff
path: root/src/transport.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-14 14:10:28 -0300
committerGabriel Schneider <[email protected]>2026-09-14 14:20:25 -0300
commit5f24c4a2284a0af84fb9d116f7a39f6a58e42ae9 (patch)
tree8679763439492361fe99ea01f5190e5204fa590c /src/transport.zig
downloadcloud9-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.zig207
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;
+}