| ... | @@ -931,7 +931,8 @@ pub const WriteFileError = PReadError || WriteError; | ... | @@ -931,7 +931,8 @@ pub const WriteFileError = PReadError || WriteError; |
| 931 | | 931 | |
| 932 | pub fn writeFileAll(self: File, in_file: File, options: BufferedWriter.WriteFileOptions) WriteFileError!void { | 932 | pub fn writeFileAll(self: File, in_file: File, options: BufferedWriter.WriteFileOptions) WriteFileError!void { |
| 933 | var file_writer = self.writer(); | 933 | var file_writer = self.writer(); |
| 934 | var bw = file_writer.interface().buffered(&.{}); | 934 | var buffer: [2000]u8 = undefined; |
| | 935 | var bw = file_writer.interface().buffered(&buffer); |
| 935 | bw.writeFileAll(in_file, options) catch |err| switch (err) { | 936 | bw.writeFileAll(in_file, options) catch |err| switch (err) { |
| 936 | error.WriteFailed => if (file_writer.err) |_| unreachable else |e| return e, | 937 | error.WriteFailed => if (file_writer.err) |_| unreachable else |e| return e, |
| 937 | else => |e| return e, | 938 | else => |e| return e, |
| ... | @@ -940,14 +941,27 @@ pub fn writeFileAll(self: File, in_file: File, options: BufferedWriter.WriteFile | ... | @@ -940,14 +941,27 @@ pub fn writeFileAll(self: File, in_file: File, options: BufferedWriter.WriteFile |
| 940 | | 941 | |
| 941 | pub const Reader = struct { | 942 | pub const Reader = struct { |
| 942 | file: File, | 943 | file: File, |
| 943 | err: ReadError!void = {}, | 944 | err: ?ReadError = null, |
| 944 | mode: Reader.Mode = .positional, | 945 | mode: Reader.Mode = .positional, |
| 945 | pos: u64 = 0, | 946 | pos: u64 = 0, |
| 946 | size: ?u64 = null, | 947 | size: ?u64 = null, |
| 947 | size_err: GetEndPosError!void = {}, | 948 | size_err: ?GetEndPosError = null, |
| 948 | seek_err: SeekError!void = {}, | 949 | seek_err: ?SeekError = null, |
| 949 | | 950 | |
| 950 | pub const Mode = enum { streaming, positional }; | 951 | pub const Mode = enum { |
| | 952 | streaming, |
| | 953 | positional, |
| | 954 | streaming_reading, |
| | 955 | positional_reading, |
| | 956 | |
| | 957 | pub fn toStreaming(m: @This()) @This() { |
| | 958 | return switch (m) { |
| | 959 | .positional => .streaming, |
| | 960 | .positional_reading => .streaming_reading, |
| | 961 | else => unreachable, |
| | 962 | }; |
| | 963 | } |
| | 964 | }; |
| 951 | | 965 | |
| 952 | pub fn interface(r: *Reader) std.io.Reader { | 966 | pub fn interface(r: *Reader) std.io.Reader { |
| 953 | return .{ | 967 | return .{ |
| ... | @@ -975,7 +989,7 @@ pub const Reader = struct { | ... | @@ -975,7 +989,7 @@ pub const Reader = struct { |
| 975 | switch (r.mode) { | 989 | switch (r.mode) { |
| 976 | .positional => { | 990 | .positional => { |
| 977 | const size = r.size orelse { | 991 | const size = r.size orelse { |
| 978 | if (r.file.getEndPos()) |size| { | 992 | if (file.getEndPos()) |size| { |
| 979 | r.size = size; | 993 | r.size = size; |
| 980 | } else |err| { | 994 | } else |err| { |
| 981 | r.size_err = err; | 995 | r.size_err = err; |
| ... | @@ -991,6 +1005,10 @@ pub const Reader = struct { | ... | @@ -991,6 +1005,10 @@ pub const Reader = struct { |
| 991 | assert(pos == 0); | 1005 | assert(pos == 0); |
| 992 | return 0; | 1006 | return 0; |
| 993 | }, | 1007 | }, |
| | 1008 | error.Unimplemented => { |
| | 1009 | r.mode = .positional_reading; |
| | 1010 | return 0; |
| | 1011 | }, |
| 994 | else => |e| { | 1012 | else => |e| { |
| 995 | r.err = e; | 1013 | r.err = e; |
| 996 | return error.ReadFailed; | 1014 | return error.ReadFailed; |
| ... | @@ -1003,12 +1021,45 @@ pub const Reader = struct { | ... | @@ -1003,12 +1021,45 @@ pub const Reader = struct { |
| 1003 | const n = bw.writeFile(file, .none, limit, &.{}, 0) catch |err| switch (err) { | 1021 | const n = bw.writeFile(file, .none, limit, &.{}, 0) catch |err| switch (err) { |
| 1004 | error.WriteFailed => return error.WriteFailed, | 1022 | error.WriteFailed => return error.WriteFailed, |
| 1005 | error.Unseekable => unreachable, // Passing `Offset.none`. | 1023 | error.Unseekable => unreachable, // Passing `Offset.none`. |
| | 1024 | error.Unimplemented => { |
| | 1025 | r.mode = .streaming_reading; |
| | 1026 | return 0; |
| | 1027 | }, |
| | 1028 | else => |e| { |
| | 1029 | r.err = e; |
| | 1030 | return error.ReadFailed; |
| | 1031 | }, |
| | 1032 | }; |
| | 1033 | r.pos = pos + n; |
| | 1034 | return n; |
| | 1035 | }, |
| | 1036 | .positional_reading => { |
| | 1037 | const dest = limit.slice(try bw.writableSliceGreedy(1)); |
| | 1038 | const n = file.pread(dest, pos) catch |err| switch (err) { |
| | 1039 | error.Unseekable => { |
| | 1040 | r.mode = .streaming_reading; |
| | 1041 | assert(pos == 0); |
| | 1042 | return 0; |
| | 1043 | }, |
| 1006 | else => |e| { | 1044 | else => |e| { |
| 1007 | r.err = e; | 1045 | r.err = e; |
| 1008 | return error.ReadFailed; | 1046 | return error.ReadFailed; |
| 1009 | }, | 1047 | }, |
| 1010 | }; | 1048 | }; |
| | 1049 | if (n == 0) return error.EndOfStream; |
| | 1050 | r.pos = pos + n; |
| | 1051 | bw.advance(n); |
| | 1052 | return n; |
| | 1053 | }, |
| | 1054 | .streaming_reading => { |
| | 1055 | const dest = limit.slice(try bw.writableSliceGreedy(1)); |
| | 1056 | const n = file.read(dest) catch |err| { |
| | 1057 | r.err = err; |
| | 1058 | return error.ReadFailed; |
| | 1059 | }; |
| | 1060 | if (n == 0) return error.EndOfStream; |
| 1011 | r.pos = pos + n; | 1061 | r.pos = pos + n; |
| | 1062 | bw.advance(n); |
| 1012 | return n; | 1063 | return n; |
| 1013 | }, | 1064 | }, |
| 1014 | } | 1065 | } |
| ... | @@ -1020,7 +1071,7 @@ pub const Reader = struct { | ... | @@ -1020,7 +1071,7 @@ pub const Reader = struct { |
| 1020 | const pos = r.pos; | 1071 | const pos = r.pos; |
| 1021 | | 1072 | |
| 1022 | switch (r.mode) { | 1073 | switch (r.mode) { |
| 1023 | .positional => { | 1074 | .positional, .positional_reading => { |
| 1024 | if (is_windows) { | 1075 | if (is_windows) { |
| 1025 | // Unfortunately, `ReadFileScatter` cannot be used since it requires | 1076 | // Unfortunately, `ReadFileScatter` cannot be used since it requires |
| 1026 | // page alignment, so we are stuck using only the first slice. | 1077 | // page alignment, so we are stuck using only the first slice. |
| ... | @@ -1053,7 +1104,7 @@ pub const Reader = struct { | ... | @@ -1053,7 +1104,7 @@ pub const Reader = struct { |
| 1053 | if (send_vecs.len == 0) return 0; // Prevent false positive end detection on empty `data`. | 1104 | if (send_vecs.len == 0) return 0; // Prevent false positive end detection on empty `data`. |
| 1054 | const n = posix.preadv(handle, send_vecs, pos) catch |err| switch (err) { | 1105 | const n = posix.preadv(handle, send_vecs, pos) catch |err| switch (err) { |
| 1055 | error.Unseekable => { | 1106 | error.Unseekable => { |
| 1056 | r.mode = .streaming; | 1107 | r.mode = r.mode.toStreaming(); |
| 1057 | assert(pos == 0); | 1108 | assert(pos == 0); |
| 1058 | return 0; | 1109 | return 0; |
| 1059 | }, | 1110 | }, |
| ... | @@ -1066,7 +1117,7 @@ pub const Reader = struct { | ... | @@ -1066,7 +1117,7 @@ pub const Reader = struct { |
| 1066 | r.pos = pos + n; | 1117 | r.pos = pos + n; |
| 1067 | return n; | 1118 | return n; |
| 1068 | }, | 1119 | }, |
| 1069 | .streaming => { | 1120 | .streaming, .streaming_reading => { |
| 1070 | if (is_windows) { | 1121 | if (is_windows) { |
| 1071 | // Unfortunately, `ReadFileScatter` cannot be used since it requires | 1122 | // Unfortunately, `ReadFileScatter` cannot be used since it requires |
| 1072 | // page alignment, so we are stuck using only the first slice. | 1123 | // page alignment, so we are stuck using only the first slice. |
| ... | @@ -1113,13 +1164,13 @@ pub const Reader = struct { | ... | @@ -1113,13 +1164,13 @@ pub const Reader = struct { |
| 1113 | const file = r.file; | 1164 | const file = r.file; |
| 1114 | const pos = r.pos; | 1165 | const pos = r.pos; |
| 1115 | switch (r.mode) { | 1166 | switch (r.mode) { |
| 1116 | .positional => { | 1167 | .positional, .positional_reading => { |
| 1117 | const size = r.size orelse { | 1168 | const size = r.size orelse { |
| 1118 | if (file.getEndPos()) |size| { | 1169 | if (file.getEndPos()) |size| { |
| 1119 | r.size = size; | 1170 | r.size = size; |
| 1120 | } else |err| { | 1171 | } else |err| { |
| 1121 | r.size_err = err; | 1172 | r.size_err = err; |
| 1122 | r.mode = .streaming; | 1173 | r.mode = r.mode.toStreaming(); |
| 1123 | } | 1174 | } |
| 1124 | return 0; | 1175 | return 0; |
| 1125 | }; | 1176 | }; |
| ... | @@ -1127,17 +1178,13 @@ pub const Reader = struct { | ... | @@ -1127,17 +1178,13 @@ pub const Reader = struct { |
| 1127 | r.pos = pos + delta; | 1178 | r.pos = pos + delta; |
| 1128 | return delta; | 1179 | return delta; |
| 1129 | }, | 1180 | }, |
| 1130 | .streaming => { | 1181 | .streaming, .streaming_reading => { |
| 1131 | // Unfortunately we can't seek forward without knowing the | 1182 | // Unfortunately we can't seek forward without knowing the |
| 1132 | // size because the seek syscalls provided to us will not | 1183 | // size because the seek syscalls provided to us will not |
| 1133 | // return the true end position if a seek would exceed the | 1184 | // return the true end position if a seek would exceed the |
| 1134 | // end. | 1185 | // end. |
| 1135 | fallback: { | 1186 | fallback: { |
| 1136 | if (r.size_err) |_| { | 1187 | if (r.size_err == null and r.seek_err == null) break :fallback; |
| 1137 | if (r.seek_err) |_| { | | |
| 1138 | break :fallback; | | |
| 1139 | } else |_| {} | | |
| 1140 | } else |_| {} | | |
| 1141 | var trash_buffer: [std.atomic.cache_line]u8 = undefined; | 1188 | var trash_buffer: [std.atomic.cache_line]u8 = undefined; |
| 1142 | const trash = &trash_buffer; | 1189 | const trash = &trash_buffer; |
| 1143 | if (is_windows) { | 1190 | if (is_windows) { |
| ... | @@ -1188,9 +1235,16 @@ pub const Writer = struct { | ... | @@ -1188,9 +1235,16 @@ pub const Writer = struct { |
| 1188 | err: WriteError!void = {}, | 1235 | err: WriteError!void = {}, |
| 1189 | mode: Writer.Mode = .positional, | 1236 | mode: Writer.Mode = .positional, |
| 1190 | pos: u64 = 0, | 1237 | pos: u64 = 0, |
| | 1238 | sendfile_err: ?SendfileError = null, |
| | 1239 | read_err: ?ReadError = null, |
| 1191 | | 1240 | |
| 1192 | pub const Mode = Reader.Mode; | 1241 | pub const Mode = Reader.Mode; |
| 1193 | | 1242 | |
| | 1243 | pub const SendfileError = error{ |
| | 1244 | UnsupportedOperation, |
| | 1245 | Unexpected, |
| | 1246 | }; |
| | 1247 | |
| 1194 | /// Number of slices to store on the stack, when trying to send as many byte | 1248 | /// Number of slices to store on the stack, when trying to send as many byte |
| 1195 | /// vectors through the underlying write calls as possible. | 1249 | /// vectors through the underlying write calls as possible. |
| 1196 | const max_buffers_len = 16; | 1250 | const max_buffers_len = 16; |
| ... | @@ -1269,75 +1323,38 @@ pub const Writer = struct { | ... | @@ -1269,75 +1323,38 @@ pub const Writer = struct { |
| 1269 | const w: *Writer = @ptrCast(@alignCast(context)); | 1323 | const w: *Writer = @ptrCast(@alignCast(context)); |
| 1270 | const out_fd = w.file.handle; | 1324 | const out_fd = w.file.handle; |
| 1271 | const in_fd = in_file.handle; | 1325 | const in_fd = in_file.handle; |
| 1272 | const len_int = switch (in_limit) { | 1326 | // TODO try using copy_file_range on Linux |
| 1273 | .nothing => return writeSplat(context, headers_and_trailers, 1), | 1327 | // TODO try using copy_file_range on FreeBSD |
| 1274 | .unlimited => 0, | 1328 | // TODO try using sendfile on macOS |
| 1275 | else => in_limit.toInt().?, | 1329 | // TODO try using sendfile on FreeBSD |
| 1276 | }; | 1330 | if (native_os == .linux and w.mode == .streaming) sf: { |
| 1277 | // TODO try using copy_file_range on linux | 1331 | // Try using sendfile on Linux. |
| 1278 | // TODO try using copy_file_range on freebsd | 1332 | if (w.sendfile_err != null) break :sf; |
| 1279 | if (native_os == .linux) sf: { | | |
| 1280 | // Linux sendfile does not support headers or trailers but it does | 1333 | // Linux sendfile does not support headers or trailers but it does |
| 1281 | // support a streaming read from in_file. | 1334 | // support a streaming read from in_file. |
| 1282 | if (headers_len > 0) return writeSplat(context, headers_and_trailers[0..headers_len], 1); | 1335 | if (headers_len > 0) return writeSplat(context, headers_and_trailers[0..headers_len], 1); |
| 1283 | const max_count = 0x7ffff000; // Avoid EINVAL. | 1336 | const max_count = 0x7ffff000; // Avoid EINVAL. |
| 1284 | const smaller_len = if (len_int == 0) max_count else @min(len_int, max_count); | 1337 | const smaller_len = in_limit.minInt(max_count); |
| 1285 | var off: std.os.linux.off_t = undefined; | 1338 | var off: std.os.linux.off_t = undefined; |
| 1286 | const off_ptr: ?*std.os.linux.off_t = if (in_offset.toInt()) |offset| b: { | 1339 | const off_ptr: ?*std.os.linux.off_t = if (in_offset.toInt()) |offset| b: { |
| 1287 | off = std.math.cast(std.os.linux.off_t, offset) orelse | 1340 | off = std.math.cast(std.os.linux.off_t, offset) orelse |
| 1288 | return writeSplat(context, headers_and_trailers, 1); | 1341 | return writeSplat(context, headers_and_trailers, 1); |
| 1289 | break :b &off; | 1342 | break :b &off; |
| 1290 | } else null; | 1343 | } else null; |
| 1291 | if (true) @panic("TODO"); | | |
| 1292 | const n = std.os.linux.wrapped.sendfile(out_fd, in_fd, off_ptr, smaller_len) catch |err| switch (err) { | 1344 | const n = std.os.linux.wrapped.sendfile(out_fd, in_fd, off_ptr, smaller_len) catch |err| switch (err) { |
| 1293 | error.UnsupportedOperation => break :sf, | 1345 | // Errors that imply sendfile should be avoided on the next write. |
| 1294 | error.Unseekable => break :sf, | 1346 | error.UnsupportedOperation, |
| 1295 | error.Unexpected => break :sf, | 1347 | error.Unexpected, |
| | 1348 | => |e| { |
| | 1349 | w.sendfile_err = e; |
| | 1350 | break :sf; |
| | 1351 | }, |
| 1296 | else => |e| return e, | 1352 | else => |e| return e, |
| 1297 | }; | 1353 | }; |
| 1298 | if (in_offset.toInt()) |offset| { | 1354 | w.pos += n; |
| 1299 | assert(n == off - offset); | | |
| 1300 | } else if (n == 0 and len_int == 0) { | | |
| 1301 | // The caller wouldn't be able to tell that the file transfer is | | |
| 1302 | // done and would incorrectly repeat the same call. | | |
| 1303 | return writeSplat(context, headers_and_trailers, 1); | | |
| 1304 | } | | |
| 1305 | return n; | 1355 | return n; |
| 1306 | } | 1356 | } |
| 1307 | var iovecs_buffer: [max_buffers_len]std.posix.iovec_const = undefined; | 1357 | return error.Unimplemented; |
| 1308 | const iovecs = iovecs_buffer[0..@min(iovecs_buffer.len, headers_and_trailers.len)]; | | |
| 1309 | for (iovecs, headers_and_trailers[0..iovecs.len]) |*v, d| v.* = .{ .base = d.ptr, .len = d.len }; | | |
| 1310 | const headers = iovecs[0..@min(headers_len, iovecs.len)]; | | |
| 1311 | const trailers = iovecs[headers.len..]; | | |
| 1312 | const flags = 0; | | |
| 1313 | return posix.sendfile(out_fd, in_fd, in_offset, len_int, headers, trailers, flags) catch |err| switch (err) { | | |
| 1314 | error.Unseekable, | | |
| 1315 | error.FastOpenAlreadyInProgress, | | |
| 1316 | error.MessageTooBig, | | |
| 1317 | error.FileDescriptorNotASocket, | | |
| 1318 | error.NetworkUnreachable, | | |
| 1319 | error.NetworkSubsystemFailed, | | |
| 1320 | => return writeFileUnseekable(out_fd, in_fd, in_offset, in_limit, headers_and_trailers, headers_len), | | |
| 1321 | | | |
| 1322 | else => |e| return e, | | |
| 1323 | }; | | |
| 1324 | } | | |
| 1325 | | | |
| 1326 | fn writeFileUnseekable( | | |
| 1327 | out_fd: Handle, | | |
| 1328 | in_fd: Handle, | | |
| 1329 | in_offset: u64, | | |
| 1330 | in_limit: std.io.Writer.Limit, | | |
| 1331 | headers_and_trailers: []const []const u8, | | |
| 1332 | headers_len: usize, | | |
| 1333 | ) std.io.Writer.FileError!usize { | | |
| 1334 | _ = out_fd; | | |
| 1335 | _ = in_fd; | | |
| 1336 | _ = in_offset; | | |
| 1337 | _ = in_limit; | | |
| 1338 | _ = headers_and_trailers; | | |
| 1339 | _ = headers_len; | | |
| 1340 | @panic("TODO writeFileUnseekable"); | | |
| 1341 | } | 1358 | } |
| 1342 | }; | 1359 | }; |
| 1343 | | 1360 | |