| ... | @@ -777,7 +777,6 @@ pub fn io(ev: *Evented) Io { | ... | @@ -777,7 +777,6 @@ pub fn io(ev: *Evented) Io { |
| 777 | .netConnectUnix = netConnectUnixUnavailable, | 777 | .netConnectUnix = netConnectUnixUnavailable, |
| 778 | .netSocketCreatePair = netSocketCreatePairUnavailable, | 778 | .netSocketCreatePair = netSocketCreatePairUnavailable, |
| 779 | .netSend = netSendUnavailable, | 779 | .netSend = netSendUnavailable, |
| 780 | .netReceive = netReceive, | | |
| 781 | .netRead = netReadUnavailable, | 780 | .netRead = netReadUnavailable, |
| 782 | .netWrite = netWriteUnavailable, | 781 | .netWrite = netWriteUnavailable, |
| 783 | .netWriteFile = netWriteFileUnavailable, | 782 | .netWriteFile = netWriteFileUnavailable, |
| ... | @@ -2092,6 +2091,18 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper | ... | @@ -2092,6 +2091,18 @@ fn operate(userdata: ?*anyopaque, operation: Io.Operation) Io.Cancelable!Io.Oper |
| 2092 | .device_io_control => |o| .{ | 2091 | .device_io_control => |o| .{ |
| 2093 | .device_io_control = try ev.deviceIoControl(try maybe_sync.enterSync(ev), o), | 2092 | .device_io_control = try ev.deviceIoControl(try maybe_sync.enterSync(ev), o), |
| 2094 | }, | 2093 | }, |
| | 2094 | .net_receive => |o| .{ |
| | 2095 | .net_receive = r: { |
| | 2096 | const opt_err, const n = ev.netReceive(&maybe_sync.cancel_region, o.socket_handle, o.message_buffer, o.data_buffer, o.flags); |
| | 2097 | break :r .{ |
| | 2098 | if (opt_err) |err| switch (err) { |
| | 2099 | error.Canceled => |e| return e, |
| | 2100 | else => |e| e, |
| | 2101 | } else null, |
| | 2102 | n, |
| | 2103 | }; |
| | 2104 | }, |
| | 2105 | }, |
| 2095 | }; | 2106 | }; |
| 2096 | } | 2107 | } |
| 2097 | | 2108 | |
| ... | @@ -2375,6 +2386,10 @@ fn batchDrainSubmitted( | ... | @@ -2375,6 +2386,10 @@ fn batchDrainSubmitted( |
| 2375 | return error.ConcurrencyUnavailable | 2386 | return error.ConcurrencyUnavailable |
| 2376 | else | 2387 | else |
| 2377 | .{ .device_io_control = try ev.deviceIoControl(try maybe_sync.enterSync(ev), o) }, | 2388 | .{ .device_io_control = try ev.deviceIoControl(try maybe_sync.enterSync(ev), o) }, |
| | 2389 | .net_receive => |o| { |
| | 2390 | _ = o; |
| | 2391 | @panic("TODO implement batchDrainSubmitted for net_receive"); |
| | 2392 | }, |
| 2378 | })) |result| { | 2393 | })) |result| { |
| 2379 | switch (batch.completed.tail) { | 2394 | switch (batch.completed.tail) { |
| 2380 | .none => batch.completed.head = index, | 2395 | .none => batch.completed.head = index, |
| ... | @@ -2475,6 +2490,7 @@ fn batchDrainReady(batch: *Io.Batch) Io.Timeout.Error!void { | ... | @@ -2475,6 +2490,7 @@ fn batchDrainReady(batch: *Io.Batch) Io.Timeout.Error!void { |
| 2475 | }, | 2490 | }, |
| 2476 | }, | 2491 | }, |
| 2477 | .device_io_control => unreachable, | 2492 | .device_io_control => unreachable, |
| | 2493 | .net_receive => @panic("TODO"), |
| 2478 | })) |result| { | 2494 | })) |result| { |
| 2479 | switch (batch.completed.tail) { | 2495 | switch (batch.completed.tail) { |
| 2480 | .none => batch.completed.head = index, | 2496 | .none => batch.completed.head = index, |
| ... | @@ -5035,37 +5051,16 @@ fn netSendUnavailable( | ... | @@ -5035,37 +5051,16 @@ fn netSendUnavailable( |
| 5035 | } | 5051 | } |
| 5036 | | 5052 | |
| 5037 | fn netReceive( | 5053 | fn netReceive( |
| 5038 | userdata: ?*anyopaque, | 5054 | ev: *Evented, |
| | 5055 | cancel_region: *CancelRegion, |
| 5039 | handle: net.Socket.Handle, | 5056 | handle: net.Socket.Handle, |
| 5040 | message_buffer: []net.IncomingMessage, | 5057 | message_buffer: []net.IncomingMessage, |
| 5041 | data_buffer: []u8, | 5058 | data_buffer: []u8, |
| 5042 | flags: net.ReceiveFlags, | 5059 | flags: net.ReceiveFlags, |
| 5043 | timeout: Io.Timeout, | 5060 | ) struct { ?net.Socket.ReceiveError, usize } { |
| 5044 | ) struct { ?net.Socket.ReceiveTimeoutError, usize } { | | |
| 5045 | const ev: *Evented = @ptrCast(@alignCast(userdata)); | | |
| 5046 | const ev_io = ev.io(); | | |
| 5047 | | | |
| 5048 | var message_i: usize = 0; | 5061 | var message_i: usize = 0; |
| 5049 | var data_i: usize = 0; | 5062 | var data_i: usize = 0; |
| 5050 | | 5063 | |
| 5051 | const deadline: ?struct { | | |
| 5052 | raw: Io.Timestamp, | | |
| 5053 | timespec: linux.kernel_timespec, | | |
| 5054 | clock: Io.Clock, | | |
| 5055 | } = if (timeout.toTimestamp(ev_io)) |deadline| deadline: { | | |
| 5056 | const ns = deadline.raw.toNanoseconds(); | | |
| 5057 | break :deadline .{ | | |
| 5058 | .raw = deadline.raw, | | |
| 5059 | .timespec = .{ | | |
| 5060 | .sec = @intCast(@divFloor(ns, std.time.ns_per_s)), | | |
| 5061 | .nsec = @intCast(@mod(ns, std.time.ns_per_s)), | | |
| 5062 | }, | | |
| 5063 | .clock = deadline.clock, | | |
| 5064 | }; | | |
| 5065 | } else null; | | |
| 5066 | | | |
| 5067 | var cancel_region: CancelRegion = .init(); | | |
| 5068 | defer cancel_region.deinit(); | | |
| 5069 | while (true) { | 5064 | while (true) { |
| 5070 | if (message_buffer.len - message_i == 0) return .{ null, message_i }; | 5065 | if (message_buffer.len - message_i == 0) return .{ null, message_i }; |
| 5071 | const message = &message_buffer[message_i]; | 5066 | const message = &message_buffer[message_i]; |
| ... | @@ -5085,7 +5080,7 @@ fn netReceive( | ... | @@ -5085,7 +5080,7 @@ fn netReceive( |
| 5085 | const thread = cancel_region.awaitIoUring() catch |err| return .{ err, message_i }; | 5080 | const thread = cancel_region.awaitIoUring() catch |err| return .{ err, message_i }; |
| 5086 | thread.enqueue().* = .{ | 5081 | thread.enqueue().* = .{ |
| 5087 | .opcode = .RECVMSG, | 5082 | .opcode = .RECVMSG, |
| 5088 | .flags = if (deadline) |_| linux.IOSQE_IO_LINK else 0, | 5083 | .flags = 0, |
| 5089 | .ioprio = 0, | 5084 | .ioprio = 0, |
| 5090 | .fd = handle, | 5085 | .fd = handle, |
| 5091 | .off = 0, | 5086 | .off = 0, |
| ... | @@ -5102,26 +5097,6 @@ fn netReceive( | ... | @@ -5102,26 +5097,6 @@ fn netReceive( |
| 5102 | .addr3 = 0, | 5097 | .addr3 = 0, |
| 5103 | .resv = 0, | 5098 | .resv = 0, |
| 5104 | }; | 5099 | }; |
| 5105 | if (deadline) |*deadline_ptr| thread.enqueue().* = .{ | | |
| 5106 | .opcode = .LINK_TIMEOUT, | | |
| 5107 | .flags = linux.IOSQE_CQE_SKIP_SUCCESS, | | |
| 5108 | .ioprio = 0, | | |
| 5109 | .fd = 0, | | |
| 5110 | .off = 0, | | |
| 5111 | .addr = @intFromPtr(&deadline_ptr.timespec), | | |
| 5112 | .len = 1, | | |
| 5113 | .rw_flags = linux.IORING_TIMEOUT_ABS | @as(u32, switch (deadline_ptr.clock) { | | |
| 5114 | .real => linux.IORING_TIMEOUT_REALTIME, | | |
| 5115 | else => 0, | | |
| 5116 | .boot => linux.IORING_TIMEOUT_BOOTTIME, | | |
| 5117 | }), | | |
| 5118 | .user_data = @intFromEnum(Completion.Userdata.wakeup), | | |
| 5119 | .buf_index = 0, | | |
| 5120 | .personality = 0, | | |
| 5121 | .splice_fd_in = 0, | | |
| 5122 | .addr3 = 0, | | |
| 5123 | .resv = 0, | | |
| 5124 | }; | | |
| 5125 | ev.yield(null, .nothing); | 5100 | ev.yield(null, .nothing); |
| 5126 | const completion = cancel_region.completion(); | 5101 | const completion = cancel_region.completion(); |
| 5127 | switch (completion.errno()) { | 5102 | switch (completion.errno()) { |
| ... | @@ -5144,9 +5119,7 @@ fn netReceive( | ... | @@ -5144,9 +5119,7 @@ fn netReceive( |
| 5144 | continue; | 5119 | continue; |
| 5145 | }, | 5120 | }, |
| 5146 | .AGAIN => unreachable, | 5121 | .AGAIN => unreachable, |
| 5147 | .INTR, .CANCELED => if (deadline) |d| if (now(ev, d.clock).nanoseconds >= d.raw.nanoseconds) | 5122 | .INTR, .CANCELED => {}, |
| 5148 | return .{ error.Timeout, message_i }, | | |
| 5149 | | | |
| 5150 | .BADF => |err| return .{ errnoBug(err), message_i }, | 5123 | .BADF => |err| return .{ errnoBug(err), message_i }, |
| 5151 | .NFILE => return .{ error.SystemFdQuotaExceeded, message_i }, | 5124 | .NFILE => return .{ error.SystemFdQuotaExceeded, message_i }, |
| 5152 | .MFILE => return .{ error.ProcessFdQuotaExceeded, message_i }, | 5125 | .MFILE => return .{ error.ProcessFdQuotaExceeded, message_i }, |