File
Blob: src/worker/do/repo/catalog/receive.ts
| 1 | import type { Logger } from "@/worker/common/logger"; |
| 2 | import type { RepoStateSchema } from "../repoState"; |
| 3 | |
| 4 | import { asTypedStorage } from "../repoState"; |
| 5 | import { |
| 6 | applyReceiveCommands, |
| 7 | isValidRefName, |
| 8 | type ReceiveCommand, |
| 9 | type ReceiveStatus, |
| 10 | validateReceiveCommands, |
| 11 | } from "@/worker/git/operations/validation"; |
| 12 | import { getDb, listActivePackCatalog, upsertPackCatalogRow } from "../db"; |
| 13 | import { DEFAULT_HEAD, bumpPacksetVersion, ensureRepoMetadataDefaults } from "./shared"; |
| 14 | import { scheduleCompactionWake, selectCompactionWork } from "./compaction/plan"; |
| 15 | |
| 16 | export 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 | |
| 34 | function 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 | |
| 52 | export 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 | } |