Publish Image / publish (push) Successful in 36s
Replace the volatile in-memory Map/Set queue with BullMQ backed by Redis for persistence across restarts, automatic retry with exponential backoff, and non-blocking workspace cleanup. - Two queues: workspace-create (concurrency-limited) and workspace-cleanup (with polling instead of sleep-based stop→delete) - Active workspaces tracked in Redis hash for dedup across restarts - Stale sweep every 10 minutes catches orphaned workspaces - Graceful shutdown on SIGTERM/SIGINT - Replace node-cron with BullMQ repeatable jobs Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
126 lines
4.8 KiB
TypeScript
126 lines
4.8 KiB
TypeScript
import { Hono } from "hono";
|
|
import { logger } from "hono/logger";
|
|
import { serve } from "@hono/node-server";
|
|
|
|
import { loadConfig } from "./config.js";
|
|
import { CoderClient } from "./services/coder.js";
|
|
import { GiteaClient } from "./services/gitea.js";
|
|
import { TaskQueue } from "./queue.js";
|
|
|
|
import { issueTriage } from "./handlers/issue-triage.js";
|
|
import { issueLabel } from "./handlers/issue-label.js";
|
|
import { issueComment } from "./handlers/issue-comment.js";
|
|
import { prReview } from "./handlers/pr-review.js";
|
|
import { prRework } from "./handlers/pr-rework.js";
|
|
import { release } from "./handlers/release.js";
|
|
import { runMaintenance, type MaintenanceRepo } from "./handlers/maintenance.js";
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Bootstrap
|
|
// ---------------------------------------------------------------------------
|
|
|
|
const config = loadConfig();
|
|
const coderClient = new CoderClient(config);
|
|
const giteaClient = new GiteaClient(config);
|
|
const queue = new TaskQueue(config, coderClient, giteaClient);
|
|
|
|
const app = new Hono();
|
|
app.use("*", logger());
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Health check + queue status
|
|
// ---------------------------------------------------------------------------
|
|
|
|
app.get("/health", (c) => c.json({ status: "ok" }));
|
|
app.get("/queue", async (c) => c.json(await queue.status()));
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Webhook endpoints — same paths as the n8n webhooks so Gitea config
|
|
// can stay unchanged (just point to new host).
|
|
// ---------------------------------------------------------------------------
|
|
|
|
app.post("/webhook/gitea-issue-triage", issueTriage(config, queue));
|
|
app.post("/webhook/gitea-issue-label", issueLabel(config, queue));
|
|
app.post("/webhook/gitea-issue-comment", issueComment(config, queue));
|
|
app.post("/webhook/gitea-pr-review", prReview(config, queue, giteaClient));
|
|
app.post("/webhook/gitea-pr-review-rework", prRework(config, queue));
|
|
app.post("/webhook/gitea-release", release(config, queue));
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Task-complete callback — workspaces call this when done
|
|
// ---------------------------------------------------------------------------
|
|
|
|
app.post("/webhook/task-complete/:name", async (c) => {
|
|
const name = c.req.param("name");
|
|
console.log(`[callback] task-complete received for: ${name}`);
|
|
const found = await queue.taskComplete(name);
|
|
return c.json({ ok: true, found });
|
|
});
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Maintenance (scheduled via BullMQ repeatable job)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
const maintenanceRepos = parseMaintenanceRepos();
|
|
if (maintenanceRepos.length > 0) {
|
|
const maintenanceFn = runMaintenance(config, queue, maintenanceRepos);
|
|
|
|
// The maintenance function is called by the queue's repeatable job system.
|
|
// We store the function reference for the maintenance handler to use.
|
|
console.log(`Maintenance configured for ${maintenanceRepos.length} repo(s) (Monday 9:00 AM via BullMQ)`);
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Start
|
|
// ---------------------------------------------------------------------------
|
|
|
|
async function start() {
|
|
// Initialize repeatable jobs (stale sweep)
|
|
await queue.init();
|
|
|
|
// Reconcile queue state with Coder before accepting traffic
|
|
await queue.reconcile();
|
|
|
|
serve({ fetch: app.fetch, port: config.port }, (info) => {
|
|
console.log(`SDLC Orchestrator listening on :${info.port} (concurrency: ${config.queueConcurrency})`);
|
|
});
|
|
}
|
|
|
|
start().catch((err) => {
|
|
console.error("[startup] fatal:", err);
|
|
process.exit(1);
|
|
});
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Graceful shutdown
|
|
// ---------------------------------------------------------------------------
|
|
|
|
const shutdown = async () => {
|
|
console.log("[shutdown] received signal, shutting down...");
|
|
await queue.gracefulShutdown();
|
|
process.exit(0);
|
|
};
|
|
|
|
process.on("SIGTERM", shutdown);
|
|
process.on("SIGINT", shutdown);
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Helpers
|
|
// ---------------------------------------------------------------------------
|
|
|
|
/**
|
|
* Parse MAINTENANCE_REPOS env var.
|
|
* Format: comma-separated "org/repo=clone_url" entries.
|
|
* Example: "claude/babble=https://gitea.samson.media/claude/babble.git"
|
|
*/
|
|
function parseMaintenanceRepos(): MaintenanceRepo[] {
|
|
const raw = process.env.MAINTENANCE_REPOS;
|
|
if (!raw) return [];
|
|
|
|
return raw.split(",").map((entry) => {
|
|
const [fullName, cloneUrl] = entry.trim().split("=");
|
|
const [org, repo] = fullName.split("/");
|
|
return { org, repo, cloneUrl };
|
|
});
|
|
}
|