| ... | @@ -24,6 +24,7 @@ const idle_stack_size = 256 * 1024; | ... | @@ -24,6 +24,7 @@ const idle_stack_size = 256 * 1024; |
| 24 | | 24 | |
| 25 | const max_idle_search = 4; | 25 | const max_idle_search = 4; |
| 26 | const max_steal_ready_search = 4; | 26 | const max_steal_ready_search = 4; |
| | 27 | const max_iovecs_len = 8; |
| 27 | | 28 | |
| 28 | const changes_buffer_len = 64; | 29 | const changes_buffer_len = 64; |
| 29 | | 30 | |
| ... | @@ -1431,13 +1432,47 @@ fn netReceive( | ... | @@ -1431,13 +1432,47 @@ fn netReceive( |
| 1431 | _ = timeout; | 1432 | _ = timeout; |
| 1432 | @panic("TODO"); | 1433 | @panic("TODO"); |
| 1433 | } | 1434 | } |
| 1434 | fn netRead(userdata: ?*anyopaque, src: net.Socket.Handle, data: [][]u8) net.Stream.Reader.Error!usize { | 1435 | |
| | 1436 | fn netRead(userdata: ?*anyopaque, fd: net.Socket.Handle, data: [][]u8) net.Stream.Reader.Error!usize { |
| 1435 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 1437 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); |
| 1436 | _ = k; | 1438 | |
| 1437 | _ = src; | 1439 | var iovecs_buffer: [max_iovecs_len]posix.iovec = undefined; |
| 1438 | _ = data; | 1440 | var i: usize = 0; |
| 1439 | @panic("TODO"); | 1441 | for (data) |buf| { |
| | 1442 | if (iovecs_buffer.len - i == 0) break; |
| | 1443 | if (buf.len != 0) { |
| | 1444 | iovecs_buffer[i] = .{ .base = buf.ptr, .len = buf.len }; |
| | 1445 | i += 1; |
| | 1446 | } |
| | 1447 | } |
| | 1448 | const dest = iovecs_buffer[0..i]; |
| | 1449 | assert(dest[0].len > 0); |
| | 1450 | |
| | 1451 | while (true) { |
| | 1452 | try k.checkCancel(); |
| | 1453 | std.debug.print("calling readv\n", .{}); |
| | 1454 | const rc = posix.system.readv(fd, dest.ptr, @intCast(dest.len)); |
| | 1455 | switch (posix.errno(rc)) { |
| | 1456 | .SUCCESS => return @intCast(rc), |
| | 1457 | .INTR => continue, |
| | 1458 | .CANCELED => return error.Canceled, |
| | 1459 | .AGAIN => @panic("TODO"), |
| | 1460 | |
| | 1461 | .INVAL => |err| return errnoBug(err), |
| | 1462 | .FAULT => |err| return errnoBug(err), |
| | 1463 | .BADF => |err| return errnoBug(err), // File descriptor used after closed. |
| | 1464 | .NOBUFS => return error.SystemResources, |
| | 1465 | .NOMEM => return error.SystemResources, |
| | 1466 | .NOTCONN => return error.SocketUnconnected, |
| | 1467 | .CONNRESET => return error.ConnectionResetByPeer, |
| | 1468 | .TIMEDOUT => return error.Timeout, |
| | 1469 | .PIPE => return error.SocketUnconnected, |
| | 1470 | .NETDOWN => return error.NetworkDown, |
| | 1471 | else => |err| return posix.unexpectedErrno(err), |
| | 1472 | } |
| | 1473 | } |
| 1440 | } | 1474 | } |
| | 1475 | |
| 1441 | fn netWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize { | 1476 | fn netWrite(userdata: ?*anyopaque, dest: net.Socket.Handle, header: []const u8, data: []const []const u8, splat: usize) net.Stream.Writer.Error!usize { |
| 1442 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 1477 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); |
| 1443 | _ = k; | 1478 | _ = k; |
| ... | @@ -1529,7 +1564,7 @@ fn openSocketPosix( | ... | @@ -1529,7 +1564,7 @@ fn openSocketPosix( |
| 1529 | else => |err| return posix.unexpectedErrno(err), | 1564 | else => |err| return posix.unexpectedErrno(err), |
| 1530 | } | 1565 | } |
| 1531 | }; | 1566 | }; |
| 1532 | fl_flags &= ~@as(usize, 1 << @bitOffsetOf(posix.O, "NONBLOCK")); | 1567 | fl_flags |= @as(usize, 1 << @bitOffsetOf(posix.O, "NONBLOCK")); |
| 1533 | while (true) { | 1568 | while (true) { |
| 1534 | try k.checkCancel(); | 1569 | try k.checkCancel(); |
| 1535 | switch (posix.errno(posix.system.fcntl(fd, posix.F.SETFL, fl_flags))) { | 1570 | switch (posix.errno(posix.system.fcntl(fd, posix.F.SETFL, fl_flags))) { |