| ... | ... | @@ -339,8 +339,8 @@ const Group = struct { |
| 339 | 339 | .canceled => true, |
| 340 | 340 | .parked => unreachable, |
| 341 | 341 | .blocked => unreachable, |
| 342 | | .blocked_alertable => unreachable, |
| 343 | | .blocked_alertable_canceling => unreachable, |
| 342 | .blocked_apc => unreachable, |
| 343 | .blocked_windows_dns => unreachable, |
| 344 | 344 | .blocked_canceling => unreachable, |
| 345 | 345 | }; |
| 346 | 346 | if (result) { |
| ... | ... | @@ -379,7 +379,10 @@ const Group = struct { |
| 379 | 379 | while (it) |thread| : (it = thread.next) { |
| 380 | 380 | // This non-mutating RMW exists for ordering reasons: see comment in `Group.Task.start` for reasons. |
| 381 | 381 | _ = thread.status.fetchOr(.{ .cancelation = @enumFromInt(0), .awaitable = .null }, .release); |
| 382 | | if (thread.cancelAwaitable(.fromGroup(g.ptr))) any_blocked = true; |
| 382 | if (thread.cancelAwaitable(.fromGroup(g.ptr))) |method| { |
| 383 | thread.interrupt_method = method; |
| 384 | any_blocked = true; |
| 385 | } |
| 383 | 386 | } |
| 384 | 387 | return any_blocked; |
| 385 | 388 | } |
| ... | ... | @@ -391,7 +394,7 @@ const Group = struct { |
| 391 | 394 | var any_signaled = false; |
| 392 | 395 | var it = t.worker_threads.load(.acquire); // acquire `Thread` values |
| 393 | 396 | while (it) |thread| : (it = thread.next) { |
| 394 | | if (thread.signalCanceledSyscall(t, .fromGroup(g.ptr))) any_signaled = true; |
| 397 | if (thread.signalCanceledSyscall(t, .fromGroup(g.ptr), thread.interrupt_method)) any_signaled = true; |
| 395 | 398 | } |
| 396 | 399 | return any_signaled; |
| 397 | 400 | } |
| ... | ... | @@ -543,8 +546,8 @@ const Future = struct { |
| 543 | 546 | .canceled => true, |
| 544 | 547 | .parked => unreachable, |
| 545 | 548 | .blocked => unreachable, |
| 546 | | .blocked_alertable => unreachable, |
| 547 | | .blocked_alertable_canceling => unreachable, |
| 549 | .blocked_apc => unreachable, |
| 550 | .blocked_windows_dns => unreachable, |
| 548 | 551 | .blocked_canceling => unreachable, |
| 549 | 552 | }; |
| 550 | 553 | thread.status.store(.{ .cancelation = .none, .awaitable = .null }, .monotonic); |
| ... | ... | @@ -573,11 +576,15 @@ const Future = struct { |
| 573 | 576 | num_completed: *std.atomic.Value(u32), |
| 574 | 577 | thread: ?*Thread, |
| 575 | 578 | ) void { |
| 576 | | var need_signal: bool = if (thread) |th| th.cancelAwaitable(.fromFuture(future)) else false; |
| 579 | var interrupt_method: ?Thread.InterruptMethod = |
| 580 | if (thread) |th| th.cancelAwaitable(.fromFuture(future)) else null; |
| 577 | 581 | var timeout_ns: u64 = 1 << 10; |
| 578 | 582 | while (true) { |
| 579 | | need_signal = need_signal and thread.?.signalCanceledSyscall(t, .fromFuture(future)); |
| 580 | | Thread.futexWaitUncancelable(&num_completed.raw, 0, if (need_signal) timeout_ns else null); |
| 583 | if (interrupt_method) |method| { |
| 584 | if (!thread.?.signalCanceledSyscall(t, .fromFuture(future), method)) |
| 585 | interrupt_method = null; |
| 586 | } |
| 587 | Thread.futexWaitUncancelable(&num_completed.raw, 0, if (interrupt_method != null) timeout_ns else null); |
| 581 | 588 | switch (num_completed.load(.acquire)) { // acquire task results |
| 582 | 589 | 0 => {}, |
| 583 | 590 | 1 => break, |
| ... | ... | @@ -625,6 +632,9 @@ const Thread = struct { |
| 625 | 632 | cancel_protection: Io.CancelProtection, |
| 626 | 633 | /// Always released when `Status.cancelation` is set to `.parked`. |
| 627 | 634 | futex_waiter: if (use_parking_futex) ?*parking_futex.Waiter else ?noreturn, |
| 635 | apc_context: if (is_windows) ?*anyopaque else void, |
| 636 | /// Used only by group cancelation code for temporary storage. |
| 637 | interrupt_method: InterruptMethod, |
| 628 | 638 | |
| 629 | 639 | csprng: Csprng, |
| 630 | 640 | |
| ... | ... | @@ -652,11 +662,13 @@ const Thread = struct { |
| 652 | 662 | /// To request cancelation, set the status to `.blocked_canceling` and repeatedly interrupt the system call until the status changes. |
| 653 | 663 | blocked = 0b011, |
| 654 | 664 | |
| 655 | | /// Windows-only: the thread is blocked in an alertable wait via |
| 656 | | /// `NtDelayExecution`. To request cancelation, set the status to |
| 657 | | /// `blocked_alertable_canceling` and repeatedly alert the thread |
| 658 | | /// until the status changes. |
| 659 | | blocked_alertable = 0b010, |
| 665 | /// Windows-only: the thread is blocked in a call to `NtDelayExecution`. |
| 666 | /// To request cancelation, set the status to `.canceling` and call `NtCancelIoFileEx`. |
| 667 | blocked_apc = 0b100, |
| 668 | |
| 669 | /// Windows-only: the thread is blocked in a call to `GetAddrInfoExW`. |
| 670 | /// To request cancelation, set the status to `.canceling` and call `GetAddrInfoExCancel`. |
| 671 | blocked_windows_dns = 0b010, |
| 660 | 672 | |
| 661 | 673 | /// The thread has an outstanding cancelation request but is not in a cancelable operation. |
| 662 | 674 | /// When it acknowledges the cancelation, it will set the status to `.canceled`. |
| ... | ... | @@ -705,8 +717,8 @@ const Thread = struct { |
| 705 | 717 | switch (status.cancelation) { |
| 706 | 718 | .parked => unreachable, |
| 707 | 719 | .blocked => unreachable, |
| 708 | | .blocked_alertable => unreachable, |
| 709 | | .blocked_alertable_canceling => unreachable, |
| 720 | .blocked_apc => unreachable, |
| 721 | .blocked_windows_dns => unreachable, |
| 710 | 722 | .blocked_canceling => unreachable, |
| 711 | 723 | .none, .canceled => {}, |
| 712 | 724 | .canceling => { |
| ... | ... | @@ -986,17 +998,17 @@ const Thread = struct { |
| 986 | 998 | /// It is possible that `thread` gets canceled by this function, but is blocked in a syscall. In |
| 987 | 999 | /// that case, the thread may need to be sent a signal to interrupt the call. This function will |
| 988 | 1000 | /// return `true` to indicate this, in which case the caller must call `signalCanceledSyscall`. |
| 989 | | fn cancelAwaitable(thread: *Thread, awaitable: AwaitableId) bool { |
| 1001 | fn cancelAwaitable(thread: *Thread, awaitable: AwaitableId) ?InterruptMethod { |
| 990 | 1002 | var status = thread.status.load(.monotonic); |
| 991 | 1003 | while (true) { |
| 992 | | if (status.awaitable != awaitable) return false; // thread is working on something else |
| 1004 | if (status.awaitable != awaitable) return null; // thread is working on something else |
| 993 | 1005 | status = switch (status.cancelation) { |
| 994 | 1006 | .none => thread.status.cmpxchgWeak( |
| 995 | 1007 | .{ .cancelation = .none, .awaitable = awaitable }, |
| 996 | 1008 | .{ .cancelation = .canceling, .awaitable = awaitable }, |
| 997 | 1009 | .monotonic, |
| 998 | 1010 | .monotonic, |
| 999 | | ) orelse return false, |
| 1011 | ) orelse return null, |
| 1000 | 1012 | |
| 1001 | 1013 | .parked => thread.status.cmpxchgWeak( |
| 1002 | 1014 | .{ .cancelation = .parked, .awaitable = awaitable }, |
| ... | ... | @@ -1009,7 +1021,7 @@ const Thread = struct { |
| 1009 | 1021 | parking_futex.removeCanceledWaiter(futex_waiter); |
| 1010 | 1022 | } |
| 1011 | 1023 | unpark(&.{thread.id}, null); |
| 1012 | | return false; |
| 1024 | return null; |
| 1013 | 1025 | }, |
| 1014 | 1026 | |
| 1015 | 1027 | .blocked => thread.status.cmpxchgWeak( |
| ... | ... | @@ -1017,7 +1029,17 @@ const Thread = struct { |
| 1017 | 1029 | .{ .cancelation = .blocked_canceling, .awaitable = awaitable }, |
| 1018 | 1030 | .monotonic, |
| 1019 | 1031 | .monotonic, |
| 1020 | | ) orelse return true, |
| 1032 | ) orelse return .sync, |
| 1033 | |
| 1034 | .blocked_apc => thread.status.cmpxchgWeak( |
| 1035 | .{ .cancelation = .blocked_apc, .awaitable = awaitable }, |
| 1036 | .{ .cancelation = .canceling, .awaitable = awaitable }, |
| 1037 | .monotonic, |
| 1038 | .monotonic, |
| 1039 | ) orelse { |
| 1040 | if (!is_windows) unreachable; |
| 1041 | return .apc; |
| 1042 | }, |
| 1021 | 1043 | |
| 1022 | 1044 | .blocked_alertable => thread.status.cmpxchgWeak( |
| 1023 | 1045 | .{ .cancelation = .blocked_alertable, .awaitable = awaitable }, |
| ... | ... | @@ -1026,14 +1048,14 @@ const Thread = struct { |
| 1026 | 1048 | .monotonic, |
| 1027 | 1049 | ) orelse { |
| 1028 | 1050 | if (!is_windows) unreachable; |
| 1029 | | return true; |
| 1051 | return .dns; |
| 1030 | 1052 | }, |
| 1031 | 1053 | |
| 1032 | 1054 | .canceling, .canceled => { |
| 1033 | 1055 | // This can happen when the task start raced with the cancelation, so the thread |
| 1034 | 1056 | // saw the cancelation on the future/group *and* we are trying to signal the |
| 1035 | 1057 | // thread here. |
| 1036 | | return false; |
| 1058 | return null; |
| 1037 | 1059 | }, |
| 1038 | 1060 | |
| 1039 | 1061 | .blocked_canceling => unreachable, // `awaitable` has not been canceled before now |
| ... | ... | @@ -1042,6 +1064,11 @@ const Thread = struct { |
| 1042 | 1064 | } |
| 1043 | 1065 | } |
| 1044 | 1066 | |
| 1067 | const InterruptMethod = switch (native_os) { |
| 1068 | .windows => enum { sync, dns, apc }, |
| 1069 | else => enum { sync }, |
| 1070 | }; |
| 1071 | |
| 1045 | 1072 | /// Sends a signal to `thread` if it is still blocked in a syscall (i.e. has not yet observed |
| 1046 | 1073 | /// the cancelation request from `cancelAwaitable`). |
| 1047 | 1074 | /// |
| ... | ... | @@ -1051,24 +1078,21 @@ const Thread = struct { |
| 1051 | 1078 | /// the thread is still blocked. For the implementation, `Future.waitForCancelWithSignaling` and |
| 1052 | 1079 | /// `Group.waitForCancelWithSignaling`: they use exponential backoff starting at a 1us delay and |
| 1053 | 1080 | /// doubling each call. In practice, it is rare to send more than one signal. |
| 1054 | | fn signalCanceledSyscall(thread: *Thread, t: *Threaded, awaitable: AwaitableId) bool { |
| 1055 | | const status = thread.status.load(.monotonic); |
| 1056 | | if (status.awaitable != awaitable) { |
| 1057 | | // The thread has moved on and is working on something totally different. |
| 1058 | | return false; |
| 1059 | | } |
| 1081 | fn signalCanceledSyscall(thread: *Thread, t: *Threaded, awaitable: AwaitableId, method: InterruptMethod) bool { |
| 1082 | const bad_status: Status = .{ .cancelation = .blocked_canceling, .awaitable = awaitable }; |
| 1083 | if (thread.status.load(.monotonic) != bad_status) return false; |
| 1060 | 1084 | |
| 1061 | 1085 | // The thread ID and/or handle can be read non-atomically because they never change and were |
| 1062 | 1086 | // released by the store that made `thread` available to us. |
| 1063 | 1087 | |
| 1064 | | switch (status.cancelation) { |
| 1065 | | .blocked_canceling => if (std.Thread.use_pthreads) { |
| 1066 | | return switch (std.c.pthread_kill(thread.handle, .IO)) { |
| 1067 | | 0 => true, |
| 1068 | | else => false, |
| 1069 | | }; |
| 1070 | | } else switch (native_os) { |
| 1071 | | .linux => { |
| 1088 | if (std.Thread.use_pthreads) switch (method) { |
| 1089 | .sync => return switch (std.c.pthread_kill(thread.handle, .IO)) { |
| 1090 | 0 => true, |
| 1091 | else => false, |
| 1092 | }, |
| 1093 | } else switch (native_os) { |
| 1094 | .linux => switch (method) { |
| 1095 | .sync => { |
| 1072 | 1096 | const pid: posix.pid_t = pid: { |
| 1073 | 1097 | const cached_pid = @atomicLoad(Pid, &t.pid, .monotonic); |
| 1074 | 1098 | if (cached_pid != .unknown) break :pid @intFromEnum(cached_pid); |
| ... | ... | @@ -1081,7 +1105,9 @@ const Thread = struct { |
| 1081 | 1105 | else => false, |
| 1082 | 1106 | }; |
| 1083 | 1107 | }, |
| 1084 | | .windows => { |
| 1108 | }, |
| 1109 | .windows => switch (method) { |
| 1110 | .sync => { |
| 1085 | 1111 | var iosb: windows.IO_STATUS_BLOCK = undefined; |
| 1086 | 1112 | return switch (windows.ntdll.NtCancelSynchronousIoFile(thread.handle, null, &iosb)) { |
| 1087 | 1113 | .NOT_FOUND => true, // this might mean the operation hasn't started yet |
| ... | ... | @@ -1089,15 +1115,8 @@ const Thread = struct { |
| 1089 | 1115 | else => false, |
| 1090 | 1116 | }; |
| 1091 | 1117 | }, |
| 1092 | | else => return false, |
| 1093 | | }, |
| 1094 | | |
| 1095 | | .blocked_alertable_canceling => { |
| 1096 | | if (!is_windows) unreachable; |
| 1097 | | return switch (windows.ntdll.NtAlertThread(thread.handle)) { |
| 1098 | | .SUCCESS => true, |
| 1099 | | else => false, |
| 1100 | | }; |
| 1118 | .dns => @panic("TODO call GetAddrInfoExCancel"), |
| 1119 | .apc => @panic("TODO call NtCancelIoFileEx"), |
| 1101 | 1120 | }, |
| 1102 | 1121 | |
| 1103 | 1122 | else => { |
| ... | ... | @@ -1145,8 +1164,8 @@ const Syscall = struct { |
| 1145 | 1164 | }, .monotonic).cancelation) { |
| 1146 | 1165 | .parked => unreachable, |
| 1147 | 1166 | .blocked => unreachable, |
| 1148 | | .blocked_alertable => unreachable, |
| 1149 | | .blocked_alertable_canceling => unreachable, |
| 1167 | .blocked_apc => unreachable, |
| 1168 | .blocked_windows_dns => unreachable, |
| 1150 | 1169 | .blocked_canceling => unreachable, |
| 1151 | 1170 | .none => return .{ .thread = thread }, // new status is `.blocked` |
| 1152 | 1171 | .canceling => return error.Canceled, // new status is `.canceled` |
| ... | ... | @@ -1165,19 +1184,14 @@ const Syscall = struct { |
| 1165 | 1184 | }, .monotonic).cancelation) { |
| 1166 | 1185 | .none => unreachable, |
| 1167 | 1186 | .parked => unreachable, |
| 1168 | | .blocked_alertable => unreachable, |
| 1169 | | .blocked_alertable_canceling => unreachable, |
| 1187 | .blocked_apc => unreachable, |
| 1188 | .blocked_windows_dns => unreachable, |
| 1170 | 1189 | .canceling => unreachable, |
| 1171 | 1190 | .canceled => unreachable, |
| 1172 | 1191 | .blocked => {}, // new status is `.blocked` (unchanged) |
| 1173 | 1192 | .blocked_canceling => return error.Canceled, // new status is `.canceled` |
| 1174 | 1193 | } |
| 1175 | 1194 | } |
| 1176 | | fn toApc(s: Syscall) Io.Cancelable!void { |
| 1177 | | // TODO set state to indicate instead of NtCancelSynchronousIoFile we |
| 1178 | | // need to use NtCancelIoFileEx |
| 1179 | | return s.checkCancel(); |
| 1180 | | } |
| 1181 | 1195 | /// Marks this syscall as finished. |
| 1182 | 1196 | fn finish(s: Syscall) void { |
| 1183 | 1197 | const thread = s.thread orelse return; |
| ... | ... | @@ -1187,8 +1201,8 @@ const Syscall = struct { |
| 1187 | 1201 | }, .monotonic).cancelation) { |
| 1188 | 1202 | .none => unreachable, |
| 1189 | 1203 | .parked => unreachable, |
| 1190 | | .blocked_alertable => unreachable, |
| 1191 | | .blocked_alertable_canceling => unreachable, |
| 1204 | .blocked_apc => unreachable, |
| 1205 | .blocked_windows_dns => unreachable, |
| 1192 | 1206 | .canceling => unreachable, |
| 1193 | 1207 | .canceled => unreachable, |
| 1194 | 1208 | .blocked => {}, // new status is `.none` |
| ... | ... | @@ -1196,25 +1210,25 @@ const Syscall = struct { |
| 1196 | 1210 | } |
| 1197 | 1211 | } |
| 1198 | 1212 | /// Indicates instead of `NtCancelSynchronousIoFile` we need to use |
| 1199 | | /// `NtAlertThread` to interrupt the wait. |
| 1213 | /// `NtCancelIoFileEx` to interrupt the wait. |
| 1200 | 1214 | /// |
| 1201 | 1215 | /// Windows only, called from blocked state only. |
| 1202 | | fn toAlertable(s: Syscall) Io.Cancelable!AlertableSyscall { |
| 1203 | | comptime assert(is_windows); |
| 1204 | | const thread = s.thread orelse return .{ .thread = null }; |
| 1216 | fn toApc(s: Syscall, apc_context: ?*anyopaque) Io.Cancelable!void { |
| 1217 | const thread = s.thread orelse return; |
| 1218 | thread.apc_context = apc_context; |
| 1205 | 1219 | var prev = thread.status.load(.monotonic); |
| 1206 | 1220 | while (true) prev = switch (prev.cancelation) { |
| 1207 | 1221 | .none => unreachable, |
| 1208 | 1222 | .parked => unreachable, |
| 1209 | | .blocked_alertable => unreachable, |
| 1210 | | .blocked_alertable_canceling => unreachable, |
| 1223 | .blocked_apc => unreachable, |
| 1224 | .blocked_windows_dns => unreachable, |
| 1211 | 1225 | .canceling => unreachable, |
| 1212 | 1226 | .canceled => unreachable, |
| 1213 | 1227 | |
| 1214 | 1228 | .blocked => thread.status.cmpxchgWeak(prev, .{ |
| 1215 | | .cancelation = .blocked_alertable, |
| 1229 | .cancelation = .blocked_apc, |
| 1216 | 1230 | .awaitable = prev.awaitable, |
| 1217 | | }, .monotonic, .monotonic) orelse return .{ .thread = thread }, |
| 1231 | }, .monotonic, .monotonic) orelse return, |
| 1218 | 1232 | |
| 1219 | 1233 | .blocked_canceling => thread.status.cmpxchgWeak(prev, .{ |
| 1220 | 1234 | .cancelation = .canceled, |
| ... | ... | @@ -1222,6 +1236,45 @@ const Syscall = struct { |
| 1222 | 1236 | }, .monotonic, .monotonic) orelse return error.Canceled, |
| 1223 | 1237 | }; |
| 1224 | 1238 | } |
| 1239 | /// Windows only, called from blocked_apc state only. |
| 1240 | fn checkCancelApc(s: Syscall) Io.Cancelable!void { |
| 1241 | const thread = s.thread orelse return; |
| 1242 | var prev = thread.status.load(.monotonic); |
| 1243 | while (true) prev = switch (prev.cancelation) { |
| 1244 | .none => unreachable, |
| 1245 | .parked => unreachable, |
| 1246 | .blocked_windows_dns => unreachable, |
| 1247 | .blocked => unreachable, |
| 1248 | .canceling => unreachable, |
| 1249 | .canceled => unreachable, |
| 1250 | .blocked_apc => return, |
| 1251 | .blocked_canceling => thread.status.cmpxchgWeak(prev, .{ |
| 1252 | .cancelation = .canceled, |
| 1253 | .awaitable = prev.awaitable, |
| 1254 | }, .monotonic, .monotonic) orelse return error.Canceled, |
| 1255 | }; |
| 1256 | } |
| 1257 | /// Windows only, called from blocked_apc state only. |
| 1258 | fn finishApc(s: Syscall) void { |
| 1259 | const thread = s.thread orelse return; |
| 1260 | var prev = thread.status.load(.monotonic); |
| 1261 | while (true) prev = switch (prev.cancelation) { |
| 1262 | .none => unreachable, |
| 1263 | .parked => unreachable, |
| 1264 | .blocked_windows_dns => unreachable, |
| 1265 | .blocked => unreachable, |
| 1266 | .canceling => unreachable, |
| 1267 | .canceled => unreachable, |
| 1268 | .blocked_apc => thread.status.cmpxchgWeak(prev, .{ |
| 1269 | .cancelation = .none, |
| 1270 | .awaitable = prev.awaitable, |
| 1271 | }, .monotonic, .monotonic) orelse return, |
| 1272 | .blocked_canceling => thread.status.cmpxchgWeak(prev, .{ |
| 1273 | .cancelation = .canceling, |
| 1274 | .awaitable = prev.awaitable, |
| 1275 | }, .monotonic, .monotonic) orelse return, |
| 1276 | }; |
| 1277 | } |
| 1225 | 1278 | /// Convenience wrapper which calls `finish`, then returns `err`. |
| 1226 | 1279 | fn fail(s: Syscall, err: anytype) @TypeOf(err) { |
| 1227 | 1280 | s.finish(); |
| ... | ... | @@ -1501,6 +1554,8 @@ fn worker(t: *Threaded) void { |
| 1501 | 1554 | .cancel_protection = .unblocked, |
| 1502 | 1555 | .futex_waiter = undefined, |
| 1503 | 1556 | .csprng = .{}, |
| 1557 | .apc_context = undefined, |
| 1558 | .interrupt_method = undefined, |
| 1504 | 1559 | }; |
| 1505 | 1560 | Thread.current = &thread; |
| 1506 | 1561 | |
| ... | ... | @@ -2109,8 +2164,8 @@ fn groupAsyncEager( |
| 2109 | 2164 | .canceled => true, |
| 2110 | 2165 | .parked => unreachable, |
| 2111 | 2166 | .blocked => unreachable, |
| 2112 | | .blocked_alertable => unreachable, |
| 2113 | | .blocked_alertable_canceling => unreachable, |
| 2167 | .blocked_apc => unreachable, |
| 2168 | .blocked_windows_dns => unreachable, |
| 2114 | 2169 | .blocked_canceling => unreachable, |
| 2115 | 2170 | }; |
| 2116 | 2171 | } else false; |
| ... | ... | @@ -2121,8 +2176,8 @@ fn groupAsyncEager( |
| 2121 | 2176 | .canceled => true, |
| 2122 | 2177 | .parked => unreachable, |
| 2123 | 2178 | .blocked => unreachable, |
| 2124 | | .blocked_alertable => unreachable, |
| 2125 | | .blocked_alertable_canceling => unreachable, |
| 2179 | .blocked_apc => unreachable, |
| 2180 | .blocked_windows_dns => unreachable, |
| 2126 | 2181 | .blocked_canceling => unreachable, |
| 2127 | 2182 | }; |
| 2128 | 2183 | } else false; |
| ... | ... | @@ -2301,8 +2356,8 @@ fn recancelInner() void { |
| 2301 | 2356 | .canceling => unreachable, // called `recancel` but cancelation was already pending |
| 2302 | 2357 | .parked => unreachable, |
| 2303 | 2358 | .blocked => unreachable, |
| 2304 | | .blocked_alertable => unreachable, |
| 2305 | | .blocked_alertable_canceling => unreachable, |
| 2359 | .blocked_apc => unreachable, |
| 2360 | .blocked_windows_dns => unreachable, |
| 2306 | 2361 | .blocked_canceling => unreachable, |
| 2307 | 2362 | } |
| 2308 | 2363 | } |
| ... | ... | @@ -8376,6 +8431,7 @@ fn fileReadStreamingWindows(file: File, data: []const []u8) File.Reader.Error!us |
| 8376 | 8431 | const buffer = data[index]; |
| 8377 | 8432 | |
| 8378 | 8433 | var io_status_block: windows.IO_STATUS_BLOCK = undefined; |
| 8434 | var done: bool = false; |
| 8379 | 8435 | |
| 8380 | 8436 | read: { |
| 8381 | 8437 | const syscall: Syscall = try .start(); |
| ... | ... | @@ -8383,8 +8439,8 @@ fn fileReadStreamingWindows(file: File, data: []const []u8) File.Reader.Error!us |
| 8383 | 8439 | switch (windows.ntdll.NtReadFile( |
| 8384 | 8440 | file.handle, |
| 8385 | 8441 | null, // event |
| 8386 | | noopApc, // apc callback |
| 8387 | | null, // apc context |
| 8442 | flagApc, // apc callback |
| 8443 | &done, // apc context |
| 8388 | 8444 | &io_status_block, |
| 8389 | 8445 | buffer.ptr, |
| 8390 | 8446 | @min(std.math.maxInt(u32), buffer.len), |
| ... | ... | @@ -8402,12 +8458,12 @@ fn fileReadStreamingWindows(file: File, data: []const []u8) File.Reader.Error!us |
| 8402 | 8458 | //else => |status| return syscall.unexpectedNtstatus(status), |
| 8403 | 8459 | } |
| 8404 | 8460 | } |
| 8405 | | try syscall.toApc(); |
| 8461 | try syscall.toApc(&done); |
| 8406 | 8462 | while (true) { |
| 8407 | 8463 | switch (windows.ntdll.NtDelayExecution(1, null)) { |
| 8408 | | .USER_APC => break syscall.finish(), |
| 8464 | .USER_APC => break syscall.finishApc(), |
| 8409 | 8465 | .SUCCESS, .CANCELLED => { |
| 8410 | | try syscall.checkCancel(); |
| 8466 | try syscall.checkCancelApc(); |
| 8411 | 8467 | continue; |
| 8412 | 8468 | }, |
| 8413 | 8469 | else => |status| std.debug.panic("fileReadStreamingWindows NtDelayExecution returned {t}", .{status}), |
| ... | ... | @@ -8425,12 +8481,13 @@ fn fileReadStreamingWindows(file: File, data: []const []u8) File.Reader.Error!us |
| 8425 | 8481 | return io_status_block.Information; |
| 8426 | 8482 | } |
| 8427 | 8483 | |
| 8428 | | fn noopApc( |
| 8484 | fn flagApc( |
| 8429 | 8485 | apc_context: ?*anyopaque, |
| 8430 | 8486 | io_status_block: *windows.IO_STATUS_BLOCK, |
| 8431 | 8487 | unused: windows.ULONG, |
| 8432 | 8488 | ) callconv(.winapi) void { |
| 8433 | | _ = apc_context; |
| 8489 | const flag: *bool = @ptrCast(apc_context); |
| 8490 | flag.* = true; |
| 8434 | 8491 | _ = io_status_block; |
| 8435 | 8492 | _ = unused; |
| 8436 | 8493 | } |
| ... | ... | @@ -12442,7 +12499,7 @@ fn netLookupFallible( |
| 12442 | 12499 | var res: *ws2_32.ADDRINFOEXW = undefined; |
| 12443 | 12500 | const timeout: ?*ws2_32.timeval = null; |
| 12444 | 12501 | while (true) { |
| 12445 | | // TODO: hook this up to cancelation with `NtDelayExecution` and APC callbacks. |
| 12502 | // TODO: hook this up to cancelation with `Thread.Status.cancelation.blocked_windows_dns`. |
| 12446 | 12503 | try Thread.checkCancel(); |
| 12447 | 12504 | // TODO make this append to the queue eagerly rather than blocking until the whole thing finishes |
| 12448 | 12505 | const rc: ws2_32.WinsockError = @enumFromInt(ws2_32.GetAddrInfoExW(name_w, port_w, .DNS, null, &hints, &res, timeout, null, null, null)); |
| ... | ... | @@ -16176,8 +16233,8 @@ const parking_futex = struct { |
| 16176 | 16233 | .canceled => break :cancelable, // status is still `.canceled` |
| 16177 | 16234 | .parked => unreachable, |
| 16178 | 16235 | .blocked => unreachable, |
| 16179 | | .blocked_alertable => unreachable, |
| 16180 | | .blocked_alertable_canceling => unreachable, |
| 16236 | .blocked_apc => unreachable, |
| 16237 | .blocked_windows_dns => unreachable, |
| 16181 | 16238 | .blocked_canceling => unreachable, |
| 16182 | 16239 | } |
| 16183 | 16240 | // We could now be unparked for a cancelation at any time! |
| ... | ... | @@ -16228,8 +16285,8 @@ const parking_futex = struct { |
| 16228 | 16285 | }, |
| 16229 | 16286 | .canceled => unreachable, |
| 16230 | 16287 | .blocked => unreachable, |
| 16231 | | .blocked_alertable => unreachable, |
| 16232 | | .blocked_alertable_canceling => unreachable, |
| 16288 | .blocked_apc => unreachable, |
| 16289 | .blocked_windows_dns => unreachable, |
| 16233 | 16290 | .blocked_canceling => unreachable, |
| 16234 | 16291 | }, |
| 16235 | 16292 | } |
| ... | ... | @@ -16270,8 +16327,8 @@ const parking_futex = struct { |
| 16270 | 16327 | .canceling => continue, // race with a canceler who hasn't called `removeCanceledWaiter` yet |
| 16271 | 16328 | .canceled => unreachable, |
| 16272 | 16329 | .blocked => unreachable, |
| 16273 | | .blocked_alertable => unreachable, |
| 16274 | | .blocked_alertable_canceling => unreachable, |
| 16330 | .blocked_apc => unreachable, |
| 16331 | .blocked_windows_dns => unreachable, |
| 16275 | 16332 | .blocked_canceling => unreachable, |
| 16276 | 16333 | } |
| 16277 | 16334 | // We're waking this waiter. Remove them from the bucket and add them to our local list. |
| ... | ... | @@ -16337,8 +16394,8 @@ const parking_sleep = struct { |
| 16337 | 16394 | .canceled => break :cancelable, // status is still `.canceled` |
| 16338 | 16395 | .parked => unreachable, |
| 16339 | 16396 | .blocked => unreachable, |
| 16340 | | .blocked_alertable => unreachable, |
| 16341 | | .blocked_alertable_canceling => unreachable, |
| 16397 | .blocked_apc => unreachable, |
| 16398 | .blocked_windows_dns => unreachable, |
| 16342 | 16399 | .blocked_canceling => unreachable, |
| 16343 | 16400 | } |
| 16344 | 16401 | while (park(deadline, null)) { |
| ... | ... | @@ -16356,8 +16413,8 @@ const parking_sleep = struct { |
| 16356 | 16413 | .none => unreachable, |
| 16357 | 16414 | .canceled => unreachable, |
| 16358 | 16415 | .blocked => unreachable, |
| 16359 | | .blocked_alertable => unreachable, |
| 16360 | | .blocked_alertable_canceling => unreachable, |
| 16416 | .blocked_apc => unreachable, |
| 16417 | .blocked_windows_dns => unreachable, |
| 16361 | 16418 | .blocked_canceling => unreachable, |
| 16362 | 16419 | } |
| 16363 | 16420 | } else |err| switch (err) { |
| ... | ... | @@ -16376,8 +16433,8 @@ const parking_sleep = struct { |
| 16376 | 16433 | .none => unreachable, |
| 16377 | 16434 | .canceled => unreachable, |
| 16378 | 16435 | .blocked => unreachable, |
| 16379 | | .blocked_alertable => unreachable, |
| 16380 | | .blocked_alertable_canceling => unreachable, |
| 16436 | .blocked_apc => unreachable, |
| 16437 | .blocked_windows_dns => unreachable, |
| 16381 | 16438 | .blocked_canceling => unreachable, |
| 16382 | 16439 | }, |
| 16383 | 16440 | } |