diff --git a/src/internal/archive.zig b/src/internal/archive.zig index e115fcb..2827598 100644 --- a/src/internal/archive.zig +++ b/src/internal/archive.zig @@ -114,60 +114,61 @@ pub const Archive = struct { const idx = parseSegmentIndex(entry.name) orelse continue; if (max_index == null or idx > max_index.?) max_index = idx; } - const last = max_index orelse return; - self.seg_index = last + 1; + var candidate = max_index orelse return; + self.seg_index = candidate + 1; + while (true) { + var name_buf: [64]u8 = undefined; + const name = formatSegmentName(&name_buf, candidate); + const bytes = self.dir.readFileAlloc(self.io, name, self.allocator, .limited(1 << 31)) catch return; + defer self.allocator.free(bytes); + if (bytes.len < segment.header_size) return; - var name_buf: [64]u8 = undefined; - const name = formatSegmentName(&name_buf, last); - const bytes = self.dir.readFileAlloc(self.io, name, self.allocator, .limited(1 << 31)) catch return; - defer self.allocator.free(bytes); - if (bytes.len < segment.header_size) return; - - if (segment.Header.decode(bytes)) |header| { - // sealed: just carry the seq counter forward - self.next_seq = @max(self.next_seq, header.max_seq + 1); - return; - } else |err| switch (err) { - error.ActiveSegment => {}, - else => { - log.warn("recovery: {s} unreadable ({s}), leaving as-is", .{ name, @errorName(err) }); + if (segment.Header.decode(bytes)) |header| { + self.next_seq = @max(self.next_seq, header.max_seq + 1); return; - }, - } - - // active file: walk frames, keep only fully-valid ones - var valid_end: usize = segment.header_size; - var max_seq: u64 = 0; - var walk: writer_mod.ActiveIterator = .{ .bytes = bytes }; - while (true) { - const before = walk.offset; - var block = (walk.next(self.allocator) catch null) orelse break; - defer block.deinit(self.allocator); - if (block.events.len == 0) break; - valid_end = walk.offset; - _ = before; - max_seq = @max(max_seq, block.events[block.events.len - 1].seq); - } - if (valid_end == segment.header_size) { - // nothing valid — drop the file and reuse the index - self.dir.deleteFile(self.io, name) catch {}; - self.seg_index = last; + } else |err| switch (err) { + error.ActiveSegment => {}, + else => { + log.warn("recovery: {s} unreadable ({s}), leaving as-is", .{ name, @errorName(err) }); + return; + }, + } + + var valid_end: usize = segment.header_size; + var max_seq: u64 = 0; + var walk: writer_mod.ActiveIterator = .{ .bytes = bytes }; + while (true) { + var block = (walk.next(self.allocator) catch null) orelse break; + defer block.deinit(self.allocator); + if (block.events.len == 0) break; + valid_end = walk.offset; + max_seq = @max(max_seq, block.events[block.events.len - 1].seq); + } + if (valid_end == segment.header_size) { + // A normal restart can leave an empty newest active segment. + // Reuse its index, then inspect the preceding sealed segment + // to recover the durable seq floor instead of resetting to 1. + try self.dir.deleteFile(self.io, name); + self.seg_index = candidate; + if (candidate == 0) return; + candidate -= 1; + continue; + } + + var w = try writer_mod.ActiveWriter.init(self.allocator); + defer w.deinit(); + w.buf.clearRetainingCapacity(); + try w.buf.appendSlice(self.allocator, bytes[0..valid_end]); + try w.seal(); + + var f = try self.dir.createFile(self.io, name, .{ .truncate = true }); + defer f.close(self.io); + try f.writeStreamingAll(self.io, w.bytes()); + try f.sync(self.io); + self.next_seq = @max(self.next_seq, max_seq + 1); + log.info("recovery: sealed interrupted segment {s} ({d} bytes valid)", .{ name, valid_end }); return; } - - // rebuild in memory (truncated to valid frames), seal, rewrite - var w = try writer_mod.ActiveWriter.init(self.allocator); - defer w.deinit(); - w.buf.clearRetainingCapacity(); - try w.buf.appendSlice(self.allocator, bytes[0..valid_end]); - try w.seal(); - - var f = try self.dir.createFile(self.io, name, .{ .truncate = true }); - defer f.close(self.io); - try f.writeStreamingAll(self.io, w.bytes()); - try f.sync(self.io); - self.next_seq = @max(self.next_seq, max_seq + 1); - log.info("recovery: sealed interrupted segment {s} ({d} bytes valid)", .{ name, valid_end }); } fn openNextSegment(self: *Archive) !void { @@ -447,6 +448,25 @@ test "archive: append, rotate, recover across restart" { a.close(); } + { + var a = try Archive.init(testing.allocator, io, data_dir); + defer a.deinit(); + // Seals seq 101 and opens the next active segment, but deliberately + // appends nothing to it before another restart. + try testing.expectEqual(@as(u64, 102), a.next_seq); + a.close(); + } + + { + var a = try Archive.init(testing.allocator, io, data_dir); + defer a.deinit(); + // The empty newest active segment must not hide the preceding sealed + // max seq and reset allocation to 1. + try testing.expectEqual(@as(u64, 102), a.next_seq); + const seq = try a.append(testEvent(1000), 1000); + try testing.expectEqual(@as(u64, 102), seq); + } + // every sealed file in the dir must parse with our reader var dir = try Io.Dir.cwd().openDir(io, data_dir, .{}); defer dir.close(io); @@ -467,7 +487,7 @@ test "archive: append, rotate, recover across restart" { total_events += sealed.header.event_count; } try testing.expect(sealed_count >= 2); - try testing.expectEqual(@as(u64, 100), total_events); // event 101 is unflushed-active + try testing.expectEqual(@as(u64, 101), total_events); // event 102 is unflushed-active } test "archive: metadata failure does not publish or strand durability" {