authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-01-08 15:07:03-08:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-01-30 22:03:13-08:00
log6a7fe61d74f80456b32cb46ef21715bbffaf49a7
treeaf4535684dfcba1505fc61dbff57508cafeeb1e2
parent93f5c99149948b104ab504eff3a171b6c1bff065

std.Io.Threaded.operate: handle poll buffer exceeded


1 files changed, 72 insertions(+), 69 deletions(-)

lib/std/Io/Threaded.zig+72-69
...@@ -2458,81 +2458,84 @@ fn operate(userdata: ?*anyopaque, operations: []Io.Operation) void {...@@ -2458,81 +2458,84 @@ fn operate(userdata: ?*anyopaque, operations: []Io.Operation) void {
24582458
2459 var poll_buffer: [poll_buffer_len]posix.pollfd = undefined;2459 var poll_buffer: [poll_buffer_len]posix.pollfd = undefined;
2460 var map_buffer: [poll_buffer_len]u8 = undefined; // poll_buffer index to operations index2460 var map_buffer: [poll_buffer_len]u8 = undefined; // poll_buffer index to operations index
2461 var poll_i: usize = 0;2461 var operation_index: usize = 0;
24622462
2463 // Put all the file reads with nonblocking enabled into the poll set.2463 while (operation_index < operations.len) {
2464 if (operations.len > poll_buffer.len) @panic("TODO");2464 var poll_i: usize = 0;
24652465 while (operation_index < operations.len) : (operation_index += 1) {
2466 for (operations, 0..) |*operation, operation_index| switch (operation.*) {2466 switch (operations[operation_index]) {
2467 .noop => continue,2467 .noop => continue,
2468 .file_read_streaming => |*o| {2468 .file_read_streaming => |*o| {
2469 if (o.nonblocking) {2469 if (o.nonblocking) {
2470 o.result = error.WouldBlock;2470 o.result = error.WouldBlock;
2471 poll_buffer[poll_i] = .{2471 poll_buffer[poll_i] = .{
2472 .fd = o.file.handle,2472 .fd = o.file.handle,
2473 .events = posix.POLL.IN,2473 .events = posix.POLL.IN,
2474 .revents = undefined,2474 .revents = 0,
2475 };2475 };
2476 map_buffer[poll_i] = @intCast(operation_index);2476 if (map_buffer.len - poll_i == 0) break;
2477 poll_i += 1;2477 map_buffer[poll_i] = @intCast(operation_index);
2478 } else {2478 poll_i += 1;
2479 o.result = fileReadStreaming(o.file, o.data) catch |err| switch (err) {2479 } else {
2480 error.Canceled => {2480 o.result = fileReadStreaming(o.file, o.data) catch |err| switch (err) {
2481 setOperationsCanceled(operations[operation_index..]);2481 error.Canceled => {
2482 return;2482 setOperationsError(operations[operation_index..], error.Canceled);
2483 },2483 return;
2484 else => err,2484 },
2485 };2485 else => err,
2486 };
2487 }
2488 },
2486 }2489 }
2487 },2490 }
2488 };
2489
2490 if (poll_i == 0) {
2491 @branchHint(.likely);
2492 return;
2493 }
24942491
2495 while (true) {2492 if (poll_i == 0) {
2496 const syscall = Syscall.start() catch |err| switch (err) {2493 @branchHint(.likely);
2497 error.Canceled => {2494 return;
2498 setAllOperationsError(operations, map_buffer[0..poll_i], error.Canceled);
2499 return;
2500 },
2501 };
2502 const poll_rc = posix.system.poll(&poll_buffer, poll_i, -1);
2503 syscall.finish();
2504 switch (posix.errno(poll_rc)) {
2505 .SUCCESS => {
2506 if (poll_rc == 0) {
2507 // Spurious timeout; handle same as INTR.
2508 continue;
2509 }
2510 break;
2511 },
2512 .INTR => continue,
2513 .NOMEM => {
2514 setAllOperationsError(operations, map_buffer[0..poll_i], error.SystemResources);
2515 return;
2516 },
2517 else => {
2518 setAllOperationsError(operations, map_buffer[0..poll_i], error.Unexpected);
2519 return;
2520 },
2521 }2495 }
2522 }
25232496
2524 for (poll_buffer[0..poll_i], map_buffer[0..poll_i]) |*poll_fd, operation_index| {2497 while (true) {
2525 if (poll_fd.revents == 0) continue;2498 const syscall = Syscall.start() catch |err| switch (err) {
2526 switch (operations[operation_index]) {2499 error.Canceled => {
2527 .noop => unreachable,2500 setPollOperationsError(operations, map_buffer[0..poll_i], error.Canceled);
2528 .file_read_streaming => |*o| {2501 setOperationsError(operations[operation_index..], error.Canceled);
2529 o.result = fileReadStreaming(o.file, o.data);2502 return;
2530 },2503 },
2504 };
2505 const poll_rc = posix.system.poll(&poll_buffer, poll_i, -1);
2506 syscall.finish();
2507 switch (posix.errno(poll_rc)) {
2508 .SUCCESS => {
2509 if (poll_rc == 0) {
2510 // Spurious timeout; handle same as INTR.
2511 continue;
2512 }
2513 for (poll_buffer[0..poll_i], map_buffer[0..poll_i]) |*poll_fd, i| {
2514 if (poll_fd.revents == 0) continue;
2515 switch (operations[i]) {
2516 .noop => unreachable,
2517 .file_read_streaming => |*o| {
2518 o.result = fileReadStreaming(o.file, o.data);
2519 },
2520 }
2521 }
2522 break;
2523 },
2524 .INTR => continue,
2525 .NOMEM => {
2526 setPollOperationsError(operations, map_buffer[0..poll_i], error.SystemResources);
2527 break;
2528 },
2529 else => {
2530 setPollOperationsError(operations, map_buffer[0..poll_i], error.Unexpected);
2531 break;
2532 },
2533 }
2531 }2534 }
2532 }2535 }
2533}2536}
25342537
2535fn setAllOperationsError(2538fn setPollOperationsError(
2536 operations: []Io.Operation,2539 operations: []Io.Operation,
2537 map: []const u8,2540 map: []const u8,
2538 err: error{ Canceled, SystemResources, Unexpected },2541 err: error{ Canceled, SystemResources, Unexpected },
...@@ -2543,10 +2546,10 @@ fn setAllOperationsError(...@@ -2543,10 +2546,10 @@ fn setAllOperationsError(
2543 };2546 };
2544}2547}
25452548
2546fn setOperationsCanceled(operations: []Io.Operation) void {2549fn setOperationsError(operations: []Io.Operation, err: error{ Canceled, SystemResources, Unexpected }) void {
2547 for (operations) |*op| switch (op.*) {2550 for (operations) |*op| switch (op.*) {
2548 .noop => unreachable,2551 .noop => unreachable,
2549 inline else => |*o| o.result = error.Canceled,2552 inline else => |*o| o.result = err,
2550 };2553 };
2551}2554}
25522555