#!/usr/bin/env node import { createServer } from "node:http"; import { pathToFileURL } from "node:url"; import { AgentStore } from "../../adapters/database/store.js"; import { GiteaClient } from "../../adapters/gitea/client/client.js"; import { actorPolicy, readSecret, repositoryParts, validateServerUrl, } from "../../core/config.js"; import { formatError, requireEnv } from "../../core/contracts.js"; import { type PublicationContext, safeFailure } from "../publication/status.js"; import { log, pumpDeliveries, pumpOutbox } from "./handlers/workers.js"; import { handleHttp, observeHttpRequest } from "./server.js"; async function main(): Promise { const serverUrl = validateServerUrl(requireEnv("GITEA_SERVER_URL")); const repository = repositoryParts(); const writeToken = await readSecret("GITEA_WRITE_TOKEN"); const webhookSecret = await readSecret("GITEA_WEBHOOK_SECRET"); const botLogin = requireEnv("CI_AGENT_BOT_LOGIN"); const policy = actorPolicy(); const store = new AgentStore( process.env.AGENT_DB_PATH || "/var/lib/gitea-agent/agent.db", ); store.recoverControllerWork(); store.purgeDeliveries(Date.now() - 7 * 24 * 60 * 60_000); const shutdown = new AbortController(); const client = new GiteaClient( serverUrl, writeToken, repository.owner, repository.repo, shutdown.signal, ); const [configuredRepository, bot] = await Promise.all([ client.getRepository(), client.getCurrentUser(), ]); if (bot.login.toLowerCase() !== botLogin.toLowerCase()) throw new Error(`Bot login ${botLogin} does not match ${bot.login}`); const context: PublicationContext = { client, botLogin, serverUrl, repository, repositoryId: configuredRepository.id, writeToken, signal: shutdown.signal, isCancelled: (jobId) => store.isCancelRequested(jobId), }; let delivering: Promise | undefined; let publishing: Promise | undefined; const pump = () => { if (shutdown.signal.aborted) return; if (!delivering) { delivering = pumpDeliveries( store, context, configuredRepository.id, configuredRepository.full_name, bot.id, policy, ) .catch((error) => console.error( log("error", "Delivery work failed", { error: safeFailure(error), }), ), ) .finally(() => { delivering = undefined; }); } if (!publishing) { publishing = pumpOutbox(store, context) .catch((error) => console.error( log("error", "Publication work failed", { error: safeFailure(error), }), ), ) .finally(() => { publishing = undefined; }); } }; const pumpTimer = setInterval(pump, 250); const retentionTimer = setInterval( () => store.purgeDeliveries(Date.now() - 7 * 24 * 60 * 60_000), 60 * 60_000, ); pumpTimer.unref(); retentionTimer.unref(); const server = createServer((request, response) => { observeHttpRequest(request, response, (entry) => { const level = entry.statusCode >= 500 ? "error" : "info"; console.log(log(level, "HTTP request", { ...entry })); }); handleHttp(request, response, { store, webhookSecret, repositoryId: configuredRepository.id, repositoryFullName: configuredRepository.full_name, }).catch((error) => { console.error( log("error", "Webhook request failed", { error: safeFailure(error), }), ); if (!response.headersSent) response.writeHead(500); response.end(); }); }); server.headersTimeout = 10_000; server.requestTimeout = 15_000; server.keepAliveTimeout = 5_000; const host = process.env.AGENT_HTTP_HOST || "0.0.0.0"; const port = parsePort(process.env.AGENT_HTTP_PORT || "8080"); await new Promise((resolve, reject) => { server.once("error", reject); server.listen(port, host, resolve); }); console.log( log("info", "Controller ready", { host, port, repository: configuredRepository.full_name, }), ); pump(); await new Promise((resolve) => { let stopping = false; const stop = () => { if (stopping) return; stopping = true; clearInterval(pumpTimer); clearInterval(retentionTimer); server.close(() => resolve()); shutdown.abort(new Error("Controller is shutting down")); }; process.once("SIGINT", stop); process.once("SIGTERM", stop); }); await Promise.all([ delivering?.catch(() => undefined), publishing?.catch(() => undefined), ]); store.close(); } function parsePort(value: string): number { const port = Number(value); if (!Number.isInteger(port) || port < 1 || port > 65_535) throw new Error(`Invalid AGENT_HTTP_PORT: ${value}`); return port; } if ( process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href ) { main().catch((error) => { console.error(formatError(error)); process.exitCode = 1; }); }