Skip to content
File

Blob: src/worker/do/repo/catalog/receive.ts

typescript172 lines
1import type { Logger } from "@/worker/common/logger";
2import type { RepoStateSchema } from "../repoState";
3 
4import { asTypedStorage } from "../repoState";
5import {
6 applyReceiveCommands,
7 isValidRefName,
8 type ReceiveCommand,
9 type ReceiveStatus,
10 validateReceiveCommands,
11} from "@/worker/git/operations/validation";
12import { getDb, listActivePackCatalog, upsertPackCatalogRow } from "../db";
13import { DEFAULT_HEAD, bumpPacksetVersion, ensureRepoMetadataDefaults } from "./shared";
14import { scheduleCompactionWake, selectCompactionWork } from "./compaction/plan";
15 
16export type FinalizeReceiveResult =
17 | {
18 status: "committed";
19 statuses: ReceiveStatus[];
20 changed: boolean;
21 empty: boolean;
22 shouldQueueCompaction: boolean;
23 }
24 | {
25 status: "ref_conflict";
26 statuses: ReceiveStatus[];
27 message: string;
28 }
29 | {
30 status: "lease_mismatch";
31 message: string;
32 };
33 
34function resolveHeadAfterReceive(args: {
35 storedHead:
36 | {
37 target: string;
38 oid?: string;
39 unborn?: boolean;
40 }
41 | undefined;
42 refs: Array<{ name: string; oid: string }>;
43}) {
44 const target = args.storedHead?.target || DEFAULT_HEAD.target;
45 const match = args.refs.find((ref) => ref.name === target);
46 if (match) {
47 return { target, oid: match.oid } as const;
48 }
49 return { target, unborn: true } as const;
50}
51 
52export async function finalizeReceiveState(args: {
53 ctx: DurableObjectState;
54 env: Env;
55 token: string;
56 commands: ReceiveCommand[];
57 stagedPack?:
58 | {
59 packKey: string;
60 packBytes: number;
61 idxBytes: number;
62 objectCount: number;
63 }
64 | undefined;
65 logger?: Logger;
66}): Promise<FinalizeReceiveResult> {
67 const store = asTypedStorage<RepoStateSchema>(args.ctx.storage);
68 await ensureRepoMetadataDefaults(store);
69 
70 const lease = await store.get("receiveLease");
71 if (!lease || lease.token !== args.token) {
72 return {
73 status: "lease_mismatch",
74 message: "Receive lease is no longer active for this request.",
75 };
76 }
77 
78 const currentRefs = (await store.get("refs")) || [];
79 const invalidStatuses = args.commands
80 .filter((command) => !isValidRefName(command.ref))
81 .map((command) => ({ ref: command.ref, ok: false, msg: "invalid" satisfies string }));
82 if (invalidStatuses.length > 0) {
83 await store.delete("receiveLease");
84 args.logger?.warn("receive:finalize-invalid-ref", {
85 invalidCount: invalidStatuses.length,
86 });
87 return {
88 status: "ref_conflict",
89 statuses: invalidStatuses,
90 message: "Receive finalization rejected invalid refs.",
91 };
92 }
93 
94 const statuses = validateReceiveCommands(currentRefs, args.commands);
95 if (!statuses.every((status) => status.ok)) {
96 await store.delete("receiveLease");
97 args.logger?.warn("receive:finalize-ref-conflict", {
98 conflictCount: statuses.filter((status) => !status.ok).length,
99 });
100 return {
101 status: "ref_conflict",
102 statuses,
103 message: "Ref expectations changed before the receive could be committed.",
104 };
105 }
106 
107 const nextRefs = applyReceiveCommands(currentRefs, args.commands);
108 const storedHead = await store.get("head");
109 const nextHead = resolveHeadAfterReceive({ storedHead, refs: nextRefs });
110 const nextRefsVersion = ((await store.get("refsVersion")) || 0) + 1;
111 
112 let shouldQueueCompaction = false;
113 if (args.stagedPack) {
114 const nextPackSeq = (await store.get("nextPackSeq")) || 1;
115 const db = getDb(args.ctx.storage);
116 await upsertPackCatalogRow(db, {
117 packKey: args.stagedPack.packKey,
118 kind: "receive",
119 state: "active",
120 tier: 0,
121 seqLo: nextPackSeq,
122 seqHi: nextPackSeq,
123 objectCount: args.stagedPack.objectCount,
124 packBytes: args.stagedPack.packBytes,
125 idxBytes: args.stagedPack.idxBytes,
126 createdAt: Date.now(),
127 supersededBy: null,
128 });
129 await store.put("nextPackSeq", nextPackSeq + 1);
130 const activeCatalog = await listActivePackCatalog(db);
131 await bumpPacksetVersion(store);
132 const compactionSelection = selectCompactionWork(activeCatalog);
133 shouldQueueCompaction = compactionSelection.status === "ready";
134 if (shouldQueueCompaction) {
135 await store.put("compactionWantedAt", Date.now());
136 await scheduleCompactionWake(args.ctx, args.env);
137 } else if (compactionSelection.status === "blocked") {
138 await store.delete("compactionWantedAt");
139 args.logger?.warn("receive:compaction-blocked", {
140 reason: compactionSelection.blocked.reason,
141 sourceTier: compactionSelection.blocked.sourceTier,
142 activePackCount: compactionSelection.blocked.activePackCount,
143 maxSourceObjects: compactionSelection.blocked.maxSourceObjects,
144 maxSourceBytes: compactionSelection.blocked.maxSourceBytes,
145 smallestWindowObjects: compactionSelection.blocked.smallestWindowObjects,
146 smallestWindowBytes: compactionSelection.blocked.smallestWindowBytes,
147 });
148 }
149 }
150 
151 await store.put("refs", nextRefs);
152 await store.put("head", nextHead);
153 await store.put("refsVersion", nextRefsVersion);
154 await store.delete("receiveLease");
155 
156 args.logger?.info("receive:finalize-committed", {
157 commandCount: args.commands.length,
158 refCount: nextRefs.length,
159 empty: nextRefs.length === 0,
160 stagedPackKey: args.stagedPack?.packKey,
161 shouldQueueCompaction,
162 });
163 
164 return {
165 status: "committed",
166 statuses,
167 changed: args.commands.length > 0,
168 empty: nextRefs.length === 0,
169 shouldQueueCompaction,
170 };
171}