lib/trace/src/store/writer/writer.zig

daab053ee43316e1809a84551d573ddd1e5bf3d2

  1 const std = @import("std");
  2 const sql = @import("sql");
  3 const sys = @import("sys");
  4 const event = @import("../../root.zig").event;
  5 const format = @import("../format/root.zig");
  6 const block = format.block;
  7 const codec = format.codec;
  8 const database = format.database;
  9 const capacity_mod = @import("capacity.zig");
 10 const storage_mod = @import("storage.zig");
 11 
 12 const Storage = storage_mod.Storage;
 13 
 14 pub const Writer = struct {
 15     storage: ?*Storage = null,
 16     root_path: []u8 = &.{},
 17     target_triple: []u8 = &.{},
 18     build_id: []u8 = &.{},
 19     mode: []u8 = &.{},
 20     endian: []u8 = &.{},
 21     manifest: format.Manifest = .{},
 22     block_builder: block.Builder = undefined,
 23     manifest_io_buffer: []u8 = &.{},
 24     manifest_buffer: []u8 = &.{},
 25     path_buffers: [capacity_mod.path_buffer_count][]u8 = .{ &.{}, &.{} },
 26     sql_allocator: std.heap.FixedBufferAllocator = undefined,
 27     database_workspace: sql.FileDatabase.Workspace = undefined,
 28     database_workspace_open: bool = false,
 29     file: sql.FileDatabase = undefined,
 30     file_open: bool = false,
 31     events: sql.Tree = undefined,
 32     chunk_event_count: u64 = 0,
 33     chunk_byte_count: u64 = 0,
 34     event_hasher: std.hash.Wyhash = format.newHasher(),
 35     initialized: bool = false,
 36     finished: bool = false,
 37 
 38     pub fn open(
 39         self: *Writer,
 40         storage: *Storage,
 41         root_path: []const u8,
 42         manifest_base: format.Manifest,
 43         trace_limits: format.Limits,
 44     ) !void {
 45         if (self.initialized) return error.TraceWriterAlreadyOpen;
 46         try trace_limits.validate();
 47         const regions = try storage.acquire(.{
 48             .root_path_bytes = root_path.len,
 49             .target_triple_bytes = manifest_base.target_triple.len,
 50             .build_id_bytes = manifest_base.build_id.len,
 51             .mode_bytes = manifest_base.mode.len,
 52             .endian_bytes = manifest_base.endian.len,
 53             .event_bytes = trace_limits.max_event_bytes,
 54         });
 55         var ownership_transferred = false;
 56         errdefer if (!ownership_transferred) storage.release();
 57         std.mem.copyForwards(u8, regions.root_path, root_path);
 58         std.mem.copyForwards(u8, regions.target_triple, manifest_base.target_triple);
 59         std.mem.copyForwards(u8, regions.build_id, manifest_base.build_id);
 60         std.mem.copyForwards(u8, regions.mode, manifest_base.mode);
 61         std.mem.copyForwards(u8, regions.endian, manifest_base.endian);
 62         self.* = .{
 63             .storage = storage,
 64             .root_path = regions.root_path,
 65             .target_triple = regions.target_triple,
 66             .build_id = regions.build_id,
 67             .mode = regions.mode,
 68             .endian = regions.endian,
 69             .manifest = manifest_base,
 70             .block_builder = block.Builder.init(regions.block),
 71             .manifest_io_buffer = regions.manifest_io,
 72             .manifest_buffer = regions.manifest,
 73             .path_buffers = regions.paths,
 74             .sql_allocator = std.heap.FixedBufferAllocator.init(regions.sql),
 75             .initialized = true,
 76         };
 77         ownership_transferred = true;
 78         errdefer self.deinit();
 79         self.manifest.format_version = event.trace_format_version;
 80         self.manifest.target_triple = self.target_triple;
 81         self.manifest.build_id = self.build_id;
 82         self.manifest.mode = self.mode;
 83         self.manifest.endian = self.endian;
 84         self.manifest.status = "open";
 85         self.manifest.limits = trace_limits;
 86         self.manifest.chunk_count = 0;
 87         self.manifest.event_count = 0;
 88         self.manifest.event_bytes = 0;
 89         self.manifest.event_checksum = 0;
 90         self.manifest.block_count = 0;
 91         self.manifest.block_bytes = 0;
 92 
 93         try sys.fs.deleteTree(self.root_path);
 94         try self.ensureLayout();
 95         try self.writeManifest();
 96         self.database_workspace = try sql.FileDatabase.Workspace.allocate(
 97             self.sql_allocator.allocator(),
 98             database.workspaceLimits(storage.capacity.database),
 99         );
100         self.database_workspace_open = true;
101         self.file = try database.openWriter(
102             self.sql_allocator.allocator(),
103             &self.database_workspace,
104             self.root_path,
105             storage.capacity.database,
106         );
107         self.file_open = true;
108         self.events = try database.events(&self.file);
109     }
110 
111     pub fn deinit(self: *Writer) void {
112         if (!self.initialized) return;
113         if (self.file_open) {
114             self.file.deinit();
115             self.file_open = false;
116         }
117         if (self.database_workspace_open) {
118             self.database_workspace.deallocate(self.sql_allocator.allocator());
119             self.database_workspace_open = false;
120         }
121         const storage = self.storage.?;
122         self.* = .{};
123         storage.release();
124     }
125 
126     pub fn sink(self: *Writer) event.Sink {
127         return .{ .context = self, .appendFn = appendFromSink };
128     }
129 
130     pub fn append(self: *Writer, item: event.Event) !void {
131         if (!self.initialized or self.finished) return error.TraceWriterNotOpen;
132         const encoded_size = codec.encodedSize(item) catch |err| switch (err) {
133             error.CapacityOverflow => return error.EventTooLarge,
134             error.InvalidEvent, error.OutputTooSmall => unreachable,
135         };
136         if (encoded_size == 0 or encoded_size > self.manifest.limits.max_event_bytes) {
137             return error.EventTooLarge;
138         }
139         const next_chunk_bytes = std.math.add(u64, self.chunk_byte_count, encoded_size) catch
140             return error.TraceTooLarge;
141         if (self.chunk_event_count != 0 and
142             next_chunk_bytes > self.manifest.limits.max_chunk_bytes)
143         {
144             try self.finishChunk();
145         }
146         if (!self.block_builder.empty() and !self.block_builder.canAppend(encoded_size)) {
147             try self.flushBlock();
148         }
149         std.debug.assert(self.block_builder.canAppend(encoded_size));
150         const next_chunk_event_count = std.math.add(u64, self.chunk_event_count, 1) catch
151             return error.TraceTooLarge;
152         const next_chunk_byte_count = std.math.add(u64, self.chunk_byte_count, encoded_size) catch
153             return error.TraceTooLarge;
154         const next_event_count = std.math.add(u64, self.manifest.event_count, 1) catch
155             return error.TraceTooLarge;
156         const next_event_bytes = std.math.add(u64, self.manifest.event_bytes, encoded_size) catch
157             return error.TraceTooLarge;
158         const encoded = self.block_builder.append(item, encoded_size) catch |err| switch (err) {
159             error.CapacityOverflow, error.OutputTooSmall => return error.EventTooLarge,
160             error.InvalidEvent => unreachable,
161         };
162         var key_buffer: [database.key_bytes]u8 = undefined;
163         const key = database.sequenceKey(&key_buffer, self.manifest.event_count);
164         database.updateChecksum(&self.event_hasher, key, encoded);
165         self.chunk_event_count = next_chunk_event_count;
166         self.chunk_byte_count = next_chunk_byte_count;
167         self.manifest.event_count = next_event_count;
168         self.manifest.event_bytes = next_event_bytes;
169     }
170 
171     pub fn finish(self: *Writer) !void {
172         if (!self.initialized or self.finished) return error.TraceWriterNotOpen;
173         if (self.chunk_event_count != 0) try self.finishChunk();
174         self.manifest.status = "closed";
175         self.manifest.event_checksum = self.event_hasher.final();
176         try self.writeManifest();
177         self.finished = true;
178     }
179 
180     fn appendFromSink(context: *anyopaque, item: event.Event) !void {
181         const self: *Writer = @ptrCast(@alignCast(context));
182         try self.append(item);
183     }
184 
185     fn finishChunk(self: *Writer) !void {
186         std.debug.assert(self.chunk_event_count != 0);
187         const next_chunk_count = std.math.add(u64, self.manifest.chunk_count, 1) catch
188             return error.TraceTooLarge;
189         try self.flushBlock();
190         try database.checkpoint(&self.file);
191         self.manifest.chunk_count = next_chunk_count;
192         self.chunk_event_count = 0;
193         self.chunk_byte_count = 0;
194     }
195 
196     fn flushBlock(self: *Writer) !void {
197         if (self.block_builder.empty()) return;
198         const encoded = self.block_builder.finish();
199         const next_block_count = std.math.add(u64, self.manifest.block_count, 1) catch
200             return error.TraceTooLarge;
201         const next_block_bytes = std.math.add(u64, self.manifest.block_bytes, encoded.len) catch
202             return error.TraceTooLarge;
203         try database.checkpointIfNeeded(&self.file);
204         var key_buffer: [database.key_bytes]u8 = undefined;
205         const key = database.sequenceKey(&key_buffer, self.manifest.block_count);
206         _ = try self.events.put(key, encoded, .{ .durability = .buffered });
207         self.manifest.block_count = next_block_count;
208         self.manifest.block_bytes = next_block_bytes;
209         self.block_builder.reset();
210     }
211 
212     fn ensureLayout(self: *Writer) !void {
213         try sys.fs.createDirPath(self.root_path);
214         try sys.fs.createDirPath(self.joinedPath(0, "snapshots"));
215     }
216 
217     fn writeManifest(self: *Writer) !void {
218         const manifest_path = self.joinedPath(0, format.manifest_path);
219         const pending_path = self.joinedPath(1, "manifest.pending");
220         try format.writeManifest(
221             self.manifest,
222             manifest_path,
223             pending_path,
224             self.manifest_buffer,
225             self.manifest_io_buffer,
226         );
227     }
228 
229     fn joinedPath(self: *Writer, buffer_index: usize, relative: []const u8) []const u8 {
230         std.debug.assert(buffer_index < self.path_buffers.len);
231         std.debug.assert(self.root_path.len != 0);
232         std.debug.assert(relative.len <= format.max_relative_path_bytes);
233         const buffer = self.path_buffers[buffer_index];
234         var length = self.root_path.len;
235         @memcpy(buffer[0..length], self.root_path);
236         if (!sys.path.isSeparator(buffer[length - 1])) {
237             buffer[length] = sys.path.separator;
238             length += 1;
239         }
240         @memcpy(buffer[length..][0..relative.len], relative);
241         length += relative.len;
242         std.debug.assert(length <= buffer.len);
243         return buffer[0..length];
244     }
245 };