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 };