| author | |
| committer | |
| log | e0e463bcf711f660cf8df5e8d31e82f855c0a967 |
| tree | c22b756f5e2c2081b8c5c3f484a8079d0dd36af1 |
| parent | 4d62f0839382a9aa611f9c898a1ef7a6050bd92a |
3 files changed, 27 insertions(+), 18 deletions(-)
lib/std/Io.zig+3-2| ... | @@ -680,8 +680,9 @@ pub const VTable = struct { | ... | @@ -680,8 +680,9 @@ pub const VTable = struct { |
| 680 | netConnectUnix: *const fn (?*anyopaque, net.UnixAddress) net.UnixAddress.ConnectError!net.Socket.Handle, | 680 | netConnectUnix: *const fn (?*anyopaque, net.UnixAddress) net.UnixAddress.ConnectError!net.Socket.Handle, |
| 681 | netSend: *const fn (?*anyopaque, net.Socket.Handle, []net.OutgoingMessage, net.SendFlags) struct { ?net.Socket.SendError, usize }, | 681 | netSend: *const fn (?*anyopaque, net.Socket.Handle, []net.OutgoingMessage, net.SendFlags) struct { ?net.Socket.SendError, usize }, |
| 682 | netReceive: *const fn (?*anyopaque, net.Socket.Handle, message_buffer: []net.IncomingMessage, data_buffer: []u8, net.ReceiveFlags, Timeout) struct { ?net.Socket.ReceiveTimeoutError, usize }, | 682 | netReceive: *const fn (?*anyopaque, net.Socket.Handle, message_buffer: []net.IncomingMessage, data_buffer: []u8, net.ReceiveFlags, Timeout) struct { ?net.Socket.ReceiveTimeoutError, usize }, |
| 683 | netRead: *const fn (?*anyopaque, src: net.Stream, data: [][]u8) net.Stream.Reader.Error!usize, | 683 | /// Returns 0 on end of stream. |
| 684 | netWrite: *const fn (?*anyopaque, dest: net.Stream, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize, | 684 | netRead: *const fn (?*anyopaque, src: net.Socket.Handle, data: [][]u8) net.Stream.Reader.Error!usize, |
| 685 | netWrite: *const fn (?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize, | ||
| 685 | netClose: *const fn (?*anyopaque, handle: net.Socket.Handle) void, | 686 | netClose: *const fn (?*anyopaque, handle: net.Socket.Handle) void, |
| 686 | netInterfaceNameResolve: *const fn (?*anyopaque, *const net.Interface.Name) net.Interface.Name.ResolveError!net.Interface, | 687 | netInterfaceNameResolve: *const fn (?*anyopaque, *const net.Interface.Name) net.Interface.Name.ResolveError!net.Interface, |
| 687 | netInterfaceName: *const fn (?*anyopaque, net.Interface) net.Interface.NameError!net.Interface.Name, | 688 | netInterfaceName: *const fn (?*anyopaque, net.Interface) net.Interface.NameError!net.Interface.Name, |
lib/std/Io/Threaded.zig+7-13| ... | @@ -1986,9 +1986,8 @@ fn netAcceptPosix(userdata: ?*anyopaque, listen_fd: Io.net.Socket.Handle) Io.net | ... | @@ -1986,9 +1986,8 @@ fn netAcceptPosix(userdata: ?*anyopaque, listen_fd: Io.net.Socket.Handle) Io.net |
| 1986 | } }; | 1986 | } }; |
| 1987 | } | 1987 | } |
| 1988 | 1988 | ||
| 1989 | fn netReadPosix(userdata: ?*anyopaque, stream: Io.net.Stream, data: [][]u8) Io.net.Stream.Reader.Error!usize { | 1989 | fn netReadPosix(userdata: ?*anyopaque, fd: Io.net.Socket.Handle, data: [][]u8) Io.net.Stream.Reader.Error!usize { |
| 1990 | const pool: *Pool = @ptrCast(@alignCast(userdata)); | 1990 | const pool: *Pool = @ptrCast(@alignCast(userdata)); |
| 1991 | const fd = stream.socket.handle; | ||
| 1992 | 1991 | ||
| 1993 | var iovecs_buffer: [max_iovecs_len]posix.iovec = undefined; | 1992 | var iovecs_buffer: [max_iovecs_len]posix.iovec = undefined; |
| 1994 | var i: usize = 0; | 1993 | var i: usize = 0; |
| ... | @@ -2006,11 +2005,9 @@ fn netReadPosix(userdata: ?*anyopaque, stream: Io.net.Stream, data: [][]u8) Io.n | ... | @@ -2006,11 +2005,9 @@ fn netReadPosix(userdata: ?*anyopaque, stream: Io.net.Stream, data: [][]u8) Io.n |
| 2006 | try pool.checkCancel(); | 2005 | try pool.checkCancel(); |
| 2007 | var n: usize = undefined; | 2006 | var n: usize = undefined; |
| 2008 | switch (std.os.wasi.fd_read(fd, dest.ptr, dest.len, &n)) { | 2007 | switch (std.os.wasi.fd_read(fd, dest.ptr, dest.len, &n)) { |
| 2009 | .SUCCESS => { | 2008 | .SUCCESS => return n, |
| 2010 | if (n == 0) return error.EndOfStream; | ||
| 2011 | return n; | ||
| 2012 | }, | ||
| 2013 | .INTR => continue, | 2009 | .INTR => continue, |
| 2010 | |||
| 2014 | .INVAL => |err| return errnoBug(err), | 2011 | .INVAL => |err| return errnoBug(err), |
| 2015 | .FAULT => |err| return errnoBug(err), | 2012 | .FAULT => |err| return errnoBug(err), |
| 2016 | .AGAIN => |err| return errnoBug(err), | 2013 | .AGAIN => |err| return errnoBug(err), |
| ... | @@ -2029,12 +2026,9 @@ fn netReadPosix(userdata: ?*anyopaque, stream: Io.net.Stream, data: [][]u8) Io.n | ... | @@ -2029,12 +2026,9 @@ fn netReadPosix(userdata: ?*anyopaque, stream: Io.net.Stream, data: [][]u8) Io.n |
| 2029 | try pool.checkCancel(); | 2026 | try pool.checkCancel(); |
| 2030 | const rc = posix.system.readv(fd, dest.ptr, @intCast(dest.len)); | 2027 | const rc = posix.system.readv(fd, dest.ptr, @intCast(dest.len)); |
| 2031 | switch (posix.errno(rc)) { | 2028 | switch (posix.errno(rc)) { |
| 2032 | .SUCCESS => { | 2029 | .SUCCESS => return @intCast(rc), |
| 2033 | const n: usize = @intCast(rc); | ||
| 2034 | if (n == 0) return error.EndOfStream; | ||
| 2035 | return n; | ||
| 2036 | }, | ||
| 2037 | .INTR => continue, | 2030 | .INTR => continue, |
| 2031 | |||
| 2038 | .INVAL => |err| return errnoBug(err), | 2032 | .INVAL => |err| return errnoBug(err), |
| 2039 | .FAULT => |err| return errnoBug(err), | 2033 | .FAULT => |err| return errnoBug(err), |
| 2040 | .AGAIN => |err| return errnoBug(err), | 2034 | .AGAIN => |err| return errnoBug(err), |
| ... | @@ -2359,7 +2353,7 @@ fn netReceive( | ... | @@ -2359,7 +2353,7 @@ fn netReceive( |
| 2359 | 2353 | ||
| 2360 | fn netWritePosix( | 2354 | fn netWritePosix( |
| 2361 | userdata: ?*anyopaque, | 2355 | userdata: ?*anyopaque, |
| 2362 | stream: Io.net.Stream, | 2356 | fd: Io.net.Socket.Handle, |
| 2363 | header: []const u8, | 2357 | header: []const u8, |
| 2364 | data: []const []const u8, | 2358 | data: []const []const u8, |
| 2365 | splat: usize, | 2359 | splat: usize, |
| ... | @@ -2406,7 +2400,7 @@ fn netWritePosix( | ... | @@ -2406,7 +2400,7 @@ fn netWritePosix( |
| 2406 | }, | 2400 | }, |
| 2407 | }; | 2401 | }; |
| 2408 | const flags = posix.MSG.NOSIGNAL; | 2402 | const flags = posix.MSG.NOSIGNAL; |
| 2409 | return posix.sendmsg(stream.socket.handle, &msg, flags); | 2403 | return posix.sendmsg(fd, &msg, flags); |
| 2410 | } | 2404 | } |
| 2411 | 2405 | ||
| 2412 | fn addBuf(v: []posix.iovec_const, i: *@FieldType(posix.msghdr_const, "iovlen"), bytes: []const u8) void { | 2406 | fn addBuf(v: []posix.iovec_const, i: *@FieldType(posix.msghdr_const, "iovlen"), bytes: []const u8) void { |
lib/std/Io/net.zig+17-3| ... | @@ -1076,6 +1076,8 @@ pub const Socket = struct { | ... | @@ -1076,6 +1076,8 @@ pub const Socket = struct { |
| 1076 | pub const Stream = struct { | 1076 | pub const Stream = struct { |
| 1077 | socket: Socket, | 1077 | socket: Socket, |
| 1078 | 1078 | ||
| 1079 | const max_iovecs_len = 8; | ||
| 1080 | |||
| 1079 | pub fn close(s: *Stream, io: Io) void { | 1081 | pub fn close(s: *Stream, io: Io) void { |
| 1080 | io.vtable.netClose(io.userdata, s.socket.handle); | 1082 | io.vtable.netClose(io.userdata, s.socket.handle); |
| 1081 | s.* = undefined; | 1083 | s.* = undefined; |
| ... | @@ -1097,7 +1099,6 @@ pub const Stream = struct { | ... | @@ -1097,7 +1099,6 @@ pub const Stream = struct { |
| 1097 | /// from it. | 1099 | /// from it. |
| 1098 | AccessDenied, | 1100 | AccessDenied, |
| 1099 | NetworkDown, | 1101 | NetworkDown, |
| 1100 | EndOfStream, | ||
| 1101 | } || Io.Cancelable || Io.UnexpectedError; | 1102 | } || Io.Cancelable || Io.UnexpectedError; |
| 1102 | 1103 | ||
| 1103 | pub fn init(stream: Stream, io: Io, buffer: []u8) Reader { | 1104 | pub fn init(stream: Stream, io: Io, buffer: []u8) Reader { |
| ... | @@ -1128,10 +1129,22 @@ pub const Stream = struct { | ... | @@ -1128,10 +1129,22 @@ pub const Stream = struct { |
| 1128 | fn readVec(io_r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize { | 1129 | fn readVec(io_r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize { |
| 1129 | const r: *Reader = @alignCast(@fieldParentPtr("interface", io_r)); | 1130 | const r: *Reader = @alignCast(@fieldParentPtr("interface", io_r)); |
| 1130 | const io = r.io; | 1131 | const io = r.io; |
| 1131 | return io.vtable.netRead(io.userdata, r.stream, data) catch |err| { | 1132 | var iovecs_buffer: [max_iovecs_len][]u8 = undefined; |
| 1133 | const dest_n, const data_size = try io_r.writableVector(&iovecs_buffer, data); | ||
| 1134 | const dest = iovecs_buffer[0..dest_n]; | ||
| 1135 | assert(dest[0].len > 0); | ||
| 1136 | const n = io.vtable.netRead(io.userdata, r.stream.socket.handle, dest) catch |err| { | ||
| 1132 | r.err = err; | 1137 | r.err = err; |
| 1133 | return error.ReadFailed; | 1138 | return error.ReadFailed; |
| 1134 | }; | 1139 | }; |
| 1140 | if (n == 0) { | ||
| 1141 | return error.EndOfStream; | ||
| 1142 | } | ||
| 1143 | if (n > data_size) { | ||
| 1144 | r.interface.end += n - data_size; | ||
| 1145 | return data_size; | ||
| 1146 | } | ||
| 1147 | return n; | ||
| 1135 | } | 1148 | } |
| 1136 | }; | 1149 | }; |
| 1137 | 1150 | ||
| ... | @@ -1166,7 +1179,8 @@ pub const Stream = struct { | ... | @@ -1166,7 +1179,8 @@ pub const Stream = struct { |
| 1166 | const w: *Writer = @alignCast(@fieldParentPtr("interface", io_w)); | 1179 | const w: *Writer = @alignCast(@fieldParentPtr("interface", io_w)); |
| 1167 | const io = w.io; | 1180 | const io = w.io; |
| 1168 | const buffered = io_w.buffered(); | 1181 | const buffered = io_w.buffered(); |
| 1169 | const n = io.vtable.netWrite(io.userdata, w.stream, buffered, data, splat) catch |err| { | 1182 | const handle = w.stream.socket.handle; |
| 1183 | const n = io.vtable.netWrite(io.userdata, handle, buffered, data, splat) catch |err| { | ||
| 1170 | w.err = err; | 1184 | w.err = err; |
| 1171 | return error.WriteFailed; | 1185 | return error.WriteFailed; |
| 1172 | }; | 1186 | }; |