Merge pull request #832 from tdurieux/fix/streamer-public-repository-context

fix: restore public file streaming and contain ZIP failures
This commit is contained in:
Thomas Durieux
2026-09-14 17:37:51 +02:00
committed by GitHub
7 changed files with 271 additions and 33 deletions
+4 -3
View File
@@ -18,6 +18,7 @@ import FileModel from "./model/files/files.model";
import { IFile } from "./model/files/files.types";
import { FilterQuery } from "mongoose";
import { createLogger, serializeError } from "./logger";
import { githubTokenForStreamer } from "./github-token-context";
const logger = createLogger("anonymized-file");
@@ -356,7 +357,7 @@ export default class AnonymizedFile {
return got.stream(join(config.STREAMER_ENTRYPOINT, "api"), {
method: "POST",
json: {
token: await this.repository.getToken(),
token: await githubTokenForStreamer(await this.repository.getToken(), this.repository.model.source.repositoryName),
repoFullName: this.repository.model.source.repositoryName,
commit: this.repository.model.source.commit,
branch: this.repository.model.source.branch,
@@ -403,7 +404,7 @@ export default class AnonymizedFile {
json: {
sha,
size,
token,
token: await githubTokenForStreamer(token, this.repository.model.source.repositoryName),
repoFullName: this.repository.model.source.repositoryName,
commit: this.repository.model.source.commit,
branch: this.repository.model.source.branch,
@@ -507,7 +508,7 @@ export default class AnonymizedFile {
resolve();
});
} catch (error) {
handleError(error, res);
reject(error);
}
});
}
+11
View File
@@ -11,6 +11,17 @@ export function githubTokenContext(token: string) {
if (!context && token.startsWith("public-read:")) throw new Error("Public repository access context expired");
return context;
}
// Public handles belong to this process. Revalidate here, then let the
// streamer fetch public bytes anonymously without forwarding the owner's token.
export async function githubTokenForStreamer(token: string, repository: string | undefined): Promise<string> {
const context = githubTokenContext(token);
if (!context?.publicRepository) return token;
if (!repository || context.publicRepository.toLowerCase() !== repository.toLowerCase()) {
throw new Error("Public repository access context mismatch");
}
await context.renew();
return "";
}
export function githubQuotaKey(token: string) {
return contexts.get(token)?.quotaKey || createHash("sha256").update(token).digest("hex").slice(0, 24);
}
+3 -3
View File
@@ -105,13 +105,13 @@ export default class GitHubStream extends GitHubBase {
const url = githubRawFileUrl(
this.data.organization,
this.data.repoName,
this.data.commit,
this.data.commit || "HEAD",
filePath
);
logger.debug("downloading via raw URL (LFS)", { url });
return got.stream(url, {
hooks: { beforeRequest: [async () => { await githubTokenContext(token)?.renew(); }] },
headers: githubTokenContext(token)?.publicRepository ? {} : { authorization: `token ${token}` },
headers: !token || githubTokenContext(token)?.publicRepository ? {} : { authorization: `token ${token}` },
followRedirect: true,
});
}
@@ -127,7 +127,7 @@ export default class GitHubStream extends GitHubBase {
): Promise<stream.Readable> {
// Public raw downloads need no bearer token and do not consume the
// unauthenticated REST API quota. GitHub also resolves LFS pointers here.
if (githubTokenContext(token)?.publicRepository) {
if (!token || githubTokenContext(token)?.publicRepository) {
return Promise.resolve(this.downloadFileViaRaw(token, filePath));
}
return new Promise<stream.Readable>((resolve) => {
+33 -26
View File
@@ -75,17 +75,23 @@ export async function streamAnonymizedZip(
on(event: string, listener: (...args: unknown[]) => void): unknown;
}
): Promise<void> {
const source = new GitHubDownload({
repoId: opt.repoId,
organization: opt.organization,
repoName: opt.repoName,
commit: opt.commit,
getToken: opt.getToken,
});
let response;
try {
response = await source.getZipUrl();
const token = await opt.getToken();
if (!token) {
// The API already checked public access. Codeload serves public archives
// directly, avoiding the streamer's shared unauthenticated REST quota.
response = { url: `https://codeload.github.com/${encodeURIComponent(opt.organization)}/${encodeURIComponent(opt.repoName)}/zip/${encodeURIComponent(opt.commit || "HEAD")}` };
} else {
const source = new GitHubDownload({
repoId: opt.repoId,
organization: opt.organization,
repoName: opt.repoName,
commit: opt.commit,
getToken: () => token,
});
response = await source.getZipUrl();
}
} catch (error) {
const code = await classifyGitHubMissError(error, {
organization: opt.organization,
@@ -124,9 +130,14 @@ export async function streamAnonymizedZip(
// opens). Destroy the response instead so the client sees a connection
// drop and knows the download failed. Same class of silent-truncation
// bug as #694.
let upstreamSucceeded = false;
let failed = false;
const parser = Parse();
const fail = (error: Error) => {
if (failed) return;
failed = true;
logger.error("upstream zipball failed", serializeError(error));
downloadStream.destroy();
downloadStream.unpipe(parser);
archive.abort();
const destroyable = res as unknown as {
destroy?: (err?: Error) => void;
@@ -138,10 +149,13 @@ export async function streamAnonymizedZip(
destroyable.end();
}
};
// pipe() returns the destination. Listen on the archive itself as well,
// including errors emitted after the upstream ZIP has finished downloading.
archive.on("error", fail);
downloadStream
.on("error", fail)
.pipe(Parse())
.pipe(parser)
.on("entry", (entry: NodeJS.ReadableStream & { type: string; path: string; autodrain: () => void }) => {
if (entry.type === "File") {
try {
@@ -161,11 +175,13 @@ export async function streamAnonymizedZip(
...opt.anonymizerOptions,
filePath: entry.path,
});
entry.on("error", fail);
anonymizer.on("error", fail);
const st = entry.pipe(anonymizer);
archive.append(st, { name: fileName });
} catch (error) {
entry.autodrain();
logger.error("entry transform failed", serializeError(error));
fail(error as Error);
}
} else {
entry.autodrain();
@@ -173,22 +189,13 @@ export async function streamAnonymizedZip(
})
.on("error", fail)
.on("finish", () => {
upstreamSucceeded = true;
if (failed) return;
try {
archive.finalize();
} catch {
/* ignored */
archive.finalize().catch(fail);
} catch (error) {
fail(error as Error);
}
});
archive.pipe(res).on("error", (error) => {
logger.error("archive pipe error", serializeError(error));
if (!upstreamSucceeded) {
// archive errored while we were still depending on upstream bytes:
// treat as failure rather than truncating.
fail(error);
return;
}
(res as { end?: () => void }).end?.();
});
archive.pipe(res).on("error", fail);
}
+2 -1
View File
@@ -12,6 +12,7 @@ import User from "../../core/User";
import { streamAnonymizedZip } from "../../core/zipStream";
import FileModel from "../../core/model/files/files.model";
import { createLogger, serializeError } from "../../core/logger";
import { githubTokenForStreamer } from "../../core/github-token-context";
import gh = require("parse-github-url");
const logger = createLogger("repository-public");
@@ -62,7 +63,7 @@ router.get(
.stream(join(config.STREAMER_ENTRYPOINT, "api/download"), {
method: "POST",
json: {
token,
token: await githubTokenForStreamer(token, repo.model.source.repositoryName),
repoFullName: repo.model.source.repositoryName,
commit: repo.model.source.commit,
branch: repo.model.source.branch,
+122
View File
@@ -0,0 +1,122 @@
const { expect } = require("chai");
const express = require("express");
const got = require("got");
const { Readable, PassThrough } = require("stream");
require("ts-node/register/transpile-only");
const config = require("../src/config").default;
const { registerGitHubToken, githubTokenForStreamer } = require("../src/core/github-token-context");
const File = require("../src/core/AnonymizedFile").default;
const GitHubStream = require("../src/core/source/GitHubStream").default;
const { AnonymizeTransformer } = require("../src/core/anonymize-utils");
const streamer = require("../src/streamer/route").default;
async function rejects(promise, message) {
try { await promise; } catch (error) { expect(error.message).to.equal(message); return; }
throw new Error("Expected rejection");
}
async function listen(app) {
return new Promise(resolve => {
const server = app.listen(0, "127.0.0.1", () => resolve(server));
});
}
describe("public repository streamer handoff", function () {
it("revalidates public access without exporting handles or owner credentials", async function () {
let valid = true;
registerGitHubToken("public-read:transport", {
quotaKey: "test", publicRepository: "owner/public",
renew: async () => { if (!valid) throw new Error("revoked"); return "private-owner-token"; },
});
expect(await githubTokenForStreamer("public-read:transport", "owner/public")).to.equal("");
expect(await githubTokenForStreamer("installation-token", "owner/private")).to.equal("installation-token");
await rejects(githubTokenForStreamer("public-read:transport", "other/repo"), "Public repository access context mismatch");
valid = false;
await rejects(githubTokenForStreamer("public-read:transport", "owner/public"), "revoked");
await rejects(githubTokenForStreamer("public-read:missing", "owner/public"), "Public repository access context expired");
});
it("rejects send when the access recheck fails before opening the streamer", async function () {
const endpoint = config.STREAMER_ENTRYPOINT;
config.STREAMER_ENTRYPOINT = "http://unused.test/";
const response = new PassThrough();
try {
registerGitHubToken("public-read:revoked-send", {
quotaKey: "test", publicRepository: "owner/public", renew: async () => { throw new Error("revoked"); },
});
const file = new File({ repository: {
options: { terms: [] }, model: { source: { repositoryName: "owner/public" } },
getToken: async () => "public-read:revoked-send",
generateAnonymizeTransformer: filePath => new AnonymizeTransformer({ terms: [], filePath }),
}, anonymizedPath: "README.md" });
file._file = { name: "README.md", path: "", sha: "sha", size: 25 };
await rejects(file.send(response), "revoked");
} finally {
config.STREAMER_ENTRYPOINT = endpoint;
response.destroy();
}
});
for (const [mode, commit] of [
["send", "abc"], ["anonymizedContent", "abc"], ["send", undefined], ["send", ""],
]) {
it(`serves a public README through ${mode} with commit ${JSON.stringify(commit)}`, async function () {
const previous = { endpoint: config.STREAMER_ENTRYPOINT, stream: got.stream, cache: GitHubStream.prototype.getFileContentCache };
const servers = [];
let payload;
const requests = [];
try {
registerGitHubToken("public-read:http-test", {
quotaKey: "test", publicRepository: "owner/public", renew: async () => "private-owner-token",
});
got.stream = (url, options) => {
if (!String(url).startsWith("https://github.com/")) return previous.stream(url, options);
requests.push({ url, options });
return Readable.from((async function* () {
for (const hook of options.hooks.beforeRequest) await hook();
yield Buffer.from("# README\nAlice wrote this.");
})());
};
GitHubStream.prototype.getFileContentCache = async function (path) {
return this.downloadWithFallback(await this.data.getToken(), "sha", path);
};
const upstream = express();
upstream.use(express.json());
upstream.use((req, _res, next) => { payload = req.body; next(); });
upstream.use("/api", streamer);
servers.push(await listen(upstream));
config.STREAMER_ENTRYPOINT = `http://127.0.0.1:${servers[0].address().port}/`;
const options = { terms: ["Alice"], image: true, link: true };
const repo = {
repoId: "test", options,
model: { source: { repositoryName: "owner/public", commit } },
getToken: async () => "public-read:http-test",
generateAnonymizeTransformer: path => new AnonymizeTransformer({ ...options, filePath: path }),
};
const file = new File({ repository: repo, anonymizedPath: "README.md" });
file._file = { name: "README.md", path: "", sha: "sha", size: 25 };
const api = express();
api.get("/file", async (_req, res) => {
try {
if (mode === "send") await file.send(res);
else (await file.anonymizedContent()).pipe(res);
} catch (error) { res.status(500).json({ error: error.message }); }
});
servers.push(await listen(api));
const response = await got(`http://127.0.0.1:${servers[1].address().port}/file`);
expect(response.body).to.equal("# README\nXXXX-1 wrote this.");
expect(payload.token).to.equal("");
expect(JSON.stringify(payload)).not.to.include("public-read:");
expect(JSON.stringify(payload)).not.to.include("private-owner-token");
expect(requests).to.have.length(1);
expect(requests[0].url).to.equal(`https://github.com/owner/public/raw/${commit || "HEAD"}/README.md`);
expect(requests[0].options.headers).not.to.have.property("authorization");
} finally {
got.stream = previous.stream;
GitHubStream.prototype.getFileContentCache = previous.cache;
config.STREAMER_ENTRYPOINT = previous.endpoint;
await Promise.all(servers.map(server => new Promise(resolve => server.close(resolve))));
}
});
}
});
+96
View File
@@ -0,0 +1,96 @@
const { expect } = require("chai");
const { Readable, PassThrough } = require("stream");
const { setImmediate } = require("timers");
const { once } = require("events");
const vm = require("vm");
const got = require("got");
const archiver = require("archiver");
require("ts-node/register/transpile-only");
const GitHubDownload = require("../src/core/source/GitHubDownload").default;
const { AnonymizeTransformer } = require("../src/core/anonymize-utils");
const { streamAnonymizedZip } = require("../src/core/zipStream");
async function fixture() {
const zip = archiver("zip");
const chunks = [];
zip.on("data", chunk => chunks.push(chunk));
const finished = once(zip, "end");
zip.append("private source content", { name: "repo/README.md" });
await zip.finalize();
await finished;
return Buffer.concat(chunks);
}
describe("ZIP stream errors", function () {
for (const commit of ["abc", undefined, ""]) {
it(`downloads a public ZIP without a REST lookup with commit ${JSON.stringify(commit)}`, async function () {
const input = await fixture();
const previous = { stream: got.stream, zip: GitHubDownload.prototype.getZipUrl };
const response = new PassThrough();
const parser = require("unzip-stream").Parse();
const files = [];
const entries = [];
parser.on("entry", entry => {
const chunks = [];
entry.on("data", chunk => chunks.push(chunk));
entries.push(once(entry, "end").then(() => files.push({ name: entry.path, content: Buffer.concat(chunks).toString() })));
});
const finished = once(parser, "finish");
response.pipe(parser);
try {
GitHubDownload.prototype.getZipUrl = async () => { throw new Error("must not use the anonymous REST quota"); };
got.stream = url => {
expect(url).to.equal(`https://codeload.github.com/owner/public/zip/${commit || "HEAD"}`);
return Readable.from([input]);
};
await streamAnonymizedZip({
repoId: "test", organization: "owner", repoName: "public", commit,
getToken: () => "", anonymizerOptions: { terms: ["private"], image: true, link: true },
}, response);
await finished;
await Promise.all(entries);
expect(files).to.deep.equal([{ name: "README.md", content: "XXXX-1 source content" }]);
} finally {
got.stream = previous.stream;
GitHubDownload.prototype.getZipUrl = previous.zip;
response.destroy();
}
});
}
it("aborts a download on an asynchronous anonymization timeout without crashing", async function () {
const input = await fixture();
const previous = { stream: got.stream, zip: GitHubDownload.prototype.getZipUrl, flush: AnonymizeTransformer.prototype._flush };
const response = new PassThrough();
const errors = [];
const chunks = [];
response.on("error", error => errors.push(error));
response.on("data", chunk => chunks.push(chunk));
const closed = new Promise(resolve => response.once("close", resolve));
const timeout = vm.runInNewContext('Object.assign(new Error("Script execution timed out after 1000ms"), { code: "ERR_SCRIPT_EXECUTION_TIMEOUT" })');
try {
GitHubDownload.prototype.getZipUrl = async () => ({ url: "https://example.test/archive.zip" });
got.stream = () => Readable.from([input]);
// VM timeout errors come from another realm and can arrive after the ZIP
// parser finishes, while archiver is still consuming the entry stream.
AnonymizeTransformer.prototype._flush = function (callback) {
setImmediate(() => callback(timeout));
};
await streamAnonymizedZip({
repoId: "test", organization: "owner", repoName: "public", commit: "abc",
getToken: () => "private-token", anonymizerOptions: { terms: ["private"], image: true, link: true },
}, response);
await closed;
await new Promise(resolve => setImmediate(resolve));
expect(response.destroyed).to.equal(true);
expect(errors).to.have.length(1);
expect(errors[0].code).to.equal("ERR_SCRIPT_EXECUTION_TIMEOUT");
expect(Buffer.concat(chunks).includes(Buffer.from("private source content"))).to.equal(false);
} finally {
got.stream = previous.stream;
GitHubDownload.prototype.getZipUrl = previous.zip;
AnonymizeTransformer.prototype._flush = previous.flush;
response.destroy();
}
});
});