From d0091f9b6c8f2d871a4b564f3c76e7987fcca0b3 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sun, 6 Sep 2026 17:18:37 -0500 Subject: [PATCH] fix: return structured dictionary rejection and request full admission --- .tangled/workflows/admission.yml | 2 +- src/internal/serve/server.zig | 59 ++++++++++++++++++++++++++++++-- 2 files changed, 58 insertions(+), 3 deletions(-) diff --git a/.tangled/workflows/admission.yml b/.tangled/workflows/admission.yml index bee6ef9..a5b1623 100644 --- a/.tangled/workflows/admission.yml +++ b/.tangled/workflows/admission.yml @@ -45,7 +45,7 @@ steps: "$PREFECT_API_URL/deployments/$deployment/create_flow_run" \ -H "Authorization: Basic $auth" \ -H 'Content-Type: application/json' \ - -d "{\"parameters\": {\"sha\": \"$TANGLED_COMMIT_SHA\", \"build_image\": true}, + -d "{\"parameters\": {\"sha\": \"$TANGLED_COMMIT_SHA\", \"only\": [], \"skip\": [], \"build_image\": true}, \"tags\": [\"spindle\", \"$TANGLED_COMMIT_SHA\"]}" \ | head -c 400 echo diff --git a/src/internal/serve/server.zig b/src/internal/serve/server.zig index c40f336..7175692 100644 --- a/src/internal/serve/server.zig +++ b/src/internal/serve/server.zig @@ -669,8 +669,8 @@ pub const Handler = struct { } if (id != zstd.v2_dictionary_id) { var buf: [224]u8 = undefined; - const body = std.fmt.bufPrint(&buf, "unknown zstd dictionary id {d}; current dictionary id is {d} (fetch it via getZstdDictionary and reconnect)\n", .{ id, zstd.v2_dictionary_id }) catch unreachable; - respond(conn, "400 Bad Request", "text/plain", body); + const body = std.fmt.bufPrint(&buf, "{{\"error\":\"UnknownZstdDictionary\",\"message\":\"dictionary {d} is no longer current; fetch dictionary {d} and reconnect\"}}\n", .{ id, zstd.v2_dictionary_id }) catch unreachable; + respond(conn, "400 Bad Request", "application/json", body); finishHttp(hub, http_started, http_handler, 400); return error.Close; } @@ -2637,3 +2637,58 @@ test "real-HTTP archive serving does not leak per request (prod RSS ratchet regr try testing.expectEqual(std.http.Status.ok, response.status); } } + +test "obsolete live dictionary returns the upstream JSON rejection" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tail = tail_mod.Tail.init(testing.allocator, io, 1 << 20); + defer tail.deinit(); + var stats: metrics.Stats = .{}; + var hub: Hub = .{ + .allocator = testing.allocator, + .io = io, + .tail = &tail, + .stats = &stats, + .serving = .init(true), + }; + defer hub.deinit(); + + var listener = (Io.net.IpAddress{ .ip4 = .loopback(0) }).listen(io, .{ .reuse_address = true }) catch unreachable; + defer listener.deinit(io); + const port = listener.socket.address.getPort(); + var ws_server = try websocket.Server(Handler).init(testing.allocator, io, .{ + .port = port, + .address = "127.0.0.1", + .max_conn = 4, + .handshake = .{ .max_size = 16 * 1024, .max_headers = 64 }, + }); + defer ws_server.deinit(); + const Fixture = struct { + fn run(server: *websocket.Server(Handler), server_listener: *Io.net.Server, fixture_hub: *Hub) !void { + server.runIo(server_listener, fixture_hub); + } + }; + var future = try io.concurrent(Fixture.run, .{ &ws_server, &listener, &hub }); + defer _ = future.cancel(io) catch {}; + + const address: Io.net.IpAddress = .{ .ip4 = .loopback(port) }; + var socket = try address.connect(io, .{ .mode = .stream }); + defer socket.close(io); + var write_buf: [1024]u8 = undefined; + var writer = socket.writer(io, &write_buf); + try writer.interface.print("GET /xrpc/network.bsky.jetstream.subscribeEvents?zstdDictionary={d} HTTP/1.1\r\nHost: localhost\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Protocol: xrpc.v1.json\r\n\r\n", .{zstd.v2_dictionary_id ^ 1}); + try writer.interface.flush(); + var read_buf: [4096]u8 = undefined; + var reader = socket.reader(io, &read_buf); + const response = try reader.interface.allocRemaining(testing.allocator, .limited(4096)); + defer testing.allocator.free(response); + try testing.expect(std.mem.startsWith(u8, response, "HTTP/1.1 400 ")); + const split = std.mem.indexOf(u8, response, "\r\n\r\n") orelse return error.MissingHeaders; + try testing.expect(std.mem.indexOf(u8, response[0..split], "application/json") != null); + const parsed = try std.json.parseFromSlice(std.json.Value, testing.allocator, response[split + 4 ..], .{}); + defer parsed.deinit(); + try testing.expectEqualStrings("UnknownZstdDictionary", parsed.value.object.get("error").?.string); + try testing.expect(std.mem.indexOf(u8, parsed.value.object.get("message").?.string, "reconnect") != null); +} -- 2.51.2