| ... | @@ -1658,31 +1658,32 @@ pub fn buildExeSubprocess( | ... | @@ -1658,31 +1658,32 @@ pub fn buildExeSubprocess( |
| 1658 | }; | 1658 | }; |
| 1659 | defer child.kill(io); | 1659 | defer child.kill(io); |
| 1660 | | 1660 | |
| 1661 | var stderr_task = io.concurrent(readStreamAlloc, .{ gpa, io, child.stderr.?, .unlimited }) catch | 1661 | var multi_reader_buffer: Io.File.MultiReader.Buffer(2) = undefined; |
| 1662 | @panic("TODO use multireader instead"); | 1662 | var multi_reader: Io.File.MultiReader = undefined; |
| 1663 | defer if (stderr_task.cancel(io)) |slice| gpa.free(slice) else |_| {}; | 1663 | multi_reader.init(gpa, io, multi_reader_buffer.toStreams(), &.{ child.stdout.?, child.stderr.? }); |
| 1664 | | 1664 | defer multi_reader.deinit(); |
| 1665 | var stdout_buffer: [512]u8 = undefined; | 1665 | const stdout = multi_reader.reader(0); |
| 1666 | var stdout_reader: Io.File.Reader = .initStreaming(child.stdout.?, io, &stdout_buffer); | 1666 | const stderr = multi_reader.reader(1); |
| 1667 | const stdout = &stdout_reader.interface; | 1667 | |
| 1668 | | 1668 | var stdin_buffer: [8]u8 = undefined; |
| 1669 | { | 1669 | var stdin_writer = child.stdin.?.writerStreaming(io, &stdin_buffer); |
| 1670 | var w = child.stdin.?.writer(io, &.{}); | 1670 | |
| 1671 | w.interface.writeStruct(Client.Message.Header{ .tag = .update, .bytes_len = 0 }, .little) catch |err| switch (err) { | 1671 | var client: Client = .{ |
| 1672 | error.WriteFailed => { | 1672 | .in = stdout, |
| 1673 | log.err("{t} writing to command: {f}", .{ w.err.?, cmd }); | 1673 | .out = &stdin_writer.interface, |
| 1674 | return error.AlreadyReported; | 1674 | }; |
| 1675 | }, | | |
| 1676 | }; | | |
| 1677 | w.interface.writeStruct(Client.Message.Header{ .tag = .exit, .bytes_len = 0 }, .little) catch |err| switch (err) { | | |
| 1678 | error.WriteFailed => { | | |
| 1679 | log.err("{t} writing to command: {f}", .{ w.err.?, cmd }); | | |
| 1680 | return error.AlreadyReported; | | |
| 1681 | }, | | |
| 1682 | }; | | |
| 1683 | } | | |
| 1684 | | 1675 | |
| 1685 | const Header = Server.Message.Header; | 1676 | (blk: { |
| | 1677 | client.serveMessageHeader(.{ .tag = .update, .bytes_len = 0 }) catch |err| break :blk err; |
| | 1678 | client.serveMessageHeader(.{ .tag = .exit, .bytes_len = 0 }) catch |err| break :blk err; |
| | 1679 | client.out.flush() catch |err| break :blk err; |
| | 1680 | }) catch |err| switch (err) { |
| | 1681 | error.WriteFailed => { |
| | 1682 | if (stdin_writer.err.? == error.Canceled) return error.Canceled; |
| | 1683 | log.err("{t} writing to command: {f}", .{ stdin_writer.err.?, cmd }); |
| | 1684 | return error.AlreadyReported; |
| | 1685 | }, |
| | 1686 | }; |
| 1686 | | 1687 | |
| 1687 | var result: ?Cache.Path = null; | 1688 | var result: ?Cache.Path = null; |
| 1688 | defer if (result) |r| gpa.free(r.sub_path); | 1689 | defer if (result) |r| gpa.free(r.sub_path); |
| ... | @@ -1690,33 +1691,29 @@ pub fn buildExeSubprocess( | ... | @@ -1690,33 +1691,29 @@ pub fn buildExeSubprocess( |
| 1690 | var result_error_bundle: ErrorBundle = .empty; | 1691 | var result_error_bundle: ErrorBundle = .empty; |
| 1691 | defer result_error_bundle.deinit(gpa); | 1692 | defer result_error_bundle.deinit(gpa); |
| 1692 | | 1693 | |
| 1693 | var body_buffer: std.ArrayList(u8) = .empty; | | |
| 1694 | defer body_buffer.deinit(gpa); | | |
| 1695 | | | |
| 1696 | var received_fs_inputs = false; | 1694 | var received_fs_inputs = false; |
| 1697 | var cache_hit = false; | 1695 | var cache_hit = false; |
| 1698 | | 1696 | |
| | 1697 | var eos_err: error{EndOfStream}!void = {}; |
| | 1698 | |
| 1699 | while (true) { | 1699 | while (true) { |
| 1700 | const header = stdout.takeStruct(Header, .little) catch |err| switch (err) { | 1700 | const header = client.receiveMessageWithMultiReader(&multi_reader, .none) catch |err| switch (err) { |
| 1701 | error.ReadFailed => { | 1701 | error.Timeout => unreachable, |
| 1702 | log.err("{t} reading from command: {f}", .{ stdout_reader.err.?, cmd }); | 1702 | error.EndOfStream => |e| { |
| 1703 | return error.AlreadyReported; | 1703 | if (client.in.bufferedLen() == 0) break; |
| 1704 | }, | 1704 | // Better to report the crash with stderr below, but we set |
| 1705 | error.EndOfStream => break, | 1705 | // this in case the child exits successfully while violating |
| 1706 | }; | 1706 | // this protocol. |
| 1707 | body_buffer.clearRetainingCapacity(); | 1707 | eos_err = e; |
| 1708 | stdout.appendExact(gpa, &body_buffer, header.bytes_len) catch |err| switch (err) { | 1708 | break; |
| 1709 | error.ReadFailed => { | | |
| 1710 | log.err("{t} reading from command: {f}", .{ stdout_reader.err.?, cmd }); | | |
| 1711 | return error.AlreadyReported; | | |
| 1712 | }, | 1709 | }, |
| 1713 | error.OutOfMemory => |e| return e, | 1710 | error.Canceled, error.OutOfMemory => |e| return e, |
| 1714 | error.EndOfStream => { | 1711 | else => |e| { |
| 1715 | log.err("unexpected end of stream from command: {f}", .{cmd}); | 1712 | log.err("{t} reading from command: {f}", .{ e, cmd }); |
| 1716 | return error.AlreadyReported; | 1713 | return error.AlreadyReported; |
| 1717 | }, | 1714 | }, |
| 1718 | }; | 1715 | }; |
| 1719 | const body = body_buffer.items; | 1716 | const body = stdout.take(header.bytes_len) catch unreachable; |
| 1720 | | 1717 | |
| 1721 | switch (header.tag) { | 1718 | switch (header.tag) { |
| 1722 | .zig_version => { | 1719 | .zig_version => { |
| ... | @@ -1767,16 +1764,15 @@ pub fn buildExeSubprocess( | ... | @@ -1767,16 +1764,15 @@ pub fn buildExeSubprocess( |
| 1767 | } | 1764 | } |
| 1768 | } | 1765 | } |
| 1769 | | 1766 | |
| 1770 | const stderr_contents = stderr_task.await(io) catch |err| switch (err) { | 1767 | const stderr_contents = stderr.buffered(); |
| 1771 | error.Canceled, error.OutOfMemory => |e| return e, | | |
| 1772 | else => |e| c: { | | |
| 1773 | log.warn("{t} reading stderr from command: {f}", .{ e, cmd }); | | |
| 1774 | break :c ""; | | |
| 1775 | }, | | |
| 1776 | }; | | |
| 1777 | if (stderr_contents.len > 0) | 1768 | if (stderr_contents.len > 0) |
| 1778 | log.warn("unexpected stderr from {s} command:\n{s}", .{ options.argv[0], stderr_contents }); | 1769 | log.warn("unexpected stderr from {s} command:\n{s}", .{ options.argv[0], stderr_contents }); |
| 1779 | | 1770 | |
| | 1771 | eos_err catch { |
| | 1772 | log.err("unexpected end of stream from command: {f}", .{cmd}); |
| | 1773 | return error.AlreadyReported; |
| | 1774 | }; |
| | 1775 | |
| 1780 | // Send EOF to stdin. | 1776 | // Send EOF to stdin. |
| 1781 | child.stdin.?.close(io); | 1777 | child.stdin.?.close(io); |
| 1782 | child.stdin = null; | 1778 | child.stdin = null; |
| ... | @@ -1834,14 +1830,6 @@ pub fn buildExeSubprocess( | ... | @@ -1834,14 +1830,6 @@ pub fn buildExeSubprocess( |
| 1834 | }; | 1830 | }; |
| 1835 | } | 1831 | } |
| 1836 | | 1832 | |
| 1837 | fn readStreamAlloc(gpa: Allocator, io: Io, file: Io.File, limit: Io.Limit) ![]u8 { | | |
| 1838 | var file_reader: Io.File.Reader = .initStreaming(file, io, &.{}); | | |
| 1839 | return file_reader.interface.allocRemaining(gpa, limit) catch |err| switch (err) { | | |
| 1840 | error.ReadFailed => return file_reader.err.?, | | |
| 1841 | else => |e| return e, | | |
| 1842 | }; | | |
| 1843 | } | | |
| 1844 | | | |
| 1845 | test { | 1833 | test { |
| 1846 | _ = Ast; | 1834 | _ = Ast; |
| 1847 | _ = AstRlAnnotate; | 1835 | _ = AstRlAnnotate; |