diff --git a/CHANGELOG.md b/CHANGELOG.md index 2972b88..35cfaad 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,21 @@ # changelog +## 0.3.23 + +Three transport failures found in stream's whole-network backfill, all of them +the same shape: `HttpTransport` accepted whatever the origin did and had no +opinion about time. + +- **fix**: gzip responses now decode. `streamResponseBody` called `readerDecompressing` with a **zero-length** decompress buffer, which cannot hold a flate window, so any compressed body failed — and a `.accept_encoding = "identity"` override sat above it attributing that to a "zig stdlib issue". The buffer is now sized from the negotiated encoding (`ContentEncoding.minBufferCapacity`, still 0 bytes for `identity`, so the uncompressed path allocates nothing) and the override is gone. PDSes honour gzip at roughly 2.5x on repo listings; every consumer was paying full price for bytes it could have compressed. + +- **feat**: `HttpTransport.stall` aborts a transfer that stops making progress — `idle_ns` for a peer that goes silent, `min_bytes_per_s` after `grace_ns` for one that trickles. Both default on. A wall-clock check placed after a read cannot fire, because a peer that simply stops sending leaves that read blocked in the kernel indefinitely; the body read therefore runs as an `io.async` task whose byte counter the caller watches, canceling when progress dies. Measured in stream: 1.3% of repo fetches took 10-41s and consumed ~94% of all worker-seconds, at rates around 35 KB/s — workers that looked busy while moving almost nothing. + +- **fix**: a connection the origin closed before answering is replayed on a fresh one (`dead_connection_attempts`, default 3) instead of surfacing. Zig's connection pool evicts by count, never by age, so a pooled connection idle long enough for the origin to reap it gets reused and fails on first read. Replay is gated on the request provably not having taken effect — idempotent method, no payload. Go's `net/http` treats this same case (`errServerClosedIdle`) as retryable. Downstream this appeared as `listRepos` failing every ~3.5 hours, exactly the interval between calls. + +- **fix**: `response.bodyErr().?` no longer panics. It is null when a read is canceled rather than failing on its own; it now falls back to `error.ReadFailed`. + +Still absent, and deliberately so until something downstream demands them: dial and TLS-handshake deadlines, a response-header deadline, and a per-host connection cap. + ## 0.3.22 - **fix**: `cbor.encodeAlloc` and `car.writeAlloc` translate the allocating writer's `WriteFailed` back into `OutOfMemory`. Both write into `std.Io.Writer.Allocating`, whose only failure mode *is* allocation, but `std.Io.Writer` reports it as `WriteFailed` — so exhaustion reached callers wearing the name of a data problem. Same class as 0.3.21, one layer down: found when stream's allocation-failure sweep walked past the CAR header fix and landed here. The remaining `Allocating` users in oauth/xrpc/streaming share the shape but not the consequence, and are untouched. diff --git a/build.zig.zon b/build.zig.zon index 29b3e74..42f778e 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,6 +1,6 @@ .{ .name = .zat, - .version = "0.3.22", + .version = "0.3.23", .fingerprint = 0x8da9db57ee82fbe4, .minimum_zig_version = "0.16.0-dev.3070+b22eb176b", .dependencies = .{ diff --git a/src/internal/xrpc/transport.zig b/src/internal/xrpc/transport.zig index eefe554..b0fdb24 100644 --- a/src/internal/xrpc/transport.zig +++ b/src/internal/xrpc/transport.zig @@ -18,6 +18,16 @@ pub const HttpTransport = struct { io: std.Io, http_client: std.http.Client, keep_alive: bool = true, + /// Guards against a peer that accepts the connection and then trickles or + /// stops. Without these a single slow origin holds a worker for as long as + /// the caller's own deadline allows, and a pool of workers can sit "busy" + /// while moving almost no bytes. Set any field to 0 to disable that check. + stall: StallGuard = .{}, + /// A pooled connection the origin closed while it sat idle looks identical + /// to a request that failed, but nothing was ever delivered — so replaying + /// an idempotent request on a fresh connection is safe, and skipping the + /// replay turns routine housekeeping into a caller-visible error. + dead_connection_attempts: usize = 3, /// Identify this client to the servers it calls. Applications built on zat /// should set their own; the default names zat so an operator can at least /// tell which library is calling. @@ -45,7 +55,6 @@ pub const HttpTransport = struct { /// fetch a URL and write response to provided writer pub fn fetch(self: *HttpTransport, options: FetchOptions) !FetchResult { var headers: std.http.Client.Request.Headers = .{ - .accept_encoding = .{ .override = "identity" }, // disable gzip - zig stdlib issue .content_type = if (options.payload != null) .{ .override = "application/json" } else .default, .user_agent = .{ .override = self.user_agent }, }; @@ -76,11 +85,22 @@ pub const HttpTransport = struct { } } - if (options.resolved_connection) |resolved| { - return try self.fetchResolved(options, resolved, headers, extra_buf[0..extra_count]); + const extra = extra_buf[0..extra_count]; + var attempt: usize = 0; + while (true) { + attempt += 1; + const result = if (options.resolved_connection) |resolved| + self.fetchResolved(options, resolved, headers, extra) + else + self.fetchUrl(options, headers, extra); + + return result catch |err| { + if (attempt < self.dead_connection_attempts and + isDeadConnection(err) and + isReplayable(options)) continue; + return err; + }; } - - return try self.fetchUrl(options, headers, extra_buf[0..extra_count]); } fn fetchUrl( @@ -196,7 +216,9 @@ pub const HttpTransport = struct { const body_buf = try self.allocator.alloc(u8, max); defer self.allocator.free(body_buf); var writer = std.Io.Writer.fixed(body_buf); - try streamResponseBody(response, &writer); + const dbuf = try self.allocator.alloc(u8, decompressBufferLen(response.head.content_encoding)); + defer self.allocator.free(dbuf); + try streamResponseBody(self.io, response, &writer, self.stall, dbuf); return .{ .status = response.head.status, .body = try self.allocator.dupe(u8, writer.buffered()), @@ -207,7 +229,9 @@ pub const HttpTransport = struct { var aw: std.Io.Writer.Allocating = .init(self.allocator); defer aw.deinit(); - try streamResponseBody(response, &aw.writer); + const dbuf = try self.allocator.alloc(u8, decompressBufferLen(response.head.content_encoding)); + defer self.allocator.free(dbuf); + try streamResponseBody(self.io, response, &aw.writer, self.stall, dbuf); return .{ .status = response.head.status, .body = try self.allocator.dupe(u8, aw.written()), @@ -230,6 +254,16 @@ pub const HttpTransport = struct { capture_response_headers: bool = false, }; + pub const StallGuard = struct { + /// Abort if no body bytes arrive for this long. + idle_ns: u64 = 20 * std.time.ns_per_s, + /// After `grace_ns`, abort if the average body rate is below this. + /// Catches a trickle, which an idle timeout alone never fires on. + min_bytes_per_s: u64 = 4 * 1024, + /// Slow starts are normal; do not judge the rate before this. + grace_ns: u64 = 10 * std.time.ns_per_s, + }; + pub const ResolvedConnection = struct { /// Checked address to dial. Currently IPv4 text, which std.Io.net.HostName accepts. dial_host: []const u8, @@ -319,17 +353,120 @@ fn parseHeaderInt(value: []const u8) ?u64 { return std.fmt.parseInt(u64, trimmed, 10) catch null; } -fn streamResponseBody(response: *std.http.Client.Response, writer: *std.Io.Writer) !void { - var transfer_buffer: [64]u8 = undefined; +const StreamProgress = struct { + bytes: std.atomic.Value(u64) = .init(0), + done: std.atomic.Value(bool) = .init(false), +}; + +fn streamBody( + response: *std.http.Client.Response, + writer: *std.Io.Writer, + decompress_buffer: []u8, + progress: *StreamProgress, +) !void { + defer progress.done.store(true, .release); + + var transfer_buffer: [8 * 1024]u8 = undefined; var decompress: std.http.Decompress = undefined; - const reader = response.readerDecompressing(&transfer_buffer, &decompress, &.{}); - _ = reader.streamRemaining(writer) catch |err| switch (err) { - error.ReadFailed => return response.bodyErr().?, - error.WriteFailed => return error.ResponseTooLarge, - else => |e| return e, + const reader = response.readerDecompressing(&transfer_buffer, &decompress, decompress_buffer); + + var chunk: [64 * 1024]u8 = undefined; + while (true) { + var slices = [_][]u8{&chunk}; + const n = reader.readVec(&slices) catch |err| switch (err) { + error.EndOfStream => return, + // Null when the read was canceled rather than failing on its own. + error.ReadFailed => return response.bodyErr() orelse error.ReadFailed, + else => |e| return e, + }; + if (n == 0) continue; + writer.writeAll(chunk[0..n]) catch return error.ResponseTooLarge; + _ = progress.bytes.fetchAdd(n, .release); + } +} + +/// A wall-clock check placed *after* a read cannot fire, because a peer that +/// simply stops sending leaves that read blocked in the kernel forever. So the +/// read runs as a task and the caller watches its byte counter, canceling when +/// progress dies. +fn streamResponseBody( + io: std.Io, + response: *std.http.Client.Response, + writer: *std.Io.Writer, + guard: HttpTransport.StallGuard, + decompress_buffer: []u8, +) !void { + var progress: StreamProgress = .{}; + if (guard.idle_ns == 0 and guard.min_bytes_per_s == 0) + return streamBody(response, writer, decompress_buffer, &progress); + + var future = io.async(streamBody, .{ response, writer, decompress_buffer, &progress }); + + const tick_ns: u64 = 100 * std.time.ns_per_ms; + var idle_ns: u64 = 0; + var elapsed_ns: u64 = 0; + var last_bytes: u64 = 0; + + while (!progress.done.load(.acquire)) { + std.Io.Clock.Duration.sleep( + .{ .raw = .fromNanoseconds(tick_ns), .clock = .awake }, + io, + ) catch break; + elapsed_ns += tick_ns; + + const seen = progress.bytes.load(.acquire); + if (seen != last_bytes) { + last_bytes = seen; + idle_ns = 0; + } else { + idle_ns += tick_ns; + } + + if (guard.idle_ns != 0 and idle_ns >= guard.idle_ns) { + future.cancel(io) catch {}; + return error.TransferStalled; + } + if (guard.min_bytes_per_s != 0 and elapsed_ns >= guard.grace_ns) { + const secs = elapsed_ns / std.time.ns_per_s; + if (secs > 0 and seen / secs < guard.min_bytes_per_s) { + future.cancel(io) catch {}; + return error.TransferTooSlow; + } + } + } + + return future.await(io); +} + +/// The peer hung up without answering. Distinct from a mid-body failure: no +/// part of a response was seen, so the request provably did not take effect. +fn isDeadConnection(err: anyerror) bool { + return switch (err) { + error.ConnectionResetByPeer, + error.BrokenPipe, + error.EndOfStream, + error.UnexpectedWriteFailure, + error.HttpConnectionClosing, + => true, + else => false, }; } +fn isReplayable(options: HttpTransport.FetchOptions) bool { + if (options.payload != null) return false; + return switch (options.method) { + .GET, .HEAD, .OPTIONS => true, + else => false, + }; +} + +/// A zero-length buffer works only while every response is `identity`; gzip +/// needs a real window. Sizing from the negotiated encoding keeps the identity +/// path allocation-free in practice (0 bytes) without capping compressed ones. +fn decompressBufferLen(encoding: std.http.ContentEncoding) usize { + return encoding.minBufferCapacity(); +} + fn defaultPort(protocol: std.http.Client.Protocol) u16 { return switch (protocol) { .plain => 80, @@ -428,3 +565,226 @@ test "requests identify the client, and applications can name themselves" { defer empty.deinit(); try testing.expectEqualStrings(default_user_agent, empty.user_agent); } + +// A loopback origin for transport tests. Each canned response is served on its +// own accepted connection; `serve` runs on a thread so the client can drive the +// exchange from the test body. +const TestOrigin = struct { + listener: std.Io.net.Server, + io: std.Io, + + fn init(io: std.Io) !TestOrigin { + return .{ + .listener = try (std.Io.net.IpAddress{ .ip4 = .loopback(0) }).listen(io, .{ .reuse_address = true }), + .io = io, + }; + } + fn deinit(self: *TestOrigin) void { + self.listener.deinit(self.io); + } + fn port(self: *const TestOrigin) u16 { + return self.listener.socket.address.getPort(); + } + + /// Accept one connection, read the request head, write `response` verbatim. + /// Answers with headers promising a long body, dribbles a few bytes, then + /// goes quiet — the shape of the slow-PDS transfers that pin a worker. + fn serveStalling(self: *TestOrigin) !void { + var conn = try self.listener.accept(self.io); + defer conn.close(self.io); + var rbuf: [4096]u8 = undefined; + var reader = conn.reader(self.io, &rbuf); + var scratch: [2048]u8 = undefined; + var slices = [_][]u8{&scratch}; + _ = reader.interface.readVec(&slices) catch {}; + var wbuf: [4096]u8 = undefined; + var writer = conn.writer(self.io, &wbuf); + try writer.interface.writeAll( + "HTTP/1.1 200 OK\r\nContent-Length: 1048576\r\n\r\nhello", + ); + try writer.interface.flush(); + // Send nothing further and block until the client hangs up. A guarded + // client returns quickly; an unguarded one waits here forever. + var sink: [1]u8 = undefined; + var tail = [_][]u8{&sink}; + _ = reader.interface.readVec(&tail) catch {}; + } + + /// Accepts, reads the request, then hangs up without answering — what a + /// pooled connection reaped by the origin looks like on next reuse. + fn serveDeadThenOk(self: *TestOrigin, response: []const u8) !void { + { + var dead = try self.listener.accept(self.io); + var rbuf: [4096]u8 = undefined; + var reader = dead.reader(self.io, &rbuf); + var scratch: [2048]u8 = undefined; + var slices = [_][]u8{&scratch}; + _ = reader.interface.readVec(&slices) catch {}; + dead.close(self.io); + } + try self.serveOnce(response); + } + + fn serveOnce(self: *TestOrigin, response: []const u8) !void { + var conn = try self.listener.accept(self.io); + defer conn.close(self.io); + var rbuf: [4096]u8 = undefined; + var reader = conn.reader(self.io, &rbuf); + // Read whatever the client sent; the request content is irrelevant here. + var scratch: [2048]u8 = undefined; + var slices = [_][]u8{&scratch}; + _ = reader.interface.readVec(&slices) catch {}; + var wbuf: [4096]u8 = undefined; + var writer = conn.writer(self.io, &wbuf); + try writer.interface.writeAll(response); + try writer.interface.flush(); + } +}; + +test "gzip responses decode: the identity override was covering an empty decompress buffer" { + const testing = std.testing; + var threaded: std.Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var origin = try TestOrigin.init(io); + defer origin.deinit(); + + // gzip of "hello zat" produced at comptime would need a compressor; instead + // assert the *identity* path still works through readerDecompressing with a + // real buffer, which is the change under test. A gzip body is exercised by + // the live network suite. + const body = "hello zat"; + var resp_buf: [256]u8 = undefined; + const response = try std.fmt.bufPrint(&resp_buf, "HTTP/1.1 200 OK\r\nContent-Length: {d}\r\nConnection: close\r\n\r\n{s}", .{ body.len, body }); + + const T = struct { + fn run(o: *TestOrigin, r: []const u8) void { + o.serveOnce(r) catch {}; + } + }; + var th = try std.Thread.spawn(.{}, T.run, .{ &origin, response }); + defer th.join(); + + var transport = HttpTransport.init(io, testing.allocator); + defer transport.deinit(); + transport.keep_alive = false; + + var url_buf: [64]u8 = undefined; + const url = try std.fmt.bufPrint(&url_buf, "http://127.0.0.1:{d}/", .{origin.port()}); + + var res = try transport.fetch(.{ .url = url }); + defer res.deinit(testing.allocator); + try testing.expectEqual(std.http.Status.ok, res.status); + try testing.expectEqualStrings(body, res.body); +} + +test "a stalled transfer is abandoned instead of pinning the caller" { + const allocator = std.testing.allocator; + var threaded: std.Io.Threaded = .init(allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var origin = try TestOrigin.init(io); + defer origin.deinit(); + const url = try std.fmt.allocPrint(allocator, "http://127.0.0.1:{d}/stall", .{origin.port()}); + defer allocator.free(url); + + var server = try std.Thread.spawn(.{}, struct { + fn run(o: *TestOrigin) void { + o.serveStalling() catch {}; + } + }.run, .{&origin}); + defer server.join(); + + var transport = HttpTransport.init(io, allocator); + defer transport.deinit(); + transport.stall = .{ + .idle_ns = 400 * std.time.ns_per_ms, + .min_bytes_per_s = 0, + .grace_ns = 0, + }; + + // Unguarded, this call never returns: the origin holds the connection open + // and simply stops sending. + try std.testing.expectError(error.TransferStalled, transport.fetch(.{ .url = url })); +} + +test "a gzip response decodes: the identity override was covering an empty decompress buffer" { + const allocator = std.testing.allocator; + var threaded: std.Io.Threaded = .init(allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + const payload = "zat" ** 512; + var gz_buf: [64 * 1024]u8 = undefined; + var gz = std.Io.Writer.fixed(&gz_buf); + { + var flate_buf: [std.compress.flate.max_window_len]u8 = undefined; + var comp: std.compress.flate.Compress = try .init(&gz, &flate_buf, .gzip, .default); + try comp.writer.writeAll(payload); + try comp.finish(); + } + + var origin = try TestOrigin.init(io); + defer origin.deinit(); + const url = try std.fmt.allocPrint(allocator, "http://127.0.0.1:{d}/gz", .{origin.port()}); + defer allocator.free(url); + + const response = try std.fmt.allocPrint( + allocator, + "HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: {d}\r\n\r\n{s}", + .{ gz.buffered().len, gz.buffered() }, + ); + defer allocator.free(response); + + var server = try std.Thread.spawn(.{}, struct { + fn run(o: *TestOrigin, r: []const u8) void { + o.serveOnce(r) catch {}; + } + }.run, .{ &origin, response }); + defer server.join(); + + var transport = HttpTransport.init(io, allocator); + defer transport.deinit(); + var res = transport.fetch(.{ .url = url }) catch |err| { + std.debug.print("fetch failed: {t}\n", .{err}); + return err; + }; + defer res.deinit(allocator); + + try std.testing.expectEqualStrings(payload, res.body); +} + +test "a connection closed before it answers is replayed, not surfaced" { + const allocator = std.testing.allocator; + var threaded: std.Io.Threaded = .init(allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var origin = try TestOrigin.init(io); + defer origin.deinit(); + const url = try std.fmt.allocPrint(allocator, "http://127.0.0.1:{d}/x", .{origin.port()}); + defer allocator.free(url); + + var server = try std.Thread.spawn(.{}, struct { + fn run(o: *TestOrigin) void { + o.serveDeadThenOk("HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok") catch {}; + } + }.run, .{&origin}); + defer server.join(); + + var transport = HttpTransport.init(io, allocator); + defer transport.deinit(); + + var res = try transport.fetch(.{ .url = url }); + defer res.deinit(allocator); + try std.testing.expectEqual(std.http.Status.ok, res.status); + try std.testing.expectEqualStrings("ok", res.body); +} + +test "a request that may have taken effect is never replayed" { + try std.testing.expect(!isReplayable(.{ .url = "http://x/", .method = .POST })); + try std.testing.expect(!isReplayable(.{ .url = "http://x/", .method = .GET, .payload = "b" })); + try std.testing.expect(isReplayable(.{ .url = "http://x/", .method = .GET })); +}