summaryrefslogtreecommitdiff
path: root/src/transport.zig
blob: d7619beeffc511015fde038900dfdc69f4768289 (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
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;
}