security: fix authz, DOS, and concurrency findings from review

- git smart-HTTP/SSH now require admin for private-repo clone, reject
  pending users in basic-auth, and rate-limit Basic auth attempts to
  block argon2-driven event-loop DOS and password spray (C1-C3)
- per-(client, kind) rate-limit buckets across all mutating routes:
  comments, reactions, uploads, repo/issue/patch/release creation (H1, H2)
- stream /raw via piped git cat-file with Range support, gate blob/diff/
  markdown rendering on MAX_RENDER_BYTES, cap MAX_RAW_DOWNLOAD_BYTES (H3)
- cap sharp limitInputPixels on avatar upload (H4)
- new contentDisposition helper for RFC 5987 filename* + safe ASCII (M2)
- CI run + release archive concurrency caps with FIFO queue, queued UI
  pill, and SSH-push CI trigger; pumpQueue is sole writer of
  runningTasks.set across triggerRun and retryRun (M3)
- replace shell-string interpolation in CI publish_zip with positional
  exec via per-call workDir (M6)
- defer H5 (Secure cookie / HSTS) with rationale captured in TODO.txt
AuthorKonata <konata@posteo.jp>
Date
Commiteb52ac0fd3fc4dfbbee9269c2efcd35cea103c8c
Parenta2f612e
23 files changed, 983 insertions(+), 156 deletions(-)
▾M.gitignore
@@ -40,3 +40,5 @@ data-test/
# Generated CSS
public/assets/main.css
.claude
▾MREADME.md
@@ -80,7 +80,7 @@ All settings are environment variables:
| `MAX_USER_UPLOAD_BYTES` | `2097152` | Max upload size for non-admin users (2 MB) |
| `INLINE_MAX_BYTES` | `524288` | Max file size rendered inline in the code view (512 KB) |
| `SSH_DISABLED` | `0` | Set to `1` to disable the embedded SSH server |
| `TRUSTED_PROXY` | `0` | Trust `X-Forwarded-For` headers |
| `TRUSTED_PROXY` | `0` | Trust `X-Forwarded-For` headers \*\* |
| `RATE_LIMIT_DISABLED` | `0` | Set to `1` to disable rate limiting |
| `HIGHLIGHT_WORKERS` | `4` | Syntax highlighting worker threads \* |
| `COMMITTER_NAME` | `$OWNER_DISPLAY_NAME` | Git committer name for merges and UI edits |
@@ -98,6 +98,8 @@ All settings are environment variables:
\* Each highlighting worker loads its own copy of the language grammars and uses ~200 MB of memory. Increase with care.
\*\* Set `TRUSTED_PROXY=1` only when Hearthforge is behind a reverse proxy that strips any incoming `X-Forwarded-For` from clients. Caddy and Traefik do this by default; nginx requires `proxy_set_header X-Forwarded-For $remote_addr;` (rather than the common `$proxy_add_x_forwarded_for`, which appends to a client-supplied value). Setting `TRUSTED_PROXY=1` in front of a proxy that does not strip means rate limits and any audit logging are spoofable per request.
## CI/CD Pipelines
Hearthforge includes a built-in CI/CD system that runs pipelines in Docker or Podman containers,
▾Msrc/app.ts
@@ -37,8 +37,9 @@ export async function createApp(port: number) {
await syncStartup();
await cancelStaleRuns();
// Recurring session cleanup — runs every 24 hours
setInterval(cleanupSessions, 24 * 60 * 60 * 1000);
// Recurring session cleanup — runs hourly so the row count tracks
// expiry instead of trailing it by up to a day.
setInterval(cleanupSessions, 60 * 60 * 1000);
return new Elysia({
serve: { maxRequestBodySize: config.MAX_UPLOAD_BYTES },
▾Msrc/config.ts
@@ -2,6 +2,13 @@ import path from "node:path";
const env = process.env;
/** parseInt with a default that distinguishes "unset" from "explicit 0". */
function intEnv(value: string | undefined, defaultValue: number): number {
if (value === undefined || value === "") return defaultValue;
const parsed = parseInt(value, 10);
return Number.isFinite(parsed) ? parsed : defaultValue;
}
const config = {
OWNER_DISPLAY_NAME: env.OWNER_DISPLAY_NAME ?? "Admin",
INLINE_MAX_BYTES: parseInt(env.INLINE_MAX_BYTES ?? "", 10) || 524288,
@@ -35,6 +42,23 @@ const config = {
CI_MAX_HISTORY: parseInt(env.CI_MAX_HISTORY ?? "", 10) || 50,
CI_MAX_CONCURRENT: parseInt(env.CI_MAX_CONCURRENT ?? "", 10) || 2,
CI_DEFAULT_TIMEOUT: parseInt(env.CI_DEFAULT_TIMEOUT ?? "", 10) || 3600,
// Cap on simultaneous source archive generation jobs (one release with
// include_source_code spawns three git-archive + compressor pipelines
// back-to-back). Without this cap, an admin firing several releases in
// quick succession can saturate CPU. Excess jobs are queued in memory.
MAX_CONCURRENT_ARCHIVE_JOBS:
parseInt(env.MAX_CONCURRENT_ARCHIVE_JOBS ?? "", 10) || 2,
// Cap on any server-side render that holds the whole content in
// RAM and runs synchronous CPU work on it (blob view, commit/patch
// diff, markdown). Above this, the UI shows a "too large to
// preview" stub and links to the raw endpoint. Setting to 0
// disables inline rendering entirely.
MAX_RENDER_BYTES: intEnv(env.MAX_RENDER_BYTES, 10 * 1024 * 1024),
// Cap on the streamed /raw download. 0 means no limit (current
// behaviour) — useful on a trusted LAN where you actually want to
// pull large blobs out of the browser. Public deployments should
// pin this to something sane.
MAX_RAW_DOWNLOAD_BYTES: intEnv(env.MAX_RAW_DOWNLOAD_BYTES, 0),
};
// Derived values that depend on other config fields
▾Msrc/constants.ts
@@ -31,6 +31,26 @@ export const LOGIN_MAX_ATTEMPTS = 10;
export const LOGIN_RATE_WINDOW_MS = 60_000;
export const REGISTRATION_MAX_ATTEMPTS = 3;
export const REGISTRATION_RATE_WINDOW_MS = 60 * 60_000;
// Bounds the per-IP cost of git smart-HTTP basic auth — every call to
// verifyBasicAuth runs argon2 (~100ms) and would otherwise be an
// unauthenticated event-loop DOS vector and a brute-force oracle.
export const GIT_AUTH_MAX_ATTEMPTS = 10;
export const GIT_AUTH_RATE_WINDOW_MS = 60_000;
// Per-user / per-IP caps on user-content writes. Numbers are deliberately
// roomy for a logged-in person clicking around but tight enough that a
// scripted client can't fill the database in seconds.
export const COMMENT_MAX_PER_MIN = 30;
export const REACTION_MAX_PER_MIN = 60;
export const ISSUE_CREATE_MAX_PER_MIN = 10;
export const PATCH_CREATE_MAX_PER_MIN = 10;
export const REPO_CREATE_MAX_PER_HOUR = 30;
export const FILE_EDIT_MAX_PER_MIN = 30;
export const RELEASE_WRITE_MAX_PER_MIN = 20;
export const LABEL_WRITE_MAX_PER_MIN = 30;
export const UPLOAD_MAX_PER_MIN = 10;
export const RATE_WINDOW_MIN_MS = 60_000;
export const RATE_WINDOW_HOUR_MS = 60 * 60_000;
// Session
export const SESSION_ID_BYTES = 32;
▾Asrc/lib/contentDisposition.ts
@@ -0,0 +1,26 @@
/**
* Build a safe `Content-Disposition` header value.
*
* The display filename is restricted to a printable ASCII subset so it can
* never break out of the quoted-string form (no `"`, `\`, CR, LF, NUL),
* and a UTF-8 `filename*` parameter is added per RFC 5987 so unicode names
* still come through to the client when possible.
*/
export function contentDisposition(
type: "inline" | "attachment",
name: string,
): string {
// Drop path components and control chars; collapse anything not
// printable-ASCII-and-safe-in-a-quoted-string into "_".
const base = name.split(/[/\\]/).pop() ?? "";
// biome-ignore lint/suspicious/noControlCharactersInRegex: control characters are exactly what we want to strip from HTTP header values
const stripped = base.replace(/[\x00-\x1f\x7f"\\]/g, "_");
const ascii = stripped.length > 0 ? stripped : "file";
// biome-ignore lint/suspicious/noControlCharactersInRegex: control characters are exactly what we want to strip from HTTP header values
const utf8 = encodeURIComponent(base.replace(/[\x00-\x1f\x7f]/g, "_"))
// RFC 5987 disallows `'` in filename* value-chars (it is the
// separator between charset, language, and value); encodeURIComponent
// does not escape it, so do it explicitly.
.replace(/'/g, "%27");
return `${type}; filename="${ascii}"; filename*=UTF-8''${utf8}`;
}
▾Msrc/lib/rateLimiter.ts
@@ -5,28 +5,53 @@ interface Bucket {
resetAt: number;
}
export type RateLimitKind =
| "login"
| "passkey"
| "git-auth"
| "comment"
| "reaction"
| "upload"
| "register"
| "repo-create"
| "issue-create"
| "patch-create"
| "label-write"
| "release-write"
| "file-edit";
const buckets = new Map<string, Bucket>();
let sweepCounter = 0;
function maybeSweep() {
if (++sweepCounter < 1000) return;
sweepCounter = 0;
function sweep() {
const now = Date.now();
for (const [key, bucket] of buckets) {
if (now > bucket.resetAt) buckets.delete(key);
}
}
// Periodic sweep so memory doesn't grow unboundedly when traffic is low
// and the request-driven sweep never reaches its threshold.
setInterval(sweep, 60 * 1000).unref();
let sweepCounter = 0;
function maybeSweep() {
if (++sweepCounter < 1000) return;
sweepCounter = 0;
sweep();
}
export function checkRateLimit(
ip: string | null,
kind: RateLimitKind,
maxRequests: number,
windowMs: number,
): boolean {
if (config.RATE_LIMIT_DISABLED || !ip) return true;
const key = `${ip}|${kind}`;
const now = Date.now();
const bucket = buckets.get(ip);
const bucket = buckets.get(key);
if (!bucket || now > bucket.resetAt) {
buckets.set(ip, { count: 1, resetAt: now + windowMs });
buckets.set(key, { count: 1, resetAt: now + windowMs });
maybeSweep();
return true;
}
@@ -47,3 +72,25 @@ export function getClientIp(
}
return server?.requestIP(request)?.address ?? null;
}
/**
* Rate-limit a mutating handler. Returns null if the request is allowed,
* or a 429 Response if it isn't. The bucket is keyed on the user id when
* available (so a single attacker can't bypass by rotating source IPs)
* and on the IP when not.
*/
export function rateLimit(
request: Request,
server: Bun.Server<unknown> | null,
userId: number | null,
kind: RateLimitKind,
maxRequests: number,
windowMs: number,
): Response | null {
const key = userId !== null ? `u${userId}` : getClientIp(request, server);
if (checkRateLimit(key, kind, maxRequests, windowMs)) return null;
return new Response("Too many requests. Please slow down.", {
status: 429,
headers: { "Content-Type": "text/plain; charset=utf-8" },
});
}
▾Msrc/routes/auth.tsx
@@ -102,7 +102,14 @@ export const authRoutes = new Elysia()
"/login",
async ({ body, request, server }) => {
const ip = getClientIp(request, server);
if (!checkRateLimit(ip, LOGIN_MAX_ATTEMPTS, LOGIN_RATE_WINDOW_MS)) {
if (
!checkRateLimit(
ip,
"login",
LOGIN_MAX_ATTEMPTS,
LOGIN_RATE_WINDOW_MS,
)
) {
return html(
<Login error="Too many login attempts. Please try again later." />,
);
@@ -154,6 +161,7 @@ export const authRoutes = new Elysia()
if (
!checkRateLimit(
ip,
"register",
REGISTRATION_MAX_ATTEMPTS,
REGISTRATION_RATE_WINDOW_MS,
)
▾Msrc/routes/avatars.ts
@@ -1,7 +1,12 @@
import { Elysia, t } from "elysia";
import config from "../config.ts";
import { YEAR_SECONDS } from "../constants.ts";
import {
RATE_WINDOW_MIN_MS,
UPLOAD_MAX_PER_MIN,
YEAR_SECONDS,
} from "../constants.ts";
import { db } from "../db/index.ts";
import { rateLimit } from "../lib/rateLimiter.ts";
import { requireAuth, resolveSession } from "../middleware/session.ts";
import {
avatarJxlPath,
@@ -52,10 +57,19 @@ export const avatarRoutes = new Elysia()
.post(
"/settings/avatar",
async ({ body, cookie }) => {
async ({ body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"upload",
UPLOAD_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
if (body.avatar.size > config.MAX_USER_UPLOAD_BYTES) {
return new Response("Avatar file too large", { status: 400 });
}
▾Msrc/routes/ci.tsx
@@ -3,10 +3,12 @@ import path from "node:path";
import { Elysia, t } from "elysia";
import { CI_RUNS_PER_PAGE, paths } from "../constants.ts";
import { db, getRepo } from "../db/index.ts";
import { contentDisposition } from "../lib/contentDisposition.ts";
import { paginate } from "../lib/pagination.ts";
import { requireAdmin, resolveSession } from "../middleware/session.ts";
import {
cancelRun,
ciQueuePosition,
parseCiConfig,
purgeRepoCaches,
retryRun,
@@ -146,6 +148,8 @@ export const ciRoutes = new Elysia()
const runsWithCounts = runs.map((r) => ({
...r,
artifact_count: artifactCountMap.get(r.id) ?? 0,
queue_position:
r.status === "queued" ? ciQueuePosition(r.id) : null,
}));
// Determine why manual trigger may be unavailable (admin-only check)
@@ -252,6 +256,9 @@ export const ciRoutes = new Elysia()
steps={steps}
artifacts={artifacts}
autoRefresh={query.refresh !== "off"}
queuePosition={
run.status === "queued" ? ciQueuePosition(run.id) : null
}
/>,
);
},
@@ -508,7 +515,10 @@ export const ciRoutes = new Elysia()
return new Response(Bun.file(filePath), {
headers: {
"Content-Disposition": `attachment; filename="${artifact.filename}"`,
"Content-Disposition": contentDisposition(
"attachment",
artifact.filename,
),
"Content-Type": "application/octet-stream",
"Content-Length": String(artifact.size),
},
▾Msrc/routes/git.ts
@@ -2,8 +2,15 @@ import { existsSync } from "node:fs";
import path from "node:path";
import * as argon2 from "argon2";
import { Elysia, t } from "elysia";
import { ADMIN_USERNAME, paths, VALID_REPO_NAME_RE } from "../constants.ts";
import {
ADMIN_USERNAME,
GIT_AUTH_MAX_ATTEMPTS,
GIT_AUTH_RATE_WINDOW_MS,
paths,
VALID_REPO_NAME_RE,
} from "../constants.ts";
import { db } from "../db";
import { checkRateLimit, getClientIp } from "../lib/rateLimiter.ts";
import {
parseCiConfig,
shouldTriggerPush,
@@ -110,6 +117,7 @@ async function verifyBasicAuth(
.selectFrom("users")
.select("password_hash")
.where("username", "=", username)
.where("is_pending", "=", 0)
.executeTakeFirst();
if (!user?.password_hash) return false;
try {
@@ -129,6 +137,37 @@ function unauthorized(): Response {
});
}
function tooManyRequests(): Response {
return new Response("Too Many Requests", {
status: 429,
headers: { "Content-Type": "text/plain" },
});
}
/**
* Rate-limit the per-IP cost of `verifyBasicAuth`. Each call costs
* ~100ms of argon2 work on the single event-loop thread, so without
* this an unauthenticated attacker can pin the CPU and use the same
* endpoint as a password-spray oracle around the /login limiter.
*
* Counts even non-Basic-header requests against the bucket: a private
* repo is only listed in the UI for admins, so a non-admin only ever
* pokes these endpoints intentionally and shouldn't get a free retry
* by omitting the header.
*/
function checkGitAuthLimit(
request: Request,
server: Bun.Server<unknown> | null,
): boolean {
const ip = getClientIp(request, server);
return checkRateLimit(
ip,
"git-auth",
GIT_AUTH_MAX_ATTEMPTS,
GIT_AUTH_RATE_WINDOW_MS,
);
}
async function getRepo(
slug: string,
): Promise<{ name: string; repoPath: string; isPrivate: boolean } | null> {
@@ -162,7 +201,7 @@ export const gitRoutes = new Elysia()
// info/refs — serves both upload-pack (clone/fetch) and receive-pack (push)
.get(
"/:repo/info/refs",
async ({ params, query, request }) => {
async ({ params, query, request, server }) => {
const service = query.service;
if (
service !== "git-upload-pack" &&
@@ -175,12 +214,18 @@ export const gitRoutes = new Elysia()
if (!repo) return new Response("Not Found", { status: 404 });
const authHeader = request.headers.get("Authorization");
if (service === "git-receive-pack") {
const requiresAuth =
service === "git-receive-pack" || repo.isPrivate;
if (requiresAuth) {
if (!checkGitAuthLimit(request, server))
return tooManyRequests();
// Match the UI: private repos are admin-only. The web UI
// returns 404 to non-admins via getRepo() in repos.tsx, so
// the smart-HTTP path must do the same — otherwise any
// logged-in user could clone a "private" repo despite the
// UI hiding it.
if (!(await verifyBasicAuth(authHeader, true)))
return unauthorized();
} else if (repo.isPrivate) {
if (!(await verifyBasicAuth(authHeader, false)))
return unauthorized();
}
const gitCmd =
@@ -211,14 +256,16 @@ export const gitRoutes = new Elysia()
)
// upload-pack POST — clone/fetch pack transfer (public for public repos)
.post("/:repo/git-upload-pack", async ({ params, request }) => {
.post("/:repo/git-upload-pack", async ({ params, request, server }) => {
const repo = await getRepo(params.repo);
if (!repo) return new Response("Not Found", { status: 404 });
if (repo.isPrivate) {
// Admin-only: see info/refs branch above.
if (!checkGitAuthLimit(request, server)) return tooManyRequests();
if (
!(await verifyBasicAuth(
request.headers.get("Authorization"),
false,
true,
))
)
return unauthorized();
@@ -237,7 +284,8 @@ export const gitRoutes = new Elysia()
})
// receive-pack POST — push pack transfer (admin only)
.post("/:repo/git-receive-pack", async ({ params, request }) => {
.post("/:repo/git-receive-pack", async ({ params, request, server }) => {
if (!checkGitAuthLimit(request, server)) return tooManyRequests();
if (
!(await verifyBasicAuth(request.headers.get("Authorization"), true))
)
▾Msrc/routes/issues.tsx
@@ -1,11 +1,20 @@
import { Elysia, t } from "elysia";
import { sql } from "kysely";
import config from "../config.ts";
import { ALLOWED_REACTIONS, ISSUES_PER_PAGE } from "../constants.ts";
import {
ALLOWED_REACTIONS,
COMMENT_MAX_PER_MIN,
ISSUE_CREATE_MAX_PER_MIN,
ISSUES_PER_PAGE,
LABEL_WRITE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
REACTION_MAX_PER_MIN,
} from "../constants.ts";
import { issuesLabelFilter } from "../db/helpers.ts";
import { db, getRepo, type LabelRow } from "../db/index.ts";
import { authorizeCommentEdit } from "../lib/commentAuth.ts";
import { paginate } from "../lib/pagination.ts";
import { rateLimit } from "../lib/rateLimiter.ts";
import {
requireAdmin,
requireAuth,
@@ -200,10 +209,19 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"issue-create",
ISSUE_CREATE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -412,10 +430,19 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues/:number/comments",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"comment",
COMMENT_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -476,10 +503,19 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues/:number/react",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"reaction",
REACTION_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -639,10 +675,19 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues/:number/edit",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"comment",
COMMENT_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -687,10 +732,19 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues/:number/comments/:id/edit",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"comment",
COMMENT_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -731,9 +785,18 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues/:number/labels/add",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
if (!user) return new Response("Unauthorized", { status: 401 });
const limited = rateLimit(
request,
server,
user.id,
"label-write",
LABEL_WRITE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user.isAdmin);
if (!repo) return new Response("Not found", { status: 404 });
@@ -783,9 +846,18 @@ export const issueRoutes = new Elysia()
.post(
"/:repo/issues/:number/labels/remove",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
if (!user) return new Response("Unauthorized", { status: 401 });
const limited = rateLimit(
request,
server,
user.id,
"label-write",
LABEL_WRITE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user.isAdmin);
if (!repo) return new Response("Not found", { status: 404 });
▾Msrc/routes/patches.tsx
@@ -1,11 +1,20 @@
import { Elysia, t } from "elysia";
import { sql } from "kysely";
import config from "../config.ts";
import { ALLOWED_REACTIONS, PATCHES_PER_PAGE } from "../constants.ts";
import {
ALLOWED_REACTIONS,
COMMENT_MAX_PER_MIN,
LABEL_WRITE_MAX_PER_MIN,
PATCH_CREATE_MAX_PER_MIN,
PATCHES_PER_PAGE,
RATE_WINDOW_MIN_MS,
REACTION_MAX_PER_MIN,
} from "../constants.ts";
import { patchesLabelFilter } from "../db/helpers.ts";
import { db, getRepo, type LabelRow } from "../db/index.ts";
import { authorizeCommentEdit } from "../lib/commentAuth.ts";
import { paginate } from "../lib/pagination.ts";
import { rateLimit } from "../lib/rateLimiter.ts";
import {
requireAdmin,
requireAuth,
@@ -235,10 +244,19 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"patch-create",
PATCH_CREATE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -649,10 +667,19 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/upload",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"patch-create",
PATCH_CREATE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -793,10 +820,19 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/comments",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"comment",
COMMENT_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -849,10 +885,19 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/react",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"reaction",
REACTION_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -926,10 +971,19 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/comments/:id/edit",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"comment",
COMMENT_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -970,10 +1024,19 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/edit",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
const deny = requireAuth(user);
if (deny) return deny;
const limited = rateLimit(
request,
server,
user?.id ?? null,
"comment",
COMMENT_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user?.isAdmin ?? false);
if (!repo) return new Response("Not found", { status: 404 });
@@ -1018,9 +1081,18 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/labels/add",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
if (!user) return new Response("Unauthorized", { status: 401 });
const limited = rateLimit(
request,
server,
user.id,
"label-write",
LABEL_WRITE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user.isAdmin);
if (!repo) return new Response("Not found", { status: 404 });
@@ -1070,9 +1142,18 @@ export const patchRoutes = new Elysia()
.post(
"/:repo/patches/:number/labels/remove",
async ({ params, body, cookie }) => {
async ({ params, body, cookie, request, server }) => {
const user = await resolveSession(cookie.session.value);
if (!user) return new Response("Unauthorized", { status: 401 });
const limited = rateLimit(
request,
server,
user.id,
"label-write",
LABEL_WRITE_MAX_PER_MIN,
RATE_WINDOW_MIN_MS,
);
if (limited) return limited;
const repo = await getRepo(params.repo, user.isAdmin);
if (!repo) return new Response("Not found", { status: 404 });
▾Msrc/routes/releases.tsx
@@ -4,6 +4,7 @@ import { Elysia, t } from "elysia";
import config from "../config.ts";
import { paths, RELEASES_PER_PAGE } from "../constants.ts";
import { db, getRepo } from "../db/index.ts";
import { contentDisposition } from "../lib/contentDisposition.ts";
import { paginate } from "../lib/pagination.ts";
import { requireAdmin, resolveSession } from "../middleware/session.ts";
import { archiveRepo, git } from "../services/git.ts";
@@ -18,6 +19,57 @@ import { html } from "../views/render.tsx";
// immediately when the corresponding release is deleted.
const archivingTasks = new Map<number, AbortController>();
// FIFO queue of release archive jobs waiting for a slot. Each entry is
// keyed by release ID so a delete can pull it out before it ever starts.
interface PendingArchive {
releaseId: number;
repoName: string;
tagName: string;
sourceDir: string;
}
const queuedArchives: PendingArchive[] = [];
function pumpArchiveQueue(): void {
while (
queuedArchives.length > 0 &&
archivingTasks.size < config.MAX_CONCURRENT_ARCHIVE_JOBS
) {
const next = queuedArchives.shift()!;
runArchiveJob(next);
}
}
function runArchiveJob(job: PendingArchive): void {
const controller = new AbortController();
archivingTasks.set(job.releaseId, controller);
void (async () => {
try {
await archiveRepo(
job.repoName,
job.tagName,
job.repoName,
job.sourceDir,
controller.signal,
);
rmSync(path.join(job.sourceDir, ".pending"), { force: true });
} catch {
// Either the release was deleted (abort) or archiving failed.
rmSync(job.sourceDir, { recursive: true, force: true });
} finally {
archivingTasks.delete(job.releaseId);
pumpArchiveQueue();
}
})();
}
function scheduleArchive(job: PendingArchive): void {
if (archivingTasks.size >= config.MAX_CONCURRENT_ARCHIVE_JOBS) {
queuedArchives.push(job);
return;
}
runArchiveJob(job);
}
function sanitizeFilename(name: string): string {
const safe = path.basename(name).replace(/[^a-zA-Z0-9._-]/g, "_");
if (!safe || /^\.+$/.test(safe)) return "_";
@@ -312,40 +364,21 @@ export const releasesRoutes = new Elysia()
// Kick off source archive generation in the background so the
// response can be sent immediately. The .pending sentinel written
// inside the transaction signals to the detail view that archives
// are still being prepared. If the release is deleted while
// generation is in progress the background task will hit errors
// (the directory will have been removed) and silently bail out;
// SQLite AUTOINCREMENT guarantees the ID is never reused, so there
// is no risk of contaminating a later release.
// are still being prepared. Concurrent jobs are capped at
// MAX_CONCURRENT_ARCHIVE_JOBS so a flurry of release creations
// can't saturate CPU; excess jobs queue in-memory.
if (includeSource) {
const sourceDir = path.join(
paths.RELEASES_DIR,
String(releaseId),
"source",
);
const controller = new AbortController();
archivingTasks.set(releaseId, controller);
(async () => {
try {
await archiveRepo(
repo.name,
tagName!,
repo.name,
sourceDir,
controller.signal,
);
rmSync(path.join(sourceDir, ".pending"), {
force: true,
});
} catch {
// Either the release was deleted (abort) or archiving
// failed. Remove the source dir so the UI shows no
// stale state.
rmSync(sourceDir, { recursive: true, force: true });
} finally {
archivingTasks.delete(releaseId);
}
})();
scheduleArchive({
releaseId,
repoName: repo.name,
tagName: tagName!,
sourceDir,
});
}
return new Response(null, {
@@ -475,9 +508,14 @@ export const releasesRoutes = new Elysia()
if (!release) return new Response("Not found", { status: 404 });
// Abort any in-progress archive generation before touching disk so
// the background task doesn't race with the rmSync below.
// the background task doesn't race with the rmSync below. Also
// pull queued (not-yet-started) archive jobs out of the queue.
archivingTasks.get(release.id)?.abort();
archivingTasks.delete(release.id);
const qIdx = queuedArchives.findIndex(
(j) => j.releaseId === release.id,
);
if (qIdx >= 0) queuedArchives.splice(qIdx, 1);
// Remove files from disk before the DB record so that a crash
// between the two leaves a broken-but-visible repo rather than a
@@ -541,7 +579,10 @@ export const releasesRoutes = new Elysia()
return new Response(file, {
headers: {
"Content-Disposition": `attachment; filename="${safeFilename}"`,
"Content-Disposition": contentDisposition(
"attachment",
safeFilename,
),
"Content-Type": "application/octet-stream",
},
});
@@ -584,7 +625,10 @@ export const releasesRoutes = new Elysia()
return new Response(file, {
headers: {
"Content-Disposition": `attachment; filename="${safeFilename}"`,
"Content-Disposition": contentDisposition(
"attachment",
safeFilename,
),
"Content-Type": "application/octet-stream",
},
});
▾Msrc/routes/repos.tsx
@@ -17,6 +17,7 @@ import {
YEAR_SECONDS,
} from "../constants.ts";
import { db } from "../db/index.ts";
import { contentDisposition } from "../lib/contentDisposition.ts";
import { redirect } from "../lib/redirect.ts";
import { requireAdmin, resolveSession } from "../middleware/session.ts";
import { git, repoPath, type TreeEntry } from "../services/git.ts";
@@ -48,6 +49,99 @@ async function getRepo(name: string, isAdmin: boolean) {
return repo;
}
/**
* Read a byte range out of a streaming source without ever holding
* the full content in memory. Used by the /raw endpoint to honour
* HTTP Range headers against `git show`'s pipe.
*/
function sliceStream(
source: ReadableStream<Uint8Array>,
start: number,
length: number,
onDone: () => void,
): ReadableStream<Uint8Array> {
const reader = source.getReader();
let skipped = 0;
let emitted = 0;
let finished = false;
const finish = () => {
if (finished) return;
finished = true;
reader.cancel().catch(() => {});
onDone();
};
return new ReadableStream<Uint8Array>({
async pull(controller) {
while (emitted < length) {
const { value, done } = await reader.read();
if (done) {
controller.close();
finish();
return;
}
let chunk = value;
if (skipped < start) {
const drop = Math.min(start - skipped, chunk.length);
skipped += drop;
chunk = chunk.subarray(drop);
if (chunk.length === 0) continue;
}
const remaining = length - emitted;
if (chunk.length > remaining)
chunk = chunk.subarray(0, remaining);
emitted += chunk.length;
controller.enqueue(chunk);
if (emitted >= length) {
controller.close();
finish();
}
return;
}
controller.close();
finish();
},
cancel() {
finish();
},
});
}
/** Wrap a process stdout stream so the underlying process is killed on
* close or cancel. Without this, a client disconnect partway through a
* large blob leaves `git cat-file` running until its pipe back-pressures. */
function streamWithKill(
source: ReadableStream<Uint8Array>,
proc: { kill: () => void },
): ReadableStream<Uint8Array> {
const reader = source.getReader();
let killed = false;
const finish = () => {
if (killed) return;
killed = true;
reader.cancel().catch(() => {});
proc.kill();
};
return new ReadableStream<Uint8Array>({
async pull(controller) {
try {
const { value, done } = await reader.read();
if (done) {
controller.close();
finish();
return;
}
controller.enqueue(value);
} catch (err) {
controller.error(err);
finish();
}
},
cancel() {
finish();
},
});
}
async function mimeForContent(
filename: string,
content: Buffer,
@@ -477,6 +571,30 @@ export const repoRoutes = new Elysia()
if (!repo) return new Response("Not found", { status: 404 });
const filePath = decodeURIComponent(params["*"]);
// Size-gate before reading the blob into memory: holding a
// huge content buffer (and then running shiki/marked over it)
// is the cheapest DOS vector against unauthenticated users on
// a public repo.
const size = await git.getFileSize(repo.name, params.ref, filePath);
const filename = path.basename(filePath);
if (size !== null && size > config.MAX_RENDER_BYTES) {
const [branches, tags] = await Promise.all([
git.branches(repo.name),
git.tags(repo.name),
]);
return html(
<FileBlob
user={user}
repo={repo}
ref={params.ref}
filePath={filePath}
view={{ type: "download", size }}
branches={branches}
tags={tags}
markdownHtml={undefined}
/>,
);
}
const [content, branches, tags, commitSHA] = await Promise.all([
git.show(repo.name, params.ref, filePath),
git.branches(repo.name),
@@ -486,7 +604,6 @@ export const repoRoutes = new Elysia()
if (!content || !commitSHA)
return new Response("Not found", { status: 404 });
const filename = path.basename(filePath);
const [view, markdownHtml] = await Promise.all([
serveFile(
content,
@@ -530,36 +647,88 @@ export const repoRoutes = new Elysia()
if (!repo) return new Response("Not found", { status: 404 });
const filePath = decodeURIComponent(params["*"]);
const content = await git.show(repo.name, params.ref, filePath);
if (!content) return new Response("Not found", { status: 404 });
// Resolve the file's blob hash and size up front. cat-file -s
// is O(1) and lets us stream the blob below without ever
// holding it whole in memory — old code did
// `await arrayBuffer()` and sliced for Range, peaking at
// file_size × concurrent_requests of RSS.
const total = await git.getFileSize(repo.name, params.ref, filePath);
if (total === null) return new Response("Not found", { status: 404 });
if (
config.MAX_RAW_DOWNLOAD_BYTES > 0 &&
total > config.MAX_RAW_DOWNLOAD_BYTES
) {
return new Response("File exceeds raw download size limit", {
status: 413,
});
}
const filename = path.basename(filePath);
const contentType = await mimeForContent(filename, content);
const total = content.length;
// We use `cat-file blob` rather than `git show` so the bytes
// streamed exactly match what `cat-file -s` reported above —
// `git show` can apply smudge filters / autocrlf, which would
// make Content-Length wrong on filtered repos.
const blobArgs = [
"git",
"-C",
repoPath(repo.name),
"cat-file",
"blob",
`${params.ref}:${filePath}`,
];
// Sniff content-type from the first bytes only — same idea as
// the old mimeForContent but without buffering the full blob.
const sniffStream = Bun.spawn(blobArgs, {
stdout: "pipe",
stderr: "ignore",
});
const sniffReader = sniffStream.stdout.getReader();
const { value: firstChunk } = await sniffReader.read();
sniffReader.cancel().catch(() => {});
sniffStream.kill();
const head = firstChunk
? Buffer.from(firstChunk.subarray(0, BINARY_DETECT_BYTES))
: Buffer.alloc(0);
const contentType = await mimeForContent(filename, head);
const rangeHeader = request.headers.get("Range");
const fullProc = Bun.spawn(blobArgs, {
stdout: "pipe",
stderr: "ignore",
});
if (rangeHeader) {
const match = rangeHeader.match(/bytes=(\d*)-(\d*)/);
if (match) {
const start = match[1] ? parseInt(match[1], 10) : 0;
const end = match[2] ? parseInt(match[2], 10) : total - 1;
const clampedEnd = Math.min(end, total - 1);
return new Response(content.subarray(start, clampedEnd + 1), {
const length = clampedEnd - start + 1;
// sliceStream kills fullProc once the slice is exhausted or
// the consumer cancels — without this, requesting a tiny
// range from a huge blob leaves `git cat-file` running.
const sliced = sliceStream(fullProc.stdout, start, length, () =>
fullProc.kill(),
);
return new Response(sliced, {
status: 206,
headers: {
"Content-Type": contentType,
"Content-Range": `bytes ${start}-${clampedEnd}/${total}`,
"Accept-Ranges": "bytes",
"Content-Length": String(clampedEnd - start + 1),
"Content-Length": String(length),
},
});
}
}
return new Response(content, {
// Wrap the full stream too so a client disconnect during a large
// download kills the underlying git process instead of leaving it
// wedged on a back-pressured pipe.
return new Response(streamWithKill(fullProc.stdout, fullProc), {
headers: {
"Content-Type": contentType,
"Content-Disposition": `inline; filename="${filename}"`,
"Content-Disposition": contentDisposition("inline", filename),
"Content-Length": String(total),
"Accept-Ranges": "bytes",
},
@@ -734,6 +903,21 @@ export const repoRoutes = new Elysia()
git.diff(repo.name, params.sha),
]);
if (!meta) return new Response("Commit not found", { status: 404 });
// Reject oversized diffs after generation but before highlighting.
// Avoids running marked / shiki / DOMPurify over a multi-megabyte
// diff which would synchronously stall the worker pool.
if (rawDiff.length > config.MAX_RENDER_BYTES) {
return html(
<CommitDetail
user={user}
repo={repo}
sha={params.sha}
meta={meta}
files={[]}
tooLarge={rawDiff.length}
/>,
);
}
const files = await prepareDiff(
rawDiff,
`commit:${repo.name}:${params.sha}`,
▾Msrc/services/avatar.ts
@@ -52,7 +52,13 @@ export async function processAndStoreAvatar(
if (!type?.mime.startsWith("image/")) {
throw new Error("Invalid image type");
}
const { data, info } = await sharp(buffer)
// Cap input pixels to defeat decompression bombs: a 2 MB
// PNG/WebP can claim 16k×16k = ~256 MP and force sharp to allocate
// ~1 GB of RGBA before the resize step. 4096*4096 = 16 MP is well
// above any plausible avatar source (the output is 128×128).
const { data, info } = await sharp(buffer, {
limitInputPixels: 4096 * 4096,
})
.resize(128, 128, { fit: "cover", position: "center" })
.ensureAlpha()
.raw()
▾Msrc/services/ci.ts
@@ -66,6 +66,44 @@ const runningTasks = new Map<
{ controller: AbortController; containerId?: string }
>();
// FIFO queue of run IDs whose DB row is `status = "queued"`. We start the
// next one in `pumpQueue` whenever `runningTasks.size` drops below
// `CI_MAX_CONCURRENT`. `pumpQueue` is the SOLE writer of `runningTasks.set`
// — callers that want to start a run push to `queuedRunIds` and call
// `pumpQueue` synchronously. This makes the size check + slot reservation
// atomic against concurrent callers (no await between).
const queuedRunIds: number[] = [];
async function promoteToPending(runId: number): Promise<void> {
await db
.updateTable("ci_runs")
.set({ status: "pending" })
.where("id", "=", runId)
.execute();
}
function pumpQueue(): void {
while (
queuedRunIds.length > 0 &&
runningTasks.size < config.CI_MAX_CONCURRENT
) {
const next = queuedRunIds.shift()!;
const controller = new AbortController();
runningTasks.set(next, { controller });
// Promote queued→pending before spawning so executeRun's failure path
// (which only matches pending/running) can still mark it failed if
// it throws very early. We await the flip inside spawnRun's wrapper
// so the order is: row=pending → executeRun starts → row=running.
spawnRun(next, controller.signal, /* needsPromote */ true);
}
}
/** Position (1-based) of this queued run within the queue, or null. */
export function ciQueuePosition(runId: number): number | null {
const idx = queuedRunIds.indexOf(runId);
return idx < 0 ? null : idx + 1;
}
// --- TOML Parsing ---
export function parseCiConfig(tomlStr: string): CiConfig | null {
@@ -515,7 +553,6 @@ async function collectArtifacts(
containerId: string,
step: CiStep,
_shell: string[],
workDir: string | undefined,
envVars: string[],
): Promise<void> {
const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
@@ -545,62 +582,43 @@ async function collectArtifacts(
}
}
// Archive commands run with `path.dirname(srcPath)` as the working
// directory and reference the source by its basename only, so the
// user-controlled path never appears as part of an interpolated shell
// string. Previously `publish_zip` used `sh -c "cd … && zip …"` with
// raw interpolation — a step author who could write the CI TOML
// could shell-inject through the source path. Today the only TOML
// author is the admin, but this removes the implicit assumption.
type ArchiveType = "tar" | "gzip" | "zip" | "zstd";
const archiveFormats: Array<{
type: ArchiveType;
paths: string[];
ext: string;
cmd: (src: string, dst: string) => string[];
cmd: (basename: string, dst: string) => string[];
}> = [
{
type: "tar",
paths: toArray(step.publish_tar),
ext: ".tar",
cmd: (src, dst) => [
"tar",
"-cf",
dst,
"-C",
path.dirname(src),
path.basename(src),
],
cmd: (basename, dst) => ["tar", "-cf", dst, basename],
},
{
type: "gzip",
paths: toArray(step.publish_gzip),
ext: ".tar.gz",
cmd: (src, dst) => [
"tar",
"-czf",
dst,
"-C",
path.dirname(src),
path.basename(src),
],
cmd: (basename, dst) => ["tar", "-czf", dst, basename],
},
{
type: "zstd",
paths: toArray(step.publish_zstd),
ext: ".tar.zst",
cmd: (src, dst) => [
"tar",
"--zstd",
"-cf",
dst,
"-C",
path.dirname(src),
path.basename(src),
],
cmd: (basename, dst) => ["tar", "--zstd", "-cf", dst, basename],
},
{
type: "zip",
paths: toArray(step.publish_zip),
ext: ".zip",
cmd: (src, dst) => [
"sh",
"-c",
`cd ${path.dirname(src)} && zip -r ${dst} ${path.basename(src)}`,
],
cmd: (basename, dst) => ["zip", "-r", dst, basename],
},
];
@@ -609,11 +627,13 @@ async function collectArtifacts(
for (const srcPath of archivePaths) {
archiveIndex++;
const tmpPath = `/tmp/hf-artifact-${runId}-${archiveIndex}${ext}`;
// Create archive inside container
// Run with the source's parent directory as the working
// directory so each archive tool can reference the source
// by its basename — no -C, no shell.
const execResult = await execInContainer(
containerId,
cmd(srcPath, tmpPath),
workDir,
cmd(path.basename(srcPath), tmpPath),
path.dirname(srcPath),
envVars,
).catch(() => null);
if (!execResult || execResult.exitCode !== 0) continue;
@@ -861,7 +881,6 @@ async function executeRun(runId: number, signal: AbortSignal): Promise<void> {
containerId,
step,
shell,
cfg.work_dir,
envArray,
).catch(() => {});
}
@@ -959,6 +978,8 @@ async function executeRun(runId: number, signal: AbortSignal): Promise<void> {
.where("id", "=", runId)
.executeTakeFirst();
if (run) await pruneHistory(run.repo_id).catch(() => {});
// A slot just freed up — start the next queued run if any.
pumpQueue();
}
}
@@ -1003,11 +1024,22 @@ export async function triggerRun(
.where("id", "=", runId.id)
.execute();
const controller = new AbortController();
runningTasks.set(runId.id, { controller });
// Fire and forget — like release archiving
spawnRun(runId.id, controller.signal);
// Flip to "queued" BEFORE pushing onto queuedRunIds. Otherwise a
// concurrent pumpQueue (from a finishing run) could observe our entry,
// promoteToPending it, and start executeRun while our UPDATE is still
// in flight — the late UPDATE would then clobber a running row back to
// "queued". Only pumpQueue (the sole writer of runningTasks.set) may
// mutate the row's status after the push. The TOCTOU concern is moot
// here because pumpQueue cannot pump a runId that hasn't been pushed.
if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
await db
.updateTable("ci_runs")
.set({ status: "queued" })
.where("id", "=", runId.id)
.execute();
}
queuedRunIds.push(runId.id);
pumpQueue();
return runId.id;
}
@@ -1015,24 +1047,40 @@ export async function triggerRun(
// Wraps the fire-and-forget executeRun so unhandled exceptions (e.g. a throw
// before/after its own try/finally) are logged and the run is reconciled to
// "failure" instead of staying pending forever.
function spawnRun(runId: number, signal: AbortSignal): void {
void executeRun(runId, signal).catch(async (err) => {
console.error(`[ci] executeRun threw for run ${runId}:`, err);
runningTasks.delete(runId);
function spawnRun(
runId: number,
signal: AbortSignal,
needsPromote = false,
): void {
void (async () => {
try {
await db
.updateTable("ci_runs")
.set({
status: "failure",
finished_at: new Date().toISOString(),
})
.where("id", "=", runId)
.where("status", "in", ["pending", "running"])
.execute();
} catch (dbErr) {
console.error(`[ci] failed to mark run ${runId} failed:`, dbErr);
if (needsPromote) await promoteToPending(runId);
await executeRun(runId, signal);
} catch (err) {
console.error(`[ci] executeRun threw for run ${runId}:`, err);
runningTasks.delete(runId);
try {
await db
.updateTable("ci_runs")
.set({
status: "failure",
finished_at: new Date().toISOString(),
})
.where("id", "=", runId)
// Include "queued" so a row that was promoted-but-not-yet-
// observed (or never promoted because promoteToPending threw)
// still gets marked failed instead of stuck.
.where("status", "in", ["pending", "running", "queued"])
.execute();
} catch (dbErr) {
console.error(
`[ci] failed to mark run ${runId} failed:`,
dbErr,
);
}
pumpQueue();
}
});
})();
}
export async function retryRun(
@@ -1056,7 +1104,12 @@ export async function retryRun(
}
await db.deleteFrom("ci_artifacts").where("run_id", "=", runId).execute();
// Reset run
// Reset to "pending" first; if a slot isn't free, flip to "queued"
// before enqueueing. Mirrors triggerRun: pumpQueue is the sole writer
// of runningTasks.set, so we never reserve a slot directly here. Two
// concurrent retryRun calls (or a retryRun racing triggerRun) all
// funnel through pumpQueue, which serializes the size check against
// slot reservation in a single synchronous turn.
await db
.updateTable("ci_runs")
.set({
@@ -1068,10 +1121,15 @@ export async function retryRun(
.where("id", "=", runId)
.execute();
const controller = new AbortController();
runningTasks.set(runId, { controller });
spawnRun(runId, controller.signal);
if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
await db
.updateTable("ci_runs")
.set({ status: "queued" })
.where("id", "=", runId)
.execute();
}
queuedRunIds.push(runId);
pumpQueue();
}
export async function cancelRun(runId: number): Promise<void> {
@@ -1083,11 +1141,13 @@ export async function cancelRun(runId: number): Promise<void> {
await removeContainer(containerId).catch(() => {});
}
}
const queueIdx = queuedRunIds.indexOf(runId);
if (queueIdx >= 0) queuedRunIds.splice(queueIdx, 1);
await db
.updateTable("ci_runs")
.set({ status: "cancelled", finished_at: new Date().toISOString() })
.where("id", "=", runId)
.where("status", "in", ["pending", "running"])
.where("status", "in", ["pending", "running", "queued"])
.execute();
}
@@ -1142,7 +1202,7 @@ export async function cancelStaleRuns(): Promise<void> {
const stale = await db
.selectFrom("ci_runs")
.select("id")
.where("status", "in", ["pending", "running"])
.where("status", "in", ["pending", "running", "queued"])
.execute();
await Promise.allSettled(
@@ -1156,7 +1216,7 @@ export async function cancelStaleRuns(): Promise<void> {
await db
.updateTable("ci_runs")
.set({ status: "cancelled", finished_at: now })
.where("status", "in", ["pending", "running"])
.where("status", "in", ["pending", "running", "queued"])
.execute();
await db
.updateTable("ci_steps")
▾Msrc/services/sshServer.ts
@@ -6,7 +6,88 @@ import { Server, utils } from "ssh2";
import config from "../config.ts";
import { ADMIN_USERNAME, paths } from "../constants.ts";
import { db } from "../db/index.ts";
import { invalidateRefCache } from "./git.ts";
import {
parseCiConfig,
shouldTriggerPush,
shouldTriggerTag,
triggerRun,
} from "./ci.ts";
import { git, invalidateRefCache } from "./git.ts";
/**
* Parse pkt-line ref updates from the leading bytes of a receive-pack
* upload. Returns oldSha/newSha/refname triples. Mirrors `parseRefUpdates`
* in `routes/git.ts` — kept duplicated to avoid coupling the route file
* to the SSH path. Both only ever look at the first ~4 KB.
*/
function parsePktLineRefUpdates(
text: string,
): Array<{ oldSha: string; newSha: string; refname: string }> {
const refs: Array<{ oldSha: string; newSha: string; refname: string }> = [];
let pos = 0;
while (pos + 4 <= text.length) {
const lenStr = text.slice(pos, pos + 4);
const len = parseInt(lenStr, 16);
if (Number.isNaN(len) || len === 0) break;
if (len < 4 || pos + len > text.length) break;
const line = text
.slice(pos + 4, pos + len)
.replace(/\0.*$/, "")
.trim();
pos += len;
const parts = line.split(" ");
if (parts.length >= 3) {
const oldSha = parts[0] ?? "";
const newSha = parts[1] ?? "";
const refname = parts[2] ?? "";
if (refname) refs.push({ oldSha, newSha, refname });
}
}
return refs;
}
async function triggerCiForPush(
repoName: string,
refUpdates: Array<{ oldSha: string; newSha: string; refname: string }>,
): Promise<void> {
for (const { newSha, refname } of refUpdates) {
if (/^0+$/.test(newSha)) continue;
const isBranch = refname.startsWith("refs/heads/");
const isTag = refname.startsWith("refs/tags/");
if (!isBranch && !isTag) continue;
const tomlBuf = await git
.show(repoName, newSha, ".hearthforge-ci.toml")
.catch(() => null);
if (!tomlBuf) continue;
const cfg = parseCiConfig(tomlBuf.toString("utf-8"));
if (!cfg) continue;
if (isBranch) {
const branch = refname.slice("refs/heads/".length);
if (shouldTriggerPush(cfg, branch)) {
triggerRun(repoName, {
triggerSource: "push",
commitSha: newSha,
commitBranch: branch,
}).catch((e) =>
console.error(`CI push trigger failed for ${repoName}:`, e),
);
}
} else if (isTag && shouldTriggerTag(cfg)) {
const tag = refname.slice("refs/tags/".length);
triggerRun(repoName, {
triggerSource: "tag",
commitSha: newSha,
commitTag: tag,
}).catch((e) =>
console.error(`CI tag trigger failed for ${repoName}:`, e),
);
}
}
}
/** Compute SHA256 fingerprint from raw SSH public key bytes (the wire-format bytes). */
function fingerprintFromBytes(keyBytes: Buffer): string {
@@ -58,6 +139,7 @@ export async function startSshServer() {
"ssh_keys.public_key",
])
.where("ssh_keys.fingerprint", "=", fingerprint)
.where("users.is_pending", "=", 0)
.executeTakeFirst();
if (!sshKey) return ctx.reject();
@@ -113,7 +195,14 @@ export async function startSshServer() {
}
}
if (repo.is_private && !authedUser) {
if (
repo.is_private &&
authedUser?.username !== ADMIN_USERNAME
) {
// Match the UI and HTTP smart-git: private repos
// are admin-only. Without this, any user with a
// registered SSH key could clone repos hidden from
// them in the web UI.
stream.stderr.write(
"error: repository access denied\n",
);
@@ -123,6 +212,24 @@ export async function startSshServer() {
}
const proc: ChildProcess = spawn(command, [repoPath]);
// For receive-pack, capture the first ~4 KB of pkt-line
// data so we can extract the ref updates after git
// finishes (mirroring the HTTP path in routes/git.ts).
// CI was previously not triggered on SSH pushes at all.
let preamble: Buffer | null = null;
if (command === "git-receive-pack") {
preamble = Buffer.alloc(0);
stream.on("data", (chunk: Buffer) => {
if (preamble && preamble.length < 4096) {
preamble = Buffer.concat([
preamble,
chunk.subarray(0, 4096 - preamble.length),
]);
}
});
}
stream.pipe(proc.stdin!);
proc.stdout?.pipe(stream, { end: false });
proc.stderr?.pipe(stream.stderr as NodeJS.WritableStream, {
@@ -132,6 +239,14 @@ export async function startSshServer() {
proc.on("close", (code: number | null) => {
if (command === "git-receive-pack") {
invalidateRefCache(repo.name);
if (code === 0 && preamble) {
const refUpdates = parsePktLineRefUpdates(
preamble.toString("utf-8"),
);
triggerCiForPush(repo.name, refUpdates).catch(
() => {},
);
}
}
stream.exit(code ?? 0);
stream.end();
▾Msrc/styles/components.css
@@ -1942,6 +1942,10 @@
background: var(--color-border);
color: var(--color-text-muted);
}
.ci-status-queued {
background: #6e7781;
color: #fff;
}
.ci-status-running {
background: #0b5cab;
color: #fff;
▾Msrc/views/ci/CiHistory.tsx
@@ -20,6 +20,7 @@ interface RunSummary {
created_at: string;
triggered_by_username: string | null;
artifact_count: number;
queue_position: number | null;
}
interface CiHistoryProps {
@@ -30,6 +31,12 @@ interface CiHistoryProps {
manualTriggerDisabledReason: string | null;
}
function queueTitle(position: number | null): string {
if (position === null) return "Waiting in the build queue";
if (position === 1) return "Waiting in the build queue — next up";
return `Waiting in the build queue — ${position - 1} run${position - 1 === 1 ? "" : "s"} ahead`;
}
function duration(start: string | null, end: string | null): string {
if (!start || !end) return "";
const ms = new Date(end).getTime() - new Date(start).getTime();
@@ -131,7 +138,10 @@ export function CiHistory({
manualTriggerDisabledReason,
}: CiHistoryProps) {
const isRunning = runs.some(
(r) => r.status === "pending" || r.status === "running",
(r) =>
r.status === "pending" ||
r.status === "running" ||
r.status === "queued",
);
return (
<Layout user={user} title={`Pipelines — ${repo.name}`}>
@@ -203,7 +213,16 @@ export function CiHistory({
href={`/${repo.name}/ci/${run.id}`}
class="release-item-title"
>
<CiStatusPill status={run.status} />
<CiStatusPill
status={run.status}
title={
run.status === "queued"
? queueTitle(
run.queue_position,
)
: undefined
}
/>
<span class="ci-run-id">
#{run.repo_run_id ?? run.id}
</span>
▾Msrc/views/ci/CiRunDetail.tsx
@@ -32,6 +32,13 @@ interface CiRunDetailProps {
steps: CiStepRow[];
artifacts: CiArtifactRow[];
autoRefresh: boolean;
queuePosition: number | null;
}
function queueTitle(position: number | null): string {
if (position === null) return "Waiting in the build queue";
if (position === 1) return "Waiting in the build queue — next up";
return `Waiting in the build queue — ${position - 1} run${position - 1 === 1 ? "" : "s"} ahead`;
}
function duration(start: string | null, end: string | null): string {
@@ -58,8 +65,14 @@ export function CiRunDetail({
steps,
artifacts,
autoRefresh,
queuePosition,
}: CiRunDetailProps) {
const isActive = run.status === "pending" || run.status === "running";
const isActive =
run.status === "pending" ||
run.status === "running" ||
run.status === "queued";
const isQueued = run.status === "queued";
const queueText = queueTitle(queuePosition);
const displayId = run.repo_run_id ?? run.id;
let variableOverrides: Record<string, string> = {};
@@ -88,8 +101,11 @@ export function CiRunDetail({
<div class="release-detail-header">
<div>
<h2 class="release-detail-title">
<CiStatusPill status={run.status} /> Pipeline #
{displayId}
<CiStatusPill
status={run.status}
title={isQueued ? queueText : undefined}
/>{" "}
Pipeline #{displayId}
</h2>
<div class="release-item-meta">
{run.commit_sha && (
@@ -191,7 +207,9 @@ export function CiRunDetail({
<h3 class="section-title">Steps</h3>
{steps.length === 0 ? (
<div class="ci-step-pending">
<span class="text-muted">Waiting to start…</span>
<span class="text-muted">
{isQueued ? queueText : "Waiting to start…"}
</span>
</div>
) : (
steps.map((step) => (
▾Msrc/views/ci/CiStatusPill.tsx
@@ -1,5 +1,6 @@
const statusStyles: Record<string, string> = {
pending: "ci-status-pending",
queued: "ci-status-queued",
running: "ci-status-running",
success: "ci-status-success",
failure: "ci-status-failure",
@@ -7,7 +8,16 @@ const statusStyles: Record<string, string> = {
skipped: "ci-status-skipped",
};
export function CiStatusPill({ status }: { status: string }) {
interface CiStatusPillProps {
status: string;
title?: string;
}
export function CiStatusPill({ status, title }: CiStatusPillProps) {
const cls = statusStyles[status] ?? "ci-status-pending";
return <span class={`ci-status-pill ${cls}`}>{status}</span>;
return (
<span class={`ci-status-pill ${cls}`} title={title}>
{status}
</span>
);
}
▾Msrc/views/repos/CommitDetail.tsx
@@ -15,6 +15,7 @@ interface CommitDetailProps {
sha: string;
meta: CommitMeta;
files: RenderedDiffFile[];
tooLarge?: number;
}
export function CommitDetail({
@@ -23,6 +24,7 @@ export function CommitDetail({
sha,
meta,
files,
tooLarge,
}: CommitDetailProps) {
return (
<Layout user={user} title={`${sha.slice(0, 7)} — ${repo.name}`}>
@@ -141,7 +143,17 @@ export function CommitDetail({
</div>
</div>
<DiffView files={files} repo={repo} sha={sha} />
{tooLarge !== undefined ? (
<div class="file-download-notice">
<p>
Diff is too large to render inline (
{(tooLarge / 1024 / 1024).toFixed(1)} MB). Browse
individual files at the tree below.
</p>
</div>
) : (
<DiffView files={files} repo={repo} sha={sha} />
)}
</div>
</Layout>
);