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 }