updated
This commit is contained in:
+50
@@ -0,0 +1,50 @@
|
||||
// Wires config + logger into the concrete dependencies (Gitea client,
|
||||
// Claude runner, job runner, queue) and returns an HTTP server ready to
|
||||
// listen. Kept separate from index.js so tests can build an app with fake
|
||||
// dependencies without touching process.env or opening real sockets.
|
||||
|
||||
import { GiteaClient } from './gitea.js';
|
||||
import { ClaudeRunner } from './claude.js';
|
||||
import { JobRunner } from './job.js';
|
||||
import { JobQueue } from './jobQueue.js';
|
||||
import { createServer } from './server.js';
|
||||
|
||||
export function buildApp(config, logger) {
|
||||
const gitea = new GiteaClient({
|
||||
baseUrl: config.giteaUrl,
|
||||
token: config.giteaToken,
|
||||
logger: logger.child({ component: 'gitea' }),
|
||||
});
|
||||
|
||||
const claude = new ClaudeRunner({
|
||||
oauthToken: config.claudeOauthToken,
|
||||
model: config.model,
|
||||
maxTurns: config.maxTurns,
|
||||
timeoutMs: config.jobTimeoutMs,
|
||||
logger: logger.child({ component: 'claude' }),
|
||||
});
|
||||
|
||||
const jobRunner = new JobRunner({
|
||||
gitea,
|
||||
claude,
|
||||
config,
|
||||
logger: logger.child({ component: 'job' }),
|
||||
});
|
||||
|
||||
const queue = new JobQueue(logger.child({ component: 'queue' }));
|
||||
|
||||
const server = createServer({
|
||||
config,
|
||||
logger: logger.child({ component: 'server' }),
|
||||
onTrigger: (trigger, log) => {
|
||||
const key = `${trigger.repo}#${trigger.issueNumber}`;
|
||||
queue.runExclusive(key, () => jobRunner.run(trigger)).catch((err) => {
|
||||
// runExclusive already catches job errors internally; this is a
|
||||
// final safety net so a bug there can never crash the process.
|
||||
log.error('unexpected error scheduling job', { error: err });
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
return { server, gitea, claude, jobRunner, queue };
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
// Wraps invoking the `claude` CLI as a single-shot, non-interactive agent
|
||||
// run. The OAuth token is passed through env only (never argv), so it
|
||||
// never appears in process listings or in logged command lines.
|
||||
|
||||
import { run } from './exec.js';
|
||||
|
||||
export class ClaudeRunner {
|
||||
constructor({ oauthToken, model, maxTurns, timeoutMs, logger }) {
|
||||
this.oauthToken = oauthToken;
|
||||
this.model = model;
|
||||
this.maxTurns = maxTurns;
|
||||
this.timeoutMs = timeoutMs;
|
||||
this.logger = logger;
|
||||
}
|
||||
|
||||
/** Runs Claude Code against `cwd` with `prompt`. Returns trimmed stdout. */
|
||||
async runOnce(cwd, prompt) {
|
||||
const env = { ...process.env, CLAUDE_CODE_OAUTH_TOKEN: this.oauthToken };
|
||||
const args = [
|
||||
'-p', prompt,
|
||||
'--permission-mode', 'acceptEdits',
|
||||
'--model', this.model,
|
||||
'--max-turns', String(this.maxTurns),
|
||||
];
|
||||
const start = Date.now();
|
||||
const { stdout, stderr } = await run('claude', args, { cwd, env, timeoutMs: this.timeoutMs });
|
||||
this.logger?.info('claude run completed', {
|
||||
durationMs: Date.now() - start,
|
||||
stderrPreview: stderr ? stderr.slice(0, 500) : undefined,
|
||||
});
|
||||
return stdout.trim();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
// All configuration is read from the environment once, at startup, and
|
||||
// validated eagerly so a misconfiguration fails fast with a clear message
|
||||
// instead of surfacing later as a confusing runtime error.
|
||||
|
||||
const REQUIRED = ['GITEA_URL', 'GITEA_TOKEN', 'CLAUDE_CODE_OAUTH_TOKEN', 'OWNER_LOGIN'];
|
||||
|
||||
/**
|
||||
* Builds and validates config from an env-like object. Exported so tests
|
||||
* can pass a fake env instead of mutating process.env.
|
||||
*/
|
||||
export function loadConfig(env = process.env) {
|
||||
const missing = REQUIRED.filter((key) => !env[key]);
|
||||
if (missing.length > 0) {
|
||||
throw new Error(`Missing required environment variable(s): ${missing.join(', ')}`);
|
||||
}
|
||||
|
||||
const giteaUrl = env.GITEA_URL.replace(/\/+$/, '');
|
||||
let parsedUrl;
|
||||
try {
|
||||
parsedUrl = new URL(giteaUrl);
|
||||
} catch {
|
||||
throw new Error(`GITEA_URL is not a valid URL: ${env.GITEA_URL}`);
|
||||
}
|
||||
if (!['http:', 'https:'].includes(parsedUrl.protocol)) {
|
||||
throw new Error(`GITEA_URL must be http or https: ${env.GITEA_URL}`);
|
||||
}
|
||||
|
||||
const port = Number.parseInt(env.PORT ?? '3000', 10);
|
||||
if (!Number.isInteger(port) || port <= 0 || port > 65535) {
|
||||
throw new Error(`PORT must be a valid port number, got: ${env.PORT}`);
|
||||
}
|
||||
|
||||
const jobTimeoutMs = Number.parseInt(env.JOB_TIMEOUT_MS ?? '', 10) || 20 * 60_000;
|
||||
const maxTurns = Number.parseInt(env.CLAUDE_MAX_TURNS ?? '', 10) || 25;
|
||||
const model = env.CLAUDE_MODEL || 'sonnet';
|
||||
const webhookSecret = env.WEBHOOK_SECRET || null;
|
||||
|
||||
return {
|
||||
giteaUrl,
|
||||
giteaToken: env.GITEA_TOKEN,
|
||||
claudeOauthToken: env.CLAUDE_CODE_OAUTH_TOKEN,
|
||||
ownerLogin: env.OWNER_LOGIN,
|
||||
port,
|
||||
jobTimeoutMs,
|
||||
maxTurns,
|
||||
model,
|
||||
webhookSecret,
|
||||
triggerPhrase: env.TRIGGER_PHRASE || '@claude',
|
||||
logLevel: (env.LOG_LEVEL || 'info').toLowerCase(),
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
// Pure logic for deciding whether an incoming Gitea webhook payload should
|
||||
// trigger a Claude run, and for extracting the fields the rest of the
|
||||
// service needs. Kept free of I/O so it's cheap to unit test exhaustively.
|
||||
|
||||
/**
|
||||
* @typedef {object} Trigger
|
||||
* @property {string} repo full_name, e.g. "owner/repo"
|
||||
* @property {number} issueNumber
|
||||
* @property {string} issueTitle
|
||||
* @property {string} issueBody
|
||||
* @property {boolean} isPullRequest
|
||||
* @property {string} defaultBranch
|
||||
* @property {string} requestText the comment or issue body containing the trigger phrase
|
||||
* @property {string} author
|
||||
*/
|
||||
|
||||
/**
|
||||
* Decides whether a webhook event should trigger a run and, if so, extracts
|
||||
* a Trigger. Returns { trigger: null, reason } when it should be ignored.
|
||||
*/
|
||||
export function classifyEvent(eventName, payload, { ownerLogin, triggerPhrase }) {
|
||||
if (!payload || typeof payload !== 'object') {
|
||||
return { trigger: null, reason: 'empty or invalid payload' };
|
||||
}
|
||||
|
||||
let author;
|
||||
let text;
|
||||
|
||||
if (eventName === 'issue_comment' && payload.action === 'created') {
|
||||
author = payload.comment?.user?.login;
|
||||
text = payload.comment?.body;
|
||||
} else if (eventName === 'issues' && ['opened', 'edited'].includes(payload.action)) {
|
||||
author = payload.issue?.user?.login;
|
||||
text = payload.issue?.body;
|
||||
} else {
|
||||
return { trigger: null, reason: `event "${eventName}" action "${payload.action}" is not handled` };
|
||||
}
|
||||
|
||||
if (author !== ownerLogin) {
|
||||
return { trigger: null, reason: `author "${author}" is not the configured owner "${ownerLogin}"` };
|
||||
}
|
||||
if (!text || !text.includes(triggerPhrase)) {
|
||||
return { trigger: null, reason: `text does not contain trigger phrase "${triggerPhrase}"` };
|
||||
}
|
||||
if (!payload.repository?.full_name || !payload.issue?.number) {
|
||||
return { trigger: null, reason: 'payload is missing repository or issue information' };
|
||||
}
|
||||
|
||||
return {
|
||||
trigger: {
|
||||
repo: payload.repository.full_name,
|
||||
issueNumber: payload.issue.number,
|
||||
issueTitle: payload.issue.title ?? '',
|
||||
issueBody: payload.issue.body ?? '',
|
||||
isPullRequest: Boolean(payload.issue.pull_request),
|
||||
defaultBranch: payload.repository.default_branch || 'main',
|
||||
requestText: text,
|
||||
author,
|
||||
},
|
||||
reason: null,
|
||||
};
|
||||
}
|
||||
+84
@@ -0,0 +1,84 @@
|
||||
// Thin wrapper around child_process, isolated in its own module so tests
|
||||
// can inject a fake runner instead of spawning real processes.
|
||||
|
||||
import { spawn } from 'node:child_process';
|
||||
|
||||
export class ProcessError extends Error {
|
||||
constructor(command, args, code, signal, stdout, stderr) {
|
||||
super(`${command} ${args.join(' ')} exited with code ${code}${signal ? ` (signal ${signal})` : ''}`);
|
||||
this.name = 'ProcessError';
|
||||
this.command = command;
|
||||
this.args = args;
|
||||
this.code = code;
|
||||
this.signal = signal;
|
||||
this.stdout = stdout;
|
||||
this.stderr = stderr;
|
||||
}
|
||||
}
|
||||
|
||||
export class ProcessTimeoutError extends Error {
|
||||
constructor(command, args, timeoutMs) {
|
||||
super(`${command} ${args.join(' ')} timed out after ${timeoutMs}ms`);
|
||||
this.name = 'ProcessTimeoutError';
|
||||
this.command = command;
|
||||
this.args = args;
|
||||
this.timeoutMs = timeoutMs;
|
||||
}
|
||||
}
|
||||
|
||||
const MAX_BUFFER = 20 * 1024 * 1024;
|
||||
|
||||
/**
|
||||
* Runs a command to completion and resolves with its output, or rejects
|
||||
* with ProcessError / ProcessTimeoutError. Truncates captured output to
|
||||
* MAX_BUFFER bytes rather than buffering without limit.
|
||||
*/
|
||||
export function run(command, args, options = {}) {
|
||||
return new Promise((resolve, reject) => {
|
||||
const child = spawn(command, args, {
|
||||
cwd: options.cwd,
|
||||
env: options.env,
|
||||
stdio: ['ignore', 'pipe', 'pipe'],
|
||||
});
|
||||
|
||||
let stdout = '';
|
||||
let stderr = '';
|
||||
let stdoutBytes = 0;
|
||||
let stderrBytes = 0;
|
||||
let settled = false;
|
||||
let timer = null;
|
||||
|
||||
const finish = (fn, value) => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
if (timer) clearTimeout(timer);
|
||||
fn(value);
|
||||
};
|
||||
|
||||
child.stdout.on('data', (chunk) => {
|
||||
stdoutBytes += chunk.length;
|
||||
if (stdoutBytes <= MAX_BUFFER) stdout += chunk;
|
||||
});
|
||||
child.stderr.on('data', (chunk) => {
|
||||
stderrBytes += chunk.length;
|
||||
if (stderrBytes <= MAX_BUFFER) stderr += chunk;
|
||||
});
|
||||
|
||||
child.on('error', (err) => finish(reject, err));
|
||||
|
||||
child.on('close', (code, signal) => {
|
||||
if (code === 0) {
|
||||
finish(resolve, { stdout, stderr });
|
||||
} else {
|
||||
finish(reject, new ProcessError(command, args, code, signal, stdout, stderr));
|
||||
}
|
||||
});
|
||||
|
||||
if (options.timeoutMs) {
|
||||
timer = setTimeout(() => {
|
||||
child.kill('SIGKILL');
|
||||
finish(reject, new ProcessTimeoutError(command, args, options.timeoutMs));
|
||||
}, options.timeoutMs);
|
||||
}
|
||||
});
|
||||
}
|
||||
+62
@@ -0,0 +1,62 @@
|
||||
// Git operations for a single job's working directory. The token is passed
|
||||
// via a git http.extraHeader env override rather than embedded in the
|
||||
// clone URL, so it never appears in argv (visible via `ps aux`) or in any
|
||||
// URL that gets logged.
|
||||
|
||||
import { run } from './exec.js';
|
||||
|
||||
export function gitAuthEnv(token, baseEnv = process.env) {
|
||||
return {
|
||||
...baseEnv,
|
||||
GIT_TERMINAL_PROMPT: '0',
|
||||
GIT_CONFIG_COUNT: '1',
|
||||
GIT_CONFIG_KEY_0: 'http.extraHeader',
|
||||
GIT_CONFIG_VALUE_0: `Authorization: token ${token}`,
|
||||
};
|
||||
}
|
||||
|
||||
export class GitRepo {
|
||||
constructor({ dir, giteaUrl, token, logger }) {
|
||||
this.dir = dir;
|
||||
this.giteaUrl = giteaUrl.replace(/\/+$/, '');
|
||||
this.env = gitAuthEnv(token);
|
||||
this.logger = logger;
|
||||
}
|
||||
|
||||
async #git(...args) {
|
||||
this.logger?.debug('git', { args });
|
||||
return run('git', args, { cwd: this.dir, env: this.env });
|
||||
}
|
||||
|
||||
async clone(repo, ref) {
|
||||
const url = `${this.giteaUrl}/${repo}.git`;
|
||||
await run('git', ['clone', '--depth', '50', '--branch', ref, url, '.'], {
|
||||
cwd: this.dir,
|
||||
env: this.env,
|
||||
});
|
||||
}
|
||||
|
||||
/** Checks out `branch`, resuming it from the remote if it already exists. */
|
||||
async checkoutWorkBranch(branch) {
|
||||
try {
|
||||
await this.#git('fetch', 'origin', branch);
|
||||
await this.#git('checkout', '-B', branch, 'FETCH_HEAD');
|
||||
} catch {
|
||||
await this.#git('checkout', '-B', branch);
|
||||
}
|
||||
}
|
||||
|
||||
async hasChanges() {
|
||||
const { stdout } = await this.#git('status', '--porcelain');
|
||||
return stdout.trim().length > 0;
|
||||
}
|
||||
|
||||
async commitAll(message, { name = 'claude', email = '[email protected]' } = {}) {
|
||||
await this.#git('add', '-A');
|
||||
await this.#git('-c', `user.name=${name}`, '-c', `user.email=${email}`, 'commit', '-m', message);
|
||||
}
|
||||
|
||||
async push(branch) {
|
||||
await this.#git('push', 'origin', `HEAD:${branch}`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
// A minimal Gitea API client covering exactly the endpoints this service
|
||||
// needs. Kept dependency-free (global fetch, available since Node 18) and
|
||||
// injectable (accepts a custom `fetchImpl` for tests).
|
||||
|
||||
export class GiteaApiError extends Error {
|
||||
constructor(method, path, status, body) {
|
||||
super(`Gitea API ${method} ${path} -> ${status}: ${body.slice(0, 300)}`);
|
||||
this.name = 'GiteaApiError';
|
||||
this.method = method;
|
||||
this.path = path;
|
||||
this.status = status;
|
||||
this.body = body;
|
||||
}
|
||||
}
|
||||
|
||||
export class GiteaClient {
|
||||
constructor({ baseUrl, token, logger, fetchImpl = fetch }) {
|
||||
this.baseUrl = baseUrl.replace(/\/+$/, '');
|
||||
this.token = token;
|
||||
this.logger = logger;
|
||||
this.fetchImpl = fetchImpl;
|
||||
}
|
||||
|
||||
async #request(method, path, body) {
|
||||
const url = `${this.baseUrl}/api/v1${path}`;
|
||||
const res = await this.fetchImpl(url, {
|
||||
method,
|
||||
headers: {
|
||||
Authorization: `token ${this.token}`,
|
||||
'Content-Type': 'application/json',
|
||||
Accept: 'application/json',
|
||||
},
|
||||
body: body !== undefined ? JSON.stringify(body) : undefined,
|
||||
});
|
||||
this.logger?.debug('gitea api call', { method, path, status: res.status });
|
||||
if (!res.ok) {
|
||||
const text = await res.text().catch(() => '');
|
||||
throw new GiteaApiError(method, path, res.status, text);
|
||||
}
|
||||
if (res.status === 204) return null;
|
||||
const text = await res.text();
|
||||
return text ? JSON.parse(text) : null;
|
||||
}
|
||||
|
||||
/** Returns the pull request, or null if `n` is an issue, not a PR. */
|
||||
async getPullRequest(repo, n) {
|
||||
try {
|
||||
return await this.#request('GET', `/repos/${repo}/pulls/${n}`);
|
||||
} catch (err) {
|
||||
if (err instanceof GiteaApiError && err.status === 404) return null;
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
async listOpenPullRequests(repo) {
|
||||
return this.#request('GET', `/repos/${repo}/pulls?state=open`);
|
||||
}
|
||||
|
||||
async createPullRequest(repo, { title, head, base, body }) {
|
||||
return this.#request('POST', `/repos/${repo}/pulls`, { title, head, base, body });
|
||||
}
|
||||
|
||||
async createIssueComment(repo, issueNumber, body) {
|
||||
return this.#request('POST', `/repos/${repo}/issues/${issueNumber}/comments`, { body });
|
||||
}
|
||||
|
||||
/** Repo clone URL with the token embedded, for git-over-http auth. */
|
||||
authenticatedCloneUrl(repo) {
|
||||
const u = new URL(`${this.baseUrl}/${repo}.git`);
|
||||
u.username = 'oauth2';
|
||||
u.password = this.token;
|
||||
return u.toString();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
#!/usr/bin/env node
|
||||
// Process entrypoint: load config, build the app, start listening, and
|
||||
// handle shutdown signals so an in-flight job's logs aren't cut off by an
|
||||
// abrupt container stop.
|
||||
|
||||
import { loadConfig } from './config.js';
|
||||
import { createLogger } from './logger.js';
|
||||
import { buildApp } from './app.js';
|
||||
|
||||
const logger = createLogger({ component: 'claude-hook' });
|
||||
|
||||
let config;
|
||||
try {
|
||||
config = loadConfig();
|
||||
} catch (err) {
|
||||
logger.error('invalid configuration, exiting', { error: err });
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
logger.info('starting', {
|
||||
giteaUrl: config.giteaUrl,
|
||||
ownerLogin: config.ownerLogin,
|
||||
model: config.model,
|
||||
maxTurns: config.maxTurns,
|
||||
jobTimeoutMs: config.jobTimeoutMs,
|
||||
webhookSignatureEnabled: Boolean(config.webhookSecret),
|
||||
port: config.port,
|
||||
});
|
||||
|
||||
const { server, queue } = buildApp(config, logger);
|
||||
|
||||
server.listen(config.port, () => {
|
||||
logger.info('listening', { port: config.port });
|
||||
});
|
||||
|
||||
function shutdown(signal) {
|
||||
logger.info('shutdown signal received', { signal });
|
||||
server.close(() => {
|
||||
logger.info('http server closed');
|
||||
process.exit(0);
|
||||
});
|
||||
// If a job is still running, give it a grace period rather than exiting
|
||||
// immediately underneath it; force-exit afterward so the container
|
||||
// doesn't hang on a stuck job forever.
|
||||
setTimeout(() => {
|
||||
logger.warn('forcing exit after shutdown grace period', {
|
||||
stillRunning: queue.running.size,
|
||||
});
|
||||
process.exit(0);
|
||||
}, 30_000).unref();
|
||||
}
|
||||
|
||||
process.on('SIGTERM', () => shutdown('SIGTERM'));
|
||||
process.on('SIGINT', () => shutdown('SIGINT'));
|
||||
|
||||
process.on('unhandledRejection', (reason) => {
|
||||
logger.error('unhandled promise rejection', { error: reason });
|
||||
});
|
||||
process.on('uncaughtException', (err) => {
|
||||
logger.error('uncaught exception', { error: err });
|
||||
process.exit(1);
|
||||
});
|
||||
+109
@@ -0,0 +1,109 @@
|
||||
// Orchestrates one end-to-end run: determine the branch, check out the
|
||||
// repo, run Claude, commit/push if anything changed, open or reuse a PR,
|
||||
// and always post a reply comment (including on failure).
|
||||
|
||||
import fs from 'node:fs/promises';
|
||||
import os from 'node:os';
|
||||
import path from 'node:path';
|
||||
import { GitRepo } from './git.js';
|
||||
import { buildPrompt } from './prompt.js';
|
||||
|
||||
const MAX_COMMENT_BODY = 60_000; // stay well under Gitea's comment size limit
|
||||
|
||||
function truncate(text, max) {
|
||||
return text.length > max ? text.slice(0, max) + '\n\n[truncated]' : text;
|
||||
}
|
||||
|
||||
export class JobRunner {
|
||||
constructor({ gitea, claude, config, logger, gitFactory = (opts) => new GitRepo(opts) }) {
|
||||
this.gitea = gitea;
|
||||
this.claude = claude;
|
||||
this.config = config;
|
||||
this.logger = logger;
|
||||
this.gitFactory = gitFactory;
|
||||
}
|
||||
|
||||
async run(trigger) {
|
||||
const key = `${trigger.repo}#${trigger.issueNumber}`;
|
||||
const log = this.logger.child({ key });
|
||||
log.info('job started', { isPullRequest: trigger.isPullRequest });
|
||||
|
||||
let workDir;
|
||||
try {
|
||||
workDir = await fs.mkdtemp(path.join(os.tmpdir(), 'claude-hook-'));
|
||||
const reply = await this.#execute(trigger, workDir, log);
|
||||
await this.gitea.createIssueComment(trigger.repo, trigger.issueNumber, truncate(reply, MAX_COMMENT_BODY));
|
||||
log.info('job finished');
|
||||
} catch (err) {
|
||||
log.error('job errored', { error: err });
|
||||
const message = err?.stderr ? String(err.stderr) : err?.message || String(err);
|
||||
await this.gitea
|
||||
.createIssueComment(
|
||||
trigger.repo,
|
||||
trigger.issueNumber,
|
||||
`Claude run failed:\n\n\`\`\`\n${truncate(message, 4000)}\n\`\`\``,
|
||||
)
|
||||
.catch((commentErr) => log.error('could not post failure comment', { error: commentErr }));
|
||||
} finally {
|
||||
if (workDir) {
|
||||
await fs.rm(workDir, { recursive: true, force: true }).catch((rmErr) =>
|
||||
log.warn('failed to clean up work dir', { workDir, error: rmErr }),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async #execute(trigger, workDir, log) {
|
||||
const { repo, issueNumber, defaultBranch, isPullRequest } = trigger;
|
||||
|
||||
let ref = defaultBranch;
|
||||
if (isPullRequest) {
|
||||
const pr = await this.gitea.getPullRequest(repo, issueNumber);
|
||||
if (!pr) throw new Error(`issue #${issueNumber} looked like a PR but the PR lookup returned nothing`);
|
||||
ref = pr.head.ref;
|
||||
}
|
||||
const branch = isPullRequest ? ref : `claude/issue-${issueNumber}`;
|
||||
log.info('resolved target branch', { ref, branch, isPullRequest });
|
||||
|
||||
const git = this.gitFactory({ dir: workDir, giteaUrl: this.gitea.baseUrl, token: this.gitea.token, logger: log });
|
||||
await git.clone(repo, ref);
|
||||
if (!isPullRequest) {
|
||||
await git.checkoutWorkBranch(branch);
|
||||
}
|
||||
|
||||
const prompt = buildPrompt(trigger);
|
||||
log.info('running claude', { promptLength: prompt.length });
|
||||
let reply = await this.claude.runOnce(workDir, prompt);
|
||||
if (!reply) reply = '(Claude produced no textual output.)';
|
||||
|
||||
if (await git.hasChanges()) {
|
||||
log.info('changes detected, committing');
|
||||
await git.commitAll(`Claude: address #${issueNumber}`);
|
||||
await git.push(branch);
|
||||
|
||||
if (!isPullRequest) {
|
||||
const prUrl = await this.#openOrReusePullRequest(trigger, branch, reply);
|
||||
reply += `\n\nPull request: ${prUrl}`;
|
||||
}
|
||||
} else {
|
||||
log.info('no changes made');
|
||||
}
|
||||
|
||||
return reply;
|
||||
}
|
||||
|
||||
async #openOrReusePullRequest(trigger, branch, reply) {
|
||||
const { repo, issueNumber, issueTitle, defaultBranch } = trigger;
|
||||
const open = await this.gitea.listOpenPullRequests(repo);
|
||||
const existing = open.find((pr) => pr.head.ref === branch);
|
||||
if (existing) return existing.html_url;
|
||||
|
||||
const pr = await this.gitea.createPullRequest(repo, {
|
||||
title: `Claude: ${issueTitle}`,
|
||||
head: branch,
|
||||
base: defaultBranch,
|
||||
body: `Closes #${issueNumber}\n\n${truncate(reply, 4000)}`,
|
||||
});
|
||||
return pr.html_url;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
// Prevents two runs for the same issue/PR from overlapping (e.g. a
|
||||
// duplicate webhook delivery, or a fast follow-up comment while a run is
|
||||
// still in progress). Keyed by "repo#issueNumber".
|
||||
|
||||
export class JobQueue {
|
||||
constructor(logger) {
|
||||
this.running = new Set();
|
||||
this.logger = logger;
|
||||
}
|
||||
|
||||
/**
|
||||
* Runs `fn()` for `key` unless a job for that key is already running, in
|
||||
* which case it logs and returns without calling `fn`. Errors thrown by
|
||||
* `fn` are caught and logged here so a bad job never crashes the process
|
||||
* or leaves the key stuck as "running".
|
||||
*/
|
||||
async runExclusive(key, fn) {
|
||||
if (this.running.has(key)) {
|
||||
this.logger.warn('job already running for key, skipping', { key });
|
||||
return;
|
||||
}
|
||||
this.running.add(key);
|
||||
try {
|
||||
await fn();
|
||||
} catch (err) {
|
||||
this.logger.error('job failed', { key, error: err });
|
||||
} finally {
|
||||
this.running.delete(key);
|
||||
}
|
||||
}
|
||||
|
||||
isRunning(key) {
|
||||
return this.running.has(key);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
// Structured JSON logger. No dependencies, stdout only — Coolify/Docker
|
||||
// captures stdout as the log stream, so there is nothing else to configure.
|
||||
// Each line is one JSON object: { time, level, msg, ...fields }.
|
||||
|
||||
const LEVELS = { debug: 10, info: 20, warn: 30, error: 40 };
|
||||
|
||||
function resolveLevel() {
|
||||
const configured = (process.env.LOG_LEVEL || 'info').toLowerCase();
|
||||
return LEVELS[configured] !== undefined ? configured : 'info';
|
||||
}
|
||||
|
||||
const minLevel = LEVELS[resolveLevel()];
|
||||
|
||||
function write(level, msg, fields) {
|
||||
if (LEVELS[level] < minLevel) return;
|
||||
const record = {
|
||||
time: new Date().toISOString(),
|
||||
level,
|
||||
msg,
|
||||
...fields,
|
||||
};
|
||||
const line = JSON.stringify(record, safeReplacer());
|
||||
if (level === 'error' || level === 'warn') {
|
||||
process.stderr.write(line + '\n');
|
||||
} else {
|
||||
process.stdout.write(line + '\n');
|
||||
}
|
||||
}
|
||||
|
||||
// Prevents JSON.stringify from throwing on circular structures (e.g. an
|
||||
// Error with a `cause` cycle, or accidentally logging a raw HTTP object).
|
||||
function safeReplacer() {
|
||||
const seen = new WeakSet();
|
||||
return (_key, value) => {
|
||||
if (value instanceof Error) {
|
||||
return { name: value.name, message: value.message, stack: value.stack };
|
||||
}
|
||||
if (typeof value === 'object' && value !== null) {
|
||||
if (seen.has(value)) return '[Circular]';
|
||||
seen.add(value);
|
||||
}
|
||||
return value;
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a logger bound to a set of fields (e.g. a request id or job key),
|
||||
* so every subsequent call carries that context without repeating it.
|
||||
*/
|
||||
export function createLogger(bindings = {}) {
|
||||
const withFields = (extra) => ({ ...bindings, ...extra });
|
||||
return {
|
||||
debug: (msg, fields) => write('debug', msg, withFields(fields)),
|
||||
info: (msg, fields) => write('info', msg, withFields(fields)),
|
||||
warn: (msg, fields) => write('warn', msg, withFields(fields)),
|
||||
error: (msg, fields) => write('error', msg, withFields(fields)),
|
||||
child: (extra) => createLogger(withFields(extra)),
|
||||
};
|
||||
}
|
||||
|
||||
export const logger = createLogger({ component: 'claude-hook' });
|
||||
@@ -0,0 +1,14 @@
|
||||
// Builds the prompt sent to Claude Code for a triggered job. Isolated so
|
||||
// its wording can be tuned and tested without touching job orchestration.
|
||||
|
||||
export function buildPrompt({ repo, issueNumber, issueTitle, issueBody, requestText }) {
|
||||
return [
|
||||
`Repository: ${repo}, issue #${issueNumber}: ${issueTitle}`,
|
||||
issueBody || '(no description)',
|
||||
`Request from the maintainer:\n${requestText}`,
|
||||
'Make the requested changes by editing files in this checkout. Do not run git commands ' +
|
||||
'(commit, push, branch) — the calling process handles version control. ' +
|
||||
'End your response with a short, plain-text summary of what you changed, ' +
|
||||
'or your answer directly if no code change was needed.',
|
||||
].join('\n\n');
|
||||
}
|
||||
+108
@@ -0,0 +1,108 @@
|
||||
// The webhook HTTP endpoint. Deliberately small: read the body, verify the
|
||||
// signature if configured, classify the event, and hand off to the job
|
||||
// queue. All decisions are logged so a silent drop is always explainable
|
||||
// from the logs alone.
|
||||
|
||||
import http from 'node:http';
|
||||
import crypto from 'node:crypto';
|
||||
import { classifyEvent } from './events.js';
|
||||
|
||||
const MAX_BODY_BYTES = 5 * 1024 * 1024;
|
||||
|
||||
function readBody(req) {
|
||||
return new Promise((resolve, reject) => {
|
||||
const chunks = [];
|
||||
let bytes = 0;
|
||||
req.on('data', (chunk) => {
|
||||
bytes += chunk.length;
|
||||
if (bytes > MAX_BODY_BYTES) {
|
||||
req.destroy();
|
||||
reject(new Error('request body too large'));
|
||||
return;
|
||||
}
|
||||
chunks.push(chunk);
|
||||
});
|
||||
req.on('end', () => resolve(Buffer.concat(chunks)));
|
||||
req.on('error', reject);
|
||||
});
|
||||
}
|
||||
|
||||
function verifySignature(secret, signatureHeader, rawBody) {
|
||||
if (!secret) return true;
|
||||
const expected = crypto.createHmac('sha256', secret).update(rawBody).digest('hex');
|
||||
const given = signatureHeader || '';
|
||||
const a = Buffer.from(expected);
|
||||
const b = Buffer.from(given);
|
||||
return a.length === b.length && crypto.timingSafeEqual(a, b);
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds the HTTP server. `onTrigger(trigger)` is called (fire-and-forget
|
||||
* from the server's point of view — the caller decides how to queue it)
|
||||
* whenever an incoming event should start a job.
|
||||
*/
|
||||
export function createServer({ config, logger, onTrigger }) {
|
||||
return http.createServer(async (req, res) => {
|
||||
const requestId = crypto.randomUUID();
|
||||
const log = logger.child({ requestId });
|
||||
|
||||
if (req.method === 'GET' && req.url === '/healthz') {
|
||||
res.writeHead(200, { 'Content-Type': 'text/plain' }).end('ok');
|
||||
return;
|
||||
}
|
||||
|
||||
if (req.method !== 'POST') {
|
||||
res.writeHead(405).end();
|
||||
return;
|
||||
}
|
||||
|
||||
let rawBody;
|
||||
try {
|
||||
rawBody = await readBody(req);
|
||||
} catch (err) {
|
||||
log.warn('failed to read request body', { error: err });
|
||||
res.writeHead(413).end();
|
||||
return;
|
||||
}
|
||||
|
||||
const eventName = req.headers['x-gitea-event'];
|
||||
const delivery = req.headers['x-gitea-delivery'];
|
||||
log.info('webhook received', { event: eventName, delivery, bytes: rawBody.length });
|
||||
|
||||
if (!verifySignature(config.webhookSecret, req.headers['x-gitea-signature'], rawBody)) {
|
||||
log.warn('rejected webhook: invalid signature', { delivery });
|
||||
res.writeHead(401).end();
|
||||
return;
|
||||
}
|
||||
|
||||
// Acknowledge immediately; Gitea only cares that we accepted delivery.
|
||||
// The job itself can run far longer than any sane webhook timeout.
|
||||
res.writeHead(202).end('accepted');
|
||||
|
||||
if (!eventName) {
|
||||
log.debug('ignored: no X-Gitea-Event header (likely a health check)');
|
||||
return;
|
||||
}
|
||||
|
||||
let payload;
|
||||
try {
|
||||
payload = JSON.parse(rawBody.toString('utf8'));
|
||||
} catch (err) {
|
||||
log.warn('ignored: invalid JSON payload', { error: err });
|
||||
return;
|
||||
}
|
||||
|
||||
const { trigger, reason } = classifyEvent(eventName, payload, config);
|
||||
if (!trigger) {
|
||||
log.info('ignored', { reason });
|
||||
return;
|
||||
}
|
||||
|
||||
log.info('trigger matched', {
|
||||
repo: trigger.repo,
|
||||
issueNumber: trigger.issueNumber,
|
||||
isPullRequest: trigger.isPullRequest,
|
||||
});
|
||||
onTrigger(trigger, log);
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user