| ... | ... | @@ -54,6 +54,10 @@ pub fn Channel(comptime T: type) type { |
| 54 | 54 | /// For a zero length buffer, use `[0]T{}`. |
| 55 | 55 | /// TODO https://github.com/ziglang/zig/issues/2765 |
| 56 | 56 | pub fn init(self: *SelfChannel, buffer: []T) void { |
| 57 | // The ring buffer implementation only works with power of 2 buffer sizes |
| 58 | // because of relying on subtracting across zero. For example (0 -% 1) % 10 == 5 |
| 59 | assert(buffer.len == 0 or @popCount(usize, buffer.len) == 1); |
| 60 | |
| 57 | 61 | self.* = SelfChannel{ |
| 58 | 62 | .buffer_len = 0, |
| 59 | 63 | .buffer_nodes = buffer, |
| ... | ... | @@ -184,11 +188,11 @@ pub fn Channel(comptime T: type) type { |
| 184 | 188 | const get_node = &self.getters.get().?.data; |
| 185 | 189 | switch (get_node.data) { |
| 186 | 190 | GetNode.Data.Normal => |info| { |
| 187 | | info.ptr.* = self.buffer_nodes[self.buffer_index -% self.buffer_len]; |
| 191 | info.ptr.* = self.buffer_nodes[(self.buffer_index -% self.buffer_len) % self.buffer_nodes.len]; |
| 188 | 192 | }, |
| 189 | 193 | GetNode.Data.OrNull => |info| { |
| 190 | 194 | _ = self.or_null_queue.remove(info.or_null); |
| 191 | | info.ptr.* = self.buffer_nodes[self.buffer_index -% self.buffer_len]; |
| 195 | info.ptr.* = self.buffer_nodes[(self.buffer_index -% self.buffer_len) % self.buffer_nodes.len]; |
| 192 | 196 | }, |
| 193 | 197 | } |
| 194 | 198 | global_event_loop.onNextTick(get_node.tick_node); |
| ... | ... | @@ -222,7 +226,7 @@ pub fn Channel(comptime T: type) type { |
| 222 | 226 | while (self.buffer_len != self.buffer_nodes.len and put_count != 0) { |
| 223 | 227 | const put_node = &self.putters.get().?.data; |
| 224 | 228 | |
| 225 | | self.buffer_nodes[self.buffer_index] = put_node.data; |
| 229 | self.buffer_nodes[self.buffer_index % self.buffer_nodes.len] = put_node.data; |
| 226 | 230 | global_event_loop.onNextTick(put_node.tick_node); |
| 227 | 231 | self.buffer_index +%= 1; |
| 228 | 232 | self.buffer_len += 1; |
| ... | ... | @@ -283,6 +287,29 @@ test "std.event.Channel" { |
| 283 | 287 | await putter; |
| 284 | 288 | } |
| 285 | 289 | |
| 290 | test "std.event.Channel wraparound" { |
| 291 | |
| 292 | // TODO provide a way to run tests in evented I/O mode |
| 293 | if (!std.io.is_async) return error.SkipZigTest; |
| 294 | |
| 295 | const channel_size = 2; |
| 296 | |
| 297 | var buf : [channel_size]i32 = undefined; |
| 298 | var channel: Channel(i32) = undefined; |
| 299 | channel.init(&buf); |
| 300 | defer channel.deinit(); |
| 301 | |
| 302 | // add items to channel and pull them out until |
| 303 | // the buffer wraps around, make sure it doesn't crash. |
| 304 | var result : i32 = undefined; |
| 305 | channel.put(5); |
| 306 | testing.expectEqual(@as(i32, 5), channel.get()); |
| 307 | channel.put(6); |
| 308 | testing.expectEqual(@as(i32, 6), channel.get()); |
| 309 | channel.put(7); |
| 310 | testing.expectEqual(@as(i32, 7), channel.get()); |
| 311 | } |
| 312 | |
| 286 | 313 | async fn testChannelGetter(channel: *Channel(i32)) void { |
| 287 | 314 | const value1 = channel.get(); |
| 288 | 315 | testing.expect(value1 == 1234); |