authorgravatar for igor.anic@gmail.comIgor Anić <igor.anic@gmail.com> 2024-02-21 20:01:45+01:00
committergravatar for igor.anic@gmail.comIgor Anić <igor.anic@gmail.com> 2024-02-21 20:01:45+01:00
logd9950298444c3a3c9d2e5ec7efbf45e722bbed02
tree7b59aa80aa28946b42b3f223b4efe7b7898696c3
parent955fd65cb1705d8279eb195bdbc69810df1b1d98

add BufferedTee

BufferedTee provides reader interface to the consumer. Data read by consumer is also written to the output. Output is hold lookahead_size bytes behind consumer. Allowing consumer to put back some bytes to be read again. On flush all consumed bytes are flushed to the output. input -> tee -> consumer | output input - underlying unbuffered reader output - writer, receives data read by consumer consumer - uses provided reader interface If lookahead_size is zero output always has same bytes as consumer.

2 files changed, 394 insertions(+), 0 deletions(-)

lib/std/io.zig+4
...@@ -411,6 +411,9 @@ pub const BufferedAtomicFile = @import("io/buffered_atomic_file.zig").BufferedAt...@@ -411,6 +411,9 @@ pub const BufferedAtomicFile = @import("io/buffered_atomic_file.zig").BufferedAt
411411
412pub const StreamSource = @import("io/stream_source.zig").StreamSource;412pub const StreamSource = @import("io/stream_source.zig").StreamSource;
413413
414pub const BufferedTee = @import("io/buffered_tee.zig").BufferedTee;
415pub const bufferedTee = @import("io/buffered_tee.zig").bufferedTee;
416
414pub const tty = @import("io/tty.zig");417pub const tty = @import("io/tty.zig");
415418
416/// A Writer that doesn't write to anything.419/// A Writer that doesn't write to anything.
...@@ -692,4 +695,5 @@ test {...@@ -692,4 +695,5 @@ test {
692 _ = @import("io/seekable_stream.zig");695 _ = @import("io/seekable_stream.zig");
693 _ = @import("io/stream_source.zig");696 _ = @import("io/stream_source.zig");
694 _ = @import("io/test.zig");697 _ = @import("io/test.zig");
698 _ = @import("io/buffered_tee.zig");
695}699}
lib/std/io/buffered_tee.zig created+390
...@@ -0,0 +1,390 @@
1//! BufferedTee provides reader interface to the consumer. Data read by consumer
2//! is also written to the output. Output is hold lookahead_size bytes behind
3//! consumer. Allowing consumer to put back some bytes to be read again. On flush
4//! all consumed bytes are flushed to the output.
5//!
6//! input -> tee -> consumer
7//! |
8//! output
9//!
10//! input - underlying unbuffered reader
11//! output - writer, receives data read by consumer
12//! consumer - uses provided reader interface
13//!
14//! If lookahead_size is zero output always has same bytes as consumer.
15//!
16
17const std = @import("std");
18const io = std.io;
19const assert = std.debug.assert;
20const testing = std.testing;
21
22pub fn BufferedTee(
23 comptime buffer_size: usize, // internal buffer size in bytes
24 comptime lookahead_size: usize, // lookahead, number of bytes to hold output behind consumer
25 comptime InputReaderType: type,
26 comptime OutputWriterType: type,
27) type {
28 comptime assert(buffer_size > lookahead_size);
29
30 return struct {
31 input: InputReaderType,
32 output: OutputWriterType,
33
34 buf: [buffer_size]u8 = undefined, // internal buffer
35 tail: usize = 0, // buffer is filled up to this position with bytes from input
36 rp: usize = 0, // reader pointer; consumer has read up to this position
37 wp: usize = 0, // writer pointer; data is sent to the output up to this position
38
39 pub const Error = InputReaderType.Error || OutputWriterType.Error;
40 pub const Reader = io.Reader(*Self, Error, read);
41
42 const Self = @This();
43
44 pub fn read(self: *Self, dest: []u8) Error!usize {
45 var dest_index: usize = 0;
46
47 while (dest_index < dest.len) {
48 const written = @min(dest.len - dest_index, self.tail - self.rp);
49 if (written == 0) {
50 try self.preserveLookahead();
51 // fill upper part of the buf
52 const n = try self.input.read(self.buf[self.tail..]);
53 if (n == 0) {
54 // reading from the unbuffered stream returned nothing
55 // so we have nothing left to read.
56 return dest_index;
57 }
58 self.tail += n;
59 } else {
60 @memcpy(dest[dest_index..][0..written], self.buf[self.rp..][0..written]);
61 self.rp += written;
62 dest_index += written;
63 try self.flush_(lookahead_size);
64 }
65 }
66 return dest.len;
67 }
68
69 /// Move lookahead_size bytes to the buffer start.
70 fn preserveLookahead(self: *Self) !void {
71 assert(self.tail == self.rp);
72 if (lookahead_size == 0) {
73 // Flush is called on each read so wp must follow rp when lookahead_size == 0.
74 assert(self.wp == self.rp);
75 // Nothing to preserve rewind pointer to the buffer start
76 self.rp = 0;
77 self.wp = 0;
78 self.tail = 0;
79 return;
80 }
81 if (self.tail <= lookahead_size) {
82 // There is still palce in the buffer, append to buffer from tail position.
83 return;
84 }
85 try self.flush_(lookahead_size);
86 const head = self.tail - lookahead_size;
87 // Preserve head..tail at the start of the buffer.
88 std.mem.copyForwards(u8, self.buf[0..lookahead_size], self.buf[head..self.tail]);
89 self.wp -= head;
90 assert(self.wp <= lookahead_size);
91 self.rp = lookahead_size;
92 self.tail = lookahead_size;
93 }
94
95 /// Flush to the output all but lookahead size bytes.
96 fn flush_(self: *Self, lookahead: usize) !void {
97 if (self.rp <= self.wp + lookahead) return;
98 const new_wp = self.rp - lookahead;
99 try self.output.writeAll(self.buf[self.wp..new_wp]);
100 self.wp = new_wp;
101 }
102
103 /// Flush to the output all consumed bytes.
104 pub fn flush(self: *Self) !void {
105 try self.flush_(0);
106 }
107
108 /// Put back some bytes to be consumed again. Usefull when we overshoot
109 /// reading and want to return that overshoot bytes. Can return maximum
110 /// of lookahead_size number of bytes.
111 pub fn putBack(self: *Self, n: usize) void {
112 assert(n <= lookahead_size and n <= self.rp);
113 self.rp -= n;
114 }
115
116 pub fn reader(self: *Self) Reader {
117 return .{ .context = self };
118 }
119 };
120}
121
122pub fn bufferedTee(
123 comptime buffer_size: usize,
124 comptime lookahead_size: usize,
125 input: anytype,
126 output: anytype,
127) BufferedTee(
128 buffer_size,
129 lookahead_size,
130 @TypeOf(input),
131 @TypeOf(output),
132) {
133 return BufferedTee(
134 buffer_size,
135 lookahead_size,
136 @TypeOf(input),
137 @TypeOf(output),
138 ){
139 .input = input,
140 .output = output,
141 };
142}
143
144// Running test from std.io.BufferedReader on BufferedTee
145// It should act as BufferedReader for consumer.
146
147fn BufferedReader(comptime buffer_size: usize, comptime ReaderType: type) type {
148 return BufferedTee(buffer_size, 0, ReaderType, @TypeOf(io.null_writer));
149}
150
151fn bufferedReader(reader: anytype) BufferedReader(4096, @TypeOf(reader)) {
152 return .{
153 .input = reader,
154 .output = io.null_writer,
155 };
156}
157
158test "io.BufferedTee io.BufferedReader OneByte" {
159 const OneByteReadReader = struct {
160 str: []const u8,
161 curr: usize,
162
163 const Error = error{NoError};
164 const Self = @This();
165 const Reader = io.Reader(*Self, Error, read);
166
167 fn init(str: []const u8) Self {
168 return Self{
169 .str = str,
170 .curr = 0,
171 };
172 }
173
174 fn read(self: *Self, dest: []u8) Error!usize {
175 if (self.str.len <= self.curr or dest.len == 0)
176 return 0;
177
178 dest[0] = self.str[self.curr];
179 self.curr += 1;
180 return 1;
181 }
182
183 fn reader(self: *Self) Reader {
184 return .{ .context = self };
185 }
186 };
187
188 const str = "This is a test";
189 var one_byte_stream = OneByteReadReader.init(str);
190 var buf_reader = bufferedReader(one_byte_stream.reader());
191 const stream = buf_reader.reader();
192
193 const res = try stream.readAllAlloc(testing.allocator, str.len + 1);
194 defer testing.allocator.free(res);
195 try testing.expectEqualSlices(u8, str, res);
196}
197
198test "io.BufferedTee io.BufferedReader Block" {
199 const BlockReader = struct {
200 block: []const u8,
201 reads_allowed: usize,
202 curr_read: usize,
203
204 const Error = error{NoError};
205 const Self = @This();
206 const Reader = io.Reader(*Self, Error, read);
207
208 fn init(block: []const u8, reads_allowed: usize) Self {
209 return Self{
210 .block = block,
211 .reads_allowed = reads_allowed,
212 .curr_read = 0,
213 };
214 }
215
216 fn read(self: *Self, dest: []u8) Error!usize {
217 if (self.curr_read >= self.reads_allowed) return 0;
218 @memcpy(dest[0..self.block.len], self.block);
219
220 self.curr_read += 1;
221 return self.block.len;
222 }
223
224 fn reader(self: *Self) Reader {
225 return .{ .context = self };
226 }
227 };
228
229 const block = "0123";
230
231 // len out == block
232 {
233 var test_buf_reader: BufferedReader(4, BlockReader) = .{
234 .input = BlockReader.init(block, 2),
235 .output = io.null_writer,
236 };
237 var out_buf: [4]u8 = undefined;
238 _ = try test_buf_reader.read(&out_buf);
239 try testing.expectEqualSlices(u8, &out_buf, block);
240 _ = try test_buf_reader.read(&out_buf);
241 try testing.expectEqualSlices(u8, &out_buf, block);
242 try testing.expectEqual(try test_buf_reader.read(&out_buf), 0);
243 }
244
245 // len out < block
246 {
247 var test_buf_reader: BufferedReader(4, BlockReader) = .{
248 .input = BlockReader.init(block, 2),
249 .output = io.null_writer,
250 };
251 var out_buf: [3]u8 = undefined;
252 _ = try test_buf_reader.read(&out_buf);
253 try testing.expectEqualSlices(u8, &out_buf, "012");
254 _ = try test_buf_reader.read(&out_buf);
255 try testing.expectEqualSlices(u8, &out_buf, "301");
256 const n = try test_buf_reader.read(&out_buf);
257 try testing.expectEqualSlices(u8, out_buf[0..n], "23");
258 try testing.expectEqual(try test_buf_reader.read(&out_buf), 0);
259 }
260
261 // len out > block
262 {
263 var test_buf_reader: BufferedReader(4, BlockReader) = .{
264 .input = BlockReader.init(block, 2),
265 .output = io.null_writer,
266 };
267 var out_buf: [5]u8 = undefined;
268 _ = try test_buf_reader.read(&out_buf);
269 try testing.expectEqualSlices(u8, &out_buf, "01230");
270 const n = try test_buf_reader.read(&out_buf);
271 try testing.expectEqualSlices(u8, out_buf[0..n], "123");
272 try testing.expectEqual(try test_buf_reader.read(&out_buf), 0);
273 }
274
275 // len out == 0
276 {
277 var test_buf_reader: BufferedReader(4, BlockReader) = .{
278 .input = BlockReader.init(block, 2),
279 .output = io.null_writer,
280 };
281 var out_buf: [0]u8 = undefined;
282 _ = try test_buf_reader.read(&out_buf);
283 try testing.expectEqualSlices(u8, &out_buf, "");
284 }
285
286 // len bufreader buf > block
287 {
288 var test_buf_reader: BufferedReader(5, BlockReader) = .{
289 .input = BlockReader.init(block, 2),
290 .output = io.null_writer,
291 };
292 var out_buf: [4]u8 = undefined;
293 _ = try test_buf_reader.read(&out_buf);
294 try testing.expectEqualSlices(u8, &out_buf, block);
295 _ = try test_buf_reader.read(&out_buf);
296 try testing.expectEqualSlices(u8, &out_buf, block);
297 try testing.expectEqual(try test_buf_reader.read(&out_buf), 0);
298 }
299}
300
301test "io.BufferedTee with zero lookahead" {
302 // output is has same bytes as reader
303 const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9 } ** 12;
304 var in = io.fixedBufferStream(&data);
305 var out = std.ArrayList(u8).init(testing.allocator);
306 defer out.deinit();
307
308 var lbr = bufferedTee(8, 0, in.reader(), out.writer());
309
310 var buf: [16]u8 = undefined;
311
312 var read_len: usize = 0;
313 for (0..buf.len) |i| {
314 const n = try lbr.read(buf[0..i]);
315 try testing.expectEqual(i, n);
316 read_len += i;
317 try testing.expectEqual(read_len, out.items.len);
318 }
319}
320
321test "io.BufferedTee with lookahead" {
322 // output is lookahead bytes behind reader
323 inline for (1..8) |lookahead| {
324 const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9 } ** 12;
325 var in = io.fixedBufferStream(&data);
326 var out = std.ArrayList(u8).init(testing.allocator);
327 defer out.deinit();
328
329 var lbr = bufferedTee(8, lookahead, in.reader(), out.writer());
330 var buf: [16]u8 = undefined;
331
332 var read_len: usize = 0;
333 for (1..buf.len) |i| {
334 const n = try lbr.read(buf[0..i]);
335 try testing.expectEqual(i, n);
336 read_len += i;
337 const out_len = if (read_len < lookahead) 0 else read_len - lookahead;
338 try testing.expectEqual(out_len, out.items.len);
339 // std.debug.print("{d} {d} {d}\n", .{ lookahead, read_len, out_len });
340 }
341 try testing.expectEqual(read_len, out.items.len + lookahead);
342 try lbr.flush();
343 try testing.expectEqual(read_len, out.items.len);
344 }
345}
346
347test "io.BufferedTee internal state" {
348 const data = [_]u8{ 0, 1, 2, 3, 4, 5, 6, 7, 8, 9 } ** 10;
349 var in = io.fixedBufferStream(&data);
350 var out = std.ArrayList(u8).init(testing.allocator);
351 defer out.deinit();
352
353 var lbr = bufferedTee(8, 4, in.reader(), out.writer());
354
355 var buf: [16]u8 = undefined;
356 var n = try lbr.read(buf[0..3]);
357 try testing.expectEqual(3, n);
358 try testing.expectEqualSlices(u8, data[0..3], buf[0..n]);
359 try testing.expectEqual(8, lbr.tail);
360 try testing.expectEqual(3, lbr.rp);
361 try testing.expectEqual(0, out.items.len);
362
363 n = try lbr.read(buf[0..6]);
364 try testing.expectEqual(6, n);
365 try testing.expectEqualSlices(u8, data[3..9], buf[0..n]);
366 try testing.expectEqual(8, lbr.tail);
367 try testing.expectEqual(5, lbr.rp);
368 try testing.expectEqualSlices(u8, data[4..12], &lbr.buf);
369 try testing.expectEqual(5, out.items.len);
370
371 n = try lbr.read(buf[0..9]);
372 try testing.expectEqual(9, n);
373 try testing.expectEqualSlices(u8, data[9..18], buf[0..n]);
374 try testing.expectEqual(8, lbr.tail);
375 try testing.expectEqual(6, lbr.rp);
376 try testing.expectEqualSlices(u8, data[12..20], &lbr.buf);
377 try testing.expectEqual(14, out.items.len);
378
379 try lbr.flush();
380 try testing.expectEqual(18, out.items.len);
381
382 lbr.putBack(4);
383 n = try lbr.read(buf[0..4]);
384 try testing.expectEqual(4, n);
385 try testing.expectEqualSlices(u8, data[14..18], buf[0..n]);
386
387 try testing.expectEqual(18, out.items.len);
388 try lbr.flush();
389 try testing.expectEqual(18, out.items.len);
390}