tickets
All repositories: gitoria
21.4 KB
// Shared WebSocket core — RFC 6455 framing, transport-agnostic (decision D8).// Compiled into each HTTP plugin like http_common.zig; NOT a standalone .so.//// Responsibilities:// - Sec-WebSocket-Accept computation (SHA1 + base64 of key + RFC GUID)// - Frame encoding (unmasked, server → client) via a write callback// - Frame decoding + message assembly (client-masked inbound), fragmentation,// ping/pong surfacing, close handshake with codes, UTF-8 validation//// The per-protocol handshakes stay in each plugin: h1 = Upgrade/101 (http1.zig),// h2 = RFC 8441 extended CONNECT (future), h3 = RFC 9220 (future). Everything// after the handshake — the byte stream — goes through this module unchanged.//// Transport neutrality: encoding writes through `WriteFn` callbacks and decoding// consumes raw bytes via `Decoder.feed()`; the module never touches a socket.const std = @import("std");pub const GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11";/// Write callback: return false on transport failure (peer gone).pub const WriteFn = *const fn (ctx: ?*anyopaque, data: []const u8) bool;pub const Opcode = enum(u4) {continuation = 0x0,text = 0x1,binary = 0x2,close = 0x8,ping = 0x9,pong = 0xA,_,pub fn isControl(self: Opcode) bool {return @intFromEnum(self) >= 0x8;}};// Close codes used by the corepub const CLOSE_NORMAL: u16 = 1000;pub const CLOSE_GOING_AWAY: u16 = 1001;pub const CLOSE_PROTOCOL_ERROR: u16 = 1002;pub const CLOSE_INVALID_PAYLOAD: u16 = 1007;pub const CLOSE_TOO_BIG: u16 = 1009;pub const CLOSE_NO_STATUS: u16 = 1005; // synthesized when close frame has no code — never sent on the wire/// Sec-WebSocket-Accept: base64(SHA1(key ++ GUID)). Output is always 28 chars.pub fn computeAcceptKey(key: []const u8, out: *[28]u8) []const u8 {var sha = std.crypto.hash.Sha1.init(.{});sha.update(key);sha.update(GUID);var digest: [20]u8 = undefined;sha.final(&digest);return std.base64.standard.Encoder.encode(out, &digest);}/// Encode a frame header into `buf` (needs >= 10 bytes). Returns header length./// Server frames are unmasked (RFC 6455 §5.1).pub fn encodeFrameHeader(buf: *[10]u8, opcode: Opcode, payload_len: usize, fin: bool) usize {buf[0] = (if (fin) @as(u8, 0x80) else 0) | @as(u8, @intFromEnum(opcode));if (payload_len <= 125) {buf[1] = @intCast(payload_len);return 2;} else if (payload_len <= 65535) {buf[1] = 126;std.mem.writeInt(u16, buf[2..4], @intCast(payload_len), .big);return 4;} else {buf[1] = 127;std.mem.writeInt(u64, buf[2..10], payload_len, .big);return 10;}}/// Write one complete unmasked frame through the callback. True on success.pub fn writeFrame(write: WriteFn, ctx: ?*anyopaque, opcode: Opcode, payload: []const u8) bool {var hdr: [10]u8 = undefined;const hlen = encodeFrameHeader(&hdr, opcode, payload.len, true);if (!write(ctx, hdr[0..hlen])) return false;if (payload.len > 0 and !write(ctx, payload)) return false;return true;}/// Write a close frame with a code and optional reason (reason truncated to 123).pub fn writeClose(write: WriteFn, ctx: ?*anyopaque, code: u16, reason: []const u8) bool {var payload: [125]u8 = undefined;std.mem.writeInt(u16, payload[0..2], code, .big);const rlen = @min(reason.len, 123);@memcpy(payload[2 .. 2 + rlen], reason[0..rlen]);return writeFrame(write, ctx, .close, payload[0 .. 2 + rlen]);}/// Write one complete MASKED frame — the client's half of RFC 6455 §5.1 ("a/// client MUST mask all frames"). The mask is fresh per frame from the kernel/// CSPRNG (`getrandom(2)`, not a seeded generator: a predictable mask is not/// a security property this module is willing to guess wrong about — same/// reasoning as hl:http1's session-token `randomToken`). A getrandom() short/// read is vanishingly rare for 4 bytes; the leftover is zeroed rather than/// left uninitialized, which only ever makes the mask weaker, never wrong./// Masking is applied through a small stack buffer so a large payload needs/// no allocation.pub fn writeFrameMasked(write: WriteFn, ctx: ?*anyopaque, opcode: Opcode, payload: []const u8) bool {var hdr: [10]u8 = undefined;const hlen = encodeFrameHeader(&hdr, opcode, payload.len, true);hdr[1] |= 0x80; // mask bitvar mask: [4]u8 = .{ 0, 0, 0, 0 };_ = std.os.linux.getrandom(&mask, mask.len, 0);if (!write(ctx, hdr[0..hlen])) return false;if (!write(ctx, &mask)) return false;if (payload.len == 0) return true;var buf: [4096]u8 = undefined;var i: usize = 0;while (i < payload.len) {const n = @min(buf.len, payload.len - i);for (0..n) |j| buf[j] = payload[i + j] ^ mask[(i + j) % 4];if (!write(ctx, buf[0..n])) return false;i += n;}return true;}/// Masked `writeClose` — see `writeFrameMasked`.pub fn writeCloseMasked(write: WriteFn, ctx: ?*anyopaque, code: u16, reason: []const u8) bool {var payload: [125]u8 = undefined;std.mem.writeInt(u16, payload[0..2], code, .big);const rlen = @min(reason.len, 123);@memcpy(payload[2 .. 2 + rlen], reason[0..rlen]);return writeFrameMasked(write, ctx, .close, payload[0 .. 2 + rlen]);}/// A decoded inbound event. Payload slices are allocated with the decoder's/// allocator and owned by the CALLER (free after use).pub const Event = union(enum) {text: []u8,binary: []u8,ping: []u8, // payload to echo in the pongpong: []u8,close: Close,/// Protocol violation — caller should send a close frame with this code/// and drop the connection.protocol_error: u16,pub const Close = struct {code: u16, // CLOSE_NO_STATUS when the frame carried no codereason: []u8,};pub fn deinitPayload(self: *Event, alloc: std.mem.Allocator) void {switch (self.*) {.text, .binary, .ping, .pong => |p| alloc.free(p),.close => |cl| alloc.free(cl.reason),.protocol_error => {},}}};fn validCloseCode(code: u16) bool {return switch (code) {1000...1003, 1007...1011 => true,3000...4999 => true,else => false,};}/// Streaming frame decoder + message assembler. One per connection./// feed() raw transport bytes, then loop next() until it returns null.pub const Decoder = struct {alloc: std.mem.Allocator,/// Client frames must be masked (server side). Set false for client use.require_masked: bool = true,max_message_size: usize = 16 * 1024 * 1024,buf: std.ArrayListUnmanaged(u8) = .empty, // raw incoming bytesread_pos: usize = 0,msg: std.ArrayListUnmanaged(u8) = .empty, // fragmented-message assemblymsg_opcode: ?Opcode = null,failed: bool = false, // after a protocol error the stream is deadpub fn init(alloc: std.mem.Allocator) Decoder {return .{ .alloc = alloc };}pub fn deinit(self: *Decoder) void {self.buf.deinit(self.alloc);self.msg.deinit(self.alloc);}pub fn feed(self: *Decoder, bytes: []const u8) !void {// Compact consumed prefix before growingif (self.read_pos > 0) {const remaining = self.buf.items.len - self.read_pos;std.mem.copyForwards(u8, self.buf.items[0..remaining], self.buf.items[self.read_pos..]);self.buf.shrinkRetainingCapacity(remaining);self.read_pos = 0;}try self.buf.appendSlice(self.alloc, bytes);}fn fail(self: *Decoder, code: u16) Event {self.failed = true;return .{ .protocol_error = code };}/// Pop the next complete event, or null if more bytes are needed.pub fn next(self: *Decoder) !?Event {if (self.failed) return null;const data = self.buf.items[self.read_pos..];if (data.len < 2) return null;const b0 = data[0];const b1 = data[1];const fin = (b0 & 0x80) != 0;const rsv = b0 & 0x70;const opcode: Opcode = @enumFromInt(@as(u4, @truncate(b0 & 0x0F)));const masked = (b1 & 0x80) != 0;const len7: u8 = b1 & 0x7F;if (rsv != 0) return self.fail(CLOSE_PROTOCOL_ERROR); // no extensions negotiatedswitch (@intFromEnum(opcode)) {0x0, 0x1, 0x2, 0x8, 0x9, 0xA => {},else => return self.fail(CLOSE_PROTOCOL_ERROR),}if (self.require_masked and !masked) return self.fail(CLOSE_PROTOCOL_ERROR);if (opcode.isControl() and (!fin or len7 > 125)) return self.fail(CLOSE_PROTOCOL_ERROR);var offset: usize = 2;var payload_len: u64 = len7;if (len7 == 126) {if (data.len < offset + 2) return null;payload_len = std.mem.readInt(u16, data[offset..][0..2], .big);offset += 2;if (payload_len <= 125) return self.fail(CLOSE_PROTOCOL_ERROR); // non-minimal encoding} else if (len7 == 127) {if (data.len < offset + 8) return null;payload_len = std.mem.readInt(u64, data[offset..][0..8], .big);offset += 8;if (payload_len <= 65535) return self.fail(CLOSE_PROTOCOL_ERROR);if (payload_len > (1 << 62)) return self.fail(CLOSE_PROTOCOL_ERROR); // MSB must be 0}if (payload_len > self.max_message_size orself.msg.items.len + payload_len > self.max_message_size)return self.fail(CLOSE_TOO_BIG);var mask: [4]u8 = .{ 0, 0, 0, 0 };if (masked) {if (data.len < offset + 4) return null;mask = data[offset..][0..4].*;offset += 4;}const plen: usize = @intCast(payload_len);if (data.len < offset + plen) return null; // frame incomplete// Unmask into an owned copyconst payload = try self.alloc.alloc(u8, plen);errdefer self.alloc.free(payload);for (0..plen) |i| {payload[i] = data[offset + i] ^ mask[i % 4];}self.read_pos += offset + plen;switch (opcode) {.ping => return .{ .ping = payload },.pong => return .{ .pong = payload },.close => {if (plen == 0) return .{ .close = .{ .code = CLOSE_NO_STATUS, .reason = payload } };if (plen == 1) {self.alloc.free(payload);return self.fail(CLOSE_PROTOCOL_ERROR);}const code = std.mem.readInt(u16, payload[0..2], .big);if (!validCloseCode(code)) {self.alloc.free(payload);return self.fail(CLOSE_PROTOCOL_ERROR);}if (!std.unicode.utf8ValidateSlice(payload[2..])) {self.alloc.free(payload);return self.fail(CLOSE_INVALID_PAYLOAD);}const reason = try self.alloc.alloc(u8, plen - 2);@memcpy(reason, payload[2..]);self.alloc.free(payload);return .{ .close = .{ .code = code, .reason = reason } };},.text, .binary => {if (self.msg_opcode != null) {self.alloc.free(payload);return self.fail(CLOSE_PROTOCOL_ERROR); // new data frame while assembling}if (fin) {if (opcode == .text and !std.unicode.utf8ValidateSlice(payload)) {self.alloc.free(payload);return self.fail(CLOSE_INVALID_PAYLOAD);}return if (opcode == .text) .{ .text = payload } else .{ .binary = payload };}// First fragmentself.msg_opcode = opcode;try self.msg.appendSlice(self.alloc, payload);self.alloc.free(payload);return self.next(); // more frames may already be buffered},.continuation => {const start_op = self.msg_opcode orelse {self.alloc.free(payload);return self.fail(CLOSE_PROTOCOL_ERROR); // continuation without a message};try self.msg.appendSlice(self.alloc, payload);self.alloc.free(payload);if (!fin) return self.next();// Message completeconst full = try self.msg.toOwnedSlice(self.alloc);self.msg_opcode = null;if (start_op == .text and !std.unicode.utf8ValidateSlice(full)) {self.alloc.free(full);return self.fail(CLOSE_INVALID_PAYLOAD);}return if (start_op == .text) .{ .text = full } else .{ .binary = full };},_ => return self.fail(CLOSE_PROTOCOL_ERROR), // unreachable — filtered above}}};// ===========================================================================// Tests// ===========================================================================const t = std.testing;test "computeAcceptKey RFC 6455 sample" {var out: [28]u8 = undefined;const accept = computeAcceptKey("dGhlIHNhbXBsZSBub25jZQ==", &out);try t.expectEqualStrings("s3pPLMBiTxaQ9kYGzzhZRbK+xOo=", accept);}fn maskedFrame(alloc: std.mem.Allocator, fin: bool, opcode: Opcode, payload: []const u8, mask: [4]u8) ![]u8 {var out = std.ArrayListUnmanaged(u8).empty;var hdr: [10]u8 = undefined;const hlen = encodeFrameHeader(&hdr, opcode, payload.len, fin);hdr[1] |= 0x80; // set mask bittry out.appendSlice(alloc, hdr[0..hlen]);try out.appendSlice(alloc, &mask);for (payload, 0..) |b, i| try out.append(alloc, b ^ mask[i % 4]);return out.toOwnedSlice(alloc);}test "decode masked text frame" {var dec = Decoder.init(t.allocator);defer dec.deinit();const frame = try maskedFrame(t.allocator, true, .text, "Hello", .{ 0x37, 0xfa, 0x21, 0x3d });defer t.allocator.free(frame);try dec.feed(frame);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqualStrings("Hello", ev.text);try t.expectEqual(@as(?Event, null), try dec.next());}test "decode split across feeds" {var dec = Decoder.init(t.allocator);defer dec.deinit();const frame = try maskedFrame(t.allocator, true, .text, "chunked delivery", .{ 1, 2, 3, 4 });defer t.allocator.free(frame);try dec.feed(frame[0..3]);try t.expectEqual(@as(?Event, null), try dec.next());try dec.feed(frame[3..]);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqualStrings("chunked delivery", ev.text);}test "fragmented message assembly" {var dec = Decoder.init(t.allocator);defer dec.deinit();const f1 = try maskedFrame(t.allocator, false, .text, "Hel", .{ 9, 8, 7, 6 });defer t.allocator.free(f1);const f2 = try maskedFrame(t.allocator, false, .continuation, "lo ", .{ 5, 4, 3, 2 });defer t.allocator.free(f2);const f3 = try maskedFrame(t.allocator, true, .continuation, "WS", .{ 1, 1, 1, 1 });defer t.allocator.free(f3);try dec.feed(f1);try dec.feed(f2);try t.expectEqual(@as(?Event, null), try dec.next());try dec.feed(f3);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqualStrings("Hello WS", ev.text);}test "control frame interleaved with fragments" {var dec = Decoder.init(t.allocator);defer dec.deinit();const f1 = try maskedFrame(t.allocator, false, .text, "par", .{ 2, 2, 2, 2 });defer t.allocator.free(f1);const ping = try maskedFrame(t.allocator, true, .ping, "hb", .{ 3, 3, 3, 3 });defer t.allocator.free(ping);const f2 = try maskedFrame(t.allocator, true, .continuation, "tial", .{ 4, 4, 4, 4 });defer t.allocator.free(f2);try dec.feed(f1);try dec.feed(ping);var ev1 = (try dec.next()).?;defer ev1.deinitPayload(t.allocator);try t.expectEqualStrings("hb", ev1.ping);try dec.feed(f2);var ev2 = (try dec.next()).?;defer ev2.deinitPayload(t.allocator);try t.expectEqualStrings("partial", ev2.text);}test "unmasked client frame is a protocol error" {var dec = Decoder.init(t.allocator);defer dec.deinit();var hdr: [10]u8 = undefined;const hlen = encodeFrameHeader(&hdr, .text, 2, true);try dec.feed(hdr[0..hlen]);try dec.feed("hi");const ev = (try dec.next()).?;try t.expectEqual(CLOSE_PROTOCOL_ERROR, ev.protocol_error);try t.expectEqual(@as(?Event, null), try dec.next()); // stream dead}test "invalid UTF-8 in text is 1007" {var dec = Decoder.init(t.allocator);defer dec.deinit();const bad = [_]u8{ 0xC3, 0x28 }; // invalid 2-byte sequenceconst frame = try maskedFrame(t.allocator, true, .text, &bad, .{ 7, 7, 7, 7 });defer t.allocator.free(frame);try dec.feed(frame);const ev = (try dec.next()).?;try t.expectEqual(CLOSE_INVALID_PAYLOAD, ev.protocol_error);}test "close with code and reason" {var dec = Decoder.init(t.allocator);defer dec.deinit();var payload: [7]u8 = undefined;std.mem.writeInt(u16, payload[0..2], 1000, .big);@memcpy(payload[2..], "done!");const frame = try maskedFrame(t.allocator, true, .close, &payload, .{ 6, 6, 6, 6 });defer t.allocator.free(frame);try dec.feed(frame);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqual(@as(u16, 1000), ev.close.code);try t.expectEqualStrings("done!", ev.close.reason);}test "close with invalid code is 1002" {var dec = Decoder.init(t.allocator);defer dec.deinit();var payload: [2]u8 = undefined;std.mem.writeInt(u16, payload[0..2], 1006, .big); // 1006 must never appear on the wireconst frame = try maskedFrame(t.allocator, true, .close, &payload, .{ 1, 2, 3, 4 });defer t.allocator.free(frame);try dec.feed(frame);const ev = (try dec.next()).?;try t.expectEqual(CLOSE_PROTOCOL_ERROR, ev.protocol_error);}test "fragmented control frame is 1002" {var dec = Decoder.init(t.allocator);defer dec.deinit();const frame = try maskedFrame(t.allocator, false, .ping, "x", .{ 1, 2, 3, 4 });defer t.allocator.free(frame);try dec.feed(frame);const ev = (try dec.next()).?;try t.expectEqual(CLOSE_PROTOCOL_ERROR, ev.protocol_error);}test "continuation without message is 1002" {var dec = Decoder.init(t.allocator);defer dec.deinit();const frame = try maskedFrame(t.allocator, true, .continuation, "x", .{ 1, 2, 3, 4 });defer t.allocator.free(frame);try dec.feed(frame);const ev = (try dec.next()).?;try t.expectEqual(CLOSE_PROTOCOL_ERROR, ev.protocol_error);}test "new data frame while assembling is 1002" {var dec = Decoder.init(t.allocator);defer dec.deinit();const f1 = try maskedFrame(t.allocator, false, .text, "a", .{ 1, 2, 3, 4 });defer t.allocator.free(f1);const f2 = try maskedFrame(t.allocator, true, .text, "b", .{ 1, 2, 3, 4 });defer t.allocator.free(f2);try dec.feed(f1);try dec.feed(f2);const ev = (try dec.next()).?;try t.expectEqual(CLOSE_PROTOCOL_ERROR, ev.protocol_error);}test "oversize message is 1009" {var dec = Decoder.init(t.allocator);defer dec.deinit();dec.max_message_size = 8;const frame = try maskedFrame(t.allocator, true, .text, "123456789", .{ 1, 2, 3, 4 });defer t.allocator.free(frame);try dec.feed(frame);const ev = (try dec.next()).?;try t.expectEqual(CLOSE_TOO_BIG, ev.protocol_error);}test "16-bit extended length round trip" {var dec = Decoder.init(t.allocator);defer dec.deinit();const big = try t.allocator.alloc(u8, 300);defer t.allocator.free(big);@memset(big, 'A');const frame = try maskedFrame(t.allocator, true, .binary, big, .{ 9, 9, 9, 9 });defer t.allocator.free(frame);try dec.feed(frame);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqual(@as(usize, 300), ev.binary.len);try t.expectEqual(@as(u8, 'A'), ev.binary[150]);}test "writeFrame emits parseable frame" {const Sink = struct {var buf: std.ArrayListUnmanaged(u8) = .empty;fn write(_: ?*anyopaque, data: []const u8) bool {buf.appendSlice(t.allocator, data) catch return false;return true;}};defer Sink.buf.deinit(t.allocator);try t.expect(writeFrame(&Sink.write, null, .text, "server says hi"));// Parse it back with an unmasked-tolerant decoder (client role)var dec = Decoder.init(t.allocator);dec.require_masked = false;defer dec.deinit();try dec.feed(Sink.buf.items);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqualStrings("server says hi", ev.text);}test "writeClose encodes code and reason" {const Sink = struct {var buf: std.ArrayListUnmanaged(u8) = .empty;fn write(_: ?*anyopaque, data: []const u8) bool {buf.appendSlice(t.allocator, data) catch return false;return true;}};defer Sink.buf.deinit(t.allocator);try t.expect(writeClose(&Sink.write, null, 1000, "bye"));var dec = Decoder.init(t.allocator);dec.require_masked = false;defer dec.deinit();try dec.feed(Sink.buf.items);var ev = (try dec.next()).?;defer ev.deinitPayload(t.allocator);try t.expectEqual(@as(u16, 1000), ev.close.code);try t.expectEqualStrings("bye", ev.close.reason);}
Branches
- mainmain branch
Latest commits
- 4137be0fantcolony#40: mission references point to the moved missionsmre
- 9bfba36aantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
- c7bd2645tickets: Hybriel master 73267707 (#122 fixed); compactNow workaround removed (#110 covered)mre
- 2ab91ee9tickets: gate checks rows appear once (session sync); re-vendor to ff51cf46 stopped on hybriel#122, stays 837fe120mre
- e01c2b1dtickets#24: installable app (manifest, service worker, offline list), own icon; gate waits for the hello's pongmre
- 752fbb7fdeploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
- 38bdd5e4deploy.sh: never send .git or .gitignore to Byrodinmre
- f12fa1bcState of 2026-09-27, before the move to gitoriamre