tiny.sql.Pager
Defined in pager.
API (51)
Actions
Public operations.
appendWalbaseGenerationbasePageCountbeginReadcanRestorecheckpointcommitCheckpointcommitDurableCheckpointcommitStagedWalcurrentViewdatabasePageCountdeinitframeCountinitinstallBaseinstallBaseAtGenerationmarkedPageAt: Returns the imagepageAtreturns, with the check mark stored with it.pageAtpositionprepareCheckpointreleaseDurableBasereplaceWalreplaceWalControlledreserverestorerestoreEpochsetBasePageCountstageWalPagestagedWalPagestagedWalPageIdstagedWalPageMutstorageswapStagedWalFrameswalByteswalCapacityByteswalStagingCapacity
Types and contracts
Public types and contracts.
Fields and members
Public fields and members.
allocatorbasebase_generationbase_indexbase_page_countcheckpoint_once_plancheckpoint_plancheckpoint_serialdatabase_page_countend_markjournalwal_index
Source
Source: lib/sql/src/pager.zig:716
zig
pub const Pager = struct { pub const Workspace = PagerWorkspace; allocator: Allocator, base: std.ArrayList(BaseImage) = .empty, base_index: std.AutoHashMapUnmanaged(u32, usize) = .empty, wal_index: WalIndex, journal: wal.Writer, checkpoint_plan: CheckpointPlan, checkpoint_once_plan: CheckpointPlan, base_generation: u64 = 0, base_page_count: u32 = 0, database_page_count: u32 = 0, end_mark: usize = 0, checkpoint_serial: u64 = 0, pub const Position = struct { journal: wal.Writer.Position, frames_len: usize, database_page_count: u32, end_mark: usize, }; pub const RestoreEpoch = struct { base_generation: u64, base_page_count: u32, checkpoint_serial: u64, }; pub fn init(allocator: Allocator, workspace: *Workspace, options: InitOptions) Error!Pager { const journal_limits: wal.Writer.Limits = .{ .header = options.header, .frames = options.wal_frames, }; const workspace_capacity = try Workspace.Capacity.derive(options); if (workspace.journal.len < workspace_capacity.journal.storage_bytes or workspace.checkpoint.len < workspace_capacity.checkpoint.storage_bytes or workspace.checkpoint_once.len < workspace_capacity.checkpoint_once.storage_bytes or workspace.wal_index.len < workspace_capacity.wal_index.storage_bytes) { return error.StorageTooShort; } var journal = wal.Writer.init(workspace.journal, journal_limits) catch |err| switch (err) { error.CapacityOverflow, error.StorageTooShort => unreachable, }; var checkpoint_plan = CheckpointPlan.init(workspace.checkpoint, .{ .pages = options.wal_frames, }) catch |err| switch (err) { error.CapacityOverflow, error.StorageTooShort => unreachable, }; var checkpoint_once_plan = CheckpointPlan.init(workspace.checkpoint_once, .{ .pages = options.wal_frames, }) catch |err| switch (err) { error.CapacityOverflow, error.StorageTooShort => unreachable, }; var wal_index = WalIndex.init(workspace.wal_index, .{ .frames = options.wal_frames, }) catch |err| switch (err) { error.CapacityOverflow, error.StorageTooShort => unreachable, }; journal.activate(); checkpoint_plan.activate(); checkpoint_once_plan.activate(); wal_index.activate(); workspace.* = undefined; return .{ .allocator = allocator, .wal_index = wal_index, .journal = journal, .checkpoint_plan = checkpoint_plan, .checkpoint_once_plan = checkpoint_once_plan, }; } pub fn deinit(self: *Pager) Workspace { self.base.deinit(self.allocator); self.base_index.deinit(self.allocator); const workspace: Workspace = .{ .journal = self.journal.deinit(), .checkpoint = self.checkpoint_plan.deinit(), .checkpoint_once = self.checkpoint_once_plan.deinit(), .wal_index = self.wal_index.deinit(), }; self.* = undefined; return workspace; } pub fn reserve(self: *Pager, capacity: Capacity) Error!void { const phase = trace.scope("pager.reserve"); defer phase.end(); if (capacity.wal_frames > self.journal.remainingFrames()) return error.WalFull; const wal_pages = capacity.wal_pages orelse capacity.wal_frames; try self.wal_index.reserve(capacity.wal_frames, wal_pages); try self.base.ensureTotalCapacityPrecise( self.allocator, try additionalCapacity(self.base.items.len, capacity.base_pages), ); try self.base_index.ensureUnusedCapacity(self.allocator, try hashMapSize(capacity.base_pages)); } pub fn replaceWal(self: *Pager, bytes: []const u8, committed_len: usize) Error!void { return self.replaceWalControlled(bytes, committed_len, .{}) catch |err| switch (err) { error.Interrupted => unreachable, else => return @errorCast(err), }; } pub fn replaceWalControlled( self: *Pager, bytes: []const u8, committed_len: usize, control: wal.Control, ) (Error || error{Interrupted})!void { const phase = trace.scope("pager.replace_wal"); defer phase.end(); if (committed_len < wal.header_size or (committed_len - wal.header_size) % wal.frame_size != 0) return error.InvalidWal; const frames_count = (committed_len - wal.header_size) / wal.frame_size; if (frames_count > self.journal.frameCapacity()) return error.WalFull; std.debug.assert(frames_count <= self.wal_index.frames.capacity); std.debug.assert(frames_count <= self.wal_index.pages.capacity); try self.journal.loadControlled(bytes, committed_len, control); try self.rebuildWalFramesFromJournalControlled(control); trace.progress("pager.replace_wal.complete"); } pub fn installBase(self: *Pager, page_id: u32, image: *const [page.size]u8) Error!void { const phase = trace.scope("pager.install_base"); defer phase.end(); const generation = try self.nextBaseGeneration(); try self.installBaseAtGeneration(page_id, image, generation); trace.progress("pager.install_base.complete"); } pub fn installBaseAtGeneration(self: *Pager, page_id: u32, image: *const [page.size]u8, generation: u64) Error!void { const phase = trace.scope("pager.install_base.generation"); defer phase.end(); if (self.base_index.get(page_id)) |index| { const cached = &self.base.items[index]; if (cached.id == page_id and cached.generation == generation) { cached.bytes = image.*; cached.checked = false; self.base_generation = @max(self.base_generation, generation); self.base_page_count = @max(self.base_page_count, page_id); self.database_page_count = @max(self.database_page_count, page_id); trace.progress("pager.install_base.generation.complete"); return; } } try self.base.ensureUnusedCapacity(self.allocator, 1); try self.base_index.ensureUnusedCapacity(self.allocator, 1); const index = self.base.items.len; self.base.appendAssumeCapacity(.{ .id = page_id, .generation = generation, .bytes = image.*, }); if (self.base_index.get(page_id)) |existing| { if (self.base.items[existing].generation <= generation) self.base_index.putAssumeCapacity(page_id, index); } else { self.base_index.putAssumeCapacity(page_id, index); } self.base_generation = @max(self.base_generation, generation); self.base_page_count = @max(self.base_page_count, page_id); self.database_page_count = @max(self.database_page_count, page_id); trace.progress("pager.install_base.generation.complete"); } pub fn appendWal(self: *Pager, page_id: u32, db_page_count: u32, image: *const [page.size]u8) Error!void { const phase = trace.scope("pager.append_wal"); defer phase.end(); if (self.journal.remainingFrames() == 0) return error.WalFull; const page_index = self.lowerBoundWalPage(page_id); const existing = page_index < self.wal_index.pages.items.len and self.wal_index.pages.items[page_index].page_id == page_id; try self.wal_index.reserve(1, @intFromBool(!existing)); try self.journal.append(page_id, db_page_count, image); self.indexAppendedWalFrame(page_id, db_page_count, page_index, existing); trace.progress("pager.append_wal.complete"); } pub fn walStagingCapacity(self: *const Pager, position_value: Position) usize { std.debug.assert(std.meta.eql(self.position(), position_value)); return self.journal.remainingFrames(); } pub fn stageWalPage(self: *Pager, position_value: Position, index: usize, page_id: u32, image: *const [page.size]u8) Error!void { std.debug.assert(std.meta.eql(self.position(), position_value)); try self.journal.stagePage(position_value.journal, index, page_id, image); } pub fn stagedWalPage(self: *const Pager, position_value: Position, index: usize) []const u8 { std.debug.assert(std.meta.eql(self.position(), position_value)); return self.journal.stagedPage(position_value.journal, index); } pub fn stagedWalPageMut(self: *Pager, position_value: Position, index: usize) *[page.size]u8 { std.debug.assert(std.meta.eql(self.position(), position_value)); return self.journal.stagedPageMut(position_value.journal, index); } pub fn stagedWalPageId(self: *const Pager, position_value: Position, index: usize) u32 { return self.journal.stagedPageId(position_value.journal, index); } pub fn swapStagedWalFrames(self: *Pager, position_value: Position, left_index: usize, right_index: usize) void { std.debug.assert(std.meta.eql(self.position(), position_value)); self.journal.swapStagedFrames(position_value.journal, left_index, right_index); } pub fn commitStagedWal(self: *Pager, position_value: Position, count: usize, database_page_count: u32) Error!void { const phase = trace.scope("pager.commit_staged_wal"); defer phase.end(); std.debug.assert(std.meta.eql(self.position(), position_value)); if (count > self.walStagingCapacity(position_value)) return error.WalFull; var new_pages: usize = 0; var previous_page_id: u32 = 0; for (0..count) |index| { const page_id = self.stagedWalPageId(position_value, index); std.debug.assert(page_id != 0); if (index != 0) std.debug.assert(page_id > previous_page_id); previous_page_id = page_id; const page_index = self.lowerBoundWalPage(page_id); if (page_index == self.wal_index.pages.items.len or self.wal_index.pages.items[page_index].page_id != page_id) { new_pages += 1; } } try self.wal_index.reserve(count, new_pages); for (0..count) |index| { const page_id = self.stagedWalPageId(position_value, index); const page_index = self.lowerBoundWalPage(page_id); const existing = page_index < self.wal_index.pages.items.len and self.wal_index.pages.items[page_index].page_id == page_id; self.journal.commitStagedFrame(position_value.journal, index, if (index + 1 == count) database_page_count else 0); self.indexAppendedWalFrame(page_id, if (index + 1 == count) database_page_count else 0, page_index, existing); } trace.progress("pager.commit_staged_wal.complete"); } pub fn beginRead(self: *const Pager) Error!Snapshot { const phase = trace.scope("pager.begin_read"); defer phase.end(); return .{ .pager = self, .view = try self.currentView() }; } pub fn checkpoint(self: *Pager, options: CheckpointOptions) Error!Checkpoint { const phase = trace.scope("pager.checkpoint"); defer phase.end(); var prepared = try self.prepareCheckpointWithPlan( &self.checkpoint_once_plan, options, ); defer prepared.deinit(); return try self.commitCheckpoint(prepared); } pub fn prepareCheckpoint(self: *Pager, options: CheckpointOptions) Error!PreparedCheckpoint { return try self.prepareCheckpointWithPlan(&self.checkpoint_plan, options); } fn prepareCheckpointWithPlan( self: *Pager, plan: *CheckpointPlan, options: CheckpointOptions, ) Error!PreparedCheckpoint { const phase = trace.scope("pager.checkpoint.prepare"); defer phase.end(); const current = try self.currentView(); const has_readers = switch (options.readers) { .none => false, .oldest => true, }; const target_mark = switch (options.readers) { .none => current.end_mark, .oldest => |oldest| @min(oldest.end_mark, current.end_mark), }; const can_rewrite_wal = options.restart_header != null and !has_readers; if (can_rewrite_wal) { const retained_frames = self.frameCount() - target_mark; std.debug.assert(retained_frames <= self.wal_index.frames.capacity); std.debug.assert(retained_frames <= self.wal_index.pages.capacity); } const bytes = self.walBytes(); try plan.begin(self.checkpointPageCount(target_mark, bytes.len)); errdefer plan.release(); const serial = std.math.add(u64, self.checkpoint_serial, 1) catch return error.GenerationOverflow; _ = std.math.add(u64, serial, 1) catch return error.GenerationOverflow; self.checkpoint_serial = serial; try self.collectCheckpointPages(target_mark, bytes.len, plan); const pages = plan.pages().len; var generation = self.base_generation; if (pages > 0) generation = try self.nextBaseGeneration(); const retain_generation = switch (options.readers) { .none => generation, .oldest => |oldest| @min(oldest.base_generation, generation), }; return .{ .pager = self, .plan = plan, .checkpoint = .{ .end_mark = target_mark, .pages = pages, .base_generation = generation, .restarted = can_rewrite_wal, }, .retain_generation = retain_generation, .restart_header = options.restart_header, .has_readers = has_readers, .serial = serial, .state = .{ .position = self.position(), .base_generation = self.base_generation, .base_images = self.base.items.len, .wal_pages = self.wal_index.pages.items.len, }, }; } pub fn commitCheckpoint(self: *Pager, prepared: PreparedCheckpoint) Error!Checkpoint { const phase = trace.scope("pager.checkpoint.commit_memory"); defer phase.end(); if (!self.preparedCheckpointCurrent(prepared)) return error.StaleCheckpoint; const checkpoint_value = prepared.checkpoint; if (checkpoint_value.pages > 0) { try self.installCheckpointPages(prepared, checkpoint_value.base_generation); self.base_generation = checkpoint_value.base_generation; } _ = try self.compactBaseHistory(prepared.retain_generation); self.finishCheckpointFrames(prepared); self.checkpoint_serial = prepared.serial + 1; trace.progress("pager.checkpoint.complete"); return checkpoint_value; } pub fn commitDurableCheckpoint(self: *Pager, prepared: PreparedCheckpoint) Error!Checkpoint { const phase = trace.scope("pager.checkpoint.commit_durable"); defer phase.end(); if (!self.preparedCheckpointCurrent(prepared)) return error.StaleCheckpoint; if (prepared.has_readers) return error.DurableCheckpointReaders; const checkpoint_value = prepared.checkpoint; if (checkpoint_value.pages > 0) { self.base_generation = checkpoint_value.base_generation; for (prepared.plan.pages()) |checkpoint_page| { self.base_page_count = @max(self.base_page_count, checkpoint_page.page_id); self.database_page_count = @max(self.database_page_count, checkpoint_page.page_id); } } self.releaseDurableBase(); self.finishCheckpointFrames(prepared); self.checkpoint_serial = prepared.serial + 1; trace.progress("pager.checkpoint.complete"); return checkpoint_value; } pub fn storage(self: *const Pager) Storage { return .{ .base_images = self.base.items.len, .base_capacity = self.base.capacity, .wal_frames = self.wal_index.frames.items.len, .wal_frame_capacity = self.wal_index.frames.capacity, .wal_pages = self.wal_index.pages.items.len, .wal_page_capacity = self.wal_index.pages.capacity, }; } pub fn releaseDurableBase(self: *Pager) void { self.base.deinit(self.allocator); self.base_index.deinit(self.allocator); self.base = .empty; self.base_index = .empty; } pub fn currentView(self: *const Pager) Error!View { return .{ .base_generation = self.base_generation, .end_mark = self.end_mark, }; } pub fn pageAt(self: *const Pager, page_id: u32, view: View) Error!?[]const u8 { const location = try self.locateImage(page_id, view) orelse return null; return switch (location) { .wal => |index| self.walImageBytes(index), .tail => |bytes| bytes, .base => |index| self.base.items[index].bytes[0..], }; } /// Returns the image `pageAt` returns, with the check mark stored with /// it. The pager only clears a mark, when it stores new bytes under it. /// A reader sets it after the image passes the reader's checks, so later /// readers can skip them. pub fn markedPageAt(self: *Pager, page_id: u32, view: View) Error!?MarkedImage { const location = try self.locateImage(page_id, view) orelse return null; return switch (location) { .wal => |index| .{ .bytes = self.walImageBytes(index), .checked = &self.wal_index.frames.items[index].checked, }, .tail => |bytes| .{ .bytes = bytes, .checked = null }, .base => |index| .{ .bytes = &self.base.items[index].bytes, .checked = &self.base.items[index].checked, }, }; } fn locateImage(self: *const Pager, page_id: u32, view: View) Error!?ImageLocation { const phase = trace.scope("pager.page_at"); defer phase.end(); if (self.indexedWalImage(page_id, view.end_mark)) |index| return .{ .wal = index }; if (view.end_mark > self.frameCount()) { if (try wal.pageAt(self.walBytes(), page_id, view.end_mark)) |bytes| { return .{ .tail = bytes[0..page.size] }; } } if (self.baseVisibleIndex(page_id, view.base_generation)) |index| return .{ .base = index }; return null; } fn walImageBytes(self: *const Pager, index: usize) *const [page.size]u8 { return self.walBytes()[self.wal_index.frames.items[index].offset..][0..page.size]; } pub fn frameCount(self: *const Pager) usize { return self.journal.frameCount(); } pub fn position(self: *const Pager) Position { return .{ .journal = self.journal.position(), .frames_len = self.wal_index.frames.items.len, .database_page_count = self.database_page_count, .end_mark = self.end_mark, }; } pub fn restoreEpoch(self: *const Pager) RestoreEpoch { return .{ .base_generation = self.base_generation, .base_page_count = self.base_page_count, .checkpoint_serial = self.checkpoint_serial, }; } pub fn canRestore( self: *const Pager, position_value: Position, epoch: RestoreEpoch, ) bool { if (!std.meta.eql(self.restoreEpoch(), epoch)) return false; if (self.wal_index.frames.items.len < position_value.frames_len) return false; if (self.journal.position().len < position_value.journal.len) return false; return true; } pub fn restore(self: *Pager, position_value: Position) void { self.journal.restore(position_value.journal); self.wal_index.frames.shrinkRetainingCapacity(position_value.frames_len); self.database_page_count = position_value.database_page_count; self.end_mark = position_value.end_mark; self.rebuildWalPages(); } pub fn walBytes(self: *const Pager) []const u8 { return self.journal.bytes(); } pub fn walCapacityBytes(self: *const Pager) usize { return self.journal.byteCapacity(); } pub fn baseGeneration(self: *const Pager) u64 { return self.base_generation; } pub fn setBasePageCount(self: *Pager, count: u32) void { self.base_page_count = @max(self.base_page_count, count); self.database_page_count = @max(self.database_page_count, count); if (count > 0 and self.base_generation == 0) self.base_generation = 1; } pub fn basePageCount(self: *const Pager) u32 { return self.base_page_count; } pub fn databasePageCount(self: *const Pager) u32 { return self.database_page_count; } fn lowerBoundWalPage(self: *const Pager, page_id: u32) usize { var low: usize = 0; var high = self.wal_index.pages.items.len; while (low < high) { const mid = low + (high - low) / 2; if (self.wal_index.pages.items[mid].page_id < page_id) { low = mid + 1; } else { high = mid; } } return low; } fn indexAppendedWalFrame(self: *Pager, page_id: u32, db_page_count: u32, page_index: usize, existing: bool) void { const frame = self.journal.frameCount(); if (db_page_count != 0) self.end_mark = frame; self.database_page_count = @max(self.database_page_count, @max(page_id, db_page_count)); const frame_index = self.wal_index.frames.items.len; self.wal_index.frames.appendAssumeCapacity(.{ .page_id = page_id, .frame = frame, .offset = self.walBytes().len - page.size, .previous = if (existing) self.wal_index.pages.items[page_index].frame_index else null, }); if (existing) { self.wal_index.pages.items[page_index].frame_index = frame_index; } else { self.wal_index.pages.insertAssumeCapacity(page_index, .{ .page_id = page_id, .frame_index = frame_index, }); } } fn lowerBoundBase(self: *const Pager, page_id: u32, generation: u64) usize { var low: usize = 0; var high = self.base.items.len; while (low < high) { const mid = low + (high - low) / 2; const image = self.base.items[mid]; if (image.id < page_id or (image.id == page_id and image.generation < generation)) { low = mid + 1; } else { high = mid; } } return low; } fn baseVisibleIndex(self: *const Pager, page_id: u32, generation: u64) ?usize { if (self.base_index.get(page_id)) |index| { const image = &self.base.items[index]; if (image.id == page_id and image.generation <= generation) return index; } var best_index: ?usize = null; var best_generation: u64 = 0; for (self.base.items, 0..) |*image, index| { if (image.id == page_id and image.generation <= generation and (best_index == null or image.generation > best_generation)) { best_index = index; best_generation = image.generation; } } return best_index; } fn upperBoundBase(self: *const Pager, page_id: u32, generation: u64) usize { var low: usize = 0; var high = self.base.items.len; while (low < high) { const mid = low + (high - low) / 2; const image = self.base.items[mid]; if (image.id < page_id or (image.id == page_id and image.generation <= generation)) { low = mid + 1; } else { high = mid; } } return low; } fn indexedWalImage(self: *const Pager, page_id: u32, max_frame: usize) ?usize { const page_index = self.lowerBoundWalPage(page_id); if (page_index == self.wal_index.pages.items.len or self.wal_index.pages.items[page_index].page_id != page_id) { return null; } const wal_len = self.walBytes().len; var frame_index: ?usize = self.wal_index.pages.items[page_index].frame_index; while (frame_index) |index| { const image = self.wal_index.frames.items[index]; if (image.frame <= max_frame and image.offset + page.size <= wal_len) return index; frame_index = image.previous; } return null; } fn latestCheckpointFrame(self: *const Pager, wal_page: WalPage, target_mark: usize, wal_len: usize) ?WalImage { var frame_index: ?usize = wal_page.frame_index; while (frame_index) |index| { const image = self.wal_index.frames.items[index]; if (image.frame <= target_mark and image.offset + page.size <= wal_len) return image; frame_index = image.previous; } return null; } fn checkpointPageCount(self: *const Pager, target_mark: usize, wal_len: usize) usize { var count: usize = 0; for (self.wal_index.pages.items) |wal_page| { if (self.latestCheckpointFrame(wal_page, target_mark, wal_len) != null) count += 1; } return count; } fn collectCheckpointPages( self: *const Pager, target_mark: usize, wal_len: usize, plan: *CheckpointPlan, ) error{CheckpointPlanCapacityExceeded}!void { for (self.wal_index.pages.items) |wal_page| { const image = self.latestCheckpointFrame(wal_page, target_mark, wal_len) orelse continue; try plan.append(.{ .page_id = image.page_id, .wal_offset = image.offset, }); } } fn installCheckpointPages(self: *Pager, prepared: PreparedCheckpoint, generation: u64) Error!void { const checkpoint_page_count = prepared.checkpointPageCount(); try self.base.ensureUnusedCapacity(self.allocator, checkpoint_page_count); try self.base_index.ensureUnusedCapacity(self.allocator, try hashMapSize(checkpoint_page_count)); for (0..checkpoint_page_count) |page_index| { const checkpoint_page = prepared.checkpointPage(page_index); const index = self.base.items.len; self.base.appendAssumeCapacity(.{ .id = checkpoint_page.page_id, .generation = generation, .bytes = checkpoint_page.bytes[0..page.size].*, }); if (self.base_index.get(checkpoint_page.page_id)) |existing| { if (self.base.items[existing].generation <= generation) self.base_index.putAssumeCapacity(checkpoint_page.page_id, index); } else { self.base_index.putAssumeCapacity(checkpoint_page.page_id, index); } self.base_page_count = @max(self.base_page_count, checkpoint_page.page_id); self.database_page_count = @max(self.database_page_count, checkpoint_page.page_id); } } fn preparedCheckpointCurrent(self: *const Pager, prepared: PreparedCheckpoint) bool { if (prepared.pager != self) return false; if (self.checkpoint_serial != prepared.serial) return false; if (!std.meta.eql(self.position(), prepared.state.position)) return false; if (self.base_generation != prepared.state.base_generation) return false; if (self.base.items.len != prepared.state.base_images) return false; if (self.wal_index.pages.items.len != prepared.state.wal_pages) return false; return true; } fn finishCheckpointFrames(self: *Pager, prepared: PreparedCheckpoint) void { const checkpoint_value = prepared.checkpoint; if (!prepared.has_readers and !checkpoint_value.restarted) _ = self.compactWalFrames(checkpoint_value.end_mark); if (checkpoint_value.restarted) { self.journal.rewriteTail(checkpoint_value.end_mark, prepared.restart_header.?); self.rebuildWalFramesFromJournal(); } } fn compactBaseHistory(self: *Pager, retain_generation: u64) Error!usize { const phase = trace.scope("pager.base_history.compact"); defer phase.end(); var retained_floor: std.AutoHashMapUnmanaged(u32, usize) = .empty; defer retained_floor.deinit(self.allocator); try retained_floor.ensureTotalCapacity(self.allocator, try hashMapSize(self.base.items.len)); for (self.base.items, 0..) |*image, index| { if (image.generation > retain_generation) continue; if (retained_floor.get(image.id)) |existing| { if (self.base.items[existing].generation < image.generation) retained_floor.putAssumeCapacity(image.id, index); } else { retained_floor.putAssumeCapacity(image.id, index); } } var write_index: usize = 0; for (self.base.items, 0..) |image, index| { const keep = image.generation > retain_generation or (retained_floor.get(image.id) orelse std.math.maxInt(usize)) == index; if (keep) { self.base.items[write_index] = image; write_index += 1; } } const removed = self.base.items.len - write_index; self.base.shrinkRetainingCapacity(write_index); try self.rebuildBaseIndex(); if (removed > 0) trace.progress("pager.base_history.compact.complete"); return removed; } fn compactWalFrames(self: *Pager, checkpoint_mark: usize) usize { const phase = trace.scope("pager.wal_frames.compact"); defer phase.end(); var write_index: usize = 0; for (self.wal_index.frames.items) |image| { if (image.frame > checkpoint_mark) { self.wal_index.frames.items[write_index] = image; write_index += 1; } } const removed = self.wal_index.frames.items.len - write_index; self.wal_index.frames.shrinkRetainingCapacity(write_index); if (removed > 0) { self.rebuildWalPages(); trace.progress("pager.wal_frames.compact.complete"); } return removed; } fn rebuildWalFramesFromJournal(self: *Pager) void { self.rebuildWalFramesFromJournalControlled(.{}) catch unreachable; } fn rebuildWalFramesFromJournalControlled( self: *Pager, control: wal.Control, ) error{Interrupted}!void { self.wal_index.frames.clearRetainingCapacity(); self.wal_index.pages.clearRetainingCapacity(); self.database_page_count = self.base_page_count; self.end_mark = 0; var reader = wal.Reader.initControlled(self.walBytes(), control) catch |err| switch (err) { error.Interrupted => return error.Interrupted, else => unreachable, }; const frames_max = self.journal.frameCount(); var frame_index: usize = 0; while (frame_index < frames_max) : (frame_index += 1) { const frame = (reader.nextControlled(control) catch |err| switch (err) { error.Interrupted => return error.Interrupted, else => unreachable, }) orelse unreachable; self.wal_index.frames.appendAssumeCapacity(.{ .page_id = frame.page_id, .frame = frame.index, .offset = wal.header_size + (frame.index - 1) * wal.frame_size + wal.frame_header_size, .previous = null, }); self.database_page_count = @max(self.database_page_count, @max(frame.page_id, frame.db_page_count)); if (frame.committed()) self.end_mark = frame.index; } try self.rebuildWalPagesControlled(control); } fn rebuildWalPages(self: *Pager) void { self.rebuildWalPagesControlled(.{}) catch unreachable; } fn rebuildWalPagesControlled( self: *Pager, control: wal.Control, ) error{Interrupted}!void { self.wal_index.pages.clearRetainingCapacity(); for (self.wal_index.frames.items, 0..) |*image, index| { try control.check(); const page_index = self.lowerBoundWalPage(image.page_id); if (page_index < self.wal_index.pages.items.len and self.wal_index.pages.items[page_index].page_id == image.page_id) { image.previous = self.wal_index.pages.items[page_index].frame_index; self.wal_index.pages.items[page_index].frame_index = index; } else { image.previous = null; self.wal_index.pages.insertAssumeCapacity(page_index, .{ .page_id = image.page_id, .frame_index = index, }); } } try control.check(); } fn rebuildBaseIndex(self: *Pager) Error!void { try self.base_index.ensureTotalCapacity(self.allocator, try hashMapSize(self.base.items.len)); self.base_index.clearRetainingCapacity(); for (self.base.items, 0..) |*image, index| { if (self.base_index.get(image.id)) |existing| { if (self.base.items[existing].generation <= image.generation) self.base_index.putAssumeCapacity(image.id, index); } else { self.base_index.putAssumeCapacity(image.id, index); } } } fn nextBaseGeneration(self: *const Pager) Error!u64 { if (self.base_generation == std.math.maxInt(u64)) return error.GenerationOverflow; return self.base_generation + 1; }};Source: lib/sql/src/root.zig:294
zig
pub const Pager = pager.Pager;Complete call list for Pager.commitStagedWal
8 direct calls.
lib.sql.src.pager.Pager.indexAppendedWalFrame[method] — private source atlib/sql/src/pager.zig:1231in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.lowerBoundWalPage[method] — private source atlib/sql/src/pager.zig:1217in nearest public ownertiny.sql.pagertiny.sql.Pager.position[method] atlib/sql/src/pager.zig:1155tiny.sql.Pager.stagedWalPageId[method] atlib/sql/src/pager.zig:916tiny.sql.Pager.walStagingCapacity[method] atlib/sql/src/pager.zig:896tiny.sql.pager.WalIndex.reserve[method] atlib/sql/src/pager.zig:311tiny.sql.trace.progress[function] atlib/sql/src/trace.zig:7tiny.sql.WalWriter.commitStagedFrame[method] atlib/sql/src/wal.zig:970
Complete caller list for Pager.position
14 direct callers.
lib.sql.src.file.Database.appendWal[method] — private source atlib/sql/src/file.zig:2317in nearest public ownertiny.sql.indextiny.sql.FileDatabase.beginWrite[method] atlib/sql/src/file.zig:2292tiny.sql.FileDatabase.restore[method] atlib/sql/src/file.zig:2194tiny.sql.FileDatabase.restoreWrites[method] atlib/sql/src/file.zig:2221tiny.sql.FileDatabase.savepoint[method] atlib/sql/src/file.zig:2181tiny.sql.Pager.commitStagedWal[method] atlib/sql/src/pager.zig:925lib.sql.src.pager.Pager.prepareCheckpointWithPlan[method] — private source atlib/sql/src/pager.zig:980in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.preparedCheckpointCurrent[method] — private source atlib/sql/src/pager.zig:1372in nearest public ownertiny.sql.pagertiny.sql.Pager.stageWalPage[method] atlib/sql/src/pager.zig:901tiny.sql.Pager.stagedWalPage[method] atlib/sql/src/pager.zig:906tiny.sql.Pager.stagedWalPageMut[method] atlib/sql/src/pager.zig:911tiny.sql.Pager.swapStagedWalFrames[method] atlib/sql/src/pager.zig:920tiny.sql.Pager.walStagingCapacity[method] atlib/sql/src/pager.zig:896tiny.sql.pager.PreparedCheckpoint.checkpointPage[method] atlib/sql/src/pager.zig:688
Complete caller list for Pager.walBytes
10 direct callers.
lib.sql.src.file.Database.ensureWalFrames[method] — private source atlib/sql/src/file.zig:2573in nearest public ownertiny.sql.indexlib.sql.src.file.Database.persistWal[method] — private source atlib/sql/src/file.zig:2581in nearest public ownertiny.sql.indexlib.sql.src.file.Database.rewriteWal[method] — private source atlib/sql/src/file.zig:2553in nearest public ownertiny.sql.indexlib.sql.src.pager.Pager.indexAppendedWalFrame[method] — private source atlib/sql/src/pager.zig:1231in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.indexedWalImage[method] — private source atlib/sql/src/pager.zig:1299in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.locateImage[method] — private source atlib/sql/src/pager.zig:1134in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.prepareCheckpointWithPlan[method] — private source atlib/sql/src/pager.zig:980in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.rebuildWalFramesFromJournalControlled[method] — private source atlib/sql/src/pager.zig:1444in nearest public ownertiny.sql.pagerlib.sql.src.pager.Pager.walImageBytes[method] — private source atlib/sql/src/pager.zig:1147in nearest public ownertiny.sql.pagertiny.sql.pager.PreparedCheckpoint.checkpointPage[method] atlib/sql/src/pager.zig:688
Audit
| Definitions | 40 |
|---|---|
| Public names | 80 |
| Members | 19 |
| Version | 26.7.0 |
| Revision | daab053ee433 |