lib/sql/src/repository/flow.zig
daab053ee43316e1809a84551d573ddd1e5bf3d2
1 const std = @import("std");
2 const remote_owner = @import("remote.zig");
3
4 const Allocator = std.mem.Allocator;
5 const fs_io = std.Options.debug_io;
6
7 pub fn Repository(comptime Sql: type, comptime Policy: type) type {
8 const branch = Sql.branch;
9 const history = Sql.history;
10 const sync = Sql.sync;
11 const version = Sql.version;
12
13 return struct {
14 const Database = Policy.Database;
15 const Remote = remote_owner.Owner(Sql, Policy);
16
17 pub const Transfer = struct {
18 commits: usize = 0,
19 records: usize = 0,
20 applied: bool = false,
21 merged: bool = false,
22 resolved: usize = 0,
23 };
24
25 pub const Head = sync.HeadHex;
26
27 pub const Status = struct {
28 remote: ?[]u8 = null,
29 local_head: Head,
30 remote_head: ?Head = null,
31 ahead: usize = 0,
32 behind: usize = 0,
33 diverged: bool = false,
34
35 pub fn deinit(self: *Status, allocator: Allocator) void {
36 if (self.remote) |remote_path| allocator.free(remote_path);
37 self.* = undefined;
38 }
39 };
40
41 pub const writeRemote = Remote.write;
42 pub const readRemote = Remote.read;
43
44 pub fn status(
45 allocator: Allocator,
46 db: *Database,
47 remote_store_dir: ?[]const u8,
48 ) !Status {
49 const local_ref = (try (try db.history.full()).ref(
50 Policy.branch_name,
51 )) orelse return error.RefNotFound;
52 var result = Status{
53 .local_head = sync.headHex(local_ref.target),
54 };
55 const remote_dir = remote_store_dir orelse return result;
56
57 result.remote = try allocator.dupe(u8, remote_dir);
58 errdefer result.deinit(allocator);
59
60 var opened_remote = try Remote.open(allocator, remote_dir);
61 defer opened_remote.deinit();
62
63 const related = try sync.historyRelation(
64 allocator,
65 try db.history.full(),
66 &opened_remote.history,
67 Policy.branch_name,
68 );
69 if (related.remote_head) |remote_head| {
70 result.remote_head = sync.headHex(remote_head);
71 }
72 result.ahead = related.ahead;
73 result.behind = related.behind;
74 result.diverged = related.diverged();
75 return result;
76 }
77
78 pub fn push(
79 allocator: Allocator,
80 workspace: *Database.Workspace,
81 db: *Database,
82 remote_store_dir: []const u8,
83 dry_run: bool,
84 ) !Transfer {
85 if (try db.dirty()) return error.SyncDirtyStore;
86
87 if (dry_run) {
88 var remote_dir = try Remote.openDir(remote_store_dir);
89 defer remote_dir.close(fs_io);
90 if (!Remote.historyExists(remote_dir)) {
91 var names = [_][]const u8{Policy.branch_name};
92 var pack = try sync.exportRefNames(
93 allocator,
94 try db.history.full(),
95 names[0..],
96 );
97 defer pack.deinit();
98 return .{
99 .commits = pack.commits.len,
100 .records = pack.frames.len,
101 };
102 }
103 var opened_remote = try Remote.open(
104 allocator,
105 remote_store_dir,
106 );
107 defer opened_remote.deinit();
108 const related = try sync.historyRelation(
109 allocator,
110 try db.history.full(),
111 &opened_remote.history,
112 Policy.branch_name,
113 );
114 const planned = sync.planHistoryTransfer(
115 allocator,
116 try db.history.full(),
117 &opened_remote.history,
118 Policy.branch_name,
119 related,
120 .push,
121 ) catch |err| return mapHistoryPlanError(err);
122 return transferFromPlan(planned);
123 }
124
125 var remote_dir = try Remote.openDir(remote_store_dir);
126 defer remote_dir.close(fs_io);
127 var lock = try Remote.lock(remote_dir);
128 defer lock.close(fs_io);
129
130 if (try Remote.needsBootstrap(allocator, remote_dir)) {
131 if (comptime Policy.validate_transitions) {
132 try Policy.validateSnapshot(
133 allocator,
134 db,
135 try db.headHash(),
136 );
137 }
138 var remote_history = try history.History.open(
139 allocator,
140 remote_dir,
141 .{ .path = Policy.history_name, .recovery = .reject },
142 );
143 defer remote_history.deinit();
144 const stats = sync.pushFastForward(
145 allocator,
146 try db.history.full(),
147 &remote_history,
148 Policy.branch_name,
149 ) catch |err| switch (err) {
150 error.NonFastForward => return error.SyncDiverged,
151 else => return err,
152 };
153 var remote_db = try Policy.openWrite(
154 allocator,
155 workspace,
156 remote_dir,
157 );
158 defer remote_db.deinit();
159 try remote_db.alignToBranchHead();
160 try afterRemoteAlignment(
161 allocator,
162 remote_store_dir,
163 &remote_db,
164 );
165 return .{
166 .commits = stats.commits,
167 .records = stats.records,
168 .applied = true,
169 };
170 }
171
172 var remote_db = try Policy.openWrite(
173 allocator,
174 workspace,
175 remote_dir,
176 );
177 defer remote_db.deinit();
178 if (try remote_db.dirty()) return error.SyncDirtyStore;
179
180 const related = try sync.historyRelation(
181 allocator,
182 try db.history.full(),
183 try remote_db.history.full(),
184 Policy.branch_name,
185 );
186 if (related.upToDate()) {
187 try afterRemoteAlignment(
188 allocator,
189 remote_store_dir,
190 &remote_db,
191 );
192 return .{};
193 }
194 if (related.diverged()) return error.SyncDiverged;
195 if (related.behind != 0) return error.SyncRemoteAhead;
196
197 if (comptime Policy.validate_transitions) {
198 try Policy.validateTransition(
199 allocator,
200 &remote_db,
201 try remote_db.headHash(),
202 db,
203 try db.headHash(),
204 );
205 }
206
207 const stats = sync.pushFastForward(
208 allocator,
209 try db.history.full(),
210 try remote_db.history.full(),
211 Policy.branch_name,
212 ) catch |err| switch (err) {
213 error.NonFastForward => return error.SyncDiverged,
214 else => return err,
215 };
216 try remote_db.alignToBranchHead();
217 try afterRemoteAlignment(
218 allocator,
219 remote_store_dir,
220 &remote_db,
221 );
222 return .{
223 .commits = stats.commits,
224 .records = stats.records,
225 .applied = true,
226 };
227 }
228
229 pub fn pull(
230 allocator: Allocator,
231 workspace: *Database.Workspace,
232 db: *Database,
233 remote_store_dir: []const u8,
234 dry_run: bool,
235 ) !Transfer {
236 if (try db.dirty()) return error.SyncDirtyStore;
237
238 var opened_remote = try Remote.open(
239 allocator,
240 remote_store_dir,
241 );
242 defer opened_remote.deinit();
243 var remote_db: ?Database = null;
244 if (comptime @hasDecl(Policy, "openPullDatabase")) {
245 remote_db = try Policy.openPullDatabase(
246 allocator,
247 workspace,
248 opened_remote.dir,
249 );
250 }
251 defer if (remote_db) |*db_value| db_value.deinit();
252
253 const related = try sync.historyRelation(
254 allocator,
255 try db.history.full(),
256 &opened_remote.history,
257 Policy.branch_name,
258 );
259 if (related.diverged()) {
260 if (try db.pristine()) {
261 if (dry_run) {
262 return transferFromPlan(try sync.planHistoryAdopt(
263 allocator,
264 try db.history.full(),
265 &opened_remote.history,
266 Policy.branch_name,
267 ));
268 }
269 try validateRemoteTransition(
270 allocator,
271 db,
272 if (remote_db) |*db_value| db_value else null,
273 related.remote_head.?,
274 );
275 return try adopt(
276 allocator,
277 db,
278 &opened_remote.history,
279 related.remote_head.?,
280 );
281 }
282 if (dry_run) {
283 const planned = try sync.planHistoryAdopt(
284 allocator,
285 try db.history.full(),
286 &opened_remote.history,
287 Policy.branch_name,
288 );
289 var transfer = transferFromPlan(planned);
290 transfer.merged = true;
291 return transfer;
292 }
293 return try mergePull(
294 allocator,
295 db,
296 &opened_remote.history,
297 related.remote_head.?,
298 );
299 }
300 if (dry_run) {
301 const planned = sync.planHistoryTransfer(
302 allocator,
303 try db.history.full(),
304 &opened_remote.history,
305 Policy.branch_name,
306 related,
307 .pull,
308 ) catch |err| return mapHistoryPlanError(err);
309 return transferFromPlan(planned);
310 }
311 if (related.upToDate() or related.remote_head == null) return .{};
312 if (related.behind == 0) return .{};
313
314 try validateRemoteTransition(
315 allocator,
316 db,
317 if (remote_db) |*db_value| db_value else null,
318 related.remote_head.?,
319 );
320 const stats = sync.pullFastForward(
321 allocator,
322 try db.history.full(),
323 &opened_remote.history,
324 Policy.branch_name,
325 ) catch |err| switch (err) {
326 error.NonFastForward => return error.SyncDiverged,
327 else => return err,
328 };
329 try db.alignToBranchHead();
330 return .{
331 .commits = stats.commits,
332 .records = stats.records,
333 .applied = true,
334 };
335 }
336
337 pub fn exportPack(
338 allocator: Allocator,
339 db: *Database,
340 pack_path: []const u8,
341 ) !Transfer {
342 var pack = try sync.exportAll(allocator, try db.history.full());
343 defer pack.deinit();
344 const bytes = try sync.encodePack(allocator, &pack);
345 defer allocator.free(bytes);
346
347 var file = try std.Io.Dir.createFileAbsolute(
348 fs_io,
349 pack_path,
350 .{ .truncate = true },
351 );
352 defer file.close(fs_io);
353 try file.writePositionalAll(fs_io, bytes, 0);
354 try file.setLength(fs_io, bytes.len);
355 try file.sync(fs_io);
356 return .{
357 .commits = pack.commits.len,
358 .records = pack.frames.len,
359 .applied = true,
360 };
361 }
362
363 pub fn importPack(
364 allocator: Allocator,
365 db: *Database,
366 bytes: []const u8,
367 dry_run: bool,
368 ) !Transfer {
369 if (try db.dirty()) return error.SyncDirtyStore;
370
371 var pack = try sync.decodePack(allocator, bytes);
372 defer pack.deinit();
373 const target = sync.packRefTarget(
374 &pack,
375 Policy.branch_name,
376 ) orelse return error.SyncPackBranchMissing;
377 const local_ref = (try (try db.history.full()).ref(
378 Policy.branch_name,
379 )) orelse return error.RefNotFound;
380
381 if (dry_run) {
382 const missing = try sync.missingPackCounts(
383 try db.history.full(),
384 &pack,
385 );
386 if (version.same(local_ref.target, target)) return .{};
387 return transferFromPlan(missing);
388 }
389
390 const stats = try sync.importObjects(try db.history.full(), &pack);
391 if (version.same(local_ref.target, target)) return .{};
392
393 const entries = try (try db.history.full()).commitEntries(allocator);
394 defer allocator.free(entries);
395 if (!try branch.canFastForward(
396 allocator,
397 entries,
398 local_ref.target,
399 target,
400 )) {
401 if (try branch.canFastForward(
402 allocator,
403 entries,
404 target,
405 local_ref.target,
406 )) return .{};
407 if (!try db.pristine()) {
408 var transfer = try mergeImported(allocator, db, target);
409 transfer.commits = stats.commits;
410 transfer.records = stats.records;
411 return transfer;
412 }
413 try validateImportedTransition(
414 allocator,
415 db,
416 local_ref.target,
417 target,
418 );
419 try db.adoptBranchHead(target);
420 return .{
421 .commits = stats.commits,
422 .records = stats.records,
423 .applied = true,
424 };
425 }
426
427 try validateImportedTransition(
428 allocator,
429 db,
430 local_ref.target,
431 target,
432 );
433 try db.fastForwardTo(target);
434 return .{
435 .commits = stats.commits,
436 .records = stats.records,
437 .applied = true,
438 };
439 }
440
441 fn validateRemoteTransition(
442 allocator: Allocator,
443 db: *Database,
444 remote_db: ?*Database,
445 remote_head: version.Hash,
446 ) !void {
447 if (comptime !Policy.validate_transitions) return;
448 comptime {
449 if (!@hasDecl(Policy, "openPullDatabase")) {
450 @compileError(
451 "transition validation requires openPullDatabase",
452 );
453 }
454 }
455 const target = remote_db orelse unreachable;
456 try Policy.validateTransition(
457 allocator,
458 db,
459 try db.headHash(),
460 target,
461 remote_head,
462 );
463 }
464
465 fn validateImportedTransition(
466 allocator: Allocator,
467 db: *Database,
468 current_head: version.Hash,
469 target_head: version.Hash,
470 ) !void {
471 if (comptime !Policy.validate_transitions) return;
472 try Policy.validateTransition(
473 allocator,
474 db,
475 current_head,
476 db,
477 target_head,
478 );
479 }
480
481 fn mergePull(
482 allocator: Allocator,
483 db: *Database,
484 remote_history: *const history.History,
485 remote_head: version.Hash,
486 ) !Transfer {
487 var names = [_][]const u8{Policy.branch_name};
488 var pack = try sync.exportMissingRefNames(
489 allocator,
490 remote_history,
491 try db.history.full(),
492 names[0..],
493 );
494 defer pack.deinit();
495 const stats = try sync.importObjects(try db.history.full(), &pack);
496 var transfer = try mergeImported(allocator, db, remote_head);
497 transfer.commits = stats.commits;
498 transfer.records = stats.records;
499 return transfer;
500 }
501
502 fn afterRemoteAlignment(
503 allocator: Allocator,
504 remote_store_dir: []const u8,
505 db: *Database,
506 ) !void {
507 if (comptime !Policy.reconcile_remote) return;
508 try Policy.afterRemoteAlignment(
509 allocator,
510 remote_store_dir,
511 db,
512 );
513 }
514
515 fn mergeImported(
516 allocator: Allocator,
517 db: *Database,
518 remote_head: version.Hash,
519 ) !Transfer {
520 const resolved = try Policy.mergeImported(
521 allocator,
522 db,
523 remote_head,
524 );
525 return .{
526 .applied = true,
527 .merged = true,
528 .resolved = resolved,
529 };
530 }
531
532 fn adopt(
533 allocator: Allocator,
534 db: *Database,
535 remote_history: *const history.History,
536 target: version.Hash,
537 ) !Transfer {
538 var names = [_][]const u8{Policy.branch_name};
539 var pack = try sync.exportMissingRefNames(
540 allocator,
541 remote_history,
542 try db.history.full(),
543 names[0..],
544 );
545 defer pack.deinit();
546 const stats = try sync.importObjects(try db.history.full(), &pack);
547 try db.adoptBranchHead(target);
548 return .{
549 .commits = stats.commits,
550 .records = stats.records,
551 .applied = true,
552 };
553 }
554
555 fn transferFromPlan(plan: sync.HistoryTransferPlan) Transfer {
556 return .{
557 .commits = plan.commits,
558 .records = plan.records,
559 };
560 }
561
562 fn mapHistoryPlanError(err: anyerror) anyerror {
563 return switch (err) {
564 error.HistoryDiverged => error.SyncDiverged,
565 error.HistoryRemoteAhead => error.SyncRemoteAhead,
566 else => err,
567 };
568 }
569 };
570 }