| ... | @@ -95,29 +95,56 @@ pub const TcpServer = struct { | ... | @@ -95,29 +95,56 @@ pub const TcpServer = struct { |
| 95 | | 95 | |
| 96 | pub const Loop = struct { | 96 | pub const Loop = struct { |
| 97 | allocator: *mem.Allocator, | 97 | allocator: *mem.Allocator, |
| 98 | epollfd: i32, | | |
| 99 | keep_running: bool, | 98 | keep_running: bool, |
| 100 | next_tick_queue: std.atomic.QueueMpsc(promise), | 99 | next_tick_queue: std.atomic.QueueMpsc(promise), |
| | 100 | os_data: OsData, |
| | 101 | |
| | 102 | const OsData = switch (builtin.os) { |
| | 103 | builtin.Os.linux => struct { |
| | 104 | epollfd: i32, |
| | 105 | }, |
| | 106 | else => struct {}, |
| | 107 | }; |
| 101 | | 108 | |
| 102 | pub const NextTickNode = std.atomic.QueueMpsc(promise).Node; | 109 | pub const NextTickNode = std.atomic.QueueMpsc(promise).Node; |
| 103 | | 110 | |
| 104 | /// The allocator must be thread-safe because we use it for multiplexing | 111 | /// The allocator must be thread-safe because we use it for multiplexing |
| 105 | /// coroutines onto kernel threads. | 112 | /// coroutines onto kernel threads. |
| 106 | pub fn init(allocator: *mem.Allocator) !Loop { | 113 | pub fn init(allocator: *mem.Allocator) !Loop { |
| 107 | const epollfd = try std.os.linuxEpollCreate(std.os.linux.EPOLL_CLOEXEC); | 114 | var self = Loop{ |
| 108 | errdefer std.os.close(epollfd); | | |
| 109 | | | |
| 110 | return Loop{ | | |
| 111 | .keep_running = true, | 115 | .keep_running = true, |
| 112 | .allocator = allocator, | 116 | .allocator = allocator, |
| 113 | .epollfd = epollfd, | 117 | .os_data = undefined, |
| 114 | .next_tick_queue = std.atomic.QueueMpsc(promise).init(), | 118 | .next_tick_queue = std.atomic.QueueMpsc(promise).init(), |
| 115 | }; | 119 | }; |
| | 120 | try self.initOsData(); |
| | 121 | errdefer self.deinitOsData(); |
| | 122 | |
| | 123 | return self; |
| 116 | } | 124 | } |
| 117 | | 125 | |
| 118 | /// must call stop before deinit | 126 | /// must call stop before deinit |
| 119 | pub fn deinit(self: *Loop) void { | 127 | pub fn deinit(self: *Loop) void { |
| 120 | std.os.close(self.epollfd); | 128 | self.deinitOsData(); |
| | 129 | } |
| | 130 | |
| | 131 | const InitOsDataError = std.os.LinuxEpollCreateError; |
| | 132 | |
| | 133 | fn initOsData(self: *Loop) InitOsDataError!void { |
| | 134 | switch (builtin.os) { |
| | 135 | builtin.Os.linux => { |
| | 136 | self.os_data.epollfd = try std.os.linuxEpollCreate(std.os.linux.EPOLL_CLOEXEC); |
| | 137 | errdefer std.os.close(self.os_data.epollfd); |
| | 138 | }, |
| | 139 | else => {}, |
| | 140 | } |
| | 141 | } |
| | 142 | |
| | 143 | fn deinitOsData(self: *Loop) void { |
| | 144 | switch (builtin.os) { |
| | 145 | builtin.Os.linux => std.os.close(self.os_data.epollfd), |
| | 146 | else => {}, |
| | 147 | } |
| 121 | } | 148 | } |
| 122 | | 149 | |
| 123 | pub fn addFd(self: *Loop, fd: i32, prom: promise) !void { | 150 | pub fn addFd(self: *Loop, fd: i32, prom: promise) !void { |
| ... | @@ -125,11 +152,11 @@ pub const Loop = struct { | ... | @@ -125,11 +152,11 @@ pub const Loop = struct { |
| 125 | .events = std.os.linux.EPOLLIN | std.os.linux.EPOLLOUT | std.os.linux.EPOLLET, | 152 | .events = std.os.linux.EPOLLIN | std.os.linux.EPOLLOUT | std.os.linux.EPOLLET, |
| 126 | .data = std.os.linux.epoll_data{ .ptr = @ptrToInt(prom) }, | 153 | .data = std.os.linux.epoll_data{ .ptr = @ptrToInt(prom) }, |
| 127 | }; | 154 | }; |
| 128 | try std.os.linuxEpollCtl(self.epollfd, std.os.linux.EPOLL_CTL_ADD, fd, &ev); | 155 | try std.os.linuxEpollCtl(self.os_data.epollfd, std.os.linux.EPOLL_CTL_ADD, fd, &ev); |
| 129 | } | 156 | } |
| 130 | | 157 | |
| 131 | pub fn removeFd(self: *Loop, fd: i32) void { | 158 | pub fn removeFd(self: *Loop, fd: i32) void { |
| 132 | std.os.linuxEpollCtl(self.epollfd, std.os.linux.EPOLL_CTL_DEL, fd, undefined) catch {}; | 159 | std.os.linuxEpollCtl(self.os_data.epollfd, std.os.linux.EPOLL_CTL_DEL, fd, undefined) catch {}; |
| 133 | } | 160 | } |
| 134 | async fn waitFd(self: *Loop, fd: i32) !void { | 161 | async fn waitFd(self: *Loop, fd: i32) !void { |
| 135 | defer self.removeFd(fd); | 162 | defer self.removeFd(fd); |
| ... | @@ -156,12 +183,22 @@ pub const Loop = struct { | ... | @@ -156,12 +183,22 @@ pub const Loop = struct { |
| 156 | resume node.data; | 183 | resume node.data; |
| 157 | } | 184 | } |
| 158 | if (!self.keep_running) break; | 185 | if (!self.keep_running) break; |
| 159 | var events: [16]std.os.linux.epoll_event = undefined; | 186 | |
| 160 | const count = std.os.linuxEpollWait(self.epollfd, events[0..], -1); | 187 | self.dispatchOsEvents(); |
| 161 | for (events[0..count]) |ev| { | 188 | } |
| 162 | const p = @intToPtr(promise, ev.data.ptr); | 189 | } |
| 163 | resume p; | 190 | |
| 164 | } | 191 | fn dispatchOsEvents(self: *Loop) void { |
| | 192 | switch (builtin.os) { |
| | 193 | builtin.Os.linux => { |
| | 194 | var events: [16]std.os.linux.epoll_event = undefined; |
| | 195 | const count = std.os.linuxEpollWait(self.os_data.epollfd, events[0..], -1); |
| | 196 | for (events[0..count]) |ev| { |
| | 197 | const p = @intToPtr(promise, ev.data.ptr); |
| | 198 | resume p; |
| | 199 | } |
| | 200 | }, |
| | 201 | else => {}, |
| 165 | } | 202 | } |
| 166 | } | 203 | } |
| 167 | }; | 204 | }; |