| ... | @@ -9,6 +9,9 @@ const net = std.Io.net; | ... | @@ -9,6 +9,9 @@ const net = std.Io.net; |
| 9 | const assert = std.debug.assert; | 9 | const assert = std.debug.assert; |
| 10 | const Allocator = std.mem.Allocator; | 10 | const Allocator = std.mem.Allocator; |
| 11 | const Alignment = std.mem.Alignment; | 11 | const Alignment = std.mem.Alignment; |
| | 12 | const posix = std.posix; |
| | 13 | const IpAddress = std.Io.net.IpAddress; |
| | 14 | const errnoBug = std.Io.Threaded.errnoBug; |
| 12 | | 15 | |
| 13 | /// Must be a thread-safe allocator. | 16 | /// Must be a thread-safe allocator. |
| 14 | gpa: Allocator, | 17 | gpa: Allocator, |
| ... | @@ -97,14 +100,10 @@ fn async( | ... | @@ -97,14 +100,10 @@ fn async( |
| 97 | context_alignment: std.mem.Alignment, | 100 | context_alignment: std.mem.Alignment, |
| 98 | start: *const fn (context: *const anyopaque, result: *anyopaque) void, | 101 | start: *const fn (context: *const anyopaque, result: *anyopaque) void, |
| 99 | ) ?*Io.AnyFuture { | 102 | ) ?*Io.AnyFuture { |
| 100 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 103 | return concurrent(userdata, result.len, result_alignment, context, context_alignment, start) catch { |
| 101 | _ = k; | 104 | start(context.ptr, result.ptr); |
| 102 | _ = result; | 105 | return null; |
| 103 | _ = result_alignment; | 106 | }; |
| 104 | _ = context; | | |
| 105 | _ = context_alignment; | | |
| 106 | _ = start; | | |
| 107 | @panic("TODO"); | | |
| 108 | } | 107 | } |
| 109 | | 108 | |
| 110 | fn concurrent( | 109 | fn concurrent( |
| ... | @@ -156,7 +155,7 @@ fn cancel( | ... | @@ -156,7 +155,7 @@ fn cancel( |
| 156 | fn cancelRequested(userdata: ?*anyopaque) bool { | 155 | fn cancelRequested(userdata: ?*anyopaque) bool { |
| 157 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 156 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); |
| 158 | _ = k; | 157 | _ = k; |
| 159 | @panic("TODO"); | 158 | return false; // TODO |
| 160 | } | 159 | } |
| 161 | | 160 | |
| 162 | fn groupAsync( | 161 | fn groupAsync( |
| ... | @@ -419,12 +418,23 @@ fn netAccept(userdata: ?*anyopaque, server: net.Socket.Handle) net.Server.Accept | ... | @@ -419,12 +418,23 @@ fn netAccept(userdata: ?*anyopaque, server: net.Socket.Handle) net.Server.Accept |
| 419 | _ = server; | 418 | _ = server; |
| 420 | @panic("TODO"); | 419 | @panic("TODO"); |
| 421 | } | 420 | } |
| 422 | fn netBindIp(userdata: ?*anyopaque, address: *const net.IpAddress, options: net.IpAddress.BindOptions) net.IpAddress.BindError!net.Socket { | 421 | fn netBindIp( |
| 423 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 422 | userdata: ?*anyopaque, |
| 424 | _ = k; | 423 | address: *const net.IpAddress, |
| 425 | _ = address; | 424 | options: net.IpAddress.BindOptions, |
| 426 | _ = options; | 425 | ) net.IpAddress.BindError!net.Socket { |
| 427 | @panic("TODO"); | 426 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); |
| | 427 | const family = Io.Threaded.posixAddressFamily(address); |
| | 428 | const socket_fd = try openSocketPosix(k, family, options); |
| | 429 | errdefer posix.close(socket_fd); |
| | 430 | var storage: Io.Threaded.PosixAddress = undefined; |
| | 431 | var addr_len = Io.Threaded.addressToPosix(address, &storage); |
| | 432 | try posixBind(k, socket_fd, &storage.any, addr_len); |
| | 433 | try posixGetSockName(k, socket_fd, &storage.any, &addr_len); |
| | 434 | return .{ |
| | 435 | .handle = socket_fd, |
| | 436 | .address = Io.Threaded.addressFromPosix(&storage), |
| | 437 | }; |
| 428 | } | 438 | } |
| 429 | fn netConnectIp(userdata: ?*anyopaque, address: *const net.IpAddress, options: net.IpAddress.ConnectOptions) net.IpAddress.ConnectError!net.Stream { | 439 | fn netConnectIp(userdata: ?*anyopaque, address: *const net.IpAddress, options: net.IpAddress.ConnectOptions) net.IpAddress.ConnectError!net.Stream { |
| 430 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 440 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); |
| ... | @@ -453,6 +463,7 @@ fn netConnectUnix( | ... | @@ -453,6 +463,7 @@ fn netConnectUnix( |
| 453 | _ = unix_address; | 463 | _ = unix_address; |
| 454 | @panic("TODO"); | 464 | @panic("TODO"); |
| 455 | } | 465 | } |
| | 466 | |
| 456 | fn netSend( | 467 | fn netSend( |
| 457 | userdata: ?*anyopaque, | 468 | userdata: ?*anyopaque, |
| 458 | handle: net.Socket.Handle, | 469 | handle: net.Socket.Handle, |
| ... | @@ -460,12 +471,22 @@ fn netSend( | ... | @@ -460,12 +471,22 @@ fn netSend( |
| 460 | flags: net.SendFlags, | 471 | flags: net.SendFlags, |
| 461 | ) struct { ?net.Socket.SendError, usize } { | 472 | ) struct { ?net.Socket.SendError, usize } { |
| 462 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); | 473 | const k: *Kqueue = @ptrCast(@alignCast(userdata)); |
| | 474 | |
| | 475 | const posix_flags: u32 = |
| | 476 | @as(u32, if (@hasDecl(posix.MSG, "CONFIRM") and flags.confirm) posix.MSG.CONFIRM else 0) | |
| | 477 | @as(u32, if (@hasDecl(posix.MSG, "DONTROUTE") and flags.dont_route) posix.MSG.DONTROUTE else 0) | |
| | 478 | @as(u32, if (@hasDecl(posix.MSG, "EOR") and flags.eor) posix.MSG.EOR else 0) | |
| | 479 | @as(u32, if (@hasDecl(posix.MSG, "OOB") and flags.oob) posix.MSG.OOB else 0) | |
| | 480 | @as(u32, if (@hasDecl(posix.MSG, "FASTOPEN") and flags.fastopen) posix.MSG.FASTOPEN else 0) | |
| | 481 | posix.MSG.NOSIGNAL; |
| | 482 | |
| 463 | _ = k; | 483 | _ = k; |
| | 484 | _ = posix_flags; |
| 464 | _ = handle; | 485 | _ = handle; |
| 465 | _ = outgoing_messages; | 486 | _ = outgoing_messages; |
| 466 | _ = flags; | | |
| 467 | @panic("TODO"); | 487 | @panic("TODO"); |
| 468 | } | 488 | } |
| | 489 | |
| 469 | fn netReceive( | 490 | fn netReceive( |
| 470 | userdata: ?*anyopaque, | 491 | userdata: ?*anyopaque, |
| 471 | handle: net.Socket.Handle, | 492 | handle: net.Socket.Handle, |
| ... | @@ -533,3 +554,153 @@ fn netLookup( | ... | @@ -533,3 +554,153 @@ fn netLookup( |
| 533 | _ = options; | 554 | _ = options; |
| 534 | @panic("TODO"); | 555 | @panic("TODO"); |
| 535 | } | 556 | } |
| | 557 | |
| | 558 | fn openSocketPosix( |
| | 559 | k: *Kqueue, |
| | 560 | family: posix.sa_family_t, |
| | 561 | options: IpAddress.BindOptions, |
| | 562 | ) error{ |
| | 563 | AddressFamilyUnsupported, |
| | 564 | ProtocolUnsupportedBySystem, |
| | 565 | ProcessFdQuotaExceeded, |
| | 566 | SystemFdQuotaExceeded, |
| | 567 | SystemResources, |
| | 568 | ProtocolUnsupportedByAddressFamily, |
| | 569 | SocketModeUnsupported, |
| | 570 | OptionUnsupported, |
| | 571 | Unexpected, |
| | 572 | Canceled, |
| | 573 | }!posix.socket_t { |
| | 574 | const mode = Io.Threaded.posixSocketMode(options.mode); |
| | 575 | const protocol = Io.Threaded.posixProtocol(options.protocol); |
| | 576 | const socket_fd = while (true) { |
| | 577 | try k.checkCancel(); |
| | 578 | const flags: u32 = mode | if (Io.Threaded.socket_flags_unsupported) 0 else posix.SOCK.CLOEXEC; |
| | 579 | const socket_rc = posix.system.socket(family, flags, protocol); |
| | 580 | switch (posix.errno(socket_rc)) { |
| | 581 | .SUCCESS => { |
| | 582 | const fd: posix.fd_t = @intCast(socket_rc); |
| | 583 | errdefer posix.close(fd); |
| | 584 | if (Io.Threaded.socket_flags_unsupported) { |
| | 585 | while (true) { |
| | 586 | try k.checkCancel(); |
| | 587 | switch (posix.errno(posix.system.fcntl(fd, posix.F.SETFD, @as(usize, posix.FD_CLOEXEC)))) { |
| | 588 | .SUCCESS => break, |
| | 589 | .INTR => continue, |
| | 590 | .CANCELED => return error.Canceled, |
| | 591 | else => |err| return posix.unexpectedErrno(err), |
| | 592 | } |
| | 593 | } |
| | 594 | |
| | 595 | var fl_flags: usize = while (true) { |
| | 596 | try k.checkCancel(); |
| | 597 | const rc = posix.system.fcntl(fd, posix.F.GETFL, @as(usize, 0)); |
| | 598 | switch (posix.errno(rc)) { |
| | 599 | .SUCCESS => break @intCast(rc), |
| | 600 | .INTR => continue, |
| | 601 | .CANCELED => return error.Canceled, |
| | 602 | else => |err| return posix.unexpectedErrno(err), |
| | 603 | } |
| | 604 | }; |
| | 605 | fl_flags &= ~@as(usize, 1 << @bitOffsetOf(posix.O, "NONBLOCK")); |
| | 606 | while (true) { |
| | 607 | try k.checkCancel(); |
| | 608 | switch (posix.errno(posix.system.fcntl(fd, posix.F.SETFL, fl_flags))) { |
| | 609 | .SUCCESS => break, |
| | 610 | .INTR => continue, |
| | 611 | .CANCELED => return error.Canceled, |
| | 612 | else => |err| return posix.unexpectedErrno(err), |
| | 613 | } |
| | 614 | } |
| | 615 | } |
| | 616 | break fd; |
| | 617 | }, |
| | 618 | .INTR => continue, |
| | 619 | .CANCELED => return error.Canceled, |
| | 620 | |
| | 621 | .AFNOSUPPORT => return error.AddressFamilyUnsupported, |
| | 622 | .INVAL => return error.ProtocolUnsupportedBySystem, |
| | 623 | .MFILE => return error.ProcessFdQuotaExceeded, |
| | 624 | .NFILE => return error.SystemFdQuotaExceeded, |
| | 625 | .NOBUFS => return error.SystemResources, |
| | 626 | .NOMEM => return error.SystemResources, |
| | 627 | .PROTONOSUPPORT => return error.ProtocolUnsupportedByAddressFamily, |
| | 628 | .PROTOTYPE => return error.SocketModeUnsupported, |
| | 629 | else => |err| return posix.unexpectedErrno(err), |
| | 630 | } |
| | 631 | }; |
| | 632 | errdefer posix.close(socket_fd); |
| | 633 | |
| | 634 | if (options.ip6_only) { |
| | 635 | if (posix.IPV6 == void) return error.OptionUnsupported; |
| | 636 | try setSocketOption(k, socket_fd, posix.IPPROTO.IPV6, posix.IPV6.V6ONLY, 0); |
| | 637 | } |
| | 638 | |
| | 639 | return socket_fd; |
| | 640 | } |
| | 641 | |
| | 642 | fn posixBind( |
| | 643 | k: *Kqueue, |
| | 644 | socket_fd: posix.socket_t, |
| | 645 | addr: *const posix.sockaddr, |
| | 646 | addr_len: posix.socklen_t, |
| | 647 | ) !void { |
| | 648 | while (true) { |
| | 649 | try k.checkCancel(); |
| | 650 | switch (posix.errno(posix.system.bind(socket_fd, addr, addr_len))) { |
| | 651 | .SUCCESS => break, |
| | 652 | .INTR => continue, |
| | 653 | .CANCELED => return error.Canceled, |
| | 654 | |
| | 655 | .ADDRINUSE => return error.AddressInUse, |
| | 656 | .BADF => |err| return errnoBug(err), // File descriptor used after closed. |
| | 657 | .INVAL => |err| return errnoBug(err), // invalid parameters |
| | 658 | .NOTSOCK => |err| return errnoBug(err), // invalid `sockfd` |
| | 659 | .AFNOSUPPORT => return error.AddressFamilyUnsupported, |
| | 660 | .ADDRNOTAVAIL => return error.AddressUnavailable, |
| | 661 | .FAULT => |err| return errnoBug(err), // invalid `addr` pointer |
| | 662 | .NOMEM => return error.SystemResources, |
| | 663 | else => |err| return posix.unexpectedErrno(err), |
| | 664 | } |
| | 665 | } |
| | 666 | } |
| | 667 | |
| | 668 | fn posixGetSockName(k: *Kqueue, socket_fd: posix.fd_t, addr: *posix.sockaddr, addr_len: *posix.socklen_t) !void { |
| | 669 | while (true) { |
| | 670 | try k.checkCancel(); |
| | 671 | switch (posix.errno(posix.system.getsockname(socket_fd, addr, addr_len))) { |
| | 672 | .SUCCESS => break, |
| | 673 | .INTR => continue, |
| | 674 | .CANCELED => return error.Canceled, |
| | 675 | |
| | 676 | .BADF => |err| return errnoBug(err), // File descriptor used after closed. |
| | 677 | .FAULT => |err| return errnoBug(err), |
| | 678 | .INVAL => |err| return errnoBug(err), // invalid parameters |
| | 679 | .NOTSOCK => |err| return errnoBug(err), // always a race condition |
| | 680 | .NOBUFS => return error.SystemResources, |
| | 681 | else => |err| return posix.unexpectedErrno(err), |
| | 682 | } |
| | 683 | } |
| | 684 | } |
| | 685 | |
| | 686 | fn setSocketOption(k: *Kqueue, fd: posix.fd_t, level: i32, opt_name: u32, option: u32) !void { |
| | 687 | const o: []const u8 = @ptrCast(&option); |
| | 688 | while (true) { |
| | 689 | try k.checkCancel(); |
| | 690 | switch (posix.errno(posix.system.setsockopt(fd, level, opt_name, o.ptr, @intCast(o.len)))) { |
| | 691 | .SUCCESS => return, |
| | 692 | .INTR => continue, |
| | 693 | .CANCELED => return error.Canceled, |
| | 694 | |
| | 695 | .BADF => |err| return errnoBug(err), // File descriptor used after closed. |
| | 696 | .NOTSOCK => |err| return errnoBug(err), |
| | 697 | .INVAL => |err| return errnoBug(err), |
| | 698 | .FAULT => |err| return errnoBug(err), |
| | 699 | else => |err| return posix.unexpectedErrno(err), |
| | 700 | } |
| | 701 | } |
| | 702 | } |
| | 703 | |
| | 704 | fn checkCancel(k: *Kqueue) error{Canceled}!void { |
| | 705 | if (cancelRequested(k)) return error.Canceled; |
| | 706 | } |