| ... | @@ -1091,87 +1091,7 @@ pub fn indexPack(allocator: Allocator, pack: std.fs.File, index_writer: anytype) | ... | @@ -1091,87 +1091,7 @@ pub fn indexPack(allocator: Allocator, pack: std.fs.File, index_writer: anytype) |
| 1091 | try index_writer.writeAll(&index_checksum); | 1091 | try index_writer.writeAll(&index_checksum); |
| 1092 | } | 1092 | } |
| 1093 | | 1093 | |
| 1094 | /// A reader that stores read data in a growable internal buffer. The read | 1094 | // Performs the first pass over the packfile data for index construction. |
| 1095 | /// position can be rewound to allow previously read data to be read again. | | |
| 1096 | fn AccumulatingReader(comptime ReaderType: type) type { | | |
| 1097 | return struct { | | |
| 1098 | child_reader: ReaderType, | | |
| 1099 | buffer: std.ArrayListUnmanaged(u8) = .{}, | | |
| 1100 | /// The position in `buffer` from which reads should start, returning | | |
| 1101 | /// buffered data. If this is `buffer.items.len`, data will be read from | | |
| 1102 | /// `child_reader` instead. | | |
| 1103 | read_start: usize = 0, | | |
| 1104 | allocator: Allocator, | | |
| 1105 | | | |
| 1106 | const Self = @This(); | | |
| 1107 | | | |
| 1108 | fn deinit(self: *Self) void { | | |
| 1109 | self.buffer.deinit(self.allocator); | | |
| 1110 | self.* = undefined; | | |
| 1111 | } | | |
| 1112 | | | |
| 1113 | const ReadError = ReaderType.Error || Allocator.Error; | | |
| 1114 | const Reader = std.io.Reader(*Self, ReadError, read); | | |
| 1115 | | | |
| 1116 | fn read(self: *Self, buf: []u8) ReadError!usize { | | |
| 1117 | if (self.read_start < self.buffer.items.len) { | | |
| 1118 | // Previously buffered data is available and should be used | | |
| 1119 | // before reading more data from the underlying reader. | | |
| 1120 | const available = self.buffer.items.len - self.read_start; | | |
| 1121 | const count = @min(available, buf.len); | | |
| 1122 | @memcpy(buf[0..count], self.buffer.items[self.read_start..][0..count]); | | |
| 1123 | self.read_start += count; | | |
| 1124 | return count; | | |
| 1125 | } | | |
| 1126 | | | |
| 1127 | try self.buffer.ensureUnusedCapacity(self.allocator, buf.len); | | |
| 1128 | const read_buffer = self.buffer.unusedCapacitySlice(); | | |
| 1129 | const count = try self.child_reader.read(read_buffer[0..buf.len]); | | |
| 1130 | @memcpy(buf[0..count], read_buffer[0..count]); | | |
| 1131 | self.buffer.items.len += count; | | |
| 1132 | self.read_start += count; | | |
| 1133 | return count; | | |
| 1134 | } | | |
| 1135 | | | |
| 1136 | fn reader(self: *Self) Reader { | | |
| 1137 | return .{ .context = self }; | | |
| 1138 | } | | |
| 1139 | | | |
| 1140 | /// Returns a slice of the buffered data that has already been read, | | |
| 1141 | /// except the last `count_before_end` bytes. | | |
| 1142 | fn readDataExcept(self: Self, count_before_end: usize) []const u8 { | | |
| 1143 | assert(count_before_end <= self.read_start); | | |
| 1144 | return self.buffer.items[0 .. self.read_start - count_before_end]; | | |
| 1145 | } | | |
| 1146 | | | |
| 1147 | /// Discards the first `count` bytes of buffered data. | | |
| 1148 | fn discard(self: *Self, count: usize) void { | | |
| 1149 | assert(count <= self.buffer.items.len); | | |
| 1150 | const retain = self.buffer.items.len - count; | | |
| 1151 | mem.copyForwards( | | |
| 1152 | u8, | | |
| 1153 | self.buffer.items[0..retain], | | |
| 1154 | self.buffer.items[count..][0..retain], | | |
| 1155 | ); | | |
| 1156 | self.buffer.items.len = retain; | | |
| 1157 | self.read_start -= @min(self.read_start, count); | | |
| 1158 | } | | |
| 1159 | | | |
| 1160 | /// Rewinds the read position to the beginning of buffered data. | | |
| 1161 | fn rewind(self: *Self) void { | | |
| 1162 | self.read_start = 0; | | |
| 1163 | } | | |
| 1164 | }; | | |
| 1165 | } | | |
| 1166 | | | |
| 1167 | fn accumulatingReader( | | |
| 1168 | allocator: Allocator, | | |
| 1169 | reader: anytype, | | |
| 1170 | ) AccumulatingReader(@TypeOf(reader)) { | | |
| 1171 | return .{ .child_reader = reader, .allocator = allocator }; | | |
| 1172 | } | | |
| 1173 | | | |
| 1174 | /// Performs the first pass over the packfile data for index construction. | | |
| 1175 | /// This will index all non-delta objects, queue delta objects for further | 1095 | /// This will index all non-delta objects, queue delta objects for further |
| 1176 | /// processing, and return the pack checksum (which is part of the index | 1096 | /// processing, and return the pack checksum (which is part of the index |
| 1177 | /// format). | 1097 | /// format). |
| ... | @@ -1181,102 +1101,62 @@ fn indexPackFirstPass( | ... | @@ -1181,102 +1101,62 @@ fn indexPackFirstPass( |
| 1181 | index_entries: *std.AutoHashMapUnmanaged(Oid, IndexEntry), | 1101 | index_entries: *std.AutoHashMapUnmanaged(Oid, IndexEntry), |
| 1182 | pending_deltas: *std.ArrayListUnmanaged(IndexEntry), | 1102 | pending_deltas: *std.ArrayListUnmanaged(IndexEntry), |
| 1183 | ) ![Sha1.digest_length]u8 { | 1103 | ) ![Sha1.digest_length]u8 { |
| 1184 | var pack_buffered_reader = std.io.bufferedReader(pack.reader()); | 1104 | var pack_counting_writer = std.io.countingWriter(std.io.null_writer); |
| 1185 | var pack_accumulating_reader = accumulatingReader(allocator, pack_buffered_reader.reader()); | 1105 | var pack_hashed_writer = std.compress.hashedWriter(pack_counting_writer.writer(), Sha1.init(.{})); |
| 1186 | defer pack_accumulating_reader.deinit(); | 1106 | var entry_crc32_writer = std.compress.hashedWriter(pack_hashed_writer.writer(), std.hash.Crc32.init()); |
| 1187 | var pack_position: usize = 0; | 1107 | var pack_buffered_reader = std.io.bufferedTee(4096, 8, pack.reader(), entry_crc32_writer.writer()); |
| 1188 | var pack_hash = Sha1.init(.{}); | 1108 | const pack_reader = pack_buffered_reader.reader(); |
| 1189 | const pack_reader = pack_accumulating_reader.reader(); | | |
| 1190 | | 1109 | |
| 1191 | const pack_header = try PackHeader.read(pack_reader); | 1110 | const pack_header = try PackHeader.read(pack_reader); |
| 1192 | const pack_header_bytes = pack_accumulating_reader.readDataExcept(0); | 1111 | try pack_buffered_reader.flush(); |
| 1193 | pack_position += pack_header_bytes.len; | | |
| 1194 | pack_hash.update(pack_header_bytes); | | |
| 1195 | pack_accumulating_reader.discard(pack_header_bytes.len); | | |
| 1196 | | 1112 | |
| 1197 | var current_entry: u32 = 0; | 1113 | var current_entry: u32 = 0; |
| 1198 | while (current_entry < pack_header.total_objects) : (current_entry += 1) { | 1114 | while (current_entry < pack_header.total_objects) : (current_entry += 1) { |
| 1199 | const entry_offset = pack_position; | 1115 | const entry_offset = pack_counting_writer.bytes_written; |
| 1200 | var entry_crc32 = std.hash.Crc32.init(); | 1116 | entry_crc32_writer.hasher = std.hash.Crc32.init(); // reset hasher |
| 1201 | | | |
| 1202 | const entry_header = try EntryHeader.read(pack_reader); | 1117 | const entry_header = try EntryHeader.read(pack_reader); |
| 1203 | const entry_header_bytes = pack_accumulating_reader.readDataExcept(0); | | |
| 1204 | pack_position += entry_header_bytes.len; | | |
| 1205 | pack_hash.update(entry_header_bytes); | | |
| 1206 | entry_crc32.update(entry_header_bytes); | | |
| 1207 | pack_accumulating_reader.discard(entry_header_bytes.len); | | |
| 1208 | | 1118 | |
| 1209 | switch (entry_header) { | 1119 | switch (entry_header) { |
| 1210 | .commit, .tree, .blob, .tag => |object| { | 1120 | inline .commit, .tree, .blob, .tag => |object, tag| { |
| 1211 | var entry_decompress_stream = std.compress.zlib.decompressor(pack_reader); | 1121 | var entry_decompress_stream = std.compress.zlib.decompressor(pack_reader); |
| 1212 | var entry_data_size: usize = 0; | 1122 | var entry_counting_reader = std.io.countingReader(entry_decompress_stream.reader()); |
| 1213 | var entry_hashed_writer = hashedWriter(std.io.null_writer, Sha1.init(.{})); | 1123 | var entry_hashed_writer = hashedWriter(std.io.null_writer, Sha1.init(.{})); |
| 1214 | const entry_writer = entry_hashed_writer.writer(); | 1124 | const entry_writer = entry_hashed_writer.writer(); |
| 1215 | | | |
| 1216 | // The object header is not included in the pack data but is | 1125 | // The object header is not included in the pack data but is |
| 1217 | // part of the object's ID | 1126 | // part of the object's ID |
| 1218 | try entry_writer.print("{s} {}\x00", .{ @tagName(entry_header), object.uncompressed_length }); | 1127 | try entry_writer.print("{s} {}\x00", .{ @tagName(tag), object.uncompressed_length }); |
| 1219 | | 1128 | var fifo = std.fifo.LinearFifo(u8, .{ .Static = 4096 }).init(); |
| 1220 | while (try entry_decompress_stream.next()) |decompressed_data| { | 1129 | try fifo.pump(entry_counting_reader.reader(), entry_writer); |
| 1221 | entry_data_size += decompressed_data.len; | 1130 | if (entry_counting_reader.bytes_read != object.uncompressed_length) { |
| 1222 | try entry_writer.writeAll(decompressed_data); | | |
| 1223 | | | |
| 1224 | const compressed_bytes = pack_accumulating_reader.readDataExcept(entry_decompress_stream.unreadBytes()); | | |
| 1225 | pack_position += compressed_bytes.len; | | |
| 1226 | pack_hash.update(compressed_bytes); | | |
| 1227 | entry_crc32.update(compressed_bytes); | | |
| 1228 | pack_accumulating_reader.discard(compressed_bytes.len); | | |
| 1229 | } | | |
| 1230 | const footer_bytes = pack_accumulating_reader.readDataExcept(entry_decompress_stream.unreadBytes()); | | |
| 1231 | pack_position += footer_bytes.len; | | |
| 1232 | pack_hash.update(footer_bytes); | | |
| 1233 | entry_crc32.update(footer_bytes); | | |
| 1234 | pack_accumulating_reader.discard(footer_bytes.len); | | |
| 1235 | pack_accumulating_reader.rewind(); | | |
| 1236 | | | |
| 1237 | if (entry_data_size != object.uncompressed_length) { | | |
| 1238 | return error.InvalidObject; | 1131 | return error.InvalidObject; |
| 1239 | } | 1132 | } |
| 1240 | | | |
| 1241 | const oid = entry_hashed_writer.hasher.finalResult(); | 1133 | const oid = entry_hashed_writer.hasher.finalResult(); |
| | 1134 | pack_buffered_reader.putBack(entry_decompress_stream.unreadBytes()); |
| | 1135 | try pack_buffered_reader.flush(); |
| 1242 | try index_entries.put(allocator, oid, .{ | 1136 | try index_entries.put(allocator, oid, .{ |
| 1243 | .offset = entry_offset, | 1137 | .offset = entry_offset, |
| 1244 | .crc32 = entry_crc32.final(), | 1138 | .crc32 = entry_crc32_writer.hasher.final(), |
| 1245 | }); | 1139 | }); |
| 1246 | }, | 1140 | }, |
| 1247 | inline .ofs_delta, .ref_delta => |delta| { | 1141 | inline .ofs_delta, .ref_delta => |delta| { |
| 1248 | var entry_decompress_stream = std.compress.zlib.decompressor(pack_reader); | 1142 | var entry_decompress_stream = std.compress.zlib.decompressor(pack_reader); |
| 1249 | var entry_data_size: usize = 0; | 1143 | var entry_counting_reader = std.io.countingReader(entry_decompress_stream.reader()); |
| 1250 | | 1144 | var fifo = std.fifo.LinearFifo(u8, .{ .Static = 4096 }).init(); |
| 1251 | while (try entry_decompress_stream.next()) |decompressed_data| { | 1145 | try fifo.pump(entry_counting_reader.reader(), std.io.null_writer); |
| 1252 | entry_data_size += decompressed_data.len; | 1146 | if (entry_counting_reader.bytes_read != delta.uncompressed_length) { |
| 1253 | | | |
| 1254 | const compressed_bytes = pack_accumulating_reader.readDataExcept(entry_decompress_stream.unreadBytes()); | | |
| 1255 | pack_position += compressed_bytes.len; | | |
| 1256 | pack_hash.update(compressed_bytes); | | |
| 1257 | entry_crc32.update(compressed_bytes); | | |
| 1258 | pack_accumulating_reader.discard(compressed_bytes.len); | | |
| 1259 | } | | |
| 1260 | const footer_bytes = pack_accumulating_reader.readDataExcept(entry_decompress_stream.unreadBytes()); | | |
| 1261 | pack_position += footer_bytes.len; | | |
| 1262 | pack_hash.update(footer_bytes); | | |
| 1263 | entry_crc32.update(footer_bytes); | | |
| 1264 | pack_accumulating_reader.discard(footer_bytes.len); | | |
| 1265 | pack_accumulating_reader.rewind(); | | |
| 1266 | | | |
| 1267 | if (entry_data_size != delta.uncompressed_length) { | | |
| 1268 | return error.InvalidObject; | 1147 | return error.InvalidObject; |
| 1269 | } | 1148 | } |
| 1270 | | 1149 | pack_buffered_reader.putBack(entry_decompress_stream.unreadBytes()); |
| | 1150 | try pack_buffered_reader.flush(); |
| 1271 | try pending_deltas.append(allocator, .{ | 1151 | try pending_deltas.append(allocator, .{ |
| 1272 | .offset = entry_offset, | 1152 | .offset = entry_offset, |
| 1273 | .crc32 = entry_crc32.final(), | 1153 | .crc32 = entry_crc32_writer.hasher.final(), |
| 1274 | }); | 1154 | }); |
| 1275 | }, | 1155 | }, |
| 1276 | } | 1156 | } |
| 1277 | } | 1157 | } |
| 1278 | | 1158 | |
| 1279 | const pack_checksum = pack_hash.finalResult(); | 1159 | const pack_checksum = pack_hashed_writer.hasher.finalResult(); |
| 1280 | const recorded_checksum = try pack_reader.readBytesNoEof(Sha1.digest_length); | 1160 | const recorded_checksum = try pack_reader.readBytesNoEof(Sha1.digest_length); |
| 1281 | if (!mem.eql(u8, &pack_checksum, &recorded_checksum)) { | 1161 | if (!mem.eql(u8, &pack_checksum, &recorded_checksum)) { |
| 1282 | return error.CorruptedPack; | 1162 | return error.CorruptedPack; |