authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-01-13 18:42:00-08:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2026-01-30 12:10:01-08:00
log0cca18e43c2685cfb4a72b3b265b68af1e85ca17
treeff24d7ab8d3e9f196e6ba91e0b19c56d12e3755e
parent1a168f08b51572ab136e2a8fcda35ce3c68c89bb

std: update rest of build runner to new File.MultiReader API


2 files changed, 109 insertions(+), 60 deletions(-)

lib/std/Build/Step/Run.zig+63-47
...@@ -1669,39 +1669,44 @@ fn evalZigTest(...@@ -1669,39 +1669,44 @@ fn evalZigTest(
16691669
1670 while (true) {1670 while (true) {
1671 var child = try process.spawn(io, spawn_options);1671 var child = try process.spawn(io, spawn_options);
1672 var poller = std.Io.poll(gpa, StdioPollEnum, .{1672 var multi_reader_buffer: Io.File.MultiReader.Buffer(2) = undefined;
1673 .stdout = child.stdout.?,1673 var multi_reader: Io.File.MultiReader = undefined;
1674 .stderr = child.stderr.?,1674 multi_reader.init(gpa, io, multi_reader_buffer.toStreams(), &.{ child.stdout.?, child.stderr.? });
1675 });
1676 var child_killed = false;1675 var child_killed = false;
1677 defer if (!child_killed) {1676 defer if (!child_killed) {
1678 child.kill(io);1677 child.kill(io);
1679 poller.deinit();1678 multi_reader.deinit();
1680 run.step.result_peak_rss = @max(1679 run.step.result_peak_rss = @max(
1681 run.step.result_peak_rss,1680 run.step.result_peak_rss,
1682 child.resource_usage_statistics.getMaxRss() orelse 0,1681 child.resource_usage_statistics.getMaxRss() orelse 0,
1683 );1682 );
1684 };1683 };
16851684
1686 switch (try pollZigTest(1685 switch (try waitZigTest(
1687 run,1686 run,
1688 &child,1687 &child,
1689 options,1688 options,
1690 fuzz_context,1689 fuzz_context,
1691 &poller,1690 &multi_reader,
1692 &test_metadata,1691 &test_metadata,
1693 &test_results,1692 &test_results,
1694 )) {1693 )) {
1695 .write_failed => |err| {1694 .write_failed => |err| {
1696 // The runner unexpectedly closed a stdio pipe, which means a crash. Make sure we've captured1695 // The runner unexpectedly closed a stdio pipe, which means a crash. Make sure we've captured
1697 // all available stderr to make our error output as useful as possible.1696 // all available stderr to make our error output as useful as possible.
1698 while (try poller.poll()) {}1697 const stderr_fr = multi_reader.fileReader(1);
1699 run.step.result_stderr = try arena.dupe(u8, poller.reader(.stderr).buffered());1698 while (true) {
1699 stderr_fr.interface.fillMore() catch |e| switch (e) {
1700 error.ReadFailed => return stderr_fr.err.?,
1701 error.EndOfStream => break,
1702 };
1703 }
1704 run.step.result_stderr = try arena.dupe(u8, stderr_fr.interface.buffered());
17001705
1701 // Clean up everything and wait for the child to exit.1706 // Clean up everything and wait for the child to exit.
1702 child.stdin.?.close(io);1707 child.stdin.?.close(io);
1703 child.stdin = null;1708 child.stdin = null;
1704 poller.deinit();1709 multi_reader.deinit();
1705 child_killed = true;1710 child_killed = true;
1706 const term = try child.wait(io);1711 const term = try child.wait(io);
1707 run.step.result_peak_rss = @max(1712 run.step.result_peak_rss = @max(
...@@ -1716,13 +1721,14 @@ fn evalZigTest(...@@ -1716,13 +1721,14 @@ fn evalZigTest(
1716 .no_poll => |no_poll| {1721 .no_poll => |no_poll| {
1717 // This might be a success (we requested exit and the child dutifully closed stdout) or1722 // This might be a success (we requested exit and the child dutifully closed stdout) or
1718 // a crash of some kind. Either way, the child will terminate by itself -- wait for it.1723 // a crash of some kind. Either way, the child will terminate by itself -- wait for it.
1719 const stderr_owned = try arena.dupe(u8, poller.reader(.stderr).buffered());1724 const stderr_reader = multi_reader.reader(1);
1720 poller.reader(.stderr).tossBuffered();1725 const stderr_owned = try arena.dupe(u8, stderr_reader.buffered());
1726 stderr_reader.tossBuffered();
17211727
1722 // Clean up everything and wait for the child to exit.1728 // Clean up everything and wait for the child to exit.
1723 child.stdin.?.close(io);1729 child.stdin.?.close(io);
1724 child.stdin = null;1730 child.stdin = null;
1725 poller.deinit();1731 multi_reader.deinit();
1726 child_killed = true;1732 child_killed = true;
1727 const term = try child.wait(io);1733 const term = try child.wait(io);
1728 run.step.result_peak_rss = @max(1734 run.step.result_peak_rss = @max(
...@@ -1770,8 +1776,9 @@ fn evalZigTest(...@@ -1770,8 +1776,9 @@ fn evalZigTest(
1770 return;1776 return;
1771 },1777 },
1772 .timeout => |timeout| {1778 .timeout => |timeout| {
1773 const stderr = poller.reader(.stderr).buffered();1779 const stderr_reader = multi_reader.reader(1);
1774 poller.reader(.stderr).tossBuffered();1780 const stderr = stderr_reader.buffered();
1781 stderr_reader.tossBuffered();
1775 if (timeout.active_test_index) |test_index| {1782 if (timeout.active_test_index) |test_index| {
1776 // A test was running. Report the timeout against that test, and continue on to1783 // A test was running. Report the timeout against that test, and continue on to
1777 // the next test.1784 // the next test.
...@@ -1796,16 +1803,16 @@ fn evalZigTest(...@@ -1796,16 +1803,16 @@ fn evalZigTest(
1796 }1803 }
1797}1804}
17981805
1799/// Polls stdout of a Zig test process until a termination condition is reached:1806/// Reads stdout of a Zig test process until a termination condition is reached:
1800/// * A write fails, indicating the child unexpectedly closed stdin1807/// * A write fails, indicating the child unexpectedly closed stdin
1801/// * A test (or a response from the test runner) times out1808/// * A test (or a response from the test runner) times out
1802/// * `poll` fails, indicating the child closed stdout and stderr1809/// * The wait fails, indicating the child closed stdout and stderr
1803fn pollZigTest(1810fn waitZigTest(
1804 run: *Run,1811 run: *Run,
1805 child: *process.Child,1812 child: *process.Child,
1806 options: Step.MakeOptions,1813 options: Step.MakeOptions,
1807 fuzz_context: ?FuzzContext,1814 fuzz_context: ?FuzzContext,
1808 poller: *std.Io.Poller(StdioPollEnum),1815 multi_reader: *Io.File.MultiReader,
1809 opt_metadata: *?TestMetadata,1816 opt_metadata: *?TestMetadata,
1810 results: *Step.TestResults,1817 results: *Step.TestResults,
1811) !union(enum) {1818) !union(enum) {
...@@ -1874,12 +1881,11 @@ fn pollZigTest(...@@ -1874,12 +1881,11 @@ fn pollZigTest(
1874 break :ns @max(options.unit_test_timeout_ns orelse 0, 60 * std.time.ns_per_s);1881 break :ns @max(options.unit_test_timeout_ns orelse 0, 60 * std.time.ns_per_s);
1875 };1882 };
18761883
1877 const stdout = poller.reader(.stdout);1884 const stdout = multi_reader.reader(0);
1878 const stderr = poller.reader(.stderr);1885 const stderr = multi_reader.reader(1);
1886 const Header = std.zig.Server.Message.Header;
18791887
1880 while (true) {1888 while (true) {
1881 const Header = std.zig.Server.Message.Header;
1882
1883 // This block is exited when `stdout` contains enough bytes for a `Header`.1889 // This block is exited when `stdout` contains enough bytes for a `Header`.
1884 header_ready: {1890 header_ready: {
1885 if (stdout.buffered().len >= @sizeOf(Header)) {1891 if (stdout.buffered().len >= @sizeOf(Header)) {
...@@ -1894,18 +1900,22 @@ fn pollZigTest(...@@ -1894,18 +1900,22 @@ fn pollZigTest(
1894 break :ns options.unit_test_timeout_ns;1900 break :ns options.unit_test_timeout_ns;
1895 };1901 };
18961902
1897 if (opt_timeout_ns) |timeout_ns| {1903 const timeout: Io.Timeout = if (opt_timeout_ns) |timeout_ns| .{ .duration = .{
1898 const remaining_ns = timeout_ns -| timer.?.read();1904 .raw = .fromNanoseconds(timeout_ns -| timer.?.read()),
1899 if (!try poller.pollTimeout(remaining_ns)) return .{ .no_poll = .{1905 .clock = .awake,
1900 .active_test_index = active_test_index,1906 } } else .none;
1901 .ns_elapsed = if (timer) |*t| t.read() else 0,1907
1902 } };1908 multi_reader.fill(timeout) catch |err| switch (err) {
1903 } else {1909 error.Timeout, error.EndOfStream => return .{ .no_poll = .{
1904 if (!try poller.poll()) return .{ .no_poll = .{
1905 .active_test_index = active_test_index,1910 .active_test_index = active_test_index,
1906 .ns_elapsed = if (timer) |*t| t.read() else 0,1911 .ns_elapsed = if (timer) |*t| t.read() else 0,
1907 } };1912 } },
1908 }1913 error.UnsupportedClock => {
1914 timer = null;
1915 continue;
1916 },
1917 else => |e| return e,
1918 };
19091919
1910 if (stdout.buffered().len >= @sizeOf(Header)) {1920 if (stdout.buffered().len >= @sizeOf(Header)) {
1911 // There wasn't a header before, but there is one after the `poll`.1921 // There wasn't a header before, but there is one after the `poll`.
...@@ -1923,11 +1933,8 @@ fn pollZigTest(...@@ -1923,11 +1933,8 @@ fn pollZigTest(
1923 }1933 }
1924 // There is definitely a header available now -- read it.1934 // There is definitely a header available now -- read it.
1925 const header = stdout.takeStruct(Header, .little) catch unreachable;1935 const header = stdout.takeStruct(Header, .little) catch unreachable;
1936 try stdout.fill(header.bytes_len);
19261937
1927 while (stdout.buffered().len < header.bytes_len) if (!try poller.poll()) return .{ .no_poll = .{
1928 .active_test_index = active_test_index,
1929 .ns_elapsed = if (timer) |*t| t.read() else 0,
1930 } };
1931 const body = stdout.take(header.bytes_len) catch unreachable;1938 const body = stdout.take(header.bytes_len) catch unreachable;
1932 var body_r: std.Io.Reader = .fixed(body);1939 var body_r: std.Io.Reader = .fixed(body);
1933 switch (header.tag) {1940 switch (header.tag) {
...@@ -2164,6 +2171,7 @@ fn evalGeneric(run: *Run, spawn_options: process.SpawnOptions) !EvalGenericResul...@@ -2164,6 +2171,7 @@ fn evalGeneric(run: *Run, spawn_options: process.SpawnOptions) !EvalGenericResul
2164 const b = run.step.owner;2171 const b = run.step.owner;
2165 const io = b.graph.io;2172 const io = b.graph.io;
2166 const arena = b.allocator;2173 const arena = b.allocator;
2174 const gpa = b.allocator;
21672175
2168 var child = try process.spawn(io, spawn_options);2176 var child = try process.spawn(io, spawn_options);
2169 defer child.kill(io);2177 defer child.kill(io);
...@@ -2211,23 +2219,31 @@ fn evalGeneric(run: *Run, spawn_options: process.SpawnOptions) !EvalGenericResul...@@ -2211,23 +2219,31 @@ fn evalGeneric(run: *Run, spawn_options: process.SpawnOptions) !EvalGenericResul
22112219
2212 if (child.stdout) |stdout| {2220 if (child.stdout) |stdout| {
2213 if (child.stderr) |stderr| {2221 if (child.stderr) |stderr| {
2214 var poller = std.Io.poll(arena, enum { stdout, stderr }, .{2222 var multi_reader_buffer: Io.File.MultiReader.Buffer(2) = undefined;
2215 .stdout = stdout,2223 var multi_reader: Io.File.MultiReader = undefined;
2216 .stderr = stderr,2224 multi_reader.init(gpa, io, multi_reader_buffer.toStreams(), &.{ stdout, stderr });
2217 });2225 defer multi_reader.deinit();
2218 defer poller.deinit();
22192226
2220 while (try poller.poll()) {2227 const stdout_reader = multi_reader.reader(0);
2228 const stderr_reader = multi_reader.reader(1);
2229
2230 while (multi_reader.fill(.none)) |_| {
2221 if (run.stdio_limit.toInt()) |limit| {2231 if (run.stdio_limit.toInt()) |limit| {
2222 if (poller.reader(.stderr).buffered().len > limit)2232 if (stdout_reader.buffered().len > limit)
2223 return error.StdoutStreamTooLong;2233 return error.StdoutStreamTooLong;
2224 if (poller.reader(.stderr).buffered().len > limit)2234 if (stderr_reader.buffered().len > limit)
2225 return error.StderrStreamTooLong;2235 return error.StderrStreamTooLong;
2226 }2236 }
2237 } else |err| switch (err) {
2238 error.UnsupportedClock, error.Timeout => unreachable,
2239 error.EndOfStream => {},
2240 else => |e| return e,
2227 }2241 }
22282242
2229 stdout_bytes = try poller.toOwnedSlice(.stdout);2243 try multi_reader.checkAnyError();
2230 stderr_bytes = try poller.toOwnedSlice(.stderr);2244
2245 stdout_bytes = try multi_reader.toOwnedSlice(0);
2246 stderr_bytes = try multi_reader.toOwnedSlice(1);
2231 } else {2247 } else {
2232 var stdout_reader = stdout.readerStreaming(io, &.{});2248 var stdout_reader = stdout.readerStreaming(io, &.{});
2233 stdout_bytes = stdout_reader.interface.allocRemaining(arena, run.stdio_limit) catch |err| switch (err) {2249 stdout_bytes = stdout_reader.interface.allocRemaining(arena, run.stdio_limit) catch |err| switch (err) {
lib/std/Io/File/MultiReader.zig+46-13
...@@ -113,10 +113,28 @@ pub fn deinit(mr: *MultiReader) void {...@@ -113,10 +113,28 @@ pub fn deinit(mr: *MultiReader) void {
113 }113 }
114}114}
115115
116pub fn fileReader(mr: *MultiReader, index: usize) *File.Reader {
117 return &mr.streams.contexts()[index].fr;
118}
119
116pub fn reader(mr: *MultiReader, index: usize) *Io.Reader {120pub fn reader(mr: *MultiReader, index: usize) *Io.Reader {
117 return &mr.streams.contexts()[index].fr.interface;121 return &mr.streams.contexts()[index].fr.interface;
118}122}
119123
124/// Checks for errors in all streams, prioritizing `error.Canceled` if it
125/// occurred anywhere.
126pub fn checkAnyError(mr: *const MultiReader) Error!void {
127 const contexts = mr.streams.contexts();
128 var other: Error!void = {};
129 for (contexts) |*context| {
130 if (context.err) |err| switch (err) {
131 error.Canceled => |e| return e,
132 else => |e| other = e,
133 };
134 }
135 return other;
136}
137
120pub fn toOwnedSlice(mr: *MultiReader, index: usize) Allocator.Error![]u8 {138pub fn toOwnedSlice(mr: *MultiReader, index: usize) Allocator.Error![]u8 {
121 const gpa = mr.gpa;139 const gpa = mr.gpa;
122 const r: *Io.Reader = reader(mr, index);140 const r: *Io.Reader = reader(mr, index);
...@@ -140,7 +158,7 @@ fn stream(r: *Io.Reader, w: *Io.Writer, limit: Io.Limit) Io.Reader.StreamError!u...@@ -140,7 +158,7 @@ fn stream(r: *Io.Reader, w: *Io.Writer, limit: Io.Limit) Io.Reader.StreamError!u
140 const fr: *File.Reader = @alignCast(@fieldParentPtr("interface", r));158 const fr: *File.Reader = @alignCast(@fieldParentPtr("interface", r));
141 const context: *Context = @fieldParentPtr("fr", fr);159 const context: *Context = @fieldParentPtr("fr", fr);
142 const mr = context.mr;160 const mr = context.mr;
143 return fill(mr, context);161 return fillUntimed(mr, context);
144}162}
145163
146fn discard(r: *Io.Reader, limit: Io.Limit) Io.Reader.Error!usize {164fn discard(r: *Io.Reader, limit: Io.Limit) Io.Reader.Error!usize {
...@@ -148,7 +166,7 @@ fn discard(r: *Io.Reader, limit: Io.Limit) Io.Reader.Error!usize {...@@ -148,7 +166,7 @@ fn discard(r: *Io.Reader, limit: Io.Limit) Io.Reader.Error!usize {
148 const fr: *File.Reader = @alignCast(@fieldParentPtr("interface", r));166 const fr: *File.Reader = @alignCast(@fieldParentPtr("interface", r));
149 const context: *Context = @fieldParentPtr("fr", fr);167 const context: *Context = @fieldParentPtr("fr", fr);
150 const mr = context.mr;168 const mr = context.mr;
151 return fill(mr, context);169 return fillUntimed(mr, context);
152}170}
153171
154fn readVec(r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize {172fn readVec(r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize {
...@@ -156,7 +174,7 @@ fn readVec(r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize {...@@ -156,7 +174,7 @@ fn readVec(r: *Io.Reader, data: [][]u8) Io.Reader.Error!usize {
156 const fr: *File.Reader = @alignCast(@fieldParentPtr("interface", r));174 const fr: *File.Reader = @alignCast(@fieldParentPtr("interface", r));
157 const context: *Context = @fieldParentPtr("fr", fr);175 const context: *Context = @fieldParentPtr("fr", fr);
158 const mr = context.mr;176 const mr = context.mr;
159 return fill(mr, context);177 return fillUntimed(mr, context);
160}178}
161179
162fn rebase(r: *Io.Reader, capacity: usize) Io.Reader.RebaseError!void {180fn rebase(r: *Io.Reader, capacity: usize) Io.Reader.RebaseError!void {
...@@ -196,20 +214,23 @@ fn rebaseGrowing(mr: *MultiReader, context: *Context, capacity: usize) Allocator...@@ -196,20 +214,23 @@ fn rebaseGrowing(mr: *MultiReader, context: *Context, capacity: usize) Allocator
196 }214 }
197}215}
198216
199fn fill(mr: *MultiReader, original_context: *Context) Io.Reader.Error!usize {217pub const FillError = Io.Batch.WaitError || error{
218 /// `fill` was called when all streams already have failed or reached the
219 /// end.
220 EndOfStream,
221};
222
223/// Wait until at least one stream receives more data.
224pub fn fill(mr: *MultiReader, timeout: Io.Timeout) FillError!void {
200 const contexts = mr.streams.contexts();225 const contexts = mr.streams.contexts();
201 const operations = mr.streams.operations();226 const operations = mr.streams.operations();
202 const io = contexts[0].fr.io;227 const io = contexts[0].fr.io;
228 var any_completed = false;
203229
204 mr.batch.wait(io, .none) catch |err| switch (err) {230 try mr.batch.wait(io, timeout);
205 error.Timeout, error.UnsupportedClock => unreachable,
206 else => |e| {
207 original_context.err = e;
208 return error.ReadFailed;
209 },
210 };
211231
212 while (mr.batch.next()) |i| {232 while (mr.batch.next()) |i| {
233 any_completed = true;
213 const context = &contexts[i];234 const context = &contexts[i];
214 const operation = &operations[i];235 const operation = &operations[i];
215 const n = operation.file_read_streaming.status.result catch |err| {236 const n = operation.file_read_streaming.status.result catch |err| {
...@@ -234,7 +255,19 @@ fn fill(mr: *MultiReader, original_context: *Context) Io.Reader.Error!usize {...@@ -234,7 +255,19 @@ fn fill(mr: *MultiReader, original_context: *Context) Io.Reader.Error!usize {
234 mr.batch.add(i);255 mr.batch.add(i);
235 }256 }
236257
237 if (original_context.err != null) return error.ReadFailed;258 if (!any_completed) return error.EndOfStream;
238 if (original_context.eos) return error.EndOfStream;259}
260
261fn fillUntimed(mr: *MultiReader, context: *Context) Io.Reader.Error!usize {
262 fill(mr, .none) catch |err| switch (err) {
263 error.Timeout, error.UnsupportedClock => unreachable,
264 error.Canceled, error.ConcurrencyUnavailable => |e| {
265 context.err = e;
266 return error.ReadFailed;
267 },
268 error.EndOfStream => |e| return e,
269 };
270 if (context.err != null) return error.ReadFailed;
271 if (context.eos) return error.EndOfStream;
239 return 0;272 return 0;
240}273}