From 5b2fe2235ea604bc35061219c2919897fb7ffe13 Mon Sep 17 00:00:00 2001 From: Daniel Roe Date: Thu, 11 Jun 2026 20:04:12 +0000 Subject: [PATCH] refactor: stream pushes via git protocol splice instead of clone --- nuxt.config.ts | 1 + package.json | 1 - pnpm-lock.yaml | 94 ---------------------------------------------------------------------------------------------- server/utils/git.ts | 55 ------------------------------------------------------- server/utils/splice.ts | 152 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ server/utils/ssh-cmd.ts | 27 +++++++++------------------ server/utils/sync-push-host.ts | 14 ++++++++++++++ server/utils/sync-push.ts | 116 ++++++++++++++++++++++++++++++++++++-------------------------------------------------------------------------------- server/utils/sync-ref.ts | 152 ++++++++++++++++++++++++++++++++++++++++++++++++++------------------------------------------------------------------------------------------------------ test/unit/git-wire-refs.spec.ts | 76 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ test/unit/pkt-line.spec.ts | 96 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ test/unit/receive-pack.spec.ts | 147 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ test/unit/splice.spec.ts | 113 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ test/unit/sync-push.spec.ts | 131 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ test/unit/sync-ref.spec.ts | 227 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------------------------------------------------------------------------------------------------------------------------------------- test/unit/upload-pack.spec.ts | 92 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ test/utils/git-wire.ts | 121 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ server/utils/git-wire/errors.ts | 76 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ server/utils/git-wire/pkt-line.ts | 158 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ server/utils/git-wire/receive-pack.ts | 210 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ server/utils/git-wire/refs.ts | 69 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ server/utils/git-wire/upload-pack.ts | 134 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 22 file(s) changed, 1763 insertion(s)(+), 499 deletion(s)(-) diff --git a/nuxt.config.ts b/nuxt.config.ts --- a/nuxt.config.ts +++ b/nuxt.config.ts @@ -19,6 +19,7 @@ githubWebhookSecret: '', cronSecret: '', workerBudgetMs: '', + maxPackBytes: '', encryptionKey: '', atprotoPrivateJwk: '', sessionPassword: '', diff --git a/package.json b/package.json --- a/package.json +++ b/package.json @@ -49,7 +49,6 @@ "@octokit/auth-app": "^8.2.0", "@octokit/webhooks-methods": "^6.0.0", "drizzle-orm": "^0.45.2", - "execa": "^9.6.1", "nuxt": "^4.4.4", "nuxt-og-image": "^6.4.11", "rolldown": "^1.0.0-rc.18", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -54,9 +54,6 @@ drizzle-orm: specifier: ^0.45.2 version: 0.45.2(@electric-sql/pglite@0.4.5)(@neondatabase/serverless@1.1.0) - execa: - specifier: ^9.6.1 - version: 9.6.1 nuxt: specifier: ^4.4.4 version: 4.4.4(@babel/core@7.29.0)(@babel/plugin-syntax-jsx@7.28.6(@babel/core@7.29.0))(@electric-sql/pglite@0.4.5)(@parcel/watcher@2.5.6)(@types/node@25.6.0)(@vue/compiler-sfc@3.5.33)(cac@6.7.14)(db0@0.3.4(@electric-sql/pglite@0.4.5)(drizzle-orm@0.45.2(@electric-sql/pglite@0.4.5)(@neondatabase/serverless@1.1.0)))(drizzle-orm@0.45.2(@electric-sql/pglite@0.4.5)(@neondatabase/serverless@1.1.0))(esbuild@0.28.0)(eslint@10.3.0(jiti@2.6.1))(ioredis@5.10.1)(magicast@0.5.2)(optionator@0.9.4)(oxlint@1.61.0(oxlint-tsgolint@0.22.0))(rolldown@1.0.0-rc.18)(rollup-plugin-visualizer@7.0.1(rolldown@1.0.0-rc.18)(rollup@4.60.2))(rollup@4.60.2)(srvx@0.11.15)(terser@5.46.2)(tsx@4.21.0)(typescript@6.0.3)(vite@7.3.2(@types/node@25.6.0)(jiti@2.6.1)(lightningcss@1.32.0)(terser@5.46.2)(tsx@4.21.0)(yaml@2.8.4))(vue-tsc@3.2.7(typescript@6.0.3))(yaml@2.8.4) @@ -2709,9 +2706,6 @@ cpu: [x64] os: [win32] - '@sec-ant/readable-stream@0.4.1': - resolution: {integrity: sha512-831qok9r2t8AlxLko40y2ebgSDhenenCatLVeW/uBtnHPyhHOvG0C7TvfgecV+wHzIm5KUICgzmVpWS+IMEAeg==} - '@sidvind/better-ajv-errors@3.0.1': resolution: {integrity: sha512-++1mEYIeozfnwWI9P1ECvOPoacy+CgDASrmGvXPMCcqgx0YUzB01vZ78uHdQ443V6sTY+e9MzHqmN9DOls02aw==} engines: {node: '>= 16.14'} @@ -3851,10 +3845,6 @@ resolution: {integrity: sha512-VyhnebXciFV2DESc+p6B+y0LjSm0krU4OgJN44qFAhBY0TJ+1V61tYD2+wHusZ6F9n5K+vl8k0sTy7PEfV4qpg==} engines: {node: '>=16.17'} - execa@9.6.1: - resolution: {integrity: sha512-9Be3ZoN4LmYR90tUoVu2te2BsbzHfhJyfEiAVfz7N5/zv+jduIfLrV2xdQXOHbaD6KgpGdO9PRPM1Y4Q9QkPkA==} - engines: {node: ^18.19.0 || >=20.5.0} - exsolve@1.0.8: resolution: {integrity: sha512-LmDxfWXwcTArk8fUEnOfSZpHOJ6zOMUJKOtFLFqJLoKJetuQG874Uc7/Kki7zFLzYybmZhp1M7+98pfMqeX8yA==} @@ -3917,10 +3907,6 @@ peerDependenciesMeta: picomatch: optional: true - - figures@6.1.0: - resolution: {integrity: sha512-d+l3qxjSesT4V7v2fh+QnmFnUWv9lSpjarhShNTgBOfA0ttejbQUAlHLitbjkoRiDulW0OPoQPYIGhIC8ohejg==} - engines: {node: '>=18'} file-entry-cache@8.0.0: resolution: {integrity: sha512-XXTUwCvisa5oacNGRP9SfNtYBNAMi+RPwBFmblZEF7N7swHYQS6/Zfk7SRwx4D5j3CH211YNRco1DEMNVfZCnQ==} @@ -4010,10 +3996,6 @@ get-stream@8.0.1: resolution: {integrity: sha512-VaUJspBffn/LMCJVoMvSAdmscJyS1auj5Zulnn5UoYcY531UWmdwhRWkcGKnGU93m5HSXP9LP2usOryrBtQowA==} engines: {node: '>=16'} - - get-stream@9.0.1: - resolution: {integrity: sha512-kVCxPF3vQM/N0B1PmoqVUqgHP+EeVjmZSQn+1oCRPxd2P21P2F19lIgbR3HBosbB1PUhOAoctJnfEn2GbN2eZA==} - engines: {node: '>=18'} get-tsconfig@4.14.0: resolution: {integrity: sha512-yTb+8DXzDREzgvYmh6s9vHsSVCHeC0G3PI5bEXNBHtmshPnO+S5O7qgLEOn0I5QvMy6kpZN8K1NKGyilLb93wA==} @@ -4126,10 +4108,6 @@ resolution: {integrity: sha512-AXcZb6vzzrFAUE61HnN4mpLqd/cSIwNQjtNWR0euPm6y0iqx3G4gOXaIDdtdDwZmhwe82LA6+zinmW4UBWVePQ==} engines: {node: '>=16.17.0'} - human-signals@8.0.1: - resolution: {integrity: sha512-eKCa6bwnJhvxj14kZk5NCPc6Hb6BdsU9DZcOnmQKSnO1VKrfV0zCvtttPZUsBvjmNDn8rpcJfpwSYnHBjc95MQ==} - engines: {node: '>=18.18.0'} - ieee754@1.2.1: resolution: {integrity: sha512-dcyqhDvX1C46lXZcVqCpK+FtMRQVdIMN6/Df5js2zouUsqG7I6sFxitIC+7KYK29KdXOLHdu9zL4sFnoVQnqaA==} @@ -4226,10 +4204,6 @@ resolution: {integrity: sha512-lJJV/5dYS+RcL8uQdBDW9c9uWFLLBNRyFhnAKXw5tVqLlKZ4RMGZKv+YQ/IA3OhD+RpbJa1LLFM1FQPGyIXvOA==} engines: {node: '>=12'} - is-plain-obj@4.1.0: - resolution: {integrity: sha512-+Pgi+vMuUNkJyExiMBt5IlFoMyKnr5zhJ4Uspz58WOhBF5QoIZkFyNHIbBAtHwzVAgk5RtndVNsDRN61/mmDqg==} - engines: {node: '>=12'} - is-reference@1.2.1: resolution: {integrity: sha512-U82MsXXiFIrjCK4otLT+o2NA2Cd2g5MLoOVXUZjIOhLurrRxpEXzI8O0KZHr3IjLvlAH1kTPYSuqer5T9ZVBKQ==} @@ -4240,14 +4214,6 @@ is-stream@3.0.0: resolution: {integrity: sha512-LnQR4bZ9IADDRSkvpqMGvt/tEJWclzklNgSw48V5EAaAeDd6qGvN8ei6k5p0tvxSR171VmGyHuTiAOfxAbr8kA==} engines: {node: ^12.20.0 || ^14.13.1 || >=16.0.0} - - is-stream@4.0.1: - resolution: {integrity: sha512-Dnz92NInDqYckGEUJv689RbRiTSEHCQ7wOVeALbkOz999YpqT46yMRIGtSNl2iCL1waAZSx40+h59NV/EwzV/A==} - engines: {node: '>=18'} - - is-unicode-supported@2.1.0: - resolution: {integrity: sha512-mE00Gnza5EEB3Ds0HfMyllZzbBrmLOX3vfWoj9A9PEnTfratQ/BcaJOuMhnkhjXvb2+FkY3VuHqtAGpTPmglFQ==} - engines: {node: '>=18'} is-wsl@2.2.0: resolution: {integrity: sha512-fKzAra0rGJUUBwGBgNkHZuToZcn+TtXHpeCgmkMJMMYx1sQDYaCSyjJBSCa2nH1DGm7s3n1oBnohoVTBaN7Lww==} @@ -4818,10 +4784,6 @@ package-json-from-dist@1.0.1: resolution: {integrity: sha512-UEZIS3/by4OC8vL3P2dTXRETpebLI2NiI5vIrjaD/5UtrkFX/tNbwjTSRAGC/+7CAo2pIcBaRgWmcBBHcsaCIw==} - parse-ms@4.0.0: - resolution: {integrity: sha512-TXfryirbmq34y8QBwgqCVLi+8oA3oWx2eAnSn62ITyEhEYaWRlVZ2DvMM9eZbMs/RfxPu/PK/aBLyGj4IrqMHw==} - engines: {node: '>=18'} - parseurl@1.3.3: resolution: {integrity: sha512-CiyeOxFT/JZyN5m0z9PfXw4SCBJ6Sygz1Dpl0wqjlhDEGGBP1GnsUVEL0p63hoG1fcj3fHynXi9NYO4nWOL+qQ==} engines: {node: '>= 0.8'} @@ -5085,10 +5047,6 @@ pretty-bytes@7.1.0: resolution: {integrity: sha512-nODzvTiYVRGRqAOvE84Vk5JDPyyxsVk0/fbA/bq7RqlnhksGpset09XTxbpvLTIjoaF7K8Z8DG8yHtKGTPSYRw==} engines: {node: '>=20'} - - pretty-ms@9.3.0: - resolution: {integrity: sha512-gjVS5hOP+M3wMm5nmNOucbIrqudzs9v/57bWRHQWLYklXqoXKrVfYW2W9+glfGsqtPgpiz5WwyEEB+ksXIx3gQ==} - engines: {node: '>=18'} process-nextick-args@2.0.1: resolution: {integrity: sha512-3ouUOpQhtgrbOa17J7+uxOTpITYWaGP7/AhoR3+A+/1e9skrzelGi/dXzEYyvbxubEF6Wn2ypscTKiKJFFn1ag==} @@ -5373,10 +5331,6 @@ strip-final-newline@3.0.0: resolution: {integrity: sha512-dOESqjYr96iWYylGObzd39EuNTa5VJxyvVAEm5Jnh7KGo75V43Hk1odPQkNDyXNmUR6k+gEiDVXnjB8HJ3crXw==} engines: {node: '>=12'} - - strip-final-newline@4.0.0: - resolution: {integrity: sha512-aulFJcD6YK8V1G7iRB5tigAP4TsHBZZrOV8pjV++zdUwmeV8uzbY7yn6h9MswN62adStNZFuCIx4haBnRuMDaw==} - engines: {node: '>=18'} strip-literal@3.1.0: resolution: {integrity: sha512-8r3mkIM/2+PpjHoOtiAW8Rg3jJLHaV7xPwG+YRGrv6FP0wwk/toTpATxWYOW0BKdWwl82VT2tFYi5DlROa0Mxg==} @@ -5921,10 +5875,6 @@ yocto-queue@0.1.0: resolution: {integrity: sha512-rVksvsnNCdJ/ohGc6xgPwyN8eheCxsiLM8mxuE/t/mOVqJewPuO1miLpTHQiRgTKCLexL4MeAFVagts7HmNZ2Q==} engines: {node: '>=10'} - - yoctocolors@2.1.2: - resolution: {integrity: sha512-CzhO+pFNo8ajLM2d2IW/R93ipy99LWjtwblvC1RsoSUMZgyLbYFr221TnSNT7GjGdYui6P459mw9JH/g/zW2ug==} - engines: {node: '>=18'} youch-core@0.3.3: resolution: {integrity: sha512-ho7XuGjLaJ2hWHoK8yFnsUGy2Y5uDpqSTq1FkHLK4/oqKtyUU1AFbOOxY4IpC9f0fTLjwYbslUz0Po5BpD1wrA==} @@ -8179,8 +8129,6 @@ '@rollup/rollup-win32-x64-msvc@4.60.2': optional: true - '@sec-ant/readable-stream@0.4.1': {} - '@sidvind/better-ajv-errors@3.0.1(ajv@8.20.0)': dependencies: ajv: 8.20.0 @@ -9324,21 +9272,6 @@ signal-exit: 4.1.0 strip-final-newline: 3.0.0 - execa@9.6.1: - dependencies: - '@sindresorhus/merge-streams': 4.0.0 - cross-spawn: 7.0.6 - figures: 6.1.0 - get-stream: 9.0.1 - human-signals: 8.0.1 - is-plain-obj: 4.1.0 - is-stream: 4.0.1 - npm-run-path: 6.0.0 - pretty-ms: 9.3.0 - signal-exit: 4.1.0 - strip-final-newline: 4.0.0 - yoctocolors: 2.1.2 - exsolve@1.0.8: {} fake-indexeddb@6.2.5: {} @@ -9392,10 +9325,6 @@ fdir@6.5.0(picomatch@4.0.4): optionalDependencies: picomatch: 4.0.4 - - figures@6.1.0: - dependencies: - is-unicode-supported: 2.1.0 file-entry-cache@8.0.0: dependencies: @@ -9501,11 +9430,6 @@ get-port-please@3.2.0: {} get-stream@8.0.1: {} - - get-stream@9.0.1: - dependencies: - '@sec-ant/readable-stream': 0.4.1 - is-stream: 4.0.1 get-tsconfig@4.14.0: dependencies: @@ -9632,8 +9556,6 @@ human-signals@5.0.0: {} - human-signals@8.0.1: {} - ieee754@1.2.1: {} ignore@5.3.2: {} @@ -9750,8 +9672,6 @@ is-path-inside@4.0.0: {} - is-plain-obj@4.1.0: {} - is-reference@1.2.1: dependencies: '@types/estree': 1.0.8 @@ -9759,10 +9679,6 @@ is-stream@2.0.1: {} is-stream@3.0.0: {} - - is-stream@4.0.1: {} - - is-unicode-supported@2.1.0: {} is-wsl@2.2.0: dependencies: @@ -10642,8 +10558,6 @@ package-json-from-dist@1.0.1: {} - parse-ms@4.0.0: {} - parseurl@1.3.3: {} path-browserify@1.0.1: {} @@ -10875,10 +10789,6 @@ prettier@3.8.3: {} pretty-bytes@7.1.0: {} - - pretty-ms@9.3.0: - dependencies: - parse-ms: 4.0.0 process-nextick-args@2.0.1: {} @@ -11218,8 +11128,6 @@ ansi-regex: 6.2.2 strip-final-newline@3.0.0: {} - - strip-final-newline@4.0.0: {} strip-literal@3.1.0: dependencies: @@ -11779,8 +11687,6 @@ yargs-parser: 22.0.0 yocto-queue@0.1.0: {} - - yoctocolors@2.1.2: {} youch-core@0.3.3: dependencies: diff --git a/server/utils/git.ts b/server/utils/git.ts deleted file mode 100644 --- a/server/utils/git.ts +++ /dev/null @@ -1,55 +0,0 @@ -import { execa, type Options } from 'execa' - -/** - * Thin wrapper over `execa` for invoking the system `git` binary with - * predictable defaults. - * - * - Forces non-interactive mode so a misconfigured ssh setup never hangs - * waiting for a passphrase or `yes/no` prompt. - * - Captures stderr so callers can produce useful error messages. - * - Adds a default 60s timeout; callers can override via `options.timeout`. - */ -export async function git(args: string[], options: Options = {}): Promise<{ stdout: string, stderr: string }> { - const result = await execa('git', args, { - timeout: 60_000, - ...options, - env: { - // Belt and braces against interactive prompts. `GIT_TERMINAL_PROMPT=0` - // makes git fail rather than hang if it would otherwise ask for input - // (e.g. credentials). - GIT_TERMINAL_PROMPT: '0', - // Don't pick up the running user's ssh config / known_hosts. The caller - // supplies a complete GIT_SSH_COMMAND for ssh transports. - GIT_CONFIG_NOSYSTEM: '1', - ...options.env, - }, - // Buffer (default) is fine for small operations; for very large fetches - // we'd want to stream stderr instead. - reject: true, - all: true, - }) - return { stdout: String(result.stdout), stderr: String(result.stderr) } -} - -/** - * Recognised remote rejection patterns from the knot when a repo no longer - * exists or our key has been revoked. Surfaces as a typed error so the - * worker can mark the mapping as terminally failed rather than retry forever. - */ -export class RemoteRejectedPushError extends Error { - constructor(message: string, public readonly reason: 'repo-gone' | 'auth-rejected' | 'other') { - super(message) - this.name = 'RemoteRejectedPushError' - } -} - -export function classifyPushFailure(stderr: string): RemoteRejectedPushError | null { - const lc = stderr.toLowerCase() - if (lc.includes('repository not found') || lc.includes('does not exist') || lc.includes('does not appear to be a git repository')) { - return new RemoteRejectedPushError(stderr.trim(), 'repo-gone') - } - if (lc.includes('permission denied') || lc.includes('publickey') && lc.includes('denied')) { - return new RemoteRejectedPushError(stderr.trim(), 'auth-rejected') - } - return null -} diff --git a/server/utils/splice.ts b/server/utils/splice.ts new file mode 100644 --- /dev/null +++ b/server/utils/splice.ts @@ -0,0 +1,152 @@ +import { + type ReceivePackFactory, + ReceivePackSession, + type RefUpdate, + sshReceivePackFactory, +} from './git-wire/receive-pack' +import { ZERO_SHA } from './git-wire/refs' +import { fetchAdvertisement, fetchPack } from './git-wire/upload-pack' +import { loadSshArgsForInstall } from './ssh-cmd' +import { sshEndpointForKnot } from './sync-push-host' + +const DEFAULT_MAX_PACK_BYTES = 1024 * 1024 * 1024 +/** Cap haves so a repo with thousands of refs can't bloat the negotiation. */ +const MAX_HAVES = 256 + +function maxPackBytes(): number { + const raw = process.env.NUXT_MAX_PACK_BYTES + if (!raw) return DEFAULT_MAX_PACK_BYTES + const n = Number.parseInt(raw, 10) + return Number.isNaN(n) || n <= 0 ? DEFAULT_MAX_PACK_BYTES : n +} + +async function sshFactory(installationId: number, knot: string, repoDid: string): Promise<{ + factory: ReceivePackFactory + cleanup: () => void +}> { + const { args, cleanup } = await loadSshArgsForInstall(installationId) + const { host, port } = sshEndpointForKnot(knot) + // ssh:// path form: leading slash, the knot resolves the repo by DID. + const factory = sshReceivePackFactory({ host, port, repoPath: `/${repoDid}`, sshArgs: args }) + return { factory, cleanup } +} + +export interface SplicePushParams { + installationId: number + repoFullName: string + knot: string + repoDid: string + /** Fully-qualified ref, e.g. `refs/heads/main`. */ + ref: string + /** The SHA to land on the knot. */ + want: string + /** GitHub installation token authorising the fetch. */ + token: string +} + +export interface SplicePushResult { + status: 'synced' | 'already-synced' + sha: string +} + +/** + * Stream a single ref update from GitHub to the knot without materialising a + * repository: + * + * 1. open receive-pack, read the knot's tips; + * 2. if the knot's tip for `ref` already equals `want`, no-op; + * 3. fetch a thin pack from GitHub with the knot's tips as haves; + * 4. send the compare-and-swap command and pipe the pack straight through; + * 5. read report-status. + * + * Steps 1 and 3 share one ssh session: it sits idle for the duration of the + * GitHub round-trip (receive-pack waits indefinitely for commands), which + * keeps the knot's advertised tip as the authoritative compare-and-swap base. + */ +export async function splicePush(params: SplicePushParams): Promise { + const { factory, cleanup } = await sshFactory(params.installationId, params.knot, params.repoDid) + try { + return await runSplice(factory, params) + } + finally { + cleanup() + } +} + +/** The fetch + push exchange over an open session. Split out for the wire test. */ +export async function runSplice( + factory: ReceivePackFactory, + params: { repoFullName: string, ref: string, want: string, token: string }, +): Promise { + const session = await ReceivePackSession.open(factory) + let pushStarted = false + try { + const old = session.tips.get(params.ref) ?? ZERO_SHA + if (old === params.want) { + await session.close() + return { status: 'already-synced', sha: params.want } + } + + const haves = [...new Set(session.tips.values())] + .filter(sha => sha !== ZERO_SHA) + .slice(0, MAX_HAVES) + + const { pack } = await fetchPack({ + repoFullName: params.repoFullName, + token: params.token, + want: params.want, + haves, + maxBytes: maxPackBytes(), + }) + + const update: RefUpdate = { ref: params.ref, old, next: params.want } + pushStarted = true + await session.push([update], pack) + return { status: 'synced', sha: params.want } + } + finally { + // `push` tears the session down itself; only close here if we threw before + // reaching it (e.g. the byte cap fired inside fetchPack's stream). + if (!pushStarted) await session.close() + } +} + +export interface SpliceDeleteResult { + status: 'synced' | 'already-absent' +} + +/** + * Delete a ref on the knot. No GitHub leg and no pack: read the knot's + * advertisement, and if the ref is absent we're already done (idempotent). + * Otherwise send a delete command with the advertised value as the + * compare-and-swap base. + */ +export async function spliceDelete(params: { + installationId: number + knot: string + repoDid: string + ref: string +}): Promise { + const { factory, cleanup } = await sshFactory(params.installationId, params.knot, params.repoDid) + try { + return await runSpliceDelete(factory, params.ref) + } + finally { + cleanup() + } +} + +/** The delete exchange over an open session. Split out for the wire test. */ +export async function runSpliceDelete(factory: ReceivePackFactory, ref: string): Promise { + const session = await ReceivePackSession.open(factory) + const old = session.tips.get(ref) + if (!old || old === ZERO_SHA) { + await session.close() + return { status: 'already-absent' } + } + // push() owns teardown for the success and rejection paths. + await session.push([{ ref, old, next: ZERO_SHA }], null) + return { status: 'synced' } +} + +export { fetchAdvertisement } diff --git a/server/utils/ssh-cmd.ts b/server/utils/ssh-cmd.ts --- a/server/utils/ssh-cmd.ts +++ b/server/utils/ssh-cmd.ts @@ -10,7 +10,8 @@ /** * Materialise the install's SSH private key as an OpenSSH-format file on disk * and return: - * - the `GIT_SSH_COMMAND` string to point `git` at it + * - `args`: the ssh option list (`-i -o ...`) ready to splice into a + * `spawn('ssh', [...args, target, command])` call * - a `cleanup()` callback that synchronously removes the temp dir * * The key file lives in `os.tmpdir()` with 0600 perms, has a random filename @@ -24,8 +25,8 @@ * commit can ship pinned host keys for the canonical knots once we know what * those are. */ -export async function loadSshCommandForInstall(installationId: number): Promise<{ - gitSshCommand: string +export async function loadSshArgsForInstall(installationId: number): Promise<{ + args: string[] cleanup: () => void }> { const db = useDb() @@ -54,18 +55,17 @@ chmodSync(keyPath, 0o600) writeFileSync(knownHostsPath, '', { mode: 0o600 }) - const gitSshCommand = [ - 'ssh', - '-i', shellQuote(keyPath), - '-o', `UserKnownHostsFile=${shellQuote(knownHostsPath)}`, + const args = [ + '-i', keyPath, + '-o', `UserKnownHostsFile=${knownHostsPath}`, '-o', 'StrictHostKeyChecking=accept-new', '-o', 'IdentitiesOnly=yes', '-o', 'BatchMode=yes', '-o', 'ConnectTimeout=15', - ].join(' ') + ] return { - gitSshCommand, + args, cleanup: () => { try { rmSync(dir, { recursive: true, force: true }) @@ -75,13 +75,4 @@ } }, } -} - -/** Minimal shell-quoting for paths inside GIT_SSH_COMMAND. */ -function shellQuote(s: string): string { - // GIT_SSH_COMMAND is split on whitespace by git, so escape spaces. We don't - // bother with full shell-quoting here because the paths we generate (in - // os.tmpdir()) won't contain quotes/backslashes; this is defense in depth. - if (!/[\s"'\\]/.test(s)) return s - return `"${s.replace(/(["\\])/g, '\\$1')}"` } diff --git a/server/utils/sync-push-host.ts b/server/utils/sync-push-host.ts --- a/server/utils/sync-push-host.ts +++ b/server/utils/sync-push-host.ts @@ -14,3 +14,17 @@ if (knot === 'knot1.tangled.sh') return 'tangled.org' return knot } + +/** + * Split a knot value into the ssh host and optional port. Self-hosted knots + * may carry a `:port` suffix for a non-default ssh port; the appview-hosted + * knot maps through `sshHostForKnot` and has no port. + */ +export function sshEndpointForKnot(knot: string): { host: string, port?: number } { + const mapped = sshHostForKnot(knot) + const colon = mapped.lastIndexOf(':') + if (colon === -1) return { host: mapped } + const port = Number.parseInt(mapped.slice(colon + 1), 10) + if (Number.isNaN(port)) return { host: mapped } + return { host: mapped.slice(0, colon), port } +} diff --git a/server/utils/sync-push.ts b/server/utils/sync-push.ts --- a/server/utils/sync-push.ts +++ b/server/utils/sync-push.ts @@ -1,13 +1,9 @@ -import { mkdtempSync, rmSync } from 'node:fs' -import os from 'node:os' -import path from 'node:path' import { and, eq, sql } from 'drizzle-orm' import { repoMapping } from '../db/schema' import { useDb } from './db' -import { classifyPushFailure, git, RemoteRejectedPushError } from './git' +import { RemoteRejectedError } from './git-wire/errors' import { installationOctokit } from './github-app' -import { loadSshCommandForInstall } from './ssh-cmd' -import { sshHostForKnot } from './sync-push-host' +import { splicePush } from './splice' const ZERO_SHA = '0000000000000000000000000000000000000000' @@ -30,17 +26,18 @@ * 1. Look up the repo_mapping (installationId, githubRepoId). Skip if absent * or disabled. * 2. Ref-tip dedupe: if lastSyncedRefs[ref] === after, no-op. Guards against - * GitHub redeliveries and v1.1's tangled-primary loop (PLAN.md). - * 3. Skip ref deletions (after = 0000…). Handled by github.delete in commit 13. - * 4. Bare-init /tmp scratch; fetch `after` from GitHub via smart-HTTP using - * the install token; push that ref to the knot over SSH with the - * install's key, force-with-lease against our last known tip. + * GitHub redeliveries. This is a cache only; correctness comes from the + * protocol-level compare-and-swap in the splice. + * 3. Skip ref deletions (after = 0000…). Handled by github.delete. + * 4. Splice: open receive-pack to the knot, fetch a thin pack of `after` + * from GitHub with the knot's tips as haves, pipe it straight through. + * Nothing touches disk. * 5. Update lastSyncedRefs[ref] = after. * - * On terminal failures (repo gone from knot, auth rejected) we mark the - * mapping as `status='error'` so the worker stops retrying. Transient - * failures (network blips, missing objects) re-throw and the queue retries - * with backoff. + * On terminal failures (repo gone from knot, auth rejected, pack too big) we + * mark the mapping `status='error'` so the worker stops retrying. A lost + * compare-and-swap (`stale-old-sha`) and other transient failures re-throw so + * the queue retries with backoff; the retry re-reads the knot's tip. */ export async function syncPush(payload: PushPayload): Promise { const db = useDb() @@ -62,87 +59,46 @@ const lastSynced = (row.lastSyncedRefs as Record)[payload.ref] if (lastSynced === payload.after) return { status: 'skipped', reason: 'already-synced' } - const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'synchub-push-')) - let sshCleanup: (() => void) | undefined + const octokit = await installationOctokit(payload.installationId) + const { token } = (await octokit.auth({ type: 'installation' })) as { token: string } try { - // 1. Bare init. No working tree, no objects until we fetch. - await git(['init', '--bare', '-q'], { cwd: tmpDir }) + const result = await splicePush({ + installationId: payload.installationId, + repoFullName: row.githubFullName, + knot: row.knot, + repoDid: row.tangledRepoDid, + ref: payload.ref, + want: payload.after, + token, + }) - // 2. Install-token-authed clone URL. The `x-access-token` username is - // GitHub's convention for installation tokens. - const octokit = await installationOctokit(payload.installationId) - const { token } = (await octokit.auth({ type: 'installation' })) as { token: string } - const githubUrl = `https://x-access-token:${token}@github.com/${row.githubFullName}.git` - - // 3. Fetch exactly the new ref. The `:` refspec asks git to - // fetch the object reachable from `after` and store it under our - // local refs/heads/... or refs/tags/... at the same name. - await git( - ['fetch', '--no-tags', '-q', githubUrl, `+${payload.after}:${payload.ref}`], - { cwd: tmpDir, timeout: 120_000 }, - ) - - // 4. Push to the knot. `force-with-lease` means "only update the ref if - // its current tip on the knot still matches what we last saw". Without - // a lease value we fall back to plain `--force` because we have no - // way to know the knot's current tip otherwise (we don't `ls-remote`). - // The lease is `` when we have one; on first - // sync we use plain force. - const { gitSshCommand, cleanup } = await loadSshCommandForInstall(payload.installationId) - sshCleanup = cleanup - - const knotUrl = `ssh://git@${sshHostForKnot(row.knot)}/${row.tangledRepoDid}` - const pushRefspec = lastSynced - ? `--force-with-lease=${payload.ref}:${lastSynced} ${payload.after}:${payload.ref}` - : `+${payload.after}:${payload.ref}` - - try { - await git( - ['push', '-q', knotUrl, ...pushRefspec.split(' ')], - { - cwd: tmpDir, - env: { GIT_SSH_COMMAND: gitSshCommand }, - timeout: 120_000, - }, - ) - } - catch (err) { - const stderr = err instanceof Error && 'stderr' in err ? String((err as { stderr: unknown }).stderr) : '' - const classified = classifyPushFailure(stderr) - if (classified?.reason === 'repo-gone') { - await markMappingError(row.id, 'knot reports repo no longer exists; stopping sync') - return { status: 'skipped', reason: 'repo-gone' } - } - throw classified ?? err - } - - // 5. Update last-synced tip for this ref. Use jsonb_set to leave other - // refs untouched. await db.update(repoMapping) .set({ - lastSyncedRefs: sql`jsonb_set(${repoMapping.lastSyncedRefs}, ${`{${jsonbPath(payload.ref)}}`}::text[], ${`"${payload.after}"`}::jsonb, true)`, + lastSyncedRefs: sql`jsonb_set(${repoMapping.lastSyncedRefs}, ${`{${jsonbPath(payload.ref)}}`}::text[], ${`"${result.sha}"`}::jsonb, true)`, updatedAt: new Date(), }) .where(eq(repoMapping.id, row.id)) return { status: 'synced' } } - finally { - sshCleanup?.() - try { - rmSync(tmpDir, { recursive: true, force: true }) + catch (err) { + if (err instanceof RemoteRejectedError && (err.reason === 'repo-gone' || err.reason === 'auth-rejected' || err.reason === 'too-big')) { + await markMappingError(row.id, terminalMessage(err)) + return { status: 'skipped', reason: 'repo-gone' } } - catch { - // best-effort - } + throw err } +} + +function terminalMessage(err: RemoteRejectedError): string { + if (err.reason === 'too-big') return `pack exceeded the configured size limit; stopping sync (${err.message})` + if (err.reason === 'auth-rejected') return 'knot rejected our ssh key; stopping sync' + return 'knot reports repo no longer exists; stopping sync' } /** jsonb_set path argument: `refs/heads/main` becomes a single text array element. */ function jsonbPath(ref: string): string { - // Escape any double-quotes inside the ref. We only support standard git ref - // names which never contain quotes, but be defensive. return `"${ref.replaceAll('"', '\\"')}"` } @@ -153,4 +109,4 @@ .where(eq(repoMapping.id, mappingId)) } -export { RemoteRejectedPushError } +export { RemoteRejectedError } diff --git a/server/utils/sync-ref.ts b/server/utils/sync-ref.ts --- a/server/utils/sync-ref.ts +++ b/server/utils/sync-ref.ts @@ -1,13 +1,9 @@ -import { mkdtempSync, rmSync } from 'node:fs' -import os from 'node:os' -import path from 'node:path' import { and, eq, sql } from 'drizzle-orm' import { repoMapping } from '../db/schema' import { useDb } from './db' -import { classifyPushFailure, git } from './git' +import { RemoteRejectedError, WireError } from './git-wire/errors' import { installationOctokit } from './github-app' -import { loadSshCommandForInstall } from './ssh-cmd' -import { sshHostForKnot } from './sync-push-host' +import { fetchAdvertisement, spliceDelete, splicePush } from './splice' export type RefType = 'branch' | 'tag' @@ -15,7 +11,7 @@ installationId: number githubRepoId: number refType: RefType - /** Short ref name as GitHub delivers it (e.g. `v1.0`, `feature-x`) \u2014 NOT + /** Short ref name as GitHub delivers it (e.g. `v1.0`, `feature-x`) — NOT * the `refs/...` qualified form. */ ref: string } @@ -32,10 +28,13 @@ * * Triggered by GitHub's `create` webhook event. For branches, GitHub also * sends a parallel `push` event (with `before = 0000…`), so the branch will - * usually have been created already by the time this fires \u2014 the push to - * knot is then a no-op via ref-tip dedupe. For lightweight and annotated - * tags, no `push` event is sent, so this is the only path that creates them - * on the knot. + * usually have been created already by the time this fires — the splice is + * then a no-op via the knot tip already matching. For lightweight and + * annotated tags, no `push` event is sent, so this is the only path that + * creates them on the knot. + * + * We resolve the ref name to a SHA from GitHub's advertisement (annotated tags + * resolve to the tag object), then splice that SHA to the knot. */ export async function syncCreateRef(payload: CreateRefPayload): Promise { if (payload.refType !== 'branch' && payload.refType !== 'tag') { @@ -46,60 +45,37 @@ if ('skip' in mapping) return mapping.skip const fullRef = qualifyRef(payload.refType, payload.ref) - const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'synchub-create-')) - let sshCleanup: (() => void) | undefined + + const octokit = await installationOctokit(payload.installationId) + const { token } = (await octokit.auth({ type: 'installation' })) as { token: string } + + const adv = await fetchAdvertisement(mapping.githubFullName, token) + const want = adv.refs.get(fullRef) + if (!want) { + // GitHub can deliver the create webhook before its replicas advertise the + // ref. Transient: re-throw so the queue retries with backoff. + throw new WireError(`github does not yet advertise ${fullRef} for ${mapping.githubFullName}`) + } try { - await git(['init', '--bare', '-q'], { cwd: tmpDir }) - - const octokit = await installationOctokit(payload.installationId) - const { token } = (await octokit.auth({ type: 'installation' })) as { token: string } - const githubUrl = `https://x-access-token:${token}@github.com/${mapping.githubFullName}.git` - - // Fetch the ref by name. Tags carry whatever object git stores at the - // ref (commit for lightweight; tag object for annotated); fetch gives - // us all the reachable objects either way. - await git( - ['fetch', '--no-tags', '-q', githubUrl, `+${fullRef}:${fullRef}`], - { cwd: tmpDir, timeout: 120_000 }, - ) - - const { gitSshCommand, cleanup } = await loadSshCommandForInstall(payload.installationId) - sshCleanup = cleanup - const knotUrl = `ssh://git@${sshHostForKnot(mapping.knot)}/${mapping.tangledRepoDid}` - - try { - await git( - ['push', '-q', knotUrl, `+${fullRef}:${fullRef}`], - { cwd: tmpDir, env: { GIT_SSH_COMMAND: gitSshCommand }, timeout: 120_000 }, - ) - } - catch (err) { - const stderr = err instanceof Error && 'stderr' in err ? String((err as { stderr: unknown }).stderr) : '' - const classified = classifyPushFailure(stderr) - if (classified?.reason === 'repo-gone') { - await markMappingError(mapping.id, 'knot reports repo no longer exists; stopping sync') - return { status: 'skipped', reason: 'repo-gone' } - } - throw classified ?? err - } - - // For branches we get the SHA from the local ref after fetch; for tags - // we still update lastSyncedRefs so a subsequent push event with the - // same SHA short-circuits via ref-tip dedupe. - const { stdout: sha } = await git(['rev-parse', fullRef], { cwd: tmpDir }) - await updateLastSyncedRef(mapping.id, fullRef, sha.trim()) - + const result = await splicePush({ + installationId: payload.installationId, + repoFullName: mapping.githubFullName, + knot: mapping.knot, + repoDid: mapping.tangledRepoDid, + ref: fullRef, + want, + token, + }) + await updateLastSyncedRef(mapping.id, fullRef, result.sha) return { status: 'synced' } } - finally { - sshCleanup?.() - try { - rmSync(tmpDir, { recursive: true, force: true }) + catch (err) { + if (err instanceof RemoteRejectedError && (err.reason === 'repo-gone' || err.reason === 'auth-rejected' || err.reason === 'too-big')) { + await markMappingError(mapping.id, 'knot reports repo no longer exists; stopping sync') + return { status: 'skipped', reason: 'repo-gone' } } - catch { - // best-effort - } + throw err } } @@ -107,11 +83,11 @@ * Mirror a branch or tag deletion from GitHub to the configured knot. * * Triggered by GitHub's `delete` webhook event. For branches, GitHub also - * sends a parallel `push` event with `after = 0000\u2026`, which `syncPush` - * currently skips (`reason: 'deletion'`) \u2014 this is the path that actually + * sends a parallel `push` event with `after = 0000…`, which `syncPush` + * currently skips (`reason: 'deletion'`) — this is the path that actually * removes the ref on the knot. Tag deletion arrives only via this event. * - * Deletion is idempotent: if the ref doesn't exist on the knot we treat it + * Deletion is idempotent: if the ref is already absent on the knot we treat it * as success. We wanted it gone, it's gone. */ export async function syncDeleteRef(payload: DeleteRefPayload): Promise { @@ -123,51 +99,23 @@ if ('skip' in mapping) return mapping.skip const fullRef = qualifyRef(payload.refType, payload.ref) - const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'synchub-delete-')) - let sshCleanup: (() => void) | undefined try { - // No fetch needed; we're only telling the remote to drop a ref. - await git(['init', '--bare', '-q'], { cwd: tmpDir }) - - const { gitSshCommand, cleanup } = await loadSshCommandForInstall(payload.installationId) - sshCleanup = cleanup - const knotUrl = `ssh://git@${sshHostForKnot(mapping.knot)}/${mapping.tangledRepoDid}` - - // The `:` (empty source) refspec means "delete on the remote". - try { - await git( - ['push', '-q', knotUrl, `:${fullRef}`], - { cwd: tmpDir, env: { GIT_SSH_COMMAND: gitSshCommand }, timeout: 60_000 }, - ) - } - catch (err) { - const stderr = err instanceof Error && 'stderr' in err ? String((err as { stderr: unknown }).stderr) : '' - // "remote ref does not exist" is success for our purposes \u2014 the ref is - // gone, which is what we wanted. - if (/remote ref does not exist|unable to delete.*does not exist/i.test(stderr)) { - await clearLastSyncedRef(mapping.id, fullRef) - return { status: 'synced' } - } - const classified = classifyPushFailure(stderr) - if (classified?.reason === 'repo-gone') { - await markMappingError(mapping.id, 'knot reports repo no longer exists; stopping sync') - return { status: 'skipped', reason: 'repo-gone' } - } - throw classified ?? err - } - + await spliceDelete({ + installationId: payload.installationId, + knot: mapping.knot, + repoDid: mapping.tangledRepoDid, + ref: fullRef, + }) await clearLastSyncedRef(mapping.id, fullRef) return { status: 'synced' } } - finally { - sshCleanup?.() - try { - rmSync(tmpDir, { recursive: true, force: true }) + catch (err) { + if (err instanceof RemoteRejectedError && err.reason === 'repo-gone') { + await markMappingError(mapping.id, 'knot reports repo no longer exists; stopping sync') + return { status: 'skipped', reason: 'repo-gone' } } - catch { - // best-effort - } + throw err } } diff --git a/test/unit/git-wire-refs.spec.ts b/test/unit/git-wire-refs.spec.ts new file mode 100644 --- /dev/null +++ b/test/unit/git-wire-refs.spec.ts @@ -0,0 +1,76 @@ +import { describe, expect, it } from 'vitest' +import { classifyNgReason, classifySshStderr } from '../../server/utils/git-wire/errors' +import { parseAdvertisement, ZERO_SHA } from '../../server/utils/git-wire/refs' + +function lines(...ls: string[]): Buffer[] { + return ls.map(l => Buffer.from(l)) +} + +describe('parseAdvertisement', () => { + it('parses a populated repo with capabilities on the first ref line', () => { + const adv = parseAdvertisement(lines( + 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa refs/heads/main\0report-status delete-refs thin-pack agent=git/2.39\n', + 'bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb refs/heads/dev\n', + )) + expect(adv.refs.get('refs/heads/main')).toBe('a'.repeat(40)) + expect(adv.refs.get('refs/heads/dev')).toBe('b'.repeat(40)) + expect(adv.capabilities.has('thin-pack')).toBe(true) + expect(adv.capabilities.has('delete-refs')).toBe(true) + }) + + it('skips the smart-HTTP service prelude', () => { + const adv = parseAdvertisement(lines( + '# service=git-upload-pack\n', + 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa refs/heads/main\0thin-pack\n', + )) + expect(adv.refs.get('refs/heads/main')).toBe('a'.repeat(40)) + expect(adv.capabilities.has('thin-pack')).toBe(true) + }) + + it('records peeled annotated-tag lines separately', () => { + const adv = parseAdvertisement(lines( + 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa refs/tags/v1\0thin-pack\n', + 'cccccccccccccccccccccccccccccccccccccccc refs/tags/v1^{}\n', + )) + expect(adv.refs.get('refs/tags/v1')).toBe('a'.repeat(40)) + expect(adv.peeled.get('refs/tags/v1')).toBe('c'.repeat(40)) + }) + + it('handles the empty-repo capabilities sentinel without inventing a ref', () => { + const adv = parseAdvertisement(lines( + `${ZERO_SHA} capabilities^{}\0report-status delete-refs\n`, + )) + expect(adv.refs.size).toBe(0) + expect(adv.capabilities.has('report-status')).toBe(true) + }) +}) + +describe('classifySshStderr', () => { + it('classifies repo-gone', () => { + expect(classifySshStderr('fatal: repository not found')?.reason).toBe('repo-gone') + expect(classifySshStderr('ERROR: does not exist')?.reason).toBe('repo-gone') + }) + + it('classifies auth-rejected', () => { + expect(classifySshStderr('git@host: Permission denied (publickey).')?.reason).toBe('auth-rejected') + }) + + it('returns null for unrecognised stderr', () => { + expect(classifySshStderr('warning: something benign')).toBeNull() + }) +}) + +describe('classifyNgReason', () => { + it('classifies a stale compare-and-swap as stale-old-sha', () => { + expect(classifyNgReason('non-fast-forward').reason).toBe('stale-old-sha') + expect(classifyNgReason('stale info').reason).toBe('stale-old-sha') + }) + + it('classifies a missing repo', () => { + expect(classifyNgReason('repository does not exist').reason).toBe('repo-gone') + }) + + it('falls back to other', () => { + expect(classifyNgReason('funny business').reason).toBe('other') + }) +}) diff --git a/test/unit/pkt-line.spec.ts b/test/unit/pkt-line.spec.ts new file mode 100644 --- /dev/null +++ b/test/unit/pkt-line.spec.ts @@ -0,0 +1,96 @@ +import { describe, expect, it } from 'vitest' +import { encodePktLine, flushPkt, lineToString, PktLineReader } from '../../server/utils/git-wire/pkt-line' + +async function* chunks(...parts: (Buffer | string)[]): AsyncGenerator { + for (const p of parts) yield typeof p === 'string' ? Buffer.from(p) : p +} + +describe('pkt-line', () => { + describe('encodePktLine', () => { + it('frames a payload with a 4-hex-digit length including the prefix', () => { + // "hello\n" is 6 bytes, +4 prefix = 10 = 0x000a. + expect(encodePktLine('hello\n').toString()).toBe('000ahello\n') + }) + + it('frames an empty payload as length 4', () => { + expect(encodePktLine('').toString()).toBe('0004') + }) + + it('does not append a trailing newline of its own', () => { + expect(encodePktLine('want abc').toString()).toBe('000cwant abc') + }) + + it('rejects payloads larger than the max', () => { + expect(() => encodePktLine(Buffer.alloc(65517))).toThrow(/too large/) + }) + + it('exposes the flush-pkt as 0000', () => { + expect(flushPkt.toString()).toBe('0000') + }) + }) + + describe('PktLineReader', () => { + it('decodes consecutive lines', async () => { + const r = new PktLineReader(chunks('0006a\n0006b\n')) + expect(lineToString((await r.next() as { data: Buffer }).data)).toBe('a') + expect(lineToString((await r.next() as { data: Buffer }).data)).toBe('b') + expect(await r.next()).toBeNull() + }) + + it('returns flush-pkts as a distinct type without ending iteration', async () => { + const r = new PktLineReader(chunks('0006a\n00000006b\n')) + expect((await r.next())!.type).toBe('line') + expect((await r.next())!.type).toBe('flush') + expect((await r.next())!.type).toBe('line') + expect(await r.next()).toBeNull() + }) + + it('reassembles a line split across chunk boundaries', async () => { + const r = new PktLineReader(chunks('00', '0a', 'hel', 'lo\n')) + expect(lineToString((await r.next() as { data: Buffer }).data)).toBe('hello') + }) + + it('readUntilFlush collects line payloads and stops at flush', async () => { + const r = new PktLineReader(chunks('0006a\n0006b\n0000')) + const lines = await r.readUntilFlush() + expect(lines!.map(l => lineToString(l))).toEqual(['a', 'b']) + }) + + it('hands raw trailing bytes back via remaining(), including pre-buffered ones', async () => { + // One pkt-line "NAK\n" then raw pack bytes "PACK..." arriving in the + // same chunk: the reader must not swallow the pack head. + const r = new PktLineReader(chunks('0008NAK\nPACK\x00\x01\x02', 'more')) + const first = await r.next() + expect(lineToString((first as { data: Buffer }).data)).toBe('NAK') + + let raw = Buffer.alloc(0) + for await (const c of r.remaining()) raw = Buffer.concat([raw, c]) + expect(raw.toString('binary')).toBe('PACK\x00\x01\x02more') + }) + + it('throws on a truncated length prefix', async () => { + const r = new PktLineReader(chunks('00')) + await expect(r.next()).rejects.toThrow(/truncated pkt-line length/) + }) + + it('throws on a truncated payload', async () => { + const r = new PktLineReader(chunks('000ahel')) + await expect(r.next()).rejects.toThrow(/wanted 10 bytes/) + }) + + it('throws on a reserved length', async () => { + const r = new PktLineReader(chunks('0001')) + await expect(r.next()).rejects.toThrow(/reserved pkt-line length/) + }) + }) + + describe('lineToString', () => { + it('strips a single trailing newline', () => { + expect(lineToString(Buffer.from('abc\n'))).toBe('abc') + }) + + it('leaves a line without a trailing newline intact', () => { + expect(lineToString(Buffer.from('abc'))).toBe('abc') + }) + }) +}) diff --git a/test/unit/receive-pack.spec.ts b/test/unit/receive-pack.spec.ts new file mode 100644 --- /dev/null +++ b/test/unit/receive-pack.spec.ts @@ -0,0 +1,147 @@ +import { execFileSync } from 'node:child_process' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { RemoteRejectedError } from '../../server/utils/git-wire/errors' +import { ReceivePackSession } from '../../server/utils/git-wire/receive-pack' +import { ZERO_SHA } from '../../server/utils/git-wire/refs' +import { fakeGithubFetch, GitFixture, localReceivePackFactory } from '../utils/git-wire' +import { fetchPack } from '../../server/utils/git-wire/upload-pack' + +async function* fromBuffer(b: Buffer): AsyncGenerator { + yield b +} + +async function push(factory: ReturnType, updates: Parameters[0], pack: AsyncIterable | null) { + const session = await ReceivePackSession.open(factory) + await session.push(updates, pack) +} + +async function drain(gen: AsyncGenerator): Promise { + const parts: Buffer[] = [] + for await (const c of gen) parts.push(c) + return Buffer.concat(parts) +} + +describe('receive-pack (against real git-receive-pack)', () => { + let fx: GitFixture + let realFetch: typeof globalThis.fetch + + beforeEach(() => { + fx = new GitFixture() + realFetch = globalThis.fetch + }) + + afterEach(() => { + globalThis.fetch = realFetch + fx.cleanup() + }) + + /** Build a pack on disk for `want` and return it as a single buffer. */ + async function packFor(ghBare: string, want: string, haves: string[]): Promise { + globalThis.fetch = fakeGithubFetch(new Map([['owner/repo', ghBare]])) as unknown as typeof globalThis.fetch + const { pack } = await fetchPack({ repoFullName: 'owner/repo', token: 't', want, haves, maxBytes: 1 << 30 }) + return drain(pack) + } + + it('pushes a new ref into an empty knot repo', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + + const pack = await packFor(gh, sha, []) + await push( + localReceivePackFactory(knot), + [{ ref: 'refs/heads/main', old: ZERO_SHA, next: sha }], + fromBuffer(pack), + ) + + expect(fx.revParse(knot, 'refs/heads/main')).toBe(sha) + }) + + it('fast-forwards an existing ref with a thin incremental pack', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const first = fx.commit(work, 'a.txt', 'one') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + await push(localReceivePackFactory(knot), [{ ref: 'refs/heads/main', old: ZERO_SHA, next: first }], fromBuffer(await packFor(gh, first, []))) + + const second = fx.commit(work, 'b.txt', 'two') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + await push(localReceivePackFactory(knot), [{ ref: 'refs/heads/main', old: first, next: second }], fromBuffer(await packFor(gh, second, [first]))) + + expect(fx.revParse(knot, 'refs/heads/main')).toBe(second) + }) + + it('rejects a stale compare-and-swap as stale-old-sha', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const first = fx.commit(work, 'a.txt', 'one') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + await push(localReceivePackFactory(knot), [{ ref: 'refs/heads/main', old: ZERO_SHA, next: first }], fromBuffer(await packFor(gh, first, []))) + + const second = fx.commit(work, 'b.txt', 'two') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + + // Claim the knot is still empty when it actually points at `first`. + await expect( + push(localReceivePackFactory(knot), [{ ref: 'refs/heads/main', old: ZERO_SHA, next: second }], fromBuffer(await packFor(gh, second, [first]))), + ).rejects.toMatchObject({ constructor: RemoteRejectedError, reason: 'stale-old-sha' }) + }) + + it('deletes a ref with no pack', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + await push(localReceivePackFactory(knot), [{ ref: 'refs/heads/main', old: ZERO_SHA, next: sha }], fromBuffer(await packFor(gh, sha, []))) + + await push(localReceivePackFactory(knot), [{ ref: 'refs/heads/main', old: sha, next: ZERO_SHA }], null) + + expect(() => execFileSync('git', ['rev-parse', 'refs/heads/main'], { cwd: knot })).toThrow(/unknown revision|ambiguous argument|fatal/) + }) + + it('kills a stalled session once the watchdog fires', async () => { + // A factory whose child accepts the connection but never advertises: the + // open() read would block forever without the watchdog. + const stalled = () => { + let resolveDone: (code: number | null) => void + const done = new Promise(r => { resolveDone = r }) + // eslint-disable-next-line require-yield -- models a stalled stream that blocks until killed and never emits + async function* neverYields(): AsyncGenerator { + await done + } + return { + stdin: { write: (_d: unknown, cb?: (e?: Error) => void) => cb?.(), end: () => {} } as unknown as NodeJS.WritableStream, + stdout: neverYields(), + stderr: () => '', + kill: () => resolveDone(null), + done, + } + } + await expect(ReceivePackSession.open(stalled, 50)).rejects.toThrow(/end of stream|advertisement/) + }) + + it('pushes an annotated tag', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + fx.git(['tag', '-a', 'v1', '-m', 'release'], work) + const tagSha = fx.git(['rev-parse', 'refs/tags/v1'], work) + fx.pushTo(work, gh, 'refs/tags/v1:refs/tags/v1') + const knot = fx.initBare('knot.git') + + await push( + localReceivePackFactory(knot), + [{ ref: 'refs/tags/v1', old: ZERO_SHA, next: tagSha }], + fromBuffer(await packFor(gh, tagSha, [])), + ) + + expect(fx.revParse(knot, 'refs/tags/v1')).toBe(tagSha) + expect(fx.git(['cat-file', '-t', 'refs/tags/v1'], knot)).toBe('tag') + }) +}) diff --git a/test/unit/splice.spec.ts b/test/unit/splice.spec.ts new file mode 100644 --- /dev/null +++ b/test/unit/splice.spec.ts @@ -0,0 +1,113 @@ +import { execFileSync } from 'node:child_process' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { RemoteRejectedError } from '../../server/utils/git-wire/errors' +import { runSplice, runSpliceDelete } from '../../server/utils/splice' +import { fakeGithubFetch, GitFixture, localReceivePackFactory } from '../utils/git-wire' + +describe('splice (end-to-end against real git binaries)', () => { + let fx: GitFixture + let realFetch: typeof globalThis.fetch + + beforeEach(() => { + fx = new GitFixture() + realFetch = globalThis.fetch + }) + + afterEach(() => { + globalThis.fetch = realFetch + fx.cleanup() + }) + + function wireGithub(ghBare: string) { + globalThis.fetch = fakeGithubFetch(new Map([['owner/repo', ghBare]])) as unknown as typeof globalThis.fetch + } + + it('mirrors a first push into an empty knot, transferring objects end to end', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello world') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + wireGithub(gh) + + const result = await runSplice(localReceivePackFactory(knot), { + repoFullName: 'owner/repo', + ref: 'refs/heads/main', + want: sha, + token: 'tok', + }) + + expect(result).toEqual({ status: 'synced', sha }) + expect(fx.revParse(knot, 'refs/heads/main')).toBe(sha) + expect(execFileSync('git', ['cat-file', '-p', `${sha}:a.txt`], { cwd: knot, encoding: 'utf8' })).toBe('hello world') + }) + + it('streams only the delta on a follow-up push (thin pack via knot tips as haves)', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const first = fx.commit(work, 'a.txt', 'a'.repeat(8000)) + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + wireGithub(gh) + + await runSplice(localReceivePackFactory(knot), { repoFullName: 'owner/repo', ref: 'refs/heads/main', want: first, token: 'tok' }) + + const second = fx.commit(work, 'b.txt', 'b'.repeat(8000)) + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const result = await runSplice(localReceivePackFactory(knot), { repoFullName: 'owner/repo', ref: 'refs/heads/main', want: second, token: 'tok' }) + + expect(result).toEqual({ status: 'synced', sha: second }) + expect(fx.revParse(knot, 'refs/heads/main')).toBe(second) + }) + + it('no-ops when the knot tip already equals want', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + wireGithub(gh) + + await runSplice(localReceivePackFactory(knot), { repoFullName: 'owner/repo', ref: 'refs/heads/main', want: sha, token: 'tok' }) + const again = await runSplice(localReceivePackFactory(knot), { repoFullName: 'owner/repo', ref: 'refs/heads/main', want: sha, token: 'tok' }) + + expect(again).toEqual({ status: 'already-synced', sha }) + }) + + it('aborts a push that exceeds the byte cap', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'x'.repeat(20_000)) + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + wireGithub(gh) + + const prev = process.env.NUXT_MAX_PACK_BYTES + process.env.NUXT_MAX_PACK_BYTES = '10' + try { + await expect( + runSplice(localReceivePackFactory(knot), { repoFullName: 'owner/repo', ref: 'refs/heads/main', want: sha, token: 'tok' }), + ).rejects.toMatchObject({ constructor: RemoteRejectedError, reason: 'too-big' }) + } + finally { + if (prev === undefined) delete process.env.NUXT_MAX_PACK_BYTES + else process.env.NUXT_MAX_PACK_BYTES = prev + } + }) + + it('deletes a ref and treats an already-absent ref as success', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const knot = fx.initBare('knot.git') + wireGithub(gh) + + await runSplice(localReceivePackFactory(knot), { repoFullName: 'owner/repo', ref: 'refs/heads/main', want: sha, token: 'tok' }) + + expect(await runSpliceDelete(localReceivePackFactory(knot), 'refs/heads/main')).toEqual({ status: 'synced' }) + expect(() => fx.revParse(knot, 'refs/heads/main')).toThrow(/unknown revision|ambiguous argument|fatal/) + + expect(await runSpliceDelete(localReceivePackFactory(knot), 'refs/heads/gone')).toEqual({ status: 'already-absent' }) + }) +}) diff --git a/test/unit/sync-push.spec.ts b/test/unit/sync-push.spec.ts new file mode 100644 --- /dev/null +++ b/test/unit/sync-push.spec.ts @@ -0,0 +1,131 @@ +import crypto from 'node:crypto' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { installation, repoMapping } from '../../server/db/schema' +import { clearDb, setDb, useDb } from '../../server/utils/db' +import { clearEncryptionKeyCache } from '../../server/utils/encryption' +import { RemoteRejectedError } from '../../server/utils/git-wire/errors' +import { createTestDb } from '../utils/db' + +const ORIGINAL_ENC_KEY = process.env.NUXT_ENCRYPTION_KEY +const ZERO = '0'.repeat(40) + +const splicePushMock = vi.fn<(params: Record) => Promise<{ status: string, sha: string }>>() +const octokitAuthMock = vi.fn<(input: { type: 'installation' }) => Promise<{ token: string }>>() + +vi.mock('../../server/utils/splice', () => ({ + splicePush: (params: Record) => splicePushMock(params), +})) + +vi.mock('../../server/utils/github-app', () => ({ + installationOctokit: async () => ({ auth: octokitAuthMock }), +})) + +const { syncPush } = await import('../../server/utils/sync-push') + +describe('sync-push', () => { + beforeEach(async () => { + process.env.NUXT_ENCRYPTION_KEY = crypto.randomBytes(32).toString('base64') + clearEncryptionKeyCache() + + setDb(await createTestDb()) + await useDb().insert(installation).values({ + id: 1, accountLogin: 'alice', accountId: 100, accountType: 'User', + }) + + splicePushMock.mockReset() + octokitAuthMock.mockReset() + octokitAuthMock.mockResolvedValue({ token: 'install-token' }) + splicePushMock.mockResolvedValue({ status: 'synced', sha: 'a'.repeat(40) }) + }) + + afterEach(() => { + if (ORIGINAL_ENC_KEY === undefined) delete process.env.NUXT_ENCRYPTION_KEY + else process.env.NUXT_ENCRYPTION_KEY = ORIGINAL_ENC_KEY + clearEncryptionKeyCache() + clearDb() + }) + + async function seedMapping(over: Partial = {}) { + await useDb().insert(repoMapping).values({ + installationId: 1, + githubRepoId: 9001, + githubFullName: 'alice/my-project', + tangledRepoDid: 'did:plc:repo-xyz', + tangledFullName: 'did:plc:abc/my-project', + knot: 'knot1.tangled.sh', + status: 'active', + ...over, + }) + } + + const payload = (over: Record = {}) => ({ + installationId: 1, + githubRepoId: 9001, + ref: 'refs/heads/main', + before: ZERO, + after: 'a'.repeat(40), + ...over, + }) + + it('skips when no mapping exists', async () => { + expect(await syncPush(payload())).toEqual({ status: 'skipped', reason: 'no-mapping' }) + expect(splicePushMock).not.toHaveBeenCalled() + }) + + it('skips when disabled', async () => { + await seedMapping({ disabledAt: new Date() }) + expect(await syncPush(payload())).toEqual({ status: 'skipped', reason: 'disabled' }) + }) + + it('skips ref deletions (after = zero sha)', async () => { + await seedMapping() + expect(await syncPush(payload({ after: ZERO }))).toEqual({ status: 'skipped', reason: 'deletion' }) + expect(splicePushMock).not.toHaveBeenCalled() + }) + + it('dedupes a redelivery via lastSyncedRefs', async () => { + await seedMapping({ lastSyncedRefs: { 'refs/heads/main': 'a'.repeat(40) } }) + expect(await syncPush(payload())).toEqual({ status: 'skipped', reason: 'already-synced' }) + expect(splicePushMock).not.toHaveBeenCalled() + }) + + it('splices the push and records the new tip', async () => { + await seedMapping() + const result = await syncPush(payload()) + expect(result).toEqual({ status: 'synced' }) + expect(splicePushMock).toHaveBeenCalledWith(expect.objectContaining({ + ref: 'refs/heads/main', + want: 'a'.repeat(40), + repoDid: 'did:plc:repo-xyz', + knot: 'knot1.tangled.sh', + token: 'install-token', + })) + const rows = await useDb().select().from(repoMapping) + expect((rows[0].lastSyncedRefs as Record)['refs/heads/main']).toBe('a'.repeat(40)) + }) + + it('marks mapping error and stops on a terminal too-big failure', async () => { + await seedMapping() + splicePushMock.mockRejectedValue(new RemoteRejectedError('pack exceeded', 'too-big')) + expect(await syncPush(payload())).toEqual({ status: 'skipped', reason: 'repo-gone' }) + const rows = await useDb().select().from(repoMapping) + expect(rows[0].status).toBe('error') + expect(rows[0].lastError).toMatch(/size limit/) + }) + + it('marks mapping error when the knot rejects our key', async () => { + await seedMapping() + splicePushMock.mockRejectedValue(new RemoteRejectedError('denied', 'auth-rejected')) + expect(await syncPush(payload())).toEqual({ status: 'skipped', reason: 'repo-gone' }) + const rows = await useDb().select().from(repoMapping) + expect(rows[0].status).toBe('error') + }) + + it('rethrows a transient stale-old-sha for queue retry', async () => { + await seedMapping() + splicePushMock.mockRejectedValue(new RemoteRejectedError('stale', 'stale-old-sha')) + await expect(syncPush(payload())).rejects.toMatchObject({ reason: 'stale-old-sha' }) + const rows = await useDb().select().from(repoMapping) + expect(rows[0].status).toBe('active') + }) +}) diff --git a/test/unit/sync-ref.spec.ts b/test/unit/sync-ref.spec.ts --- a/test/unit/sync-ref.spec.ts +++ b/test/unit/sync-ref.spec.ts @@ -3,32 +3,24 @@ import { installation, repoMapping } from '../../server/db/schema' import { clearDb, setDb, useDb } from '../../server/utils/db' import { clearEncryptionKeyCache } from '../../server/utils/encryption' +import { RemoteRejectedError } from '../../server/utils/git-wire/errors' import { createTestDb } from '../utils/db' const ORIGINAL_ENC_KEY = process.env.NUXT_ENCRYPTION_KEY -// Mock git + ssh + octokit so we can exercise mapping-lookup + envelope -// branches without invoking real binaries. -const gitMock = vi.fn<(args: string[], opts?: unknown) => Promise<{ stdout: string, stderr: string }>>() -const sshLoadMock = vi.fn<(installationId: number) => Promise<{ gitSshCommand: string, cleanup: () => void }>>() +const fetchAdvertisementMock = vi.fn<(repo: string, token: string) => Promise<{ refs: Map }>>() +const splicePushMock = vi.fn<(params: Record) => Promise<{ status: string, sha: string }>>() +const spliceDeleteMock = vi.fn<(params: Record) => Promise<{ status: string }>>() const octokitAuthMock = vi.fn<(input: { type: 'installation' }) => Promise<{ token: string }>>() -vi.mock('../../server/utils/git', async () => { - const actual = await vi.importActual('../../server/utils/git') - return { - ...actual, - git: (args: string[], opts: unknown) => gitMock(args, opts), - } -}) - -vi.mock('../../server/utils/ssh-cmd', () => ({ - loadSshCommandForInstall: (id: number) => sshLoadMock(id), +vi.mock('../../server/utils/splice', () => ({ + fetchAdvertisement: (repo: string, token: string) => fetchAdvertisementMock(repo, token), + splicePush: (params: Record) => splicePushMock(params), + spliceDelete: (params: Record) => spliceDeleteMock(params), })) vi.mock('../../server/utils/github-app', () => ({ - installationOctokit: async () => ({ - auth: octokitAuthMock, - }), + installationOctokit: async () => ({ auth: octokitAuthMock }), })) const { syncCreateRef, syncDeleteRef } = await import('../../server/utils/sync-ref') @@ -44,13 +36,21 @@ id: 1, accountLogin: 'alice', accountId: 100, accountType: 'User', }) - gitMock.mockReset() - sshLoadMock.mockReset() + fetchAdvertisementMock.mockReset() + splicePushMock.mockReset() + spliceDeleteMock.mockReset() octokitAuthMock.mockReset() - sshLoadMock.mockResolvedValue({ gitSshCommand: 'ssh -i /tmp/key', cleanup: () => {} }) octokitAuthMock.mockResolvedValue({ token: 'install-token' }) - gitMock.mockResolvedValue({ stdout: 'abc1234567890abc1234567890abc1234567890a', stderr: '' }) + fetchAdvertisementMock.mockResolvedValue({ + refs: new Map([ + ['refs/heads/main', 'a'.repeat(40)], + ['refs/heads/feature-x', 'b'.repeat(40)], + ['refs/tags/v1.0.0', 'c'.repeat(40)], + ]), + }) + splicePushMock.mockResolvedValue({ status: 'synced', sha: 'a'.repeat(40) }) + spliceDeleteMock.mockResolvedValue({ status: 'synced' }) }) afterEach(() => { @@ -78,176 +78,105 @@ it('skips non-branch/tag ref types', async () => { await seedMapping() const result = await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'repository' as never, - ref: 'whatever', + installationId: 1, githubRepoId: 9001, refType: 'repository' as never, ref: 'whatever', }) expect(result).toEqual({ status: 'skipped', reason: 'not-branch-or-tag' }) - expect(gitMock).not.toHaveBeenCalled() + expect(splicePushMock).not.toHaveBeenCalled() }) it('skips when no mapping exists', async () => { - const result = await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'main', - }) + const result = await syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }) expect(result).toEqual({ status: 'skipped', reason: 'no-mapping' }) }) it('skips when mapping is disabled', async () => { await seedMapping({ disabledAt: new Date() }) - const result = await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'main', - }) + const result = await syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }) expect(result).toEqual({ status: 'skipped', reason: 'disabled' }) }) - it('qualifies branch refs as refs/heads/', async () => { + it('resolves a branch ref to its SHA and splices it', async () => { await seedMapping() - await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'feature-x', - }) - - // git init, git fetch, git push, git rev-parse - const calls = gitMock.mock.calls.map(c => c[0]) - const fetch = calls.find(args => args[0] === 'fetch') - const push = calls.find(args => args[0] === 'push') - expect(fetch).toBeDefined() - expect(fetch).toContain('+refs/heads/feature-x:refs/heads/feature-x') - expect(push).toContain('+refs/heads/feature-x:refs/heads/feature-x') + await syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'feature-x' }) + expect(splicePushMock).toHaveBeenCalledWith(expect.objectContaining({ + ref: 'refs/heads/feature-x', + want: 'b'.repeat(40), + repoDid: 'did:plc:repo-xyz', + })) }) - it('qualifies tag refs as refs/tags/', async () => { + it('resolves a tag ref to its SHA', async () => { await seedMapping() - await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'tag', - ref: 'v1.0.0', - }) - const calls = gitMock.mock.calls.map(c => c[0]) - const fetch = calls.find(args => args[0] === 'fetch') - expect(fetch).toContain('+refs/tags/v1.0.0:refs/tags/v1.0.0') + await syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'tag', ref: 'v1.0.0' }) + expect(splicePushMock).toHaveBeenCalledWith(expect.objectContaining({ + ref: 'refs/tags/v1.0.0', + want: 'c'.repeat(40), + })) }) - it('routes push for knot1.tangled.sh via tangled.org', async () => { + it('retries (throws) when github does not yet advertise the ref', async () => { await seedMapping() - await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'main', - }) - const push = gitMock.mock.calls.map(c => c[0]).find(a => a[0] === 'push')! - const url = push.find(a => a.startsWith('ssh://'))! - expect(url).toContain('@tangled.org/') + fetchAdvertisementMock.mockResolvedValue({ refs: new Map() }) + await expect( + syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }), + ).rejects.toThrow(/does not yet advertise/) + expect(splicePushMock).not.toHaveBeenCalled() }) - it('updates lastSyncedRefs with the fetched SHA', async () => { + it('updates lastSyncedRefs with the synced SHA', async () => { await seedMapping() - gitMock.mockImplementation(async args => { - if (args[0] === 'rev-parse') { - return { stdout: 'deadbeef1234567890deadbeef1234567890dead', stderr: '' } - } - return { stdout: '', stderr: '' } - }) + splicePushMock.mockResolvedValue({ status: 'synced', sha: 'd'.repeat(40) }) + await syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }) - await syncCreateRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'main', - }) + const rows = await useDb().select().from(repoMapping) + expect((rows[0].lastSyncedRefs as Record)['refs/heads/main']).toBe('d'.repeat(40)) + }) - const db = useDb() - const rows = await db.select().from(repoMapping) - const refs = rows[0].lastSyncedRefs as Record - expect(refs['refs/heads/main']).toBe('deadbeef1234567890deadbeef1234567890dead') + it('marks mapping error when the knot reports repo gone', async () => { + await seedMapping() + splicePushMock.mockRejectedValue(new RemoteRejectedError('gone', 'repo-gone')) + const result = await syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }) + expect(result).toEqual({ status: 'skipped', reason: 'repo-gone' }) + const rows = await useDb().select().from(repoMapping) + expect(rows[0].status).toBe('error') + }) + + it('rethrows a transient stale-old-sha for queue retry', async () => { + await seedMapping() + splicePushMock.mockRejectedValue(new RemoteRejectedError('stale', 'stale-old-sha')) + await expect( + syncCreateRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }), + ).rejects.toMatchObject({ reason: 'stale-old-sha' }) }) }) describe('syncDeleteRef', () => { - it('uses the empty-source refspec to delete on the remote', async () => { + it('splices a delete for the qualified ref', async () => { await seedMapping() - await syncDeleteRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'old-branch', - }) - const push = gitMock.mock.calls.map(c => c[0]).find(a => a[0] === 'push')! - expect(push).toContain(':refs/heads/old-branch') - expect(push.some(s => s.startsWith(':'))).toBe(true) + await syncDeleteRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'old-branch' }) + expect(spliceDeleteMock).toHaveBeenCalledWith(expect.objectContaining({ ref: 'refs/heads/old-branch' })) }) - it('treats "remote ref does not exist" as success', async () => { + it('treats an already-absent ref as success', async () => { await seedMapping() - gitMock.mockImplementation(async args => { - if (args[0] === 'push') { - throw Object.assign(new Error('exit 1'), { - stderr: 'error: unable to delete \'refs/tags/v9\': remote ref does not exist\n', - }) - } - return { stdout: '', stderr: '' } - }) - - const result = await syncDeleteRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'tag', - ref: 'v9', - }) + spliceDeleteMock.mockResolvedValue({ status: 'already-absent' }) + const result = await syncDeleteRef({ installationId: 1, githubRepoId: 9001, refType: 'tag', ref: 'v9' }) expect(result).toEqual({ status: 'synced' }) }) it('removes the ref from lastSyncedRefs', async () => { - const db = useDb() - await seedMapping({ - lastSyncedRefs: { 'refs/heads/main': 'abc', 'refs/heads/old': 'def' }, - }) - - await syncDeleteRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'old', - }) - - const rows = await db.select().from(repoMapping) - const refs = rows[0].lastSyncedRefs as Record - expect(refs).toEqual({ 'refs/heads/main': 'abc' }) + await seedMapping({ lastSyncedRefs: { 'refs/heads/main': 'abc', 'refs/heads/old': 'def' } }) + await syncDeleteRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'old' }) + const rows = await useDb().select().from(repoMapping) + expect(rows[0].lastSyncedRefs).toEqual({ 'refs/heads/main': 'abc' }) }) it('marks mapping as error if knot reports repo gone', async () => { await seedMapping() - gitMock.mockImplementation(async args => { - if (args[0] === 'push') { - throw Object.assign(new Error('exit 1'), { - stderr: 'fatal: repository not found\n', - }) - } - return { stdout: '', stderr: '' } - }) - - const result = await syncDeleteRef({ - installationId: 1, - githubRepoId: 9001, - refType: 'branch', - ref: 'main', - }) + spliceDeleteMock.mockRejectedValue(new RemoteRejectedError('gone', 'repo-gone')) + const result = await syncDeleteRef({ installationId: 1, githubRepoId: 9001, refType: 'branch', ref: 'main' }) expect(result).toEqual({ status: 'skipped', reason: 'repo-gone' }) - - const db = useDb() - const rows = await db.select().from(repoMapping) + const rows = await useDb().select().from(repoMapping) expect(rows[0].status).toBe('error') }) }) diff --git a/test/unit/upload-pack.spec.ts b/test/unit/upload-pack.spec.ts new file mode 100644 --- /dev/null +++ b/test/unit/upload-pack.spec.ts @@ -0,0 +1,92 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { RemoteRejectedError } from '../../server/utils/git-wire/errors' +import { fetchAdvertisement, fetchPack } from '../../server/utils/git-wire/upload-pack' +import { fakeGithubFetch, GitFixture } from '../utils/git-wire' + +const PACK_MAGIC = Buffer.from('PACK') + +async function drain(gen: AsyncGenerator): Promise { + const parts: Buffer[] = [] + for await (const c of gen) parts.push(c) + return Buffer.concat(parts) +} + +describe('upload-pack (against real git-upload-pack)', () => { + let fx: GitFixture + let realFetch: typeof globalThis.fetch + + beforeEach(() => { + fx = new GitFixture() + realFetch = globalThis.fetch + }) + + afterEach(() => { + globalThis.fetch = realFetch + fx.cleanup() + vi.restoreAllMocks() + }) + + function wire(repos: Map) { + globalThis.fetch = fakeGithubFetch(repos) as unknown as typeof globalThis.fetch + } + + it('resolves a ref name to a SHA via the advertisement', async () => { + const bare = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, bare, 'HEAD:refs/heads/main') + wire(new Map([['owner/repo', bare]])) + + const adv = await fetchAdvertisement('owner/repo', 'tok') + expect(adv.refs.get('refs/heads/main')).toBe(sha) + expect(adv.capabilities.has('thin-pack')).toBe(true) + }) + + it('fetches a full pack on first sync (no haves)', async () => { + const bare = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, bare, 'HEAD:refs/heads/main') + wire(new Map([['owner/repo', bare]])) + + const { pack } = await fetchPack({ repoFullName: 'owner/repo', token: 'tok', want: sha, haves: [], maxBytes: 1 << 30 }) + const bytes = await drain(pack) + expect(bytes.subarray(0, 4)).toEqual(PACK_MAGIC) + expect(bytes.length).toBeGreaterThan(0) + }) + + it('sends a smaller pack when the knot tip is offered as a have', async () => { + const bare = fx.initBare('gh.git') + const work = fx.initWork('work') + const first = fx.commit(work, 'a.txt', 'a'.repeat(5000)) + fx.pushTo(work, bare, 'HEAD:refs/heads/main') + const second = fx.commit(work, 'b.txt', 'b'.repeat(5000)) + fx.pushTo(work, bare, 'HEAD:refs/heads/main') + wire(new Map([['owner/repo', bare]])) + + const full = await drain((await fetchPack({ repoFullName: 'owner/repo', token: 'tok', want: second, haves: [], maxBytes: 1 << 30 })).pack) + const incremental = await drain((await fetchPack({ repoFullName: 'owner/repo', token: 'tok', want: second, haves: [first], maxBytes: 1 << 30 })).pack) + + expect(incremental.subarray(0, 4)).toEqual(PACK_MAGIC) + expect(incremental.length).toBeLessThan(full.length) + }) + + it('aborts with too-big once the byte cap is exceeded', async () => { + const bare = fx.initBare('gh.git') + const work = fx.initWork('work') + const sha = fx.commit(work, 'a.txt', 'x'.repeat(10_000)) + fx.pushTo(work, bare, 'HEAD:refs/heads/main') + wire(new Map([['owner/repo', bare]])) + + const { pack } = await fetchPack({ repoFullName: 'owner/repo', token: 'tok', want: sha, haves: [], maxBytes: 10 }) + await expect(drain(pack)).rejects.toMatchObject({ + constructor: RemoteRejectedError, + reason: 'too-big', + }) + }) + + it('throws on a 404 from github', async () => { + wire(new Map()) + await expect(fetchAdvertisement('missing/repo', 'tok')).rejects.toThrow(/404/) + }) +}) diff --git a/test/utils/git-wire.ts b/test/utils/git-wire.ts new file mode 100644 --- /dev/null +++ b/test/utils/git-wire.ts @@ -0,0 +1,121 @@ +import { type ChildProcessWithoutNullStreams, execFileSync, spawn, spawnSync } from 'node:child_process' +import { mkdtempSync, rmSync } from 'node:fs' +import os from 'node:os' +import path from 'node:path' +import type { ReceivePackFactory } from '../../server/utils/git-wire/receive-pack' +import { encodePktLine, flushPkt } from '../../server/utils/git-wire/pkt-line' + +/** + * Local git fixtures for wire-protocol tests. No network, no ssh: we drive the + * same `git-upload-pack` / `git-receive-pack` binaries GitHub and the knot run + * server-side, so the bytes on the pipe are real protocol output. + */ +export class GitFixture { + readonly dir: string + + constructor() { + this.dir = mkdtempSync(path.join(os.tmpdir(), 'gitwire-test-')) + } + + git(args: string[], cwd = this.dir): string { + return execFileSync('git', args, { + cwd, + encoding: 'utf8', + env: { ...process.env, GIT_CONFIG_NOSYSTEM: '1', GIT_TERMINAL_PROMPT: '0' }, + }).trim() + } + + initBare(name: string): string { + const repo = path.join(this.dir, name) + this.git(['init', '-q', '--bare', repo]) + return repo + } + + /** Create a non-bare work repo with one commit on `main`, return its path. */ + initWork(name: string): string { + const repo = path.join(this.dir, name) + this.git(['init', '-q', '-b', 'main', repo]) + this.git(['config', 'user.email', 't@example.com'], repo) + this.git(['config', 'user.name', 'Test'], repo) + return repo + } + + commit(workRepo: string, file: string, content: string): string { + execFileSync('bash', ['-c', `printf %s ${JSON.stringify(content)} > ${JSON.stringify(path.join(workRepo, file))}`]) + this.git(['add', '.'], workRepo) + this.git(['commit', '-q', '-m', `add ${file}`], workRepo) + return this.git(['rev-parse', 'HEAD'], workRepo) + } + + pushTo(workRepo: string, bareRepo: string, refspec: string): void { + this.git(['push', '-q', bareRepo, refspec], workRepo) + } + + revParse(repo: string, ref: string): string { + return this.git(['rev-parse', ref], repo) + } + + cleanup(): void { + rmSync(this.dir, { recursive: true, force: true }) + } +} + +/** + * Stand in for GitHub's `git-upload-pack` HTTP endpoint by running the binary + * in `--stateless-rpc` mode against a local bare repo. Mirrors the request / + * response framing the real endpoint uses, so `upload-pack.ts` can be pointed + * at it through a patched `fetch`. + */ +export function fakeGithubFetch(repos: Map) { + return async function fetchImpl(input: string | URL, init?: RequestInit): Promise { + const url = typeof input === 'string' ? input : input.toString() + const match = url.match(/github\.com\/(.+?)\.git\/(info\/refs|git-upload-pack)/) + if (!match) throw new Error(`fakeGithubFetch: unexpected url ${url}`) + const repoPath = repos.get(match[1]!) + if (!repoPath) return new Response(null, { status: 404, statusText: 'Not Found' }) + + if (match[2] === 'info/refs') { + const adv = execFileSync('git-upload-pack', ['--stateless-rpc', '--advertise-refs', repoPath]) + // The HTTP transport prepends the service banner + flush; the binary does not. + const banner = Buffer.concat([ + encodePktLine('# service=git-upload-pack\n'), + flushPkt, + ]) + return new Response(new Uint8Array(Buffer.concat([banner, adv])), { status: 200 }) + } + + const reqBody = Buffer.from(await new Response(init!.body as BodyInit).arrayBuffer()) + const proc = spawnSync('git-upload-pack', ['--stateless-rpc', repoPath], { input: reqBody, maxBuffer: 1 << 30 }) + if (proc.status !== 0) { + throw new Error(`git-upload-pack exited ${proc.status}: ${proc.stderr.toString()}`) + } + return new Response(new Uint8Array(proc.stdout), { status: 200 }) + } +} + +const STDERR_CAP = 16 * 1024 + +/** + * A `ReceivePackFactory` that spawns the real `git-receive-pack` binary + * against a local bare repo, bypassing ssh. The stdio protocol is identical + * to what the knot speaks. + */ +export function localReceivePackFactory(bareRepo: string): ReceivePackFactory { + return () => { + const child: ChildProcessWithoutNullStreams = spawn('git-receive-pack', [bareRepo], { + stdio: ['pipe', 'pipe', 'pipe'], + }) + let stderrBuf = Buffer.alloc(0) + child.stderr.on('data', (chunk: Buffer) => { + stderrBuf = Buffer.concat([stderrBuf, chunk]).subarray(-STDERR_CAP) + }) + const done = new Promise(resolve => child.on('close', code => resolve(code))) + return { + stdin: child.stdin, + stdout: child.stdout, + stderr: () => stderrBuf.toString('utf8'), + kill: () => child.kill('SIGKILL'), + done, + } + } +} diff --git a/server/utils/git-wire/errors.ts b/server/utils/git-wire/errors.ts new file mode 100644 --- /dev/null +++ b/server/utils/git-wire/errors.ts @@ -0,0 +1,76 @@ +/** + * Typed failures from the git wire splice. `reason` drives the worker's + * retry-vs-give-up decision in `sync-push.ts` / `sync-ref.ts`. + * + * - repo-gone the knot no longer has the repo (or our key was revoked + * such that it reports "not found"); terminal, mark error. + * - auth-rejected ssh public-key auth refused; terminal, mark error. + * - stale-old-sha our compare-and-swap lost: the knot's ref moved between + * reading its advertisement and sending the command, or a + * concurrent worker won. Transient; retry re-reads the tip. + * - too-big the pack exceeded the configured byte cap; terminal, it + * will never fit. + * - other anything unclassified; transient, let the queue retry. + */ +export type WireFailureReason + = | 'repo-gone' + | 'auth-rejected' + | 'stale-old-sha' + | 'too-big' + | 'other' + +export class WireError extends Error { + constructor(message: string) { + super(message) + this.name = 'WireError' + } +} + +export class RemoteRejectedError extends WireError { + constructor(message: string, public readonly reason: WireFailureReason) { + super(message) + this.name = 'RemoteRejectedError' + } +} + +/** + * Classify ssh / sshd / knot stderr (the child process's stderr band, since + * we deliberately do not request side-band multiplexing). Returns null when + * nothing matches so the caller can fall back to a generic transient error. + */ +export function classifySshStderr(stderr: string): RemoteRejectedError | null { + const lc = stderr.toLowerCase() + if (lc.includes('repository not found') || lc.includes('does not exist') || lc.includes('does not appear to be a git repository')) { + return new RemoteRejectedError(stderr.trim(), 'repo-gone') + } + if (lc.includes('permission denied') || (lc.includes('publickey') && lc.includes('denied'))) { + return new RemoteRejectedError(stderr.trim(), 'auth-rejected') + } + return null +} + +/** + * Classify a receive-pack `ng ` rejection. Any rejection that + * means "the ref's current value is not what you said" maps to stale-old-sha + * so the worker retries against a fresh advertisement. git phrases this two + * ways: `non-fast-forward` / `stale info` when updating a moved ref, and + * `failed to update ref` (stderr: "reference already exists") when our command + * claimed a create but the ref already exists. + */ +export function classifyNgReason(reason: string): RemoteRejectedError { + const lc = reason.toLowerCase() + if ( + lc.includes('non-fast-forward') + || lc.includes('fetch first') + || lc.includes('stale info') + || lc.includes('not a fast forward') + || lc.includes('failed to update ref') + || lc.includes('reference already exists') + ) { + return new RemoteRejectedError(reason.trim(), 'stale-old-sha') + } + if (lc.includes('not found') || lc.includes('does not exist')) { + return new RemoteRejectedError(reason.trim(), 'repo-gone') + } + return new RemoteRejectedError(reason.trim(), 'other') +} diff --git a/server/utils/git-wire/pkt-line.ts b/server/utils/git-wire/pkt-line.ts new file mode 100644 --- /dev/null +++ b/server/utils/git-wire/pkt-line.ts @@ -0,0 +1,158 @@ +/** + * Git pkt-line framing (protocol v0). A pkt-line is a 4-hex-digit length + * prefix (counting the 4 prefix bytes themselves) followed by that many bytes + * of payload. `0000` is the flush-pkt: a section delimiter carrying no + * payload. Lengths `0001`-`0003` are reserved and invalid in v0. + * + * See `Documentation/gitprotocol-common.txt` in git.git. + */ + +const FLUSH = '0000' +const MAX_DATA = 65516 + +export const flushPkt: Buffer = Buffer.from(FLUSH, 'ascii') + +/** + * Frame a payload as a pkt-line. Accepts a string (encoded UTF-8) or raw + * bytes. Does NOT append a trailing newline; callers that want the + * conventional `\n` (command and capability lines) must include it. + */ +export function encodePktLine(data: string | Buffer): Buffer { + const payload = typeof data === 'string' ? Buffer.from(data, 'utf8') : data + if (payload.length > MAX_DATA) { + throw new RangeError(`pkt-line payload too large: ${payload.length} > ${MAX_DATA}`) + } + const len = payload.length + 4 + const prefix = len.toString(16).padStart(4, '0') + return Buffer.concat([Buffer.from(prefix, 'ascii'), payload]) +} + +export type PktLine = + | { type: 'line', data: Buffer } + | { type: 'flush' } + +/** + * Incrementally decode pkt-lines from a byte source, then hand back whatever + * raw bytes follow the section we consumed. + * + * The git smart protocol switches from pkt-line framing to a raw packfile + * stream mid-response (after the NAK/ACK line on a fetch). A naive reader that + * buffers ahead would swallow the first chunk of the pack, so this reader + * tracks exactly how much it has consumed and exposes the remainder via + * `remaining()`. + */ +export class PktLineReader { + private buf: Buffer = Buffer.alloc(0) + private done = false + private readonly iter: AsyncIterator + + constructor(source: AsyncIterable) { + this.iter = (async function* normalise() { + for await (const chunk of source) { + yield Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk) + } + })() + } + + /** + * Read the next pkt-line, or `null` at end of stream. A flush-pkt is + * returned as `{ type: 'flush' }` rather than ending iteration; the wire + * protocol uses multiple flush-delimited sections per stream. + */ + async next(): Promise { + while (this.buf.length < 4) { + // eslint-disable-next-line no-await-in-loop -- each fill extends a shared buffer the next iteration inspects; the reads are inherently sequential + if (!(await this.fill())) { + if (this.buf.length === 0) return null + throw new Error('unexpected end of stream: truncated pkt-line length') + } + } + + const len = Number.parseInt(this.buf.toString('ascii', 0, 4), 16) + if (Number.isNaN(len)) { + throw new Error(`invalid pkt-line length: ${JSON.stringify(this.buf.toString('ascii', 0, 4))}`) + } + if (len === 0) { + this.buf = this.buf.subarray(4) + return { type: 'flush' } + } + if (len < 4) { + throw new Error(`reserved pkt-line length ${len} is invalid in protocol v0`) + } + + while (this.buf.length < len) { + // eslint-disable-next-line no-await-in-loop -- sequential read; see next() + if (!(await this.fill())) { + throw new Error(`unexpected end of stream: pkt-line wanted ${len} bytes, had ${this.buf.length}`) + } + } + + const data = this.buf.subarray(4, len) + this.buf = this.buf.subarray(len) + return { type: 'line', data } + } + + /** + * Read pkt-lines up to and including the next flush-pkt, returning the line + * payloads (flush excluded). Returns `null` if the stream ends before any + * line is read. + */ + async readUntilFlush(): Promise { + const lines: Buffer[] = [] + for (;;) { + // eslint-disable-next-line no-await-in-loop -- sequential read; see next() + const pkt = await this.next() + if (pkt === null) return lines.length > 0 ? lines : null + if (pkt.type === 'flush') return lines + lines.push(pkt.data) + } + } + + /** + * The bytes already buffered past the last consumed pkt-line. Used to seed + * the raw packfile stream once negotiation framing ends. + */ + buffered(): Buffer { + return this.buf + } + + /** + * Yield the remainder of the source as a raw byte stream: first any bytes + * already buffered, then the rest of the underlying iterator verbatim. After + * calling this, do not call `next()` again. + */ + async *remaining(): AsyncGenerator { + if (this.buf.length > 0) { + yield this.buf + this.buf = Buffer.alloc(0) + } + if (this.done) return + for (;;) { + // eslint-disable-next-line no-await-in-loop -- sequential drain of the source iterator + const { value, done } = await this.iter.next() + if (done) { + this.done = true + return + } + yield value + } + } + + private async fill(): Promise { + if (this.done) return false + const { value, done } = await this.iter.next() + if (done) { + this.done = true + return false + } + this.buf = this.buf.length === 0 ? value : Buffer.concat([this.buf, value]) + return true + } +} + +/** Decode a single line's payload as a UTF-8 string with any trailing `\n` removed. */ +export function lineToString(data: Buffer): string { + return data.length > 0 && data[data.length - 1] === 0x0A + ? data.toString('utf8', 0, data.length - 1) + : data.toString('utf8') +} diff --git a/server/utils/git-wire/receive-pack.ts b/server/utils/git-wire/receive-pack.ts new file mode 100644 --- /dev/null +++ b/server/utils/git-wire/receive-pack.ts @@ -0,0 +1,210 @@ +import { type ChildProcessWithoutNullStreams, spawn } from 'node:child_process' +import { Readable } from 'node:stream' +import { classifyNgReason, classifySshStderr, RemoteRejectedError, WireError } from './errors' +import { encodePktLine, flushPkt, lineToString, PktLineReader } from './pkt-line' +import { type Advertisement, parseAdvertisement, ZERO_SHA } from './refs' + +const AGENT = 'synchub.to' +const STDERR_CAP = 16 * 1024 +/** + * Whole-session budget. receive-pack blocks indefinitely waiting for commands, + * so without this a knot that accepts the connection but stalls mid-protocol + * would hang a worker until its job lease expires. On expiry we SIGKILL the + * child; the in-flight read sees the stream end and throws, surfacing as a + * transient failure the queue retries. + */ +const SESSION_TIMEOUT_MS = 120_000 + +export interface RefUpdate { + ref: string + /** Current value on the knot, or the zero SHA to create. The compare-and-swap. */ + old: string + /** New value, or the zero SHA to delete. */ + next: string +} + +/** + * A spawned process exposing `git-receive-pack`'s stdio. The default factory + * runs ssh to the knot; tests inject a factory that spawns the binary against + * a local bare repo, so the stdio protocol is identical either way. + */ +export interface ReceivePackProcess { + stdin: NodeJS.WritableStream + stdout: AsyncIterable + /** Last bytes of stderr, for diagnostics (we don't request side-band). */ + stderr(): string + kill(): void + /** Resolves with the exit code once the process ends. */ + done: Promise +} + +export type ReceivePackFactory = () => ReceivePackProcess + +export interface SshTarget { + host: string + port?: number + repoPath: string + sshArgs: string[] +} + +/** Default transport: ssh to the knot and invoke its `git-receive-pack`. */ +export function sshReceivePackFactory(target: SshTarget): ReceivePackFactory { + return () => { + const portArgs = target.port ? ['-p', String(target.port)] : [] + // ssh:// transports invoke the remote command with the path including its + // leading slash, single-quoted. The knot resolves repos by that path. + const remoteCmd = `git-receive-pack '${target.repoPath}'` + const child = spawn('ssh', [...target.sshArgs, ...portArgs, `git@${target.host}`, remoteCmd], { + stdio: ['pipe', 'pipe', 'pipe'], + }) + return wrapChild(child) + } +} + +function wrapChild(child: ChildProcessWithoutNullStreams): ReceivePackProcess { + let stderrBuf = Buffer.alloc(0) + child.stderr.on('data', (chunk: Buffer) => { + stderrBuf = Buffer.concat([stderrBuf, chunk]).subarray(-STDERR_CAP) + }) + const done = new Promise(resolve => child.on('close', resolve)) + return { + stdin: child.stdin, + stdout: child.stdout, + stderr: () => stderrBuf.toString('utf8'), + kill: () => child.kill('SIGKILL'), + done, + } +} + +/** + * An open receive-pack session. Read `tips` after construction to learn the + * knot's current refs (needed as haves and as the compare-and-swap base), + * then call `push` once with the commands and packfile. + */ +export class ReceivePackSession { + readonly tips: Map + readonly capabilities: Set + + private readonly watchdog: NodeJS.Timeout + + private constructor( + private readonly proc: ReceivePackProcess, + private readonly reader: PktLineReader, + adv: Advertisement, + watchdog: NodeJS.Timeout, + ) { + this.tips = adv.refs + this.capabilities = adv.capabilities + this.watchdog = watchdog + } + + /** Open the session and read the advertisement. */ + static async open(factory: ReceivePackFactory, timeoutMs = SESSION_TIMEOUT_MS): Promise { + const proc = factory() + const watchdog = setTimeout(() => proc.kill(), timeoutMs) + try { + const reader = new PktLineReader(proc.stdout) + const advLines = await reader.readUntilFlush() + if (advLines === null) { + const err = classifySshStderr(proc.stderr()) + throw err ?? new WireError(`receive-pack: no advertisement (stderr: ${proc.stderr().trim() || 'empty'})`) + } + const adv = parseAdvertisement(advLines) + assertCapabilities(adv) + return new ReceivePackSession(proc, reader, adv, watchdog) + } + catch (err) { + clearTimeout(watchdog) + proc.kill() + await proc.done.catch(() => null) + throw err + } + } + + /** + * Send the ref update commands plus (for non-deletions) the packfile, then + * read and validate report-status. The packfile is streamed straight from + * `packStream` into stdin; nothing is buffered. Pass `null` for pure + * deletions. + */ + async push(updates: RefUpdate[], packStream: AsyncIterable | null): Promise { + if (updates.length === 0) throw new WireError('receive-pack: no updates') + try { + await writeAll(this.proc.stdin, buildCommandList(updates)) + if (packStream) await pipePack(this.proc.stdin, packStream) + else this.proc.stdin.end() + + const report = await this.reader.readUntilFlush() + parseReportStatus(report ?? [], updates, this.proc.stderr()) + } + finally { + clearTimeout(this.watchdog) + this.proc.kill() + await this.proc.done.catch(() => null) + } + } + + /** Close the session without pushing (advertisement-only use). */ + async close(): Promise { + clearTimeout(this.watchdog) + this.proc.stdin.end() + this.proc.kill() + await this.proc.done.catch(() => null) + } +} + +function assertCapabilities(adv: Advertisement): void { + if (!adv.capabilities.has('report-status')) { + throw new WireError('knot receive-pack does not advertise report-status') + } +} + +function buildCommandList(updates: RefUpdate[]): Buffer { + const parts: Buffer[] = [] + updates.forEach((u, i) => { + const caps = i === 0 ? `\0report-status agent=${AGENT}/1` : '' + parts.push(encodePktLine(`${u.old} ${u.next} ${u.ref}${caps}\n`)) + }) + parts.push(flushPkt) + return Buffer.concat(parts) +} + +function parseReportStatus(lines: Buffer[], updates: RefUpdate[], stderr: string): void { + if (lines.length === 0) { + const err = classifySshStderr(stderr) + throw err ?? new WireError(`receive-pack: empty report-status (stderr: ${stderr.trim() || 'empty'})`) + } + const unpack = lineToString(lines[0]!) + if (unpack !== 'unpack ok') { + throw new WireError(`receive-pack: ${unpack}`) + } + for (const raw of lines.slice(1)) { + const line = lineToString(raw) + if (line.startsWith('ng ')) { + // `ng ` + const rest = line.slice(3) + const sp = rest.indexOf(' ') + const reason = sp === -1 ? rest : rest.slice(sp + 1) + throw classifyNgReason(reason) + } + } +} + +async function writeAll(stream: NodeJS.WritableStream, data: Buffer): Promise { + await new Promise((resolve, reject) => { + stream.write(data, err => (err ? reject(err) : resolve())) + }) +} + +async function pipePack(stdin: NodeJS.WritableStream, packStream: AsyncIterable): Promise { + const src = Readable.from(packStream) + await new Promise((resolve, reject) => { + src.on('error', reject) + stdin.on('error', reject) + src.pipe(stdin, { end: true }) + stdin.on('finish', resolve) + stdin.on('close', resolve) + }) +} + +export { RemoteRejectedError, ZERO_SHA } diff --git a/server/utils/git-wire/refs.ts b/server/utils/git-wire/refs.ts new file mode 100644 --- /dev/null +++ b/server/utils/git-wire/refs.ts @@ -0,0 +1,69 @@ +import { lineToString } from './pkt-line' + +const ZERO_SHA = '0000000000000000000000000000000000000000' + +export interface Advertisement { + /** Ref name -> object SHA (unpeeled). For annotated tags this is the tag object. */ + refs: Map + /** For annotated tags, the `^{}` peeled line: ref name -> commit SHA. */ + peeled: Map + capabilities: Set +} + +/** + * Parse a git ref advertisement (protocol v0) from a list of pkt-line + * payloads (flush-pkts already stripped by the reader). + * + * Handles three shapes that occur in practice: + * - smart-HTTP prelude: a leading `# service=git-upload-pack` line, which + * the ssh transport omits; + * - a populated repo: ` \0` on the first ref line, + * ` ` thereafter, with ` ^{}` peeled lines + * for annotated tags; + * - an empty repo: a single ` capabilities^{}\0` line that + * carries capabilities but advertises no usable ref. + */ +export function parseAdvertisement(lines: Buffer[]): Advertisement { + const refs = new Map() + const peeled = new Map() + const capabilities = new Set() + + let first = true + for (const raw of lines) { + const line = lineToString(raw) + if (line.startsWith('# service=')) continue + + let sha: string + let rest: string + if (first) { + const nul = line.indexOf('\0') + const head = nul === -1 ? line : line.slice(0, nul) + const caps = nul === -1 ? '' : line.slice(nul + 1) + for (const cap of caps.split(' ')) { + if (cap) capabilities.add(cap) + } + first = false + const sp = head.indexOf(' ') + sha = head.slice(0, sp) + rest = head.slice(sp + 1) + // Empty-repo sentinel: zero SHA + the literal "capabilities^{}" name. + if (sha === ZERO_SHA && rest === 'capabilities^{}') continue + } + else { + const sp = line.indexOf(' ') + sha = line.slice(0, sp) + rest = line.slice(sp + 1) + } + + if (rest.endsWith('^{}')) { + peeled.set(rest.slice(0, -3), sha) + } + else { + refs.set(rest, sha) + } + } + + return { refs, peeled, capabilities } +} + +export { ZERO_SHA } diff --git a/server/utils/git-wire/upload-pack.ts b/server/utils/git-wire/upload-pack.ts new file mode 100644 --- /dev/null +++ b/server/utils/git-wire/upload-pack.ts @@ -0,0 +1,134 @@ +import { Buffer } from 'node:buffer' +import { RemoteRejectedError, WireError } from './errors' +import { encodePktLine, flushPkt, lineToString, PktLineReader } from './pkt-line' +import { type Advertisement, parseAdvertisement } from './refs' + +const AGENT = 'synchub.to' +const ADVERTISEMENT_TIMEOUT_MS = 30_000 + +function repoUrl(repoFullName: string): string { + return `https://github.com/${repoFullName}.git` +} + +function authHeader(token: string): string { + return `Basic ${Buffer.from(`x-access-token:${token}`).toString('base64')}` +} + +async function* streamBytes(body: ReadableStream): AsyncGenerator { + const reader = body.getReader() + try { + for (;;) { + // eslint-disable-next-line no-await-in-loop -- sequential drain of the response body + const { value, done } = await reader.read() + if (done) return + if (value) yield Buffer.from(value) + } + } + finally { + reader.releaseLock() + } +} + +/** + * Fetch GitHub's `git-upload-pack` ref advertisement over smart HTTP. We need + * this both to resolve a ref name to a SHA (create-ref path) and, more + * generally, to learn the capability set before negotiating. + */ +export async function fetchAdvertisement(repoFullName: string, token: string): Promise { + const url = `${repoUrl(repoFullName)}/info/refs?service=git-upload-pack` + const res = await fetch(url, { + headers: { + Authorization: authHeader(token), + // Pin protocol v0; v2 would frame the advertisement differently. + 'Git-Protocol': 'version=0', + }, + signal: AbortSignal.timeout(ADVERTISEMENT_TIMEOUT_MS), + }) + if (!res.ok || !res.body) { + throw new WireError(`github info/refs failed: ${res.status} ${res.statusText}`) + } + const reader = new PktLineReader(streamBytes(res.body)) + const lines = await reader.readUntilFlush() + // The first flush ends the `# service` banner; the advertisement follows. + const adv = await reader.readUntilFlush() + return parseAdvertisement([...(lines ?? []), ...(adv ?? [])]) +} + +export interface FetchPackOptions { + repoFullName: string + token: string + /** SHA we want fetched. Requires GitHub's allow-reachable-sha1-in-want. */ + want: string + /** Knot ref tips to advertise as haves so GitHub sends a thin delta. */ + haves: string[] + /** Abort and throw `too-big` once the pack exceeds this many bytes. */ + maxBytes: number +} + +export interface FetchPackResult { + /** Raw packfile bytes. Pipe straight into receive-pack; do not buffer. */ + pack: AsyncGenerator +} + +/** + * Negotiate a thin pack from GitHub for `want`, advertising `haves` so the + * server deltas against objects the knot already holds. Returns a streaming + * generator of the raw packfile bytes; the caller pipes them into + * receive-pack and never materialises them. + * + * Protocol v0, no side-band: after the single NAK/ACK pkt-line the response + * body is the raw packfile to EOF, which is exactly what we forward. + */ +export async function fetchPack(opts: FetchPackOptions): Promise { + const { repoFullName, token, want, haves, maxBytes } = opts + + const wantLine = `want ${want} thin-pack ofs-delta agent=${AGENT}/1\n` + const body: Buffer[] = [encodePktLine(wantLine), flushPkt] + for (const have of haves) { + body.push(encodePktLine(`have ${have}\n`)) + } + body.push(encodePktLine('done\n')) + + const res = await fetch(`${repoUrl(repoFullName)}/git-upload-pack`, { + method: 'POST', + headers: { + Authorization: authHeader(token), + 'Content-Type': 'application/x-git-upload-pack-request', + 'Accept': 'application/x-git-upload-pack-result', + 'Git-Protocol': 'version=0', + }, + body: Buffer.concat(body), + }) + if (!res.ok || !res.body) { + throw new WireError(`github git-upload-pack failed: ${res.status} ${res.statusText}`) + } + + const reader = new PktLineReader(streamBytes(res.body)) + // Read the negotiation result: one ACK/NAK line, or an ERR line on failure. + const ack = await reader.next() + if (ack === null || ack.type === 'flush') { + throw new WireError('github git-upload-pack: empty negotiation response') + } + const ackStr = lineToString(ack.data) + if (ackStr.startsWith('ERR ')) { + // `ERR upload-pack: not our ref` is a propagation race on GitHub's side; + // surface as a plain WireError so the queue retries with backoff. + throw new WireError(`github git-upload-pack: ${ackStr.slice(4)}`) + } + if (!ackStr.startsWith('ACK') && !ackStr.startsWith('NAK')) { + throw new WireError(`github git-upload-pack: unexpected negotiation line ${JSON.stringify(ackStr)}`) + } + + async function* capped(): AsyncGenerator { + let total = 0 + for await (const chunk of reader.remaining()) { + total += chunk.length + if (total > maxBytes) { + throw new RemoteRejectedError(`pack exceeded ${maxBytes} bytes`, 'too-big') + } + yield chunk + } + } + + return { pack: capped() } +} -- tangled.sh