File
Blob: test/util/streaming-helpers.ts
| 1 | import { expect } from "vitest"; |
| 2 | import { exports as workerExports } from "cloudflare:workers"; |
| 3 | |
| 4 | import { concatChunks, flushPkt, pktLine, decodePktLines } from "@/worker/git/core"; |
| 5 | import { encodeGitObject } from "@/worker/git/core/objects"; |
| 6 | import { buildPack } from "./git-pack"; |
| 7 | import { buildTreePayload } from "./packed-repo"; |
| 8 | import { lookupPushAuth } from "./repoSeed"; |
| 9 | |
| 10 | /** |
| 11 | * No-op: all repos are now implicitly streaming. |
| 12 | * Kept as a stub so existing test call-sites don't need updating. |
| 13 | */ |
| 14 | export async function promoteToStreaming(_owner: string, _repo: string) { |
| 15 | // Streaming is the only mode after the closure release — nothing to do. |
| 16 | } |
| 17 | |
| 18 | export function decodeReportStatus(responseBody: Uint8Array): string[] { |
| 19 | return decodePktLines(responseBody) |
| 20 | .filter((item) => item.type === "line") |
| 21 | .map((item: any) => String(item.text).trim()); |
| 22 | } |
| 23 | |
| 24 | export function decodeReceiveSideband(responseBody: Uint8Array): { |
| 25 | progress: string[]; |
| 26 | fatal: string[]; |
| 27 | reportStatus: string[]; |
| 28 | } { |
| 29 | const progress: string[] = []; |
| 30 | const fatal: string[] = []; |
| 31 | const reportStatusChunks: Uint8Array[] = []; |
| 32 | |
| 33 | for (const item of decodePktLines(responseBody)) { |
| 34 | if (item.type !== "line") continue; |
| 35 | const raw = item.raw; |
| 36 | const band = raw[0]; |
| 37 | const payload = raw.subarray(1); |
| 38 | |
| 39 | if (band === 1) { |
| 40 | reportStatusChunks.push(payload); |
| 41 | continue; |
| 42 | } |
| 43 | if (band === 2) { |
| 44 | progress.push(new TextDecoder().decode(payload)); |
| 45 | continue; |
| 46 | } |
| 47 | if (band === 3) { |
| 48 | fatal.push(new TextDecoder().decode(payload)); |
| 49 | } |
| 50 | } |
| 51 | |
| 52 | const reportStatus = |
| 53 | reportStatusChunks.length > 0 ? decodeReportStatus(concatChunks(reportStatusChunks)) : []; |
| 54 | |
| 55 | return { progress, fatal, reportStatus }; |
| 56 | } |
| 57 | |
| 58 | export async function buildStreamingReceiveBody(args: { |
| 59 | parentOid: string; |
| 60 | nextText: string; |
| 61 | commitMessage: string; |
| 62 | capabilities: string; |
| 63 | }): Promise<{ |
| 64 | body: Uint8Array; |
| 65 | blob: { oid: string }; |
| 66 | tree: { oid: string }; |
| 67 | commit: { oid: string }; |
| 68 | }> { |
| 69 | const author = "You <you@example.com> 0 +0000"; |
| 70 | const blobPayload = new TextEncoder().encode(args.nextText); |
| 71 | const blob = await encodeGitObject("blob", blobPayload); |
| 72 | const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blob.oid }]); |
| 73 | const tree = await encodeGitObject("tree", treePayload); |
| 74 | const commitPayload = new TextEncoder().encode( |
| 75 | `tree ${tree.oid}\n` + |
| 76 | `parent ${args.parentOid}\n` + |
| 77 | `author ${author}\n` + |
| 78 | `committer ${author}\n\n` + |
| 79 | `${args.commitMessage}\n` |
| 80 | ); |
| 81 | const commit = await encodeGitObject("commit", commitPayload); |
| 82 | const pack = await buildPack([ |
| 83 | { type: "blob", payload: blobPayload }, |
| 84 | { type: "tree", payload: treePayload }, |
| 85 | { type: "commit", payload: commitPayload }, |
| 86 | ]); |
| 87 | |
| 88 | return { |
| 89 | body: concatChunks([ |
| 90 | pktLine(`${args.parentOid} ${commit.oid} refs/heads/main\0 ${args.capabilities}\n`), |
| 91 | flushPkt(), |
| 92 | pack, |
| 93 | ]), |
| 94 | blob, |
| 95 | tree, |
| 96 | commit, |
| 97 | }; |
| 98 | } |
| 99 | |
| 100 | export async function pushStreamingUpdate( |
| 101 | owner: string, |
| 102 | repo: string, |
| 103 | parentOid: string, |
| 104 | nextText: string, |
| 105 | options?: { authHeader?: string } |
| 106 | ): Promise<{ |
| 107 | commitOid: string; |
| 108 | blob: { oid: string }; |
| 109 | tree: { oid: string }; |
| 110 | commit: { oid: string }; |
| 111 | }> { |
| 112 | const built = await buildStreamingReceiveBody({ |
| 113 | parentOid, |
| 114 | nextText, |
| 115 | commitMessage: "streaming update", |
| 116 | capabilities: "report-status ofs-delta agent=test", |
| 117 | }); |
| 118 | |
| 119 | const headers: Record<string, string> = { |
| 120 | "Content-Type": "application/x-git-receive-pack-request", |
| 121 | }; |
| 122 | const auth = options?.authHeader ?? lookupPushAuth(owner, repo); |
| 123 | if (auth) headers.Authorization = auth; |
| 124 | const response = await workerExports.default.fetch( |
| 125 | `https://example.com/${owner}/${repo}/git-receive-pack`, |
| 126 | { |
| 127 | method: "POST", |
| 128 | headers, |
| 129 | body: built.body, |
| 130 | } as any |
| 131 | ); |
| 132 | expect(response.status).toBe(200); |
| 133 | expect(decodeReportStatus(new Uint8Array(await response.arrayBuffer()))).toContain( |
| 134 | "ok refs/heads/main" |
| 135 | ); |
| 136 | return { commitOid: built.commit.oid, blob: built.blob, tree: built.tree, commit: built.commit }; |
| 137 | } |