authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2024-05-30 12:53:53-07:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2024-05-31 04:06:12-04:00
logc564a16a01811c18113fe304dbc3e8255af9f172
tree79779ef0cdb54747ff5eb1fc256b77d227a607cf
parent30a35a897f2220d1b252d211c389cf0319170a94

std.Progress: IPC fixes

Reduce node_storage_buffer_len from 200 to 83. This makes messages over the pipe fit in a single packet (4096 bytes). There is now a comptime assert to ensure this. In practice this is plenty of storage because typical terminal heights are significantly less than 83 rows. Handling of split reads is fixed; instead of using a global `remaining_read_trash_bytes`, the value is stored in the "saved metadata" for the IPC node. Saved metadata is split into two arrays so that the "find" operation can quickly scan over fds for a match, looking at 332 bytes maximum, and only reading the memory for the other data upon match. More typical number of bytes read for this operation would be 0 (no child processes), 4 (1 child process), or 64 (16 child processes reporting progress). Removed an align(4) that was leftover from an older design. This also includes part of Jacob Young's not-yet-landed patch that implements `writevNonblock`.

1 files changed, 83 insertions(+), 35 deletions(-)

lib/std/Progress.zig+83-35
...@@ -324,7 +324,7 @@ var global_progress: Progress = .{...@@ -324,7 +324,7 @@ var global_progress: Progress = .{
324 .node_end_index = 0,324 .node_end_index = 0,
325};325};
326326
327const node_storage_buffer_len = 200;327const node_storage_buffer_len = 83;
328var node_parents_buffer: [node_storage_buffer_len]Node.Parent = undefined;328var node_parents_buffer: [node_storage_buffer_len]Node.Parent = undefined;
329var node_storage_buffer: [node_storage_buffer_len]Node.Storage = undefined;329var node_storage_buffer: [node_storage_buffer_len]Node.Storage = undefined;
330var node_freelist_buffer: [node_storage_buffer_len]Node.OptionalIndex = undefined;330var node_freelist_buffer: [node_storage_buffer_len]Node.OptionalIndex = undefined;
...@@ -755,8 +755,10 @@ const Serialized = struct {...@@ -755,8 +755,10 @@ const Serialized = struct {
755755
756 parents_copy: [node_storage_buffer_len]Node.Parent,756 parents_copy: [node_storage_buffer_len]Node.Parent,
757 storage_copy: [node_storage_buffer_len]Node.Storage,757 storage_copy: [node_storage_buffer_len]Node.Storage,
758 ipc_metadata_fds_copy: [node_storage_buffer_len]Fd,
758 ipc_metadata_copy: [node_storage_buffer_len]SavedMetadata,759 ipc_metadata_copy: [node_storage_buffer_len]SavedMetadata,
759760
761 ipc_metadata_fds: [node_storage_buffer_len]Fd,
760 ipc_metadata: [node_storage_buffer_len]SavedMetadata,762 ipc_metadata: [node_storage_buffer_len]SavedMetadata,
761 };763 };
762};764};
...@@ -810,36 +812,39 @@ fn serialize(serialized_buffer: *Serialized.Buffer) Serialized {...@@ -810,36 +812,39 @@ fn serialize(serialized_buffer: *Serialized.Buffer) Serialized {
810}812}
811813
812const SavedMetadata = struct {814const SavedMetadata = struct {
813 ipc_fd: u16,815 remaining_read_trash_bytes: u16,
814 main_index: u8,816 main_index: u8,
815 start_index: u8,817 start_index: u8,
816 nodes_len: u8,818 nodes_len: u8,
819};
817820
818 fn getIpcFd(metadata: SavedMetadata) posix.fd_t {821const Fd = enum(i32) {
819 return if (is_windows)822 _,
820 @ptrFromInt(@as(usize, metadata.ipc_fd) << 2)823
821 else824 fn init(fd: posix.fd_t) Fd {
822 metadata.ipc_fd;825 return @enumFromInt(if (is_windows) @as(isize, @bitCast(@intFromPtr(fd))) else fd);
823 }826 }
824827
825 fn setIpcFd(fd: posix.fd_t) u16 {828 fn get(fd: Fd) posix.fd_t {
826 return @intCast(if (is_windows)829 return if (is_windows)
827 @shrExact(@intFromPtr(fd), 2)830 @ptrFromInt(@as(usize, @bitCast(@as(isize, @intFromEnum(fd)))))
828 else831 else
829 fd);832 @intFromEnum(fd);
830 }833 }
831};834};
832835
833var ipc_metadata_len: u8 = 0;836var ipc_metadata_len: u8 = 0;
834var remaining_read_trash_bytes: usize = 0;
835837
836fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buffer) usize {838fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buffer) usize {
839 const ipc_metadata_fds_copy = &serialized_buffer.ipc_metadata_fds_copy;
837 const ipc_metadata_copy = &serialized_buffer.ipc_metadata_copy;840 const ipc_metadata_copy = &serialized_buffer.ipc_metadata_copy;
841 const ipc_metadata_fds = &serialized_buffer.ipc_metadata_fds;
838 const ipc_metadata = &serialized_buffer.ipc_metadata;842 const ipc_metadata = &serialized_buffer.ipc_metadata;
839843
840 var serialized_len = start_serialized_len;844 var serialized_len = start_serialized_len;
841 var pipe_buf: [2 * 4096]u8 align(4) = undefined;845 var pipe_buf: [2 * 4096]u8 = undefined;
842846
847 const old_ipc_metadata_fds = ipc_metadata_fds_copy[0..ipc_metadata_len];
843 const old_ipc_metadata = ipc_metadata_copy[0..ipc_metadata_len];848 const old_ipc_metadata = ipc_metadata_copy[0..ipc_metadata_len];
844 ipc_metadata_len = 0;849 ipc_metadata_len = 0;
845850
...@@ -850,6 +855,7 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff...@@ -850,6 +855,7 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff
850 ) |main_parent, *main_storage, main_index| {855 ) |main_parent, *main_storage, main_index| {
851 if (main_parent == .unused) continue;856 if (main_parent == .unused) continue;
852 const fd = main_storage.getIpcFd() orelse continue;857 const fd = main_storage.getIpcFd() orelse continue;
858 const opt_saved_metadata = findOld(fd, old_ipc_metadata_fds, old_ipc_metadata);
853 var bytes_read: usize = 0;859 var bytes_read: usize = 0;
854 while (true) {860 while (true) {
855 const n = posix.read(fd, pipe_buf[bytes_read..]) catch |err| switch (err) {861 const n = posix.read(fd, pipe_buf[bytes_read..]) catch |err| switch (err) {
...@@ -862,24 +868,26 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff...@@ -862,24 +868,26 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff
862 },868 },
863 };869 };
864 if (n == 0) break;870 if (n == 0) break;
865 if (remaining_read_trash_bytes > 0) {871 if (opt_saved_metadata) |m| {
866 assert(bytes_read == 0);872 if (m.remaining_read_trash_bytes > 0) {
867 if (remaining_read_trash_bytes >= n) {873 assert(bytes_read == 0);
868 remaining_read_trash_bytes -= n;874 if (m.remaining_read_trash_bytes >= n) {
875 m.remaining_read_trash_bytes = @intCast(m.remaining_read_trash_bytes - n);
876 continue;
877 }
878 const src = pipe_buf[m.remaining_read_trash_bytes..n];
879 std.mem.copyForwards(u8, &pipe_buf, src);
880 m.remaining_read_trash_bytes = 0;
881 bytes_read = src.len;
869 continue;882 continue;
870 }883 }
871 const src = pipe_buf[remaining_read_trash_bytes..n];
872 std.mem.copyForwards(u8, &pipe_buf, src);
873 remaining_read_trash_bytes = 0;
874 bytes_read = src.len;
875 continue;
876 }884 }
877 bytes_read += n;885 bytes_read += n;
878 }886 }
879 // Ignore all but the last message on the pipe.887 // Ignore all but the last message on the pipe.
880 var input: []u8 = pipe_buf[0..bytes_read];888 var input: []u8 = pipe_buf[0..bytes_read];
881 if (input.len == 0) {889 if (input.len == 0) {
882 serialized_len = useSavedIpcData(serialized_len, serialized_buffer, main_storage, main_index, old_ipc_metadata);890 serialized_len = useSavedIpcData(serialized_len, serialized_buffer, main_storage, main_index, opt_saved_metadata, 0, fd);
883 continue;891 continue;
884 }892 }
885893
...@@ -888,9 +896,8 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff...@@ -888,9 +896,8 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff
888 const expected_bytes = 1 + subtree_len * (@sizeOf(Node.Storage) + @sizeOf(Node.Parent));896 const expected_bytes = 1 + subtree_len * (@sizeOf(Node.Storage) + @sizeOf(Node.Parent));
889 if (input.len < expected_bytes) {897 if (input.len < expected_bytes) {
890 // Ignore short reads. We'll handle the next full message when it comes instead.898 // Ignore short reads. We'll handle the next full message when it comes instead.
891 assert(remaining_read_trash_bytes == 0);899 const remaining_read_trash_bytes: u16 = @intCast(expected_bytes - input.len);
892 remaining_read_trash_bytes = expected_bytes - input.len;900 serialized_len = useSavedIpcData(serialized_len, serialized_buffer, main_storage, main_index, opt_saved_metadata, remaining_read_trash_bytes, fd);
893 serialized_len = useSavedIpcData(serialized_len, serialized_buffer, main_storage, main_index, old_ipc_metadata);
894 continue :main_loop;901 continue :main_loop;
895 }902 }
896 if (input.len > expected_bytes) {903 if (input.len > expected_bytes) {
...@@ -908,8 +915,9 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff...@@ -908,8 +915,9 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff
908 const nodes_len: u8 = @intCast(@min(parents.len - 1, serialized_buffer.storage.len - serialized_len));915 const nodes_len: u8 = @intCast(@min(parents.len - 1, serialized_buffer.storage.len - serialized_len));
909916
910 // Remember in case the pipe is empty on next update.917 // Remember in case the pipe is empty on next update.
918 ipc_metadata_fds[ipc_metadata_len] = Fd.init(fd);
911 ipc_metadata[ipc_metadata_len] = .{919 ipc_metadata[ipc_metadata_len] = .{
912 .ipc_fd = SavedMetadata.setIpcFd(fd),920 .remaining_read_trash_bytes = 0,
913 .start_index = @intCast(serialized_len),921 .start_index = @intCast(serialized_len),
914 .nodes_len = nodes_len,922 .nodes_len = nodes_len,
915 .main_index = @intCast(main_index),923 .main_index = @intCast(main_index),
...@@ -950,6 +958,7 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff...@@ -950,6 +958,7 @@ fn serializeIpc(start_serialized_len: usize, serialized_buffer: *Serialized.Buff
950 // Save a copy in case any pipes are empty on the next update.958 // Save a copy in case any pipes are empty on the next update.
951 @memcpy(serialized_buffer.parents_copy[0..serialized_len], serialized_buffer.parents[0..serialized_len]);959 @memcpy(serialized_buffer.parents_copy[0..serialized_len], serialized_buffer.parents[0..serialized_len]);
952 @memcpy(serialized_buffer.storage_copy[0..serialized_len], serialized_buffer.storage[0..serialized_len]);960 @memcpy(serialized_buffer.storage_copy[0..serialized_len], serialized_buffer.storage[0..serialized_len]);
961 @memcpy(ipc_metadata_fds_copy[0..ipc_metadata_len], ipc_metadata_fds[0..ipc_metadata_len]);
953 @memcpy(ipc_metadata_copy[0..ipc_metadata_len], ipc_metadata[0..ipc_metadata_len]);962 @memcpy(ipc_metadata_copy[0..ipc_metadata_len], ipc_metadata[0..ipc_metadata_len]);
954963
955 return serialized_len;964 return serialized_len;
...@@ -963,9 +972,13 @@ fn copyRoot(dest: *Node.Storage, src: *align(1) Node.Storage) void {...@@ -963,9 +972,13 @@ fn copyRoot(dest: *Node.Storage, src: *align(1) Node.Storage) void {
963 };972 };
964}973}
965974
966fn findOld(ipc_fd: posix.fd_t, old_metadata: []const SavedMetadata) ?*const SavedMetadata {975fn findOld(
967 for (old_metadata) |*m| {976 ipc_fd: posix.fd_t,
968 if (m.getIpcFd() == ipc_fd)977 old_metadata_fds: []Fd,
978 old_metadata: []SavedMetadata,
979) ?*SavedMetadata {
980 for (old_metadata_fds, old_metadata) |fd, *m| {
981 if (fd.get() == ipc_fd)
969 return m;982 return m;
970 }983 }
971 return null;984 return null;
...@@ -976,16 +989,28 @@ fn useSavedIpcData(...@@ -976,16 +989,28 @@ fn useSavedIpcData(
976 serialized_buffer: *Serialized.Buffer,989 serialized_buffer: *Serialized.Buffer,
977 main_storage: *Node.Storage,990 main_storage: *Node.Storage,
978 main_index: usize,991 main_index: usize,
979 old_metadata: []const SavedMetadata,992 opt_saved_metadata: ?*SavedMetadata,
993 remaining_read_trash_bytes: u16,
994 fd: posix.fd_t,
980) usize {995) usize {
981 const parents_copy = &serialized_buffer.parents_copy;996 const parents_copy = &serialized_buffer.parents_copy;
982 const storage_copy = &serialized_buffer.storage_copy;997 const storage_copy = &serialized_buffer.storage_copy;
998 const ipc_metadata_fds = &serialized_buffer.ipc_metadata_fds;
983 const ipc_metadata = &serialized_buffer.ipc_metadata;999 const ipc_metadata = &serialized_buffer.ipc_metadata;
9841000
985 const ipc_fd = main_storage.getIpcFd().?;1001 const saved_metadata = opt_saved_metadata orelse {
986 const saved_metadata = findOld(ipc_fd, old_metadata) orelse {
987 main_storage.completed_count = 0;1002 main_storage.completed_count = 0;
988 main_storage.estimated_total_count = 0;1003 main_storage.estimated_total_count = 0;
1004 if (remaining_read_trash_bytes > 0) {
1005 ipc_metadata_fds[ipc_metadata_len] = Fd.init(fd);
1006 ipc_metadata[ipc_metadata_len] = .{
1007 .remaining_read_trash_bytes = remaining_read_trash_bytes,
1008 .start_index = @intCast(start_serialized_len),
1009 .nodes_len = 0,
1010 .main_index = @intCast(main_index),
1011 };
1012 ipc_metadata_len += 1;
1013 }
989 return start_serialized_len;1014 return start_serialized_len;
990 };1015 };
9911016
...@@ -993,8 +1018,9 @@ fn useSavedIpcData(...@@ -993,8 +1018,9 @@ fn useSavedIpcData(
993 const nodes_len = @min(saved_metadata.nodes_len, serialized_buffer.storage.len - start_serialized_len);1018 const nodes_len = @min(saved_metadata.nodes_len, serialized_buffer.storage.len - start_serialized_len);
994 const old_main_index = saved_metadata.main_index;1019 const old_main_index = saved_metadata.main_index;
9951020
1021 ipc_metadata_fds[ipc_metadata_len] = Fd.init(fd);
996 ipc_metadata[ipc_metadata_len] = .{1022 ipc_metadata[ipc_metadata_len] = .{
997 .ipc_fd = SavedMetadata.setIpcFd(ipc_fd),1023 .remaining_read_trash_bytes = remaining_read_trash_bytes,
998 .start_index = @intCast(start_serialized_len),1024 .start_index = @intCast(start_serialized_len),
999 .nodes_len = nodes_len,1025 .nodes_len = nodes_len,
1000 .main_index = @intCast(main_index),1026 .main_index = @intCast(main_index),
...@@ -1209,6 +1235,11 @@ fn writeIpc(fd: posix.fd_t, serialized: Serialized) error{BrokenPipe}!void {...@@ -1209,6 +1235,11 @@ fn writeIpc(fd: posix.fd_t, serialized: Serialized) error{BrokenPipe}!void {
1209 .{ .base = parents.ptr, .len = parents.len },1235 .{ .base = parents.ptr, .len = parents.len },
1210 };1236 };
12111237
1238 // Ensures the packet can fit in the pipe buffer.
1239 const upper_bound_msg_len = 1 + node_storage_buffer_len * @sizeOf(Node.Storage) +
1240 node_storage_buffer_len * @sizeOf(Node.OptionalIndex);
1241 comptime assert(upper_bound_msg_len <= 4096);
1242
1212 while (remaining_write_trash_bytes > 0) {1243 while (remaining_write_trash_bytes > 0) {
1213 // We do this in a separate write call to give a better chance for the1244 // We do this in a separate write call to give a better chance for the
1214 // writev below to be in a single packet.1245 // writev below to be in a single packet.
...@@ -1228,7 +1259,7 @@ fn writeIpc(fd: posix.fd_t, serialized: Serialized) error{BrokenPipe}!void {...@@ -1228,7 +1259,7 @@ fn writeIpc(fd: posix.fd_t, serialized: Serialized) error{BrokenPipe}!void {
12281259
1229 // If this write would block we do not want to keep trying, but we need to1260 // If this write would block we do not want to keep trying, but we need to
1230 // know if a partial message was written.1261 // know if a partial message was written.
1231 if (posix.writev(fd, &vecs)) |written| {1262 if (writevNonblock(fd, &vecs)) |written| {
1232 const total = header.len + storage.len + parents.len;1263 const total = header.len + storage.len + parents.len;
1233 if (written < total) {1264 if (written < total) {
1234 remaining_write_trash_bytes = total - written;1265 remaining_write_trash_bytes = total - written;
...@@ -1243,6 +1274,23 @@ fn writeIpc(fd: posix.fd_t, serialized: Serialized) error{BrokenPipe}!void {...@@ -1243,6 +1274,23 @@ fn writeIpc(fd: posix.fd_t, serialized: Serialized) error{BrokenPipe}!void {
1243 }1274 }
1244}1275}
12451276
1277fn writevNonblock(fd: posix.fd_t, iov: []posix.iovec_const) posix.WriteError!usize {
1278 var iov_index: usize = 0;
1279 var written: usize = 0;
1280 var total_written: usize = 0;
1281 while (true) {
1282 while (if (iov_index < iov.len)
1283 written >= iov[iov_index].len
1284 else
1285 return total_written) : (iov_index += 1) written -= iov[iov_index].len;
1286 iov[iov_index].base += written;
1287 iov[iov_index].len -= written;
1288 written = try posix.writev(fd, iov[iov_index..]);
1289 if (written == 0) return total_written;
1290 total_written += written;
1291 }
1292}
1293
1246fn maybeUpdateSize(resize_flag: bool) void {1294fn maybeUpdateSize(resize_flag: bool) void {
1247 if (!resize_flag) return;1295 if (!resize_flag) return;
12481296