File
Blob: src/worker/git/pack/rewrite.ts
| 1 | import type { OrderedPackSnapshot } from "@/worker/git/operations/fetch/types"; |
| 2 | |
| 3 | import { createLogger } from "@/worker/common"; |
| 4 | import { |
| 5 | buildSelection, |
| 6 | buildOutputOrder, |
| 7 | canPassthroughSinglePack, |
| 8 | computeHeaderLengths, |
| 9 | } from "./rewrite/plan"; |
| 10 | import { |
| 11 | ensurePackReadState, |
| 12 | type RewriteFailure, |
| 13 | type RewriteFailureRecorder, |
| 14 | type RewriteOptions, |
| 15 | } from "./rewrite/shared"; |
| 16 | import { createPassthroughStream, createRewriteStream } from "./rewrite/stream"; |
| 17 | |
| 18 | export type PackRewriteResult = |
| 19 | | { status: "ok"; stream: ReadableStream<Uint8Array> } |
| 20 | | { status: "failed"; failure: RewriteFailure }; |
| 21 | |
| 22 | export async function rewritePackResult( |
| 23 | env: Env, |
| 24 | snapshot: OrderedPackSnapshot, |
| 25 | neededOids: string[], |
| 26 | options?: RewriteOptions |
| 27 | ): Promise<PackRewriteResult> { |
| 28 | const log = createLogger(env.LOG_LEVEL, { service: "PackRewrite" }); |
| 29 | const startedAt = Date.now(); |
| 30 | const warnedFlags = new Set<string>(); |
| 31 | const failure: RewriteFailureRecorder = options?.failure || {}; |
| 32 | const rewriteOptions: RewriteOptions = { ...options, failure }; |
| 33 | |
| 34 | function failed(reason: string, retryable: boolean, details?: Record<string, unknown>) { |
| 35 | return { |
| 36 | status: "failed" as const, |
| 37 | failure: failure.value || { reason, retryable, details }, |
| 38 | }; |
| 39 | } |
| 40 | |
| 41 | if (rewriteOptions.signal?.aborted) { |
| 42 | return failed("aborted", true); |
| 43 | } |
| 44 | if (!rewriteOptions.limiter) { |
| 45 | throw new Error("rewrite: limiter required"); |
| 46 | } |
| 47 | if (!rewriteOptions.countSubrequest) { |
| 48 | throw new Error("rewrite: countSubrequest required"); |
| 49 | } |
| 50 | |
| 51 | const selection = await buildSelection( |
| 52 | env, |
| 53 | snapshot, |
| 54 | neededOids, |
| 55 | log, |
| 56 | warnedFlags, |
| 57 | rewriteOptions |
| 58 | ); |
| 59 | if (!selection) { |
| 60 | return failed("selection-failed", true, { needed: neededOids.length }); |
| 61 | } |
| 62 | |
| 63 | const { table, readerStates } = selection; |
| 64 | |
| 65 | if (canPassthroughSinglePack(snapshot, table)) { |
| 66 | const readState = await ensurePackReadState( |
| 67 | env, |
| 68 | snapshot.packs[0]!, |
| 69 | 0, |
| 70 | readerStates, |
| 71 | log, |
| 72 | warnedFlags, |
| 73 | rewriteOptions |
| 74 | ); |
| 75 | |
| 76 | log.info("rewrite:passthrough", { |
| 77 | packKey: snapshot.packs[0]?.packKey, |
| 78 | objects: table.count, |
| 79 | }); |
| 80 | |
| 81 | return { |
| 82 | status: "ok", |
| 83 | stream: createPassthroughStream({ |
| 84 | env, |
| 85 | snapshotPack: snapshot.packs[0]!, |
| 86 | readState, |
| 87 | log, |
| 88 | warnedFlags, |
| 89 | options: rewriteOptions, |
| 90 | onComplete: () => { |
| 91 | log.info("rewrite:stream-complete", { |
| 92 | passthrough: true, |
| 93 | wholePackLoads: countWholePackLoads(readerStates), |
| 94 | timeMs: Date.now() - startedAt, |
| 95 | }); |
| 96 | }, |
| 97 | }), |
| 98 | }; |
| 99 | } |
| 100 | |
| 101 | if (!buildOutputOrder(table, log)) { |
| 102 | return failed("topology-incomplete", false, { selected: table.count }); |
| 103 | } |
| 104 | if (!computeHeaderLengths(table, log)) { |
| 105 | return failed("header-lengths-did-not-converge", false, { selected: table.count }); |
| 106 | } |
| 107 | |
| 108 | return { |
| 109 | status: "ok", |
| 110 | stream: createRewriteStream(table, snapshot, readerStates, log, rewriteOptions, () => { |
| 111 | log.info("rewrite:stream-complete", { |
| 112 | passthrough: false, |
| 113 | wholePackLoads: countWholePackLoads(readerStates), |
| 114 | timeMs: Date.now() - startedAt, |
| 115 | }); |
| 116 | }), |
| 117 | }; |
| 118 | } |
| 119 | |
| 120 | export async function rewritePack( |
| 121 | env: Env, |
| 122 | snapshot: OrderedPackSnapshot, |
| 123 | neededOids: string[], |
| 124 | options?: RewriteOptions |
| 125 | ): Promise<ReadableStream<Uint8Array> | undefined> { |
| 126 | const result = await rewritePackResult(env, snapshot, neededOids, options); |
| 127 | return result.status === "ok" ? result.stream : undefined; |
| 128 | } |
| 129 | |
| 130 | function countWholePackLoads(readerStates: Map<number, { wholePack?: Uint8Array }>): number { |
| 131 | let count = 0; |
| 132 | for (const state of readerStates.values()) { |
| 133 | if (state.wholePack) count++; |
| 134 | } |
| 135 | return count; |
| 136 | } |