diff --git a/lib/compiler/Maker/Step.zig b/lib/compiler/Maker/Step.zig index 812613928eaa9134dd249e4fc2ece38065cad01c..b8c5992cce243ed5b7f5fb0cdfc1b6b9d92f7fdb 100644 --- a/lib/compiler/Maker/Step.zig +++ b/lib/compiler/Maker/Step.zig @@ -561,24 +561,26 @@ fn zigProcessUpdate(step_index: Configuration.Step.Index, maker: *Maker, zp: *Zi var result: ?Path = null; var eos_err: error{EndOfStream}!void = {}; - const stdout = zp.multi_reader.fileReader(0); + var client: std.zig.Client = .{ + .in = zp.multi_reader.reader(0), + .out = undefined, + }; while (true) { - const Header = std.zig.Server.Message.Header; - const header = stdout.interface.takeStruct(Header, .little) catch |err| switch (err) { - error.EndOfStream => break, - error.ReadFailed => return stdout.err.?, - }; - const body = stdout.interface.take(header.bytes_len) catch |err| switch (err) { + const header = client.receiveMessageWithMultiReader(&zp.multi_reader, .none) catch |err| switch (err) { + error.Timeout => unreachable, error.EndOfStream => |e| { + if (client.in.bufferedLen() == 0) break; // Better to report the crash with stderr below, but we set // this in case the child exits successfully while violating // this protocol. eos_err = e; break; }, - error.ReadFailed => return stdout.err.?, + else => |e| return e, }; + const body = client.in.take(header.bytes_len) catch unreachable; + switch (header.tag) { .zig_version => { if (!std.mem.eql(u8, builtin.zig_version_string, body)) { diff --git a/lib/compiler/Maker/Step/Run.zig b/lib/compiler/Maker/Step/Run.zig index b6fc911f01f8d25a27cf2f3829118accfe2f58ea..488225aadca8deb146e8fc1a2fd4c196c3811b2f 100644 --- a/lib/compiler/Maker/Step/Run.zig +++ b/lib/compiler/Maker/Step/Run.zig @@ -384,13 +384,23 @@ fn waitZigTest( var sub_prog_node: ?std.Progress.Node = null; defer if (sub_prog_node) |n| n.end(); + const stdout = multi_reader.reader(0); + const stderr = multi_reader.reader(1); + + var stdin_writer = child.stdin.?.writerStreaming(io, &.{}); + + var client: std.zig.Client = .{ + .in = stdout, + .out = &stdin_writer.interface, + }; + if (opt_metadata.*) |*md| { // Previous unit test process died or was killed; we're continuing where it left off - requestNextTest(io, child.stdin.?, md, &sub_prog_node) catch |err| return .{ .write_failed = err }; + requestNextTest(&client, md, &sub_prog_node) catch |err| return .{ .write_failed = err }; } else { // Running unit tests normally run.fuzz_tests.clearRetainingCapacity(); - sendMessage(io, child.stdin.?, .query_test_metadata) catch |err| return .{ .write_failed = err }; + client.serveBodylessMessage(.query_test_metadata) catch |err| return .{ .write_failed = err }; } var active_test_index: ?u32 = null; @@ -410,10 +420,6 @@ fn waitZigTest( .raw = .fromNanoseconds(ns), } else null; - const stdout = multi_reader.reader(0); - const stderr = multi_reader.reader(1); - const Header = std.zig.Server.Message.Header; - while (true) { const timeout: Io.Timeout = t: { const opt_duration = if (active_test_index == null) response_timeout else test_timeout; @@ -421,46 +427,20 @@ fn waitZigTest( break :t .{ .deadline = last_update.addDuration(duration) }; }; - // This block is exited when `stdout` contains enough bytes for a `Header`. - header_ready: { - if (stdout.buffered().len >= @sizeOf(Header)) { - // We already have one, no need to poll! - break :header_ready; - } - - multi_reader.fill(64, timeout) catch |err| switch (err) { - error.Timeout => return .{ .timeout = .{ - .active_test_index = active_test_index, - .ns_elapsed = @intCast(last_update.untilNow(io).raw.nanoseconds), - } }, - error.EndOfStream => return .{ .no_poll = .{ - .active_test_index = active_test_index, - .ns_elapsed = @intCast(last_update.untilNow(io).raw.nanoseconds), - } }, - else => |e| return e, - }; - - continue; - } - // There is definitely a header available now -- read it. - const header = stdout.takeStruct(Header, .little) catch unreachable; - - while (stdout.buffered().len < header.bytes_len) { - multi_reader.fill(64, timeout) catch |err| switch (err) { - error.Timeout => return .{ .timeout = .{ - .active_test_index = active_test_index, - .ns_elapsed = @intCast(last_update.untilNow(io).raw.nanoseconds), - } }, - error.EndOfStream => return .{ .no_poll = .{ - .active_test_index = active_test_index, - .ns_elapsed = @intCast(last_update.untilNow(io).raw.nanoseconds), - } }, - else => |e| return e, - }; - } - - const body = stdout.take(header.bytes_len) catch unreachable; + const header = client.receiveMessageWithMultiReader(multi_reader, timeout) catch |err| switch (err) { + error.Timeout => return .{ .timeout = .{ + .active_test_index = active_test_index, + .ns_elapsed = @intCast(last_update.untilNow(io).raw.nanoseconds), + } }, + error.EndOfStream => return .{ .no_poll = .{ + .active_test_index = active_test_index, + .ns_elapsed = @intCast(last_update.untilNow(io).raw.nanoseconds), + } }, + else => |e| return e, + }; + const body = client.in.take(header.bytes_len) catch unreachable; var body_r: std.Io.Reader = .fixed(body); + switch (header.tag) { .zig_version => { if (!std.mem.eql(u8, builtin.zig_version_string, body)) return step.fail( @@ -500,7 +480,7 @@ fn waitZigTest( active_test_index = null; last_update = .now(io, .awake); - requestNextTest(io, child.stdin.?, &opt_metadata.*.?, &sub_prog_node) catch |err| return .{ .write_failed = err }; + requestNextTest(&client, &opt_metadata.*.?, &sub_prog_node) catch |err| return .{ .write_failed = err }; }, .test_started => { active_test_index = opt_metadata.*.?.next_index - 1; @@ -551,7 +531,7 @@ fn waitZigTest( md.ns_per_test[tr_hdr.index] = @intCast(last_update.durationTo(now).raw.nanoseconds); last_update = now; - requestNextTest(io, child.stdin.?, md, &sub_prog_node) catch |err| return .{ .write_failed = err }; + requestNextTest(&client, md, &sub_prog_node) catch |err| return .{ .write_failed = err }; }, else => {}, // ignore other messages } @@ -697,17 +677,18 @@ const FuzzTestRunner = struct { for (0.., f.instances) |id, *instance| { const id32: u32 = @intCast(id); + var writer = instance.child.stdin.?.writerStreaming(io, &.{}); + const client: std.zig.Client = .{ + .in = undefined, + .out = &writer.interface, + }; (switch (f.ctx.fuzz.mode) { - .forever => sendRunFuzzTestMessage( - io, - instance.child.stdin.?, + .forever => client.serveRunFuzzTestMessage( run.fuzz_tests.items, .forever, id32, ), - .limit => |limit| sendRunFuzzTestMessage( - io, - instance.child.stdin.?, + .limit => |limit| client.serveRunFuzzTestMessage( run.fuzz_tests.items, .iterations, limit.amount, @@ -1315,7 +1296,7 @@ pub const CachedTestMetadata = struct { } }; -fn requestNextTest(io: Io, in: Io.File, metadata: *TestMetadata, sub_prog_node: *?std.Progress.Node) !void { +fn requestNextTest(client: *std.zig.Client, metadata: *TestMetadata, sub_prog_node: *?std.Progress.Node) !void { while (metadata.next_index < metadata.names.len) { const i = metadata.next_index; metadata.next_index += 1; @@ -1326,76 +1307,11 @@ fn requestNextTest(io: Io, in: Io.File, metadata: *TestMetadata, sub_prog_node: if (sub_prog_node.*) |n| n.end(); sub_prog_node.* = metadata.prog_node.start(name, 0); - try sendRunTestMessage(io, in, .run_test, i); + try client.serveRunTest(i); return; } else { metadata.next_index = std.math.maxInt(u32); // indicate that all tests are done - try sendMessage(io, in, .exit); - } -} - -fn sendMessage(io: Io, file: Io.File, tag: std.zig.Client.Message.Tag) !void { - const header: std.zig.Client.Message.Header = .{ - .tag = tag, - .bytes_len = 0, - }; - var w = file.writerStreaming(io, &.{}); - w.interface.writeStruct(header, .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; -} - -fn sendRunTestMessage(io: Io, file: Io.File, tag: std.zig.Client.Message.Tag, index: u32) !void { - const header: std.zig.Client.Message.Header = .{ - .tag = tag, - .bytes_len = 4, - }; - var w = file.writerStreaming(io, &.{}); - w.interface.writeStruct(header, .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; - w.interface.writeInt(u32, index, .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; -} - -fn sendRunFuzzTestMessage( - io: Io, - file: Io.File, - test_names: []const []const u8, - kind: std.Build.abi.fuzz.LimitKind, - amount_or_instance: u64, -) !void { - const header: std.zig.Client.Message.Header = .{ - .tag = .start_fuzzing, - .bytes_len = 1 + 8 + 4 + count: { - var c: u32 = @intCast(test_names.len * 4); - for (test_names) |name| { - c += @intCast(name.len); - } - break :count c; - }, - }; - var w = file.writerStreaming(io, &.{}); - w.interface.writeStruct(header, .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; - w.interface.writeByte(@backingInt(kind)) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; - w.interface.writeInt(u64, amount_or_instance, .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; - w.interface.writeInt(u32, @intCast(test_names.len), .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; - for (test_names) |test_name| { - w.interface.writeInt(u32, @intCast(test_name.len), .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; - w.interface.writeAll(test_name) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; + try client.serveBodylessMessage(.exit); } } diff --git a/lib/compiler/std-docs.zig b/lib/compiler/std-docs.zig index f53123ae925a7fb4d997e1b8c1a5463624077780..0ed2c5bf186c2c91730b1aa4aa1755babdab0be2 100644 --- a/lib/compiler/std-docs.zig +++ b/lib/compiler/std-docs.zig @@ -346,29 +346,39 @@ fn buildWasmBinary( multi_reader.init(gpa, io, multi_reader_buffer.toStreams(), &.{ child.stdout.?, child.stderr.? }); defer multi_reader.deinit(); - try sendMessage(io, child.stdin.?, .update); - try sendMessage(io, child.stdin.?, .exit); + const stdout = multi_reader.reader(0); + + var stdin_buffer: [256]u8 = undefined; + var stdin_writer = child.stdin.?.writerStreaming(io, &stdin_buffer); + + var client: std.zig.Client = .{ + .in = stdout, + .out = &stdin_writer.interface, + }; + + try client.serveMessageHeader(.{ .tag = .update, .bytes_len = 0 }); + try client.serveMessageHeader(.{ .tag = .exit, .bytes_len = 0 }); + try client.out.flush(); var result: ?Cache.Path = null; var result_error_bundle = std.zig.ErrorBundle.empty; - const stdout = multi_reader.fileReader(0); - const MessageHeader = std.zig.Server.Message.Header; - var eos_err: error{EndOfStream}!void = {}; while (true) { - const header = stdout.interface.takeStruct(MessageHeader, .little) catch |err| switch (err) { - error.EndOfStream => break, - error.ReadFailed => return stdout.err.?, - }; - const body = stdout.interface.take(header.bytes_len) catch |err| switch (err) { + const header = client.receiveMessageWithMultiReader(&multi_reader, .none) catch |err| switch (err) { + error.Timeout => unreachable, error.EndOfStream => |e| { + if (client.in.bufferedLen() == 0) break; + // Better to report the crash with stderr below, but we set + // this in case the child exits successfully while violating + // this protocol. eos_err = e; break; }, - error.ReadFailed => return stdout.err.?, + else => |e| return e, }; + const body = client.in.take(header.bytes_len) catch unreachable; switch (header.tag) { .zig_version => { @@ -435,17 +445,6 @@ fn buildWasmBinary( }; } -fn sendMessage(io: Io, file: Io.File, tag: std.zig.Client.Message.Tag) !void { - const header: std.zig.Client.Message.Header = .{ - .tag = tag, - .bytes_len = 0, - }; - var w = file.writer(io, &.{}); - w.interface.writeStruct(header, .little) catch |err| switch (err) { - error.WriteFailed => return w.err.?, - }; -} - fn openBrowserTab(io: Io, url: []const u8) !void { // Until https://github.com/ziglang/zig/issues/19205 is implemented, we // spawn and then leak a concurrent task for this child process. diff --git a/lib/std/zig/Client.zig b/lib/std/zig/Client.zig index fe50f2314a0b68b5bd9d6510efec6d4146477237..df4eb067bd7867ab5898cf1d27717b9bffd0a0ed 100644 --- a/lib/std/zig/Client.zig +++ b/lib/std/zig/Client.zig @@ -1,3 +1,17 @@ +const Client = @This(); + +const std = @import("std"); +const Io = std.Io; +const Allocator = std.mem.Allocator; +const assert = std.debug.assert; +const OutMessage = std.zig.Client.Message; +const InMessage = std.zig.Server.Message; +const Reader = Io.Reader; +const Writer = Io.Writer; + +in: *Reader, +out: *Writer, + pub const Message = struct { pub const Header = extern struct { tag: Tag, @@ -50,7 +64,79 @@ pub const Message = struct { }; comptime { - const std = @import("std"); - std.debug.assert(@sizeOf(std.Build.abi.fuzz.LimitKind) == 1); + assert(@sizeOf(std.Build.abi.fuzz.LimitKind) == 1); } }; + +pub fn receiveMessage(c: *const Client) Reader.Error!InMessage.Header { + return c.in.takeStruct(InMessage.Header, .little); +} + +/// Assumes that `c.in` is a reader in `multi_reader`. +/// Guarantees that the response body will be buffered in `c.in` on success. +pub fn receiveMessageWithMultiReader( + c: *Client, + multi_reader: *Io.File.MultiReader, + timeout: Io.Timeout, +) (Io.File.MultiReader.Error || Io.Timeout.Error)!InMessage.Header { + while (c.in.bufferedLen() < @sizeOf(InMessage.Header)) { + multi_reader.fill(64, timeout) catch |err| switch (err) { + error.Canceled, + error.Timeout, + error.ConcurrencyUnavailable, + error.EndOfStream, + => |e| return e, + }; + } + const header = c.in.takeStruct(InMessage.Header, .little) catch unreachable; + while (c.in.bufferedLen() < header.bytes_len) { + try multi_reader.fill(header.bytes_len - c.in.bufferedLen(), timeout); + } + try multi_reader.checkAnyError(); + return header; +} + +/// Don't forget to flush! +pub fn serveMessageHeader(c: *const Client, header: OutMessage.Header) Writer.Error!void { + try c.out.writeStruct(header, .little); +} + +pub fn serveBodylessMessage(c: *const Client, tag: OutMessage.Tag) Writer.Error!void { + try c.serveMessageHeader(.{ .tag = tag, .bytes_len = 0 }); + try c.out.flush(); +} + +pub fn serveRunTest(c: *const Client, index: u32) !void { + try c.serveMessageHeader(.{ + .tag = .run_test, + .bytes_len = @sizeOf(u32), + }); + try c.out.writeInt(u32, index, .little); + try c.out.flush(); +} + +pub fn serveRunFuzzTestMessage( + c: *const Client, + test_names: []const []const u8, + kind: std.Build.abi.fuzz.LimitKind, + amount_or_instance: u64, +) !void { + try c.serveMessageHeader(.{ + .tag = .start_fuzzing, + .bytes_len = 1 + 8 + 4 + count: { + var bytes_len: u32 = @intCast(test_names.len * 4); + for (test_names) |name| { + bytes_len += @intCast(name.len); + } + break :count bytes_len; + }, + }); + try c.out.writeByte(@backingInt(kind)); + try c.out.writeInt(u64, amount_or_instance, .little); + try c.out.writeInt(u32, @intCast(test_names.len), .little); + for (test_names) |test_name| { + try c.out.writeInt(u32, @intCast(test_name.len), .little); + try c.out.writeAll(test_name); + } + try c.out.flush(); +} diff --git a/src/Compilation.zig b/src/Compilation.zig index 83b0f7ff2b206b23c947437fc514ca1c1cefed6a..bdb1d1a7de07b30cb03b21795f60b9f8b4bd654f 100644 --- a/src/Compilation.zig +++ b/src/Compilation.zig @@ -6006,26 +6006,30 @@ fn spawnZigRc( multi_reader.init(gpa, io, multi_reader_buffer.toStreams(), &.{ child.stdout.?, child.stderr.? }); defer multi_reader.deinit(); - const stdout = multi_reader.fileReader(0); - const MessageHeader = std.zig.Server.Message.Header; + const stdout = multi_reader.reader(0); var eos_err: error{EndOfStream}!void = {}; + var client: std.zig.Client = .{ + .in = stdout, + .out = undefined, + }; + while (true) { - const header = stdout.interface.takeStruct(MessageHeader, .little) catch |err| switch (err) { - error.EndOfStream => break, - error.ReadFailed => return stdout.err.?, - }; - const body = stdout.interface.take(header.bytes_len) catch |err| switch (err) { + const header = client.receiveMessageWithMultiReader(&multi_reader, .none) catch |err| switch (err) { + error.Timeout => unreachable, error.EndOfStream => |e| { + if (client.in.bufferedLen() == 0) break; // Better to report the crash with stderr below, but we set // this in case the child exits successfully while violating // this protocol. eos_err = e; break; }, - error.ReadFailed => return stdout.err.?, + else => |e| return e, }; + const body = client.in.take(header.bytes_len) catch unreachable; + switch (header.tag) { // We expect exactly one ErrorBundle, and if any error_bundle header is // sent then it's a fatal error. diff --git a/tools/incr-check.zig b/tools/incr-check.zig index 89c14ce1e7f60a710e863349025d3803f2dc37bb..cbc1ec659409eadefd89d516c8816ed36f92d4e5 100644 --- a/tools/incr-check.zig +++ b/tools/incr-check.zig @@ -305,21 +305,23 @@ const Eval = struct { fn check(eval: *Eval, mr: *Io.File.MultiReader, update: Case.Update, prog_node: std.Progress.Node) !void { const arena = eval.arena; - const stdout = mr.fileReader(0); - const stderr = &mr.fileReader(1).interface; - const Header = std.zig.Server.Message.Header; + const stdout = mr.reader(0); + const stderr = mr.reader(1); + + var client: std.zig.Client = .{ + .in = stdout, + .out = undefined, + }; while (true) { - const header = stdout.interface.takeStruct(Header, .little) catch |err| switch (err) { - error.EndOfStream => break, - error.ReadFailed => return stdout.err.?, - }; - const body = stdout.interface.take(header.bytes_len) catch |err| switch (err) { + const header = client.receiveMessageWithMultiReader(mr, .none) catch |err| switch (err) { + error.Timeout => unreachable, // If this panic triggers it might be helpful to rework this // code to print the stderr from the abnormally terminated child. error.EndOfStream => @panic("unexpected mid-message end of stream"), - error.ReadFailed => return stdout.err.?, + else => |e| return e, }; + const body = client.in.take(header.bytes_len) catch unreachable; switch (header.tag) { .error_bundle => { @@ -605,12 +607,13 @@ const Eval = struct { fn requestUpdate(eval: *Eval) !void { const io = eval.io; - const header: std.zig.Client.Message.Header = .{ - .tag = .update, - .bytes_len = 0, + + var w = eval.child.stdin.?.writerStreaming(io, &.{}); + var client: std.zig.Client = .{ + .in = undefined, + .out = &w.interface, }; - var w = eval.child.stdin.?.writer(io, &.{}); - w.interface.writeStruct(header, .little) catch |err| switch (err) { + client.serveBodylessMessage(.update) catch |err| switch (err) { error.WriteFailed => return w.err.?, }; } @@ -618,22 +621,23 @@ const Eval = struct { fn end(eval: *Eval, mr: *Io.File.MultiReader) !void { requestExit(eval.child, eval); - const stdout = mr.fileReader(0); - const Header = std.zig.Server.Message.Header; + var client: std.zig.Client = .{ + .in = mr.reader(0), + .out = undefined, + }; while (true) { - const header = stdout.interface.takeStruct(Header, .little) catch |err| switch (err) { - error.EndOfStream => break, - error.ReadFailed => return stdout.err.?, - }; - stdout.interface.discardAll(header.bytes_len) catch |err| switch (err) { - error.ReadFailed => return stdout.err.?, - error.EndOfStream => |e| return e, + const header = client.receiveMessageWithMultiReader(mr, .none) catch |err| switch (err) { + error.Timeout => unreachable, + error.EndOfStream => |e| { + if (client.in.bufferedLen() == 0) break; + return e; + }, + else => |e| return e, }; + try client.in.discardAll(header.bytes_len); } - try mr.fillRemaining(.none); - const stderr = mr.reader(1).buffered(); if (stderr.len > 0) eval.fatal("unexpected stderr:\n{s}", .{stderr}); } @@ -899,12 +903,12 @@ fn requestExit(child: *std.process.Child, eval: *Eval) void { if (child.stdin == null) return; const io = eval.io; - const header: std.zig.Client.Message.Header = .{ - .tag = .exit, - .bytes_len = 0, + var w = eval.child.stdin.?.writerStreaming(io, &.{}); + var client: std.zig.Client = .{ + .in = undefined, + .out = &w.interface, }; - var w = eval.child.stdin.?.writer(io, &.{}); - w.interface.writeStruct(header, .little) catch |err| switch (err) { + client.serveBodylessMessage(.exit) catch |err| switch (err) { error.WriteFailed => switch (w.err.?) { error.BrokenPipe => {}, else => |e| eval.fatal("failed to send exit: {t}", .{e}),