A websocket implementation for zig
0

Configure Feed

Select the types of activity you want to include in your feed.

zig fmt + enable client compression autobahn tests

Karl Seguin (Jul 28, 2025, 11:05 PM +0800) ab9afcf4 d6179ab5

+39 -42
+3 -3
src/client/client.zig
··· 173 173 try sendHandshake(path, key, buf, &opts, self._compression_opts, stream); 174 174 175 175 const res = try HandShakeReply.read(buf, key, &opts, self._compression_opts, stream); 176 - errdefer self.close(.{.code = 1001}) catch unreachable; 176 + errdefer self.close(.{ .code = 1001 }) catch unreachable; 177 177 178 178 // Set up compression with agreed-on parameters 179 179 try self.setupCompression(res.compression); ··· 378 378 var fbs = std.io.fixedBufferStream(data); 379 379 _ = try compressor.compress(fbs.reader()); 380 380 try compressor.flush(); 381 - payload = writer.items[0..writer.items.len - 4]; 381 + payload = writer.items[0 .. writer.items.len - 4]; 382 382 383 383 if (c.reset) { 384 384 c.compressor = try Compression.Type.init(writer.writer(), .{}); ··· 736 736 737 737 if (client_max_bits != 15) { 738 738 // We don't offer client window, so if the server asks for one, that's an error 739 - return error.InvalidExtensionHeader; 739 + return error.InvalidExtensionHeader; 740 740 } 741 741 742 742 return .{
+3 -3
src/proto.zig
··· 317 317 if (is_continuation) { 318 318 if (self.fragment) |*f| { 319 319 if (f.compressed) { 320 - return . {more, .{.data = try self.decompress(try f.last(payload)), .type = f.type}}; 320 + return .{ more, .{ .data = try self.decompress(try f.last(payload)), .type = f.type } }; 321 321 } 322 322 return .{ more, .{ .type = f.type, .data = try f.last(payload) } }; 323 323 } ··· 330 330 } 331 331 332 332 if (compressed) { 333 - return . {more, .{.data = try self.decompress(payload), .type = message_type}}; 333 + return .{ more, .{ .data = try self.decompress(payload), .type = message_type } }; 334 334 } 335 335 336 336 // just a normal single-fragment message (most common case) ··· 428 428 .buf = try provider.pool.acquireOrCreate(), 429 429 }; 430 430 } else { 431 - writer = .{ 431 + writer = .{ 432 432 .pooled = false, 433 433 .provider = provider, 434 434 .buf = try provider.allocator.alloc(u8, @intFromFloat(@as(f64, @floatFromInt(compressed.len)) * 1.25)),
+10 -17
src/server/handshake.zig
··· 69 69 headers.add(name, value); 70 70 switch (std.meta.stringToEnum(SpecialHeader, name) orelse .none) { 71 71 .upgrade => { 72 - if (!ascii.eqlIgnoreCase("websocket", value)) { 72 + if (!ascii.eqlIgnoreCase("websocket", value)) { 73 73 return error.InvalidUpgrade; 74 74 } 75 75 required_headers |= 1; ··· 127 127 "HTTP/1.1 101 Switching Protocols\r\n" ++ 128 128 "Upgrade: websocket\r\n" ++ 129 129 "Connection: upgrade\r\n" ++ 130 - "Sec-Websocket-Accept: " 131 - ; 130 + "Sec-Websocket-Accept: "; 132 131 133 132 @memcpy(buf[0..HEADER.len], HEADER); 134 133 var pos = HEADER.len; ··· 168 167 } 169 168 170 169 for (headers.keys[0..headers.len], headers.values[0..headers.len]) |k, v| { 171 - pos += (try std.fmt.bufPrint(buf[pos..], "\r\n{s}: {s}", .{k, v})).len; 170 + pos += (try std.fmt.bufPrint(buf[pos..], "\r\n{s}: {s}", .{ k, v })).len; 172 171 } 173 172 174 173 const end = pos + 4; ··· 373 372 max_res_headers: usize, 374 373 states: []*Handshake.State, 375 374 376 - pub fn init(allocator: Allocator, count: usize, buffer_size: usize, max_req_headers: usize, max_res_headers: usize) !*Pool { 375 + pub fn init(allocator: Allocator, count: usize, buffer_size: usize, max_req_headers: usize, max_res_headers: usize) !*Pool { 377 376 const states = try allocator.alloc(*Handshake.State, count); 378 377 errdefer allocator.free(states); 379 378 ··· 535 534 "HTTP/1.1 101 Switching Protocols\r\n" ++ 536 535 "Upgrade: websocket\r\n" ++ 537 536 "Connection: upgrade\r\n" ++ 538 - "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n\r\n" 539 - ; 537 + "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n\r\n"; 540 538 try t.expectString(expected, try Handshake.createReply("this is my key", &res_headers, null, &buf)); 541 539 } 542 540 ··· 547 545 "Upgrade: websocket\r\n" ++ 548 546 "Connection: upgrade\r\n" ++ 549 547 "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n" ++ 550 - "Sec-WebSocket-Extensions: permessage-deflate\r\n\r\n" 551 - ; 548 + "Sec-WebSocket-Extensions: permessage-deflate\r\n\r\n"; 552 549 try t.expectString(expected, try Handshake.createReply("this is my key", &res_headers, .{}, &buf)); 553 550 } 554 551 ··· 559 556 "Upgrade: websocket\r\n" ++ 560 557 "Connection: upgrade\r\n" ++ 561 558 "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n" ++ 562 - "Sec-WebSocket-Extensions: permessage-deflate; server_no_context_takeover; client_no_context_takeover\r\n\r\n" 563 - ; 559 + "Sec-WebSocket-Extensions: permessage-deflate; server_no_context_takeover; client_no_context_takeover\r\n\r\n"; 564 560 try t.expectString(expected, try Handshake.createReply("this is my key", &res_headers, .{ 565 561 .client_no_context_takeover = true, 566 562 .server_no_context_takeover = true, ··· 576 572 "Upgrade: websocket\r\n" ++ 577 573 "Connection: upgrade\r\n" ++ 578 574 "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n" ++ 579 - "Set-Cookie: Yummy!\r\n\r\n" 580 - ; 575 + "Set-Cookie: Yummy!\r\n\r\n"; 581 576 try t.expectString(expected, try Handshake.createReply("this is my key", &res_headers, null, &buf)); 582 577 } 583 578 ··· 589 584 "Connection: upgrade\r\n" ++ 590 585 "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n" ++ 591 586 "Sec-WebSocket-Extensions: permessage-deflate\r\n" ++ 592 - "Set-Cookie: Yummy!\r\n\r\n" 593 - ; 587 + "Set-Cookie: Yummy!\r\n\r\n"; 594 588 try t.expectString(expected, try Handshake.createReply("this is my key", &res_headers, .{}, &buf)); 595 589 } 596 590 ··· 602 596 "Connection: upgrade\r\n" ++ 603 597 "Sec-Websocket-Accept: flzHu2DevQ2dSCSVqKSii5e9C2o=\r\n" ++ 604 598 "Sec-WebSocket-Extensions: permessage-deflate; server_no_context_takeover; client_no_context_takeover\r\n" ++ 605 - "Set-Cookie: Yummy!\r\n\r\n" 606 - ; 599 + "Set-Cookie: Yummy!\r\n\r\n"; 607 600 try t.expectString(expected, try Handshake.createReply("this is my key", &res_headers, .{ 608 601 .client_no_context_takeover = true, 609 602 .server_no_context_takeover = true,
+14 -15
src/server/server.zig
··· 382 382 } 383 383 if (hc.handler != null) { 384 384 // if we have a handler, the our handshake completed 385 - try conn_manager.setupCompression(hc, compression); 385 + try conn_manager.setupCompression(hc, compression); 386 386 break; 387 387 } 388 388 if (timestamp() > deadline) { ··· 681 681 conn_manager.inactive(hc); 682 682 } 683 683 684 - try conn_manager.setupCompression(hc, compression); 684 + try conn_manager.setupCompression(hc, compression); 685 685 return true; 686 686 } 687 687 }; ··· 1382 1382 var fbs = std.io.fixedBufferStream(data); 1383 1383 _ = try compressor.compress(fbs.reader()); 1384 1384 try compressor.flush(); 1385 - payload = writer.items[0..writer.items.len - 4]; 1385 + payload = writer.items[0 .. writer.items.len - 4]; 1386 1386 1387 1387 if (c.reset) { 1388 1388 c.compressor = try Conn.Compression.Type.init(writer.writer(), .{}); ··· 1485 1485 }; 1486 1486 }; 1487 1487 1488 - fn handleHandshake(comptime H: type, worker: anytype, hc: *HandlerConn(H), ctx: anytype) struct{?Compression, bool} { 1488 + fn handleHandshake(comptime H: type, worker: anytype, hc: *HandlerConn(H), ctx: anytype) struct { ?Compression, bool } { 1489 1489 return _handleHandshake(H, worker, hc, ctx) catch |err| { 1490 1490 log.warn("({}) uncaugh error processing handshake: {}", .{ hc.conn.address, err }); 1491 - return .{null, false}; 1491 + return .{ null, false }; 1492 1492 }; 1493 1493 } 1494 1494 1495 - fn _handleHandshake(comptime H: type, worker: anytype, hc: *HandlerConn(H), ctx: anytype) !struct{?Compression, bool} { 1495 + fn _handleHandshake(comptime H: type, worker: anytype, hc: *HandlerConn(H), ctx: anytype) !struct { ?Compression, bool } { 1496 1496 std.debug.assert(hc.handler == null); 1497 1497 1498 1498 var state = hc.handshake orelse blk: { ··· 1507 1507 1508 1508 if (len == buf.len) { 1509 1509 log.warn("({}) handshake request exceeded maximum configured size ({d})", .{ conn.address, buf.len }); 1510 - return .{null, false}; 1510 + return .{ null, false }; 1511 1511 } 1512 1512 1513 1513 const n = posix.read(hc.socket, buf[len..]) catch |err| { ··· 1519 1519 }, 1520 1520 else => log.warn("({}) handshake error reading from socket: {}", .{ conn.address, err }), 1521 1521 } 1522 - return .{null, false}; 1522 + return .{ null, false }; 1523 1523 }; 1524 1524 1525 1525 if (n == 0) { 1526 1526 log.debug("({}) handshake connection closed", .{conn.address}); 1527 - return .{null, false}; 1527 + return .{ null, false }; 1528 1528 } 1529 1529 1530 1530 state.len = len + n; 1531 1531 var handshake = Handshake.parse(state) catch |err| { 1532 1532 log.debug("({}) error parsing handshake: {}", .{ conn.address, err }); 1533 1533 respondToHandshakeError(conn, err); 1534 - return .{null, false}; 1534 + return .{ null, false }; 1535 1535 } orelse { 1536 1536 // we need more data 1537 - return .{null, true}; 1537 + return .{ null, true }; 1538 1538 }; 1539 1539 1540 1540 var agreed_compression: ?Compression = null; ··· 1550 1550 defer state.release(); 1551 1551 hc.handshake = null; 1552 1552 1553 - 1554 1553 // After this, the app has access to &hc.conn, so any access to the 1555 1554 // conn has to be synchronized (which the conn does internally). 1556 1555 ··· 1561 1560 respondToHandshakeError(conn, err); 1562 1561 } 1563 1562 log.debug("({}) " ++ @typeName(H) ++ ".init rejected request {}", .{ conn.address, err }); 1564 - return .{null, false}; 1563 + return .{ null, false }; 1565 1564 }; 1566 1565 1567 1566 hc.handler = handler; ··· 1575 1574 const res = if (params.len == 1) hc.handler.?.afterInit() else hc.handler.?.afterInit(ctx); 1576 1575 res catch |err| { 1577 1576 log.debug("({}) " ++ @typeName(H) ++ ".afterInit error: {}", .{ conn.address, err }); 1578 - return .{null, false}; 1577 + return .{ null, false }; 1579 1578 }; 1580 1579 } 1581 1580 1582 1581 log.debug("({}) connection successfully upgraded", .{conn.address}); 1583 - return .{agreed_compression, true}; 1582 + return .{ agreed_compression, true }; 1584 1583 } 1585 1584 1586 1585 fn handleClientData(comptime H: type, hc: *HandlerConn(H), allocator: Allocator, fba: *FixedBufferAllocator) bool {
+1 -1
src/testing.zig
··· 27 27 errdefer arena.deinit(); 28 28 29 29 const port = opts.port orelse 0; 30 - const pair = t.SocketPair.init(.{.port = port}); 30 + const pair = t.SocketPair.init(.{ .port = port }); 31 31 const timeout = std.mem.toBytes(std.posix.timeval{ 32 32 .sec = 0, 33 33 .usec = 50_000,
+1 -1
src/websocket.zig
··· 46 46 test "frameText" { 47 47 { 48 48 const framed = frameText(""); 49 - try t.expectString(&[_]u8{ 129, 0}, &framed); 49 + try t.expectString(&[_]u8{ 129, 0 }, &framed); 50 50 } 51 51 52 52 {
+6 -1
support/autobahn/client/main.zig
··· 29 29 "9.2.5", "9.2.6", "9.3.1", "9.3.2", "9.3.3", "9.3.4", "9.3.5", "9.3.6", "9.3.7", "9.3.8", "9.3.9", "9.4.1", "9.4.2", "9.4.3", "9.4.4", "9.4.5", 30 30 "9.4.6", "9.4.7", "9.4.8", "9.4.9", "9.5.1", "9.5.2", "9.5.3", "9.5.4", "9.5.5", "9.5.6", "9.6.1", "9.6.2", "9.6.3", "9.6.4", "9.6.5", "9.6.6", 31 31 "9.7.1", "9.7.2", "9.7.3", "9.7.4", "9.7.5", "9.7.6", "9.8.1", "9.8.2", "9.8.3", "9.8.4", "9.8.5", "9.8.6", "10.1.1", 32 + "12.1.1", "12.1.2", "12.1.3", "12.1.4", "12.1.5", "12.1.6", "12.1.7", "12.1.8", "12.1.9", "12.1.10", "12.1.11", "12.1.12", "12.1.13", "12.1.14", "12.1.15", "12.1.16", "12.1.17", "12.1.18", "12.2.1", "12.2.2", "12.2.3", "12.2.4", "12.2.5", "12.2.6", "12.2.7", "12.2.8", "12.2.9", "12.2.10", "12.2.11", "12.2.12", "12.2.13", "12.2.14", "12.2.15", "12.2.16", "12.2.17", "12.2.18", "12.3.1", "12.3.2", "12.3.3", "12.3.4", "12.3.5", "12.3.6", "12.3.7", "12.3.8", "12.3.9", "12.3.10", "12.3.11", "12.3.12", "12.3.13", "12.3.14", "12.3.15", "12.3.16", "12.3.17", "12.3.18", "12.4.1", "12.4.2", "12.4.3", "12.4.4", "12.4.5", "12.4.6", "12.4.7", "12.4.8", "12.4.9", "12.4.10", "12.4.11", "12.4.12", "12.4.13", "12.4.14", "12.4.15", "12.4.16", "12.4.17", "12.4.18", "12.5.1", "12.5.2", "12.5.3", "12.5.4", "12.5.5", "12.5.6", "12.5.7", "12.5.8", "12.5.9", "12.5.10", "12.5.11", "12.5.12", "12.5.13", "12.5.14", "12.5.15", "12.5.16", "12.5.17", "12.5.18", 33 + "13.1.1", "13.1.2", "13.1.3", "13.1.4", "13.1.5", "13.1.6", "13.1.7", "13.1.8", "13.1.9", "13.1.10", "13.1.11", "13.1.12", "13.1.13", "13.1.14", "13.1.15", "13.1.16", "13.1.17", "13.1.18", "13.2.1", "13.2.2", "13.2.3", "13.2.4", "13.2.5", "13.2.6", "13.2.7", "13.2.8", "13.2.9", "13.2.10", "13.2.11", "13.2.12", "13.2.13", "13.2.14", "13.2.15", "13.2.16", "13.2.17", "13.2.18", "13.3.1", "13.3.2", "13.3.3", "13.3.4", "13.3.5", "13.3.6", "13.3.7", "13.3.8", "13.3.9", "13.3.10", "13.3.11", "13.3.12", "13.3.13", "13.3.14", "13.3.15", "13.3.16", "13.3.17", "13.3.18", "13.4.1", "13.4.2", "13.4.3", "13.4.4", "13.4.5", "13.4.6", "13.4.7", "13.4.8", "13.4.9", "13.4.10", "13.4.11", "13.4.12", "13.4.13", "13.4.14", "13.4.15", "13.4.16", "13.4.17", "13.4.18", "13.5.1", "13.5.2", "13.5.3", "13.5.4", "13.5.5", "13.5.6", "13.5.7", "13.5.8", "13.5.9", "13.5.10", "13.5.11", "13.5.12", "13.5.13", "13.5.14", "13.5.15", "13.5.16", "13.5.17", "13.5.18", "13.6.1", "13.6.2", "13.6.3", "13.6.4", "13.6.5", "13.6.6", "13.6.7", "13.6.8", "13.6.9", "13.6.10", "13.6.11", "13.6.12", "13.6.13", "13.6.14", "13.6.15", "13.6.16", "13.6.17", "13.6.18", "13.7.1", "13.7.2", "13.7.3", "13.7.4", "13.7.5", "13.7.6", "13.7.7", "13.7.8", "13.7.9", "13.7.10", "13.7.11", "13.7.12", "13.7.13", "13.7.14", "13.7.15", "13.7.16", "13.7.17", "13.7.18", 32 34 }; 33 35 34 36 // wait 5 seconds for autobanh server to be up ··· 38 40 defer buffer_provider.deinit(); 39 41 40 42 for (cases, 0..) |case, i| { 41 - // if (!std.mem.eql(u8, case, "6.3.1")) continue; 43 + // if (!std.mem.eql(u8, case, "12.1.1")) continue; 42 44 std.debug.print("running case: {s}\n", .{case}); 43 45 44 46 { ··· 90 92 .port = 9001, 91 93 .host = "localhost", 92 94 .buffer_provider = buffer_provider, 95 + .compression = .{ 96 + .write_threshold = 0, 97 + }, 93 98 }); 94 99 errdefer client.deinit(); 95 100 try client.handshake(path, .{
+1 -1
support/autobahn/client/run.sh
··· 15 15 crossbario/autobahn-testsuite \ 16 16 /opt/pypy/bin/wstest --mode fuzzingserver --spec /ab/config.json & 17 17 18 - cd support/autobahn/client/ && zig build run 18 + cd support/autobahn/client/ && zig build -Doptimize=ReleaseFast run 19 19 20 20 if grep FAILED support/autobahn/client/reports/index.json*; then 21 21 exit 1