| ... | @@ -2547,14 +2547,13 @@ fn batchCancel(userdata: ?*anyopaque, b: *Io.Batch) void { | ... | @@ -2547,14 +2547,13 @@ fn batchCancel(userdata: ?*anyopaque, b: *Io.Batch) void { |
| 2547 | | 2547 | |
| 2548 | fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!void { | 2548 | fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!void { |
| 2549 | const t: *Threaded = @ptrCast(@alignCast(userdata)); | 2549 | const t: *Threaded = @ptrCast(@alignCast(userdata)); |
| 2550 | _ = t; | | |
| 2551 | | 2550 | |
| 2552 | if (operations.len == 1) { | 2551 | if (operations.len == 1) { |
| 2553 | @branchHint(.likely); | 2552 | @branchHint(.likely); |
| 2554 | return operate(&operations[0]); | 2553 | return operate(&operations[0]); |
| 2555 | } | 2554 | } |
| 2556 | | 2555 | |
| 2557 | if (is_windows) @panic("TODO"); | 2556 | if (is_windows) return batchWindows(t, operations); |
| 2558 | | 2557 | |
| 2559 | var poll_buffer: [poll_buffer_len]posix.pollfd = undefined; | 2558 | var poll_buffer: [poll_buffer_len]posix.pollfd = undefined; |
| 2560 | var map_buffer: [poll_buffer_len]u8 = undefined; // poll_buffer index to operations index | 2559 | var map_buffer: [poll_buffer_len]u8 = undefined; // poll_buffer index to operations index |
| ... | @@ -2578,7 +2577,7 @@ fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!v | ... | @@ -2578,7 +2577,7 @@ fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!v |
| 2578 | const map = map_buffer[0..poll_i]; | 2577 | const map = map_buffer[0..poll_i]; |
| 2579 | | 2578 | |
| 2580 | var pending = poll_i; | 2579 | var pending = poll_i; |
| 2581 | while (pending > 1) { | 2580 | while (pending > 0) { |
| 2582 | const syscall = Syscall.start() catch |err| switch (err) { | 2581 | const syscall = Syscall.start() catch |err| switch (err) { |
| 2583 | error.Canceled => { | 2582 | error.Canceled => { |
| 2584 | if (!setOperationsError(operations, polls, map, error.Canceled)) | 2583 | if (!setOperationsError(operations, polls, map, error.Canceled)) |
| ... | @@ -2589,17 +2588,11 @@ fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!v | ... | @@ -2589,17 +2588,11 @@ fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!v |
| 2589 | const rc = posix.system.poll(polls.ptr, polls.len, -1); | 2588 | const rc = posix.system.poll(polls.ptr, polls.len, -1); |
| 2590 | syscall.finish(); | 2589 | syscall.finish(); |
| 2591 | switch (posix.errno(rc)) { | 2590 | switch (posix.errno(rc)) { |
| 2592 | .SUCCESS => { | 2591 | .SUCCESS => for (polls, map) |*poll_fd, i| { |
| 2593 | if (rc == 0) { | 2592 | if (poll_fd.revents == 0) continue; |
| 2594 | // Spurious timeout; handle the same as INTR. | 2593 | poll_fd.fd = -1; |
| 2595 | continue; | 2594 | pending -= 1; |
| 2596 | } | 2595 | operate(&operations[i]); |
| 2597 | for (polls, map) |*poll_fd, i| { | | |
| 2598 | if (poll_fd.revents == 0) continue; | | |
| 2599 | poll_fd.fd = -1; | | |
| 2600 | pending -= 1; | | |
| 2601 | operate(&operations[i]); | | |
| 2602 | } | | |
| 2603 | }, | 2596 | }, |
| 2604 | .INTR => continue, | 2597 | .INTR => continue, |
| 2605 | .NOMEM => { | 2598 | .NOMEM => { |
| ... | @@ -2612,11 +2605,67 @@ fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!v | ... | @@ -2612,11 +2605,67 @@ fn batch(userdata: ?*anyopaque, operations: []Io.Operation) Io.ConcurrentError!v |
| 2612 | }, | 2605 | }, |
| 2613 | } | 2606 | } |
| 2614 | } | 2607 | } |
| | 2608 | } |
| 2615 | | 2609 | |
| 2616 | if (pending == 1) for (poll_buffer[0..poll_i], map_buffer[0..poll_i]) |*poll_fd, i| { | 2610 | fn batchWindows(t: *Threaded, operations: []Io.Operation) Io.ConcurrentError!void { |
| 2617 | if (poll_fd.fd == -1) continue; | 2611 | _ = t; |
| 2618 | operate(&operations[i]); | 2612 | var overlapped_buffer: [poll_buffer_len]windows.OVERLAPPED = undefined; |
| | 2613 | var handles_buffer: [poll_buffer_len]windows.HANDLE = undefined; |
| | 2614 | var map_buffer: [poll_buffer_len]u8 = undefined; // handles_buffer index to operations index |
| | 2615 | var buffer_i: usize = 0; |
| | 2616 | |
| | 2617 | for (operations, 0..) |*op, operation_index| switch (op.*) { |
| | 2618 | .noop => continue, |
| | 2619 | .file_read_streaming => |*o| { |
| | 2620 | if (handles_buffer.len - buffer_i == 0) return error.ConcurrencyUnavailable; |
| | 2621 | |
| | 2622 | const overlapped = &overlapped_buffer[buffer_i]; |
| | 2623 | overlapped.* = .{ |
| | 2624 | .Internal = 0, |
| | 2625 | .InternalHigh = 0, |
| | 2626 | .DUMMYUNIONNAME = .{ |
| | 2627 | .DUMMYSTRUCTNAME = .{ |
| | 2628 | .Offset = 0, |
| | 2629 | .OffsetHigh = 0, |
| | 2630 | }, |
| | 2631 | .Pointer = null, |
| | 2632 | }, |
| | 2633 | .hEvent = null, |
| | 2634 | }; |
| | 2635 | var n: windows.DWORD = undefined; |
| | 2636 | const buf = o.data[0]; |
| | 2637 | if (windows.kernel32.ReadFile(o.file.handle, buf.ptr, buf.len, &n, overlapped) == 0) { |
| | 2638 | @panic("TODO"); |
| | 2639 | } |
| | 2640 | handles_buffer[buffer_i] = o.file.handle; |
| | 2641 | map_buffer[buffer_i] = @intCast(operation_index); |
| | 2642 | buffer_i += 1; |
| | 2643 | }, |
| 2619 | }; | 2644 | }; |
| | 2645 | |
| | 2646 | const handles = handles_buffer[0..buffer_i]; |
| | 2647 | const map = map_buffer[0..buffer_i]; |
| | 2648 | var pending = buffer_i; |
| | 2649 | |
| | 2650 | while (pending > 0) { |
| | 2651 | const syscall: Syscall = try .start(); |
| | 2652 | const index = windows.WaitForMultipleObjectsEx(handles, false, windows.INFINITE, true); |
| | 2653 | syscall.finish(); |
| | 2654 | var n: windows.DWORD = undefined; |
| | 2655 | if (0 == windows.kernel32.GetOverlappedResult(handles[index], overlapped_buffer[index], &n, 0)) { |
| | 2656 | switch (windows.GetLastError()) { |
| | 2657 | .BROKEN_PIPE => @panic("TODO"), |
| | 2658 | .OPERATION_ABORTED => @panic("TODO"), |
| | 2659 | else => @panic("TODO"), |
| | 2660 | } |
| | 2661 | } else switch (operations[map[index]]) { |
| | 2662 | .noop => unreachable, |
| | 2663 | .file_read_streaming => |*o| { |
| | 2664 | o.status = .{ .result = n }; |
| | 2665 | pending -= 1; |
| | 2666 | }, |
| | 2667 | } |
| | 2668 | } |
| 2620 | } | 2669 | } |
| 2621 | | 2670 | |
| 2622 | fn setOperationsError( | 2671 | fn setOperationsError( |