tiny.http.ClientWebSocket
Defined in tiny.http.
API (13)
Actions
Public operations.
Fields and members
Public fields and members.
allocatorconnectionframe_payload_bytesopenpending_consumedread_lengthread_storagewrite_storage
Source
Source: lib/http/src/client/runtime.zig:710
zig
pub const ClientWebSocket = struct { allocator: Allocator, connection: PersistentConnection, read_storage: []u8, write_storage: []u8, frame_payload_bytes: usize, read_length: usize = 0, pending_consumed: usize = 0, open: bool = true, fn init( allocator: Allocator, connection: PersistentConnection, scratch: WebSocketScratch, ) !ClientWebSocket { const expected = try alloc_phase.capacity.add( usize, client_websocket.frame_header_bytes, scratch.frame_payload_bytes, ); if (scratch.read.len != expected or scratch.write.len != expected) { return error.InvalidWebSocketScratch; } return .{ .allocator = allocator, .connection = connection, .read_storage = scratch.read, .write_storage = scratch.write, .frame_payload_bytes = scratch.frame_payload_bytes, }; } pub fn deinit(self: *ClientWebSocket) void { destroyConnection( self.allocator, self.connection, true, ); std.crypto.secureZero(u8, self.read_storage); std.crypto.secureZero(u8, self.write_storage); self.* = undefined; } pub fn pollReadable( self: *ClientWebSocket, timeout_ms: i32, ) !bool { if (timeout_ms < 0) return error.InvalidTimeout; if (self.read_length > self.pending_consumed or connectionBuffered(self.connection)) { return true; } return try sys.pollReadable( connectionSocket(self.connection), timeout_ms, ); } pub fn sendBinary( self: *ClientWebSocket, payload: []const u8, ) !void { try self.sendFrame(.binary, payload); } pub fn receive( self: *ClientWebSocket, ) !?WebSocketMessage { self.releaseBorrowed(); while (self.open) { if (self.read_length == 0) { try self.readFrame(); } const parsed = http.Frame.parse( self.wire()[0..self.read_length], self.frame_payload_bytes, ) catch |err| return err; if (parsed.frame.mask != null) { return error.MaskedServerFrame; } if (!parsed.frame.fin or parsed.frame.opcode == .continuation) { return error.FragmentedServerMessage; } switch (parsed.frame.opcode) { .binary, .text => { self.pending_consumed = parsed.consumed; return .{ .opcode = parsed.frame.opcode, .payload = parsed.frame.payload, }; }, .ping => { try self.sendFrame( .pong, parsed.frame.payload, ); self.consume(parsed.consumed); }, .pong => self.consume(parsed.consumed), .close => { self.sendFrame( .close, parsed.frame.payload, ) catch {}; self.consume(parsed.consumed); self.open = false; return null; }, .continuation => unreachable, } } return null; } pub fn close(self: *ClientWebSocket) !void { if (!self.open) return; try self.sendFrame(.close, &.{}); self.open = false; } fn sendFrame( self: *ClientWebSocket, opcode: http.Opcode, payload: []const u8, ) !void { if (!self.open and opcode != .close) { return error.ConnectionClosed; } var mask: [4]u8 = undefined; try sys_root.random.secureBytes(&mask); const encoded = try (http.Frame{ .fin = true, .opcode = opcode, .mask = mask, .payload = payload, }).serializeInto( self.output(), self.frame_payload_bytes, ); defer std.crypto.secureZero(u8, encoded); const writer = connectionWriter( self.connection, ); try writer.writeAll(encoded); try flushConnection(self.connection); } fn readFrame(self: *ClientWebSocket) !void { const reader = connectionReader(self.connection); const target = self.wire(); try readClientWebSocketExact( reader, target[0..2], ); if ((target[1] & 0x80) != 0) { return error.MaskedServerFrame; } const short_length = target[1] & 0x7f; var header_bytes: usize = 2; const payload_bytes: usize = switch (short_length) { 0...125 => short_length, 126 => length: { try readClientWebSocketExact( reader, target[2..4], ); header_bytes = 4; break :length std.mem.readInt( u16, target[2..4], .big, ); }, 127 => length: { try readClientWebSocketExact( reader, target[2..10], ); header_bytes = 10; const wide = std.mem.readInt( u64, target[2..10], .big, ); if (wide > self.frame_payload_bytes or wide > std.math.maxInt(usize)) { return error.PayloadTooLarge; } break :length @intCast(wide); }, else => unreachable, }; if (payload_bytes > self.frame_payload_bytes) { return error.PayloadTooLarge; } const total = std.math.add( usize, header_bytes, payload_bytes, ) catch return error.PayloadTooLarge; if (total > target.len) { return error.FrameCapacity; } try readClientWebSocketExact( reader, target[header_bytes..total], ); self.read_length = total; } fn releaseBorrowed(self: *ClientWebSocket) void { if (self.pending_consumed == 0) return; self.consume(self.pending_consumed); self.pending_consumed = 0; } fn consume( self: *ClientWebSocket, consumed: usize, ) void { const remaining = self.read_length - consumed; const previous_length = self.read_length; std.mem.copyForwards( u8, self.wire()[0..remaining], self.wire()[consumed..self.read_length], ); std.crypto.secureZero( u8, self.wire()[remaining..previous_length], ); self.read_length = remaining; } fn wire(self: *ClientWebSocket) []u8 { return self.read_storage; } fn output(self: *ClientWebSocket) []u8 { return self.write_storage; }};Source: lib/http/src/root.zig:79
zig
pub const ClientWebSocket = client.ClientWebSocket;Audit
| Definitions | 6 |
|---|---|
| Public names | 6 |
| Members | 8 |
| Version | 26.7.0 |
| Revision | daab053ee433 |