diff --git a/CHANGELOG.md b/CHANGELOG.md index 001c376..6dbb507 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,9 @@ # changelog +## 0.2.11 + +- **fix**: enable TCP keepalive on websocket connections — detect dead peers in ~20s instead of blocking forever + ## 0.2.10 - **deps**: bump websocket.zig to fork commit `9e6d732` — TCP split guard for HTTP body reads behind reverse proxies diff --git a/build.zig.zon b/build.zig.zon index 0722172..89faa63 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,6 +1,6 @@ .{ .name = .zat, - .version = "0.2.10", + .version = "0.2.11", .fingerprint = 0x8da9db57ee82fbe4, .minimum_zig_version = "0.15.0", .dependencies = .{ diff --git a/src/internal/streaming/firehose.zig b/src/internal/streaming/firehose.zig index 3bf8c60..97c54ea 100644 --- a/src/internal/streaming/firehose.zig +++ b/src/internal/streaming/firehose.zig @@ -485,6 +485,7 @@ pub const FirehoseClient = struct { const host_header = std.fmt.bufPrint(&host_header_buf, "Host: {s}\r\n", .{host}) catch host; try client.handshake(path, .{ .headers = host_header }); + configureKeepalive(&client); log.info("firehose connected to {s}", .{host}); @@ -527,6 +528,23 @@ fn WsHandler(comptime H: type) type { }; } +/// enable TCP keepalive so reads don't block forever when a peer +/// disappears without FIN/RST (network partition, crash, power loss). +/// detection time: 10s idle + 5s × 2 probes = 20s. +fn configureKeepalive(client: *websocket.Client) void { + const fd = client.stream.stream.handle; + const builtin = @import("builtin"); + posix.setsockopt(fd, posix.SOL.SOCKET, posix.SO.KEEPALIVE, &std.mem.toBytes(@as(i32, 1))) catch return; + const tcp: i32 = @intCast(posix.IPPROTO.TCP); + if (builtin.os.tag == .linux) { + posix.setsockopt(fd, tcp, posix.TCP.KEEPIDLE, &std.mem.toBytes(@as(i32, 10))) catch return; + } else if (builtin.os.tag == .macos) { + posix.setsockopt(fd, tcp, posix.TCP.KEEPALIVE, &std.mem.toBytes(@as(i32, 10))) catch return; + } + posix.setsockopt(fd, tcp, posix.TCP.KEEPINTVL, &std.mem.toBytes(@as(i32, 5))) catch return; + posix.setsockopt(fd, tcp, posix.TCP.KEEPCNT, &std.mem.toBytes(@as(i32, 2))) catch return; +} + // === tests === test "decode frame header" { diff --git a/src/internal/streaming/jetstream.zig b/src/internal/streaming/jetstream.zig index 6801934..fe7a927 100644 --- a/src/internal/streaming/jetstream.zig +++ b/src/internal/streaming/jetstream.zig @@ -202,6 +202,7 @@ pub const JetstreamClient = struct { const host_header = std.fmt.bufPrint(&host_header_buf, "Host: {s}\r\n", .{host}) catch host; try client.handshake(path, .{ .headers = host_header }); + configureKeepalive(&client); log.info("jetstream connected to {s}", .{host}); @@ -270,6 +271,23 @@ fn WsHandler(comptime H: type) type { }; } +/// enable TCP keepalive so reads don't block forever when a peer +/// disappears without FIN/RST (network partition, crash, power loss). +/// detection time: 10s idle + 5s × 2 probes = 20s. +fn configureKeepalive(client: *websocket.Client) void { + const fd = client.stream.stream.handle; + const builtin = @import("builtin"); + posix.setsockopt(fd, posix.SOL.SOCKET, posix.SO.KEEPALIVE, &std.mem.toBytes(@as(i32, 1))) catch return; + const tcp: i32 = @intCast(posix.IPPROTO.TCP); + if (builtin.os.tag == .linux) { + posix.setsockopt(fd, tcp, posix.TCP.KEEPIDLE, &std.mem.toBytes(@as(i32, 10))) catch return; + } else if (builtin.os.tag == .macos) { + posix.setsockopt(fd, tcp, posix.TCP.KEEPALIVE, &std.mem.toBytes(@as(i32, 10))) catch return; + } + posix.setsockopt(fd, tcp, posix.TCP.KEEPINTVL, &std.mem.toBytes(@as(i32, 5))) catch return; + posix.setsockopt(fd, tcp, posix.TCP.KEEPCNT, &std.mem.toBytes(@as(i32, 2))) catch return; +} + // === tests === test "parse commit event" {