import type { Logger } from "@/worker/common/logger"; import type { RepoStateSchema } from "../repoState"; import { asTypedStorage } from "../repoState"; import { applyReceiveCommands, isValidRefName, type ReceiveCommand, type ReceiveStatus, validateReceiveCommands, } from "@/worker/git/operations/validation"; import { getDb, listActivePackCatalog, upsertPackCatalogRow } from "../db"; import { DEFAULT_HEAD, bumpPacksetVersion, ensureRepoMetadataDefaults } from "./shared"; import { scheduleCompactionWake, selectCompactionWork } from "./compaction/plan"; export type FinalizeReceiveResult = | { status: "committed"; statuses: ReceiveStatus[]; changed: boolean; empty: boolean; shouldQueueCompaction: boolean; } | { status: "ref_conflict"; statuses: ReceiveStatus[]; message: string; } | { status: "lease_mismatch"; message: string; }; function resolveHeadAfterReceive(args: { storedHead: | { target: string; oid?: string; unborn?: boolean; } | undefined; refs: Array<{ name: string; oid: string }>; }) { const target = args.storedHead?.target || DEFAULT_HEAD.target; const match = args.refs.find((ref) => ref.name === target); if (match) { return { target, oid: match.oid } as const; } return { target, unborn: true } as const; } export async function finalizeReceiveState(args: { ctx: DurableObjectState; env: Env; token: string; commands: ReceiveCommand[]; stagedPack?: | { packKey: string; packBytes: number; idxBytes: number; objectCount: number; } | undefined; logger?: Logger; }): Promise { const store = asTypedStorage(args.ctx.storage); await ensureRepoMetadataDefaults(store); const lease = await store.get("receiveLease"); if (!lease || lease.token !== args.token) { return { status: "lease_mismatch", message: "Receive lease is no longer active for this request.", }; } const currentRefs = (await store.get("refs")) || []; const invalidStatuses = args.commands .filter((command) => !isValidRefName(command.ref)) .map((command) => ({ ref: command.ref, ok: false, msg: "invalid" satisfies string })); if (invalidStatuses.length > 0) { await store.delete("receiveLease"); args.logger?.warn("receive:finalize-invalid-ref", { invalidCount: invalidStatuses.length, }); return { status: "ref_conflict", statuses: invalidStatuses, message: "Receive finalization rejected invalid refs.", }; } const statuses = validateReceiveCommands(currentRefs, args.commands); if (!statuses.every((status) => status.ok)) { await store.delete("receiveLease"); args.logger?.warn("receive:finalize-ref-conflict", { conflictCount: statuses.filter((status) => !status.ok).length, }); return { status: "ref_conflict", statuses, message: "Ref expectations changed before the receive could be committed.", }; } const nextRefs = applyReceiveCommands(currentRefs, args.commands); const storedHead = await store.get("head"); const nextHead = resolveHeadAfterReceive({ storedHead, refs: nextRefs }); const nextRefsVersion = ((await store.get("refsVersion")) || 0) + 1; let shouldQueueCompaction = false; if (args.stagedPack) { const nextPackSeq = (await store.get("nextPackSeq")) || 1; const db = getDb(args.ctx.storage); await upsertPackCatalogRow(db, { packKey: args.stagedPack.packKey, kind: "receive", state: "active", tier: 0, seqLo: nextPackSeq, seqHi: nextPackSeq, objectCount: args.stagedPack.objectCount, packBytes: args.stagedPack.packBytes, idxBytes: args.stagedPack.idxBytes, createdAt: Date.now(), supersededBy: null, }); await store.put("nextPackSeq", nextPackSeq + 1); const activeCatalog = await listActivePackCatalog(db); await bumpPacksetVersion(store); const compactionSelection = selectCompactionWork(activeCatalog); shouldQueueCompaction = compactionSelection.status === "ready"; if (shouldQueueCompaction) { await store.put("compactionWantedAt", Date.now()); await scheduleCompactionWake(args.ctx, args.env); } else if (compactionSelection.status === "blocked") { await store.delete("compactionWantedAt"); args.logger?.warn("receive:compaction-blocked", { reason: compactionSelection.blocked.reason, sourceTier: compactionSelection.blocked.sourceTier, activePackCount: compactionSelection.blocked.activePackCount, maxSourceObjects: compactionSelection.blocked.maxSourceObjects, maxSourceBytes: compactionSelection.blocked.maxSourceBytes, smallestWindowObjects: compactionSelection.blocked.smallestWindowObjects, smallestWindowBytes: compactionSelection.blocked.smallestWindowBytes, }); } } await store.put("refs", nextRefs); await store.put("head", nextHead); await store.put("refsVersion", nextRefsVersion); await store.delete("receiveLease"); args.logger?.info("receive:finalize-committed", { commandCount: args.commands.length, refCount: nextRefs.length, empty: nextRefs.length === 0, stagedPackKey: args.stagedPack?.packKey, shouldQueueCompaction, }); return { status: "committed", statuses, changed: args.commands.length > 0, empty: nextRefs.length === 0, shouldQueueCompaction, }; }