diff --git a/server/utils/git-wire/receive-pack.ts b/server/utils/git-wire/receive-pack.ts index 5cfa977..3db41cd 100644 --- a/server/utils/git-wire/receive-pack.ts +++ b/server/utils/git-wire/receive-pack.ts @@ -95,6 +95,12 @@ export function ssh2ReceivePackFactory(target: SshTarget): ReceivePackFactory { stdin.pipe(channel) channel.pipe(stdout) channel.stderr.on('data', appendStderr) + // A channel-level error (peer reset, transport dying mid-stream) would + // otherwise be lost: the write side is a PassThrough whose 'error' we + // swallow, and a half-sent pack reaches the knot as a corrupt stream + // it reports back as an `unpack` failure. Fold it into the stderr band + // so push() can tell a truncated push apart from a clean one. + channel.on('error', (err: Error) => { if (!killed) connError = err }) channel.on('exit', code => settle(typeof code === 'number' ? code : null)) channel.on('close', () => { client.end(); stdout.end() }) }) @@ -128,6 +134,13 @@ export function ssh2ReceivePackFactory(target: SshTarget): ReceivePackFactory { username: 'git', privateKey: target.privateKey, readyTimeout: 8_000, + // The channel sits idle while we fetch the thin pack from GitHub (we need + // the knot's advertised tips as haves before we can ask for the delta). + // Keepalives stop the knot or an intermediary from dropping that idle + // channel out from under us, which otherwise surfaces as a truncated pack + // (`unpack` failure) once we resume streaming. + keepaliveInterval: 5_000, + keepaliveCountMax: 6, hostVerifier: () => true, }) @@ -208,6 +221,18 @@ export class ReceivePackSession { else this.proc.stdin.end() const report = await this.reader.readUntilFlush() + // If the transport died while we were streaming the pack, the knot saw a + // truncated stream and reports it as an `unpack` failure (e.g. "pack + // signature mismatch"). That reads like a terminal protocol error but is + // really a transient network fault, and a naive retry replays the same + // broken pipe. Check the channel/exit state first and surface it as a + // plain transient WireError so the queue retries with a fresh session. + const exitCode = await this.proc.done.catch(() => null) + const transport = classifySshStderr(this.proc.stderr()) + if (transport) throw transport + if (exitCode !== null && exitCode !== 0) { + throw new WireError(`receive-pack exited ${exitCode} (stderr: ${this.proc.stderr().trim() || 'empty'})`) + } parseReportStatus(report ?? [], updates, this.proc.stderr()) } finally { @@ -272,10 +297,21 @@ async function writeAll(stream: NodeJS.WritableStream, data: Buffer): Promise): Promise { const src = Readable.from(packStream) await new Promise((resolve, reject) => { - src.on('error', reject) - stdin.on('error', reject) + let settled = false + const done = (err?: Error) => { + if (settled) return + settled = true + if (err) reject(err) + else resolve() + } + // A `close` before `finish` means the write side went away before the pack + // was fully flushed (the knot's channel reset mid-stream). Resolving there + // would hand a truncated pack to report-status parsing; reject so push() + // treats it as the transient transport failure it is. + src.on('error', done) + stdin.on('error', done) + stdin.on('finish', () => done()) + stdin.on('close', () => done(new WireError('receive-pack: stdin closed before pack finished streaming'))) src.pipe(stdin, { end: true }) - stdin.on('finish', resolve) - stdin.on('close', resolve) }) } diff --git a/test/unit/receive-pack.spec.ts b/test/unit/receive-pack.spec.ts index c61606d..faf6e1c 100644 --- a/test/unit/receive-pack.spec.ts +++ b/test/unit/receive-pack.spec.ts @@ -1,7 +1,10 @@ import { execFileSync } from 'node:child_process' +import { PassThrough } from 'node:stream' +import { Buffer } from 'node:buffer' 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 { RemoteRejectedError, WireError } from '../../server/utils/git-wire/errors' +import { ReceivePackSession, type ReceivePackProcess } from '../../server/utils/git-wire/receive-pack' +import { encodePktLine, flushPkt } from '../../server/utils/git-wire/pkt-line' 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' @@ -125,6 +128,50 @@ describe('receive-pack (against real git-receive-pack)', () => { await expect(ReceivePackSession.open(stalled, 50)).rejects.toThrow(/end of stream|advertisement/) }) + it('surfaces a mid-push channel death as a transient WireError, not a report parse', async () => { + // Advertise one ref so open() succeeds, then model the transport dying + // while the pack streams: stdin closes early and the process exits non-zero + // with no report-status. The knot in production would report a truncated + // pack as an `unpack` failure; the session must instead throw a plain + // (transient, retryable) WireError. + const factory = () => { + const stdin = new PassThrough() + const advertisement = Buffer.concat([ + encodePktLine(`${ZERO_SHA} refs/heads/main\0report-status\n`), + flushPkt, + ]) + const stdout = new PassThrough() + stdout.write(advertisement) + let resolveDone: (code: number | null) => void + const done = new Promise(r => { resolveDone = r }) + // Let stdin drain (so writeAll + pipePack complete their writes), but end + // stdout without a report-status flush and exit non-zero, as a reset ssh + // channel would after eating a partial pack. + stdin.resume() + stdin.on('finish', () => { + stdout.end() + resolveDone(128) + }) + const proc: ReceivePackProcess = { + stdin, + stdout, + stderr: () => '', + kill: () => { stdout.end(); resolveDone(128) }, + done, + } + return proc + } + + async function* pack(): AsyncGenerator { + yield Buffer.from('PACK') + yield Buffer.alloc(1024) + } + + const session = await ReceivePackSession.open(factory) + await expect(session.push([{ ref: 'refs/heads/main', old: ZERO_SHA, next: 'a'.repeat(40) }], pack())) + .rejects.toBeInstanceOf(WireError) + }) + it('pushes an annotated tag', async () => { const gh = fx.initBare('gh.git') const work = fx.initWork('work') diff --git a/test/unit/splice-negotiation.spec.ts b/test/unit/splice-negotiation.spec.ts new file mode 100644 index 0000000..fc3470b --- /dev/null +++ b/test/unit/splice-negotiation.spec.ts @@ -0,0 +1,88 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { fetchPack } from '../../server/utils/git-wire/upload-pack' +import { fakeGithubFetch, GitFixture, localReceivePackFactory } from '../utils/git-wire' +import { ReceivePackSession } from '../../server/utils/git-wire/receive-pack' +import { ZERO_SHA } from '../../server/utils/git-wire/refs' + +async function drain(gen: AsyncGenerator): Promise { + const parts: Buffer[] = [] + for await (const c of gen) parts.push(c) + return Buffer.concat(parts) +} + +describe('fetchPack negotiation edge cases', () => { + let fx: GitFixture + let realFetch: typeof globalThis.fetch + + beforeEach(() => { + fx = new GitFixture() + realFetch = globalThis.fetch + }) + + afterEach(() => { + globalThis.fetch = realFetch + fx.cleanup() + }) + + function firstBytes(buf: Buffer, n = 4): string { + return buf.subarray(0, n).toString('ascii') + } + + it('pack begins with the PACK signature when a have is a real common commit', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const base = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const head = fx.commit(work, 'b.txt', 'world') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + + globalThis.fetch = fakeGithubFetch(new Map([['owner/repo', gh]])) as unknown as typeof globalThis.fetch + const { pack } = await fetchPack({ repoFullName: 'owner/repo', token: 't', want: head, haves: [base], maxBytes: 1 << 30 }) + expect(firstBytes(await drain(pack))).toBe('PACK') + }) + + it('pack begins with PACK when a have is unknown to GitHub (diverged knot ref)', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const base = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const head = fx.commit(work, 'b.txt', 'world') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + + // A commit GitHub has never seen (exists only on a throwaway repo). The + // knot could advertise such a tip; we must not let it corrupt the stream. + const other = fx.initWork('other') + const stranger = fx.commit(other, 'x.txt', 'stranger') + + globalThis.fetch = fakeGithubFetch(new Map([['owner/repo', gh]])) as unknown as typeof globalThis.fetch + const { pack } = await fetchPack({ repoFullName: 'owner/repo', token: 't', want: head, haves: [stranger, base], maxBytes: 1 << 30 }) + expect(firstBytes(await drain(pack))).toBe('PACK') + }) + + it('pushes into a knot that advertises a stranger tip alongside the base', async () => { + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const base = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const head = fx.commit(work, 'b.txt', 'world') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + + // Knot has the base on main plus an unrelated tip on another ref, so its + // advertisement (and therefore our have-set) mixes a real common commit + // with one GitHub can't resolve. + const knot = fx.initBare('knot.git') + fx.pushTo(work, knot, `${base}:refs/heads/main`) + const other = fx.initWork('other') + const stranger = fx.commit(other, 'x.txt', 'stranger') + fx.pushTo(other, knot, `${stranger}:refs/heads/wip`) + + globalThis.fetch = fakeGithubFetch(new Map([['owner/repo', gh]])) as unknown as typeof globalThis.fetch + + const session = await ReceivePackSession.open(localReceivePackFactory(knot)) + const haves = [...new Set(session.tips.values())].filter(s => s !== ZERO_SHA) + const { pack } = await fetchPack({ repoFullName: 'owner/repo', token: 't', want: head, haves, maxBytes: 1 << 30 }) + await session.push([{ ref: 'refs/heads/main', old: base, next: head }], pack) + + expect(fx.revParse(knot, 'refs/heads/main')).toBe(head) + }) +}) diff --git a/test/unit/splice-thin-pack.spec.ts b/test/unit/splice-thin-pack.spec.ts new file mode 100644 index 0000000..4857ce1 --- /dev/null +++ b/test/unit/splice-thin-pack.spec.ts @@ -0,0 +1,44 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { runSplice } from '../../server/utils/splice' +import { fakeGithubFetch, GitFixture, localReceivePackFactory } from '../utils/git-wire' + +describe('runSplice thin-pack negotiation (knot already holds the base)', () => { + let fx: GitFixture + let realFetch: typeof globalThis.fetch + + beforeEach(() => { + fx = new GitFixture() + realFetch = globalThis.fetch + }) + + afterEach(() => { + globalThis.fetch = realFetch + fx.cleanup() + }) + + it('pushes an incremental update where the knot advertises the base commit as a have', async () => { + // GitHub side: two commits on main (base -> head). + const gh = fx.initBare('gh.git') + const work = fx.initWork('work') + const base = fx.commit(work, 'a.txt', 'hello') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + const head = fx.commit(work, 'b.txt', 'world') + fx.pushTo(work, gh, 'HEAD:refs/heads/main') + + // Knot side: already has `base` on main (the state after the first sync). + const knot = fx.initBare('knot.git') + fx.pushTo(work, knot, `${base}:refs/heads/main`) + + globalThis.fetch = fakeGithubFetch(new Map([['owner/repo', gh]])) as unknown as typeof globalThis.fetch + + const result = await runSplice(localReceivePackFactory(knot), { + repoFullName: 'owner/repo', + ref: 'refs/heads/main', + want: head, + token: 't', + }) + + expect(result.status).toBe('synced') + expect(fx.revParse(knot, 'refs/heads/main')).toBe(head) + }) +})