Compare commits

...
7 Commits
Author SHA1 Message Date
Daniel SamsonandClaude Opus 4.6 36c80458e0 Safe 409 handling: check workspace age/status before deleting
Publish Image / publish (push) Successful in 20s
Instead of blindly deleting on 409 conflict, now checks:
- Stopped/failed → safe to delete and recreate
- Running but >30 min old → stale, delete and recreate
- Running and <30 min old → adopt existing workspace (mid-task)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:59:59 +01:00
Daniel SamsonandClaude Opus 4.6 06cf4d1b9e Handle stale workspaces: 409 conflict auto-cleanup + verbose logging
Publish Image / publish (push) Successful in 20s
When the orchestrator restarts, it loses in-memory queue state. If a
stale workspace still exists, createWorkspace now auto-deletes it and
retries instead of failing with 409. Also adds --verbose to Claude Code
invocation for better task log debugging.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:58:08 +01:00
Daniel SamsonandClaude Opus 4.6 4f1d306fdf Fix ImagePullBackOff: use multi-arch base image tag
The tag `ubuntu-arm64` doesn't exist on Docker Hub. Switch to `latest`
which includes arm64 in its multi-arch manifest.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:44:44 +01:00
Daniel SamsonandClaude Opus 4.6 7145fb536a Add workspace lifecycle management via task-complete callback
Publish Image / publish (push) Successful in 20s
- Workspaces curl POST /webhook/task-complete/{name} when done
- Orchestrator deletes the workspace via Coder API on callback
- Active workspaces count toward concurrency limit
- GET /queue now shows active workspaces with elapsed time
- Coder template gains callback_url parameter
- TTL becomes a safety net, not the primary shutdown mechanism

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:37:08 +01:00
Daniel SamsonandClaude Opus 4.6 f23fb4e2e0 Add task queue with dedup and concurrency control
Publish Image / publish (push) Successful in 54s
- Tasks keyed by {taskType}-{repo}-{issue} for deduplication
- Duplicate requests (pending or running) are dropped
- Configurable concurrency via QUEUE_CONCURRENCY env (default 2)
- Workspace names now {taskType}-{repo}-{issue} for observability
- GET /queue endpoint shows pending/running tasks
- All handlers route through the queue instead of calling Coder directly

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:30:27 +01:00
Daniel SamsonandClaude Opus 4.6 6062d1b518 Trigger analyst on edited comments, not just new ones
Publish Image / publish (push) Successful in 21s
An edited comment during analysis is equally meaningful as a new one
for continuing the conversation with the analyst.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:13:39 +01:00
Daniel SamsonandClaude Opus 4.6 4e5993f489 Fix cert-manager ClusterIssuer name to letsencrypt-prod
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 13:02:08 +01:00
13 changed files with 455 additions and 113 deletions
+96 -3
View File
@@ -208,6 +208,15 @@ data "coder_parameter" "anthropic_api_key" {
mutable = false mutable = false
} }
data "coder_parameter" "callback_url" {
name = "callback_url"
display_name = "Task Complete Callback URL"
description = "URL to POST when task finishes (set by orchestrator)"
type = "string"
default = ""
mutable = false
}
data "coder_parameter" "deploy_env" { data "coder_parameter" "deploy_env" {
name = "deploy_env" name = "deploy_env"
display_name = "Deploy Environment" display_name = "Deploy Environment"
@@ -433,6 +442,25 @@ resource "coder_agent" "main" {
web_terminal = true web_terminal = true
} }
metadata {
display_name = "Task Status"
key = "task_status"
script = <<-EOS
if [ -f ~/task-output.log ]; then
if grep -q "Task completed" ~/task-output.log 2>/dev/null; then
echo "✅ Complete"
else
echo "⏳ Running ($(wc -l < ~/task-output.log) lines)"
fi
elif [ -n "$TASK_TYPE" ]; then
echo "⏳ Starting..."
else
echo "Interactive"
fi
EOS
interval = 10
}
env = { env = {
ANTHROPIC_API_KEY = data.coder_parameter.anthropic_api_key.value ANTHROPIC_API_KEY = data.coder_parameter.anthropic_api_key.value
GITHUB_TOKEN = data.coder_external_auth.github.access_token GITHUB_TOKEN = data.coder_external_auth.github.access_token
@@ -443,6 +471,7 @@ resource "coder_agent" "main" {
ISSUE_NUMBER = data.coder_parameter.issue_number.value ISSUE_NUMBER = data.coder_parameter.issue_number.value
REPO_CLONE_URL = data.coder_parameter.repo_clone_url.value REPO_CLONE_URL = data.coder_parameter.repo_clone_url.value
DEPLOY_ENV = data.coder_parameter.deploy_env.value DEPLOY_ENV = data.coder_parameter.deploy_env.value
CALLBACK_URL = data.coder_parameter.callback_url.value
} }
startup_script = <<-EOT startup_script = <<-EOT
@@ -506,6 +535,41 @@ resource "coder_agent" "main" {
npm run db:migrate:local npm run db:migrate:local
fi fi
# --- Start log viewer on port 13338 ---
if [ -n "$TASK_TYPE" ]; then
touch ~/task-output.log
cat > ~/log-server.sh << 'LOGEOF'
#!/bin/bash
while true; do
{
echo "HTTP/1.1 200 OK"
echo "Content-Type: text/html; charset=utf-8"
echo "Connection: close"
echo ""
echo "<html><head><title>Task Log</title>"
echo "<meta http-equiv='refresh' content='5'>"
echo "<style>body{background:#1e1e1e;color:#d4d4d4;font-family:monospace;font-size:13px;padding:16px;white-space:pre-wrap;}"
echo "h2{color:#569cd6;margin:0 0 8px}.meta{color:#6a9955;margin-bottom:16px;display:block}</style></head><body>"
echo "<h2>$TASK_TYPE #$ISSUE_NUMBER — $(hostname)</h2>"
if [ -f ~/task-output.log ]; then
LINES=$(wc -l < ~/task-output.log)
if grep -q "Task completed" ~/task-output.log 2>/dev/null; then
echo "<span class='meta'>Status: ✅ Complete ($LINES lines)</span>"
else
echo "<span class='meta'>Status: ⏳ Running ($LINES lines) — auto-refreshing every 5s</span>"
fi
sed 's/&/\&amp;/g; s/</\&lt;/g; s/>/\&gt;/g' ~/task-output.log
else
echo "Waiting for task to start..."
fi
echo "</body></html>"
} | nc -l -p 13338 -q 1 2>/dev/null || true
done
LOGEOF
chmod +x ~/log-server.sh
nohup bash ~/log-server.sh &>/dev/null &
fi
# --- Automated task execution --- # --- Automated task execution ---
if [ -n "$TASK_TYPE" ] && [ -d ~/project ]; then if [ -n "$TASK_TYPE" ] && [ -d ~/project ]; then
cd ~/project cd ~/project
@@ -515,14 +579,38 @@ resource "coder_agent" "main" {
ARGS="$DEPLOY_ENV" ARGS="$DEPLOY_ENV"
fi fi
# Run Claude Code in non-interactive mode # Run Claude Code in non-interactive mode with tool access
claude --print --dangerously-skip-permissions "/project:$TASK_TYPE $ARGS" 2>&1 | tee ~/task-output.log claude -p --dangerously-skip-permissions --verbose "/project:$TASK_TYPE $ARGS" 2>&1 | tee ~/task-output.log
EXIT_CODE=$?
echo "Claude exited with code: $EXIT_CODE" | tee -a ~/task-output.log
echo "Task completed. Output saved to ~/task-output.log" echo "Task completed. Output saved to ~/task-output.log"
# Notify orchestrator that the task is done
if [ -n "$CALLBACK_URL" ]; then
curl -s -X POST -H "Content-Type: application/json" \
-d '{"status":"complete"}' \
"$CALLBACK_URL" || echo "Callback failed (non-fatal)"
fi
fi fi
EOT EOT
} }
resource "coder_app" "task_log" {
agent_id = coder_agent.main.id
slug = "task-log"
display_name = "Task Log"
icon = "/icon/document.svg"
url = "http://localhost:13338"
share = "owner"
healthcheck {
url = "http://localhost:13338"
interval = 10
threshold = 3
}
}
# ============================================================================= # =============================================================================
# Main Workspace Deployment # Main Workspace Deployment
# ============================================================================= # =============================================================================
@@ -603,7 +691,7 @@ resource "kubernetes_deployment_v1" "workspace" {
container { container {
name = "coder-agent" name = "coder-agent"
image = "codercom/enterprise-base:ubuntu-arm64" image = "codercom/enterprise-base:latest"
command = ["sh", "-c", coder_agent.main.init_script] command = ["sh", "-c", coder_agent.main.init_script]
@@ -657,6 +745,11 @@ resource "kubernetes_deployment_v1" "workspace" {
value = data.coder_parameter.deploy_env.value value = data.coder_parameter.deploy_env.value
} }
env {
name = "CALLBACK_URL"
value = data.coder_parameter.callback_url.value
}
resources { resources {
requests = { requests = {
cpu = "${tonumber(data.coder_parameter.cpu.value) * 500}m" cpu = "${tonumber(data.coder_parameter.cpu.value) * 500}m"
+1 -1
View File
@@ -69,7 +69,7 @@ metadata:
name: sdlc-orchestrator name: sdlc-orchestrator
namespace: sdlc-orchestrator namespace: sdlc-orchestrator
annotations: annotations:
cert-manager.io/cluster-issuer: letsencrypt-production cert-manager.io/cluster-issuer: letsencrypt-prod
spec: spec:
tls: tls:
- hosts: - hosts:
+3
View File
@@ -19,6 +19,8 @@ export interface Config {
botReviewUsername: string; botReviewUsername: string;
/** Gitea base URL (e.g. https://gitea.samson.media) */ /** Gitea base URL (e.g. https://gitea.samson.media) */
giteaUrl: string; giteaUrl: string;
/** Base URL for task-complete callbacks (e.g. https://sdlc.samson.media) */
callbackUrl: string;
/** Optional webhook secret for verifying Gitea signatures */ /** Optional webhook secret for verifying Gitea signatures */
webhookSecret?: string; webhookSecret?: string;
} }
@@ -43,6 +45,7 @@ export function loadConfig(): Config {
botDevUsername: process.env.BOT_DEV_USERNAME || "claude-dev", botDevUsername: process.env.BOT_DEV_USERNAME || "claude-dev",
botReviewUsername: process.env.BOT_REVIEW_USERNAME || "claude-review", botReviewUsername: process.env.BOT_REVIEW_USERNAME || "claude-review",
giteaUrl: process.env.GITEA_URL || "https://gitea.samson.media", giteaUrl: process.env.GITEA_URL || "https://gitea.samson.media",
callbackUrl: required("CALLBACK_URL"),
webhookSecret: process.env.WEBHOOK_SECRET, webhookSecret: process.env.WEBHOOK_SECRET,
}; };
} }
+7 -10
View File
@@ -1,19 +1,19 @@
import type { Context } from "hono"; import type { Context } from "hono";
import type { IssueCommentEvent, TaskRequest } from "../types.js"; import type { IssueCommentEvent, TaskRequest } from "../types.js";
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
/** /**
* Issue comment created → re-trigger analyst if still in analysis. * Issue comment created or edited → re-trigger analyst if still in analysis.
* *
* Replaces: issue-comment-reply.json * Replaces: issue-comment-reply.json
*/ */
export function issueComment(config: Config, coder: CoderClient) { export function issueComment(config: Config, queue: TaskQueue) {
return async (c: Context) => { return async (c: Context) => {
const event = await c.req.json<IssueCommentEvent>(); const event = await c.req.json<IssueCommentEvent>();
if (event.action !== "created") { if (event.action !== "created" && event.action !== "edited") {
return c.text("ignored: not a created event", 200); return c.text("ignored: not a created or edited event", 200);
} }
const commenter = event.comment.user.login; const commenter = event.comment.user.login;
@@ -43,11 +43,8 @@ export function issueComment(config: Config, coder: CoderClient) {
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
console.log(`[issue-comment] #${task.issueNumber} → re-trigger analyst`); const queued = queue.enqueue(task);
const workspace = await coder.createWorkspace(task); return c.json({ ok: true, queued });
console.log(`[issue-comment] workspace created: ${workspace.name}`);
return c.json({ ok: true, workspace: workspace.name });
}; };
} }
+7 -15
View File
@@ -1,7 +1,6 @@
import type { Context } from "hono"; import type { Context } from "hono";
import type { IssueEvent, TaskRequest } from "../types.js"; import type { IssueEvent, TaskRequest } from "../types.js";
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { GiteaClient } from "../services/gitea.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
const LABEL_TO_TASK: Record<string, string> = { const LABEL_TO_TASK: Record<string, string> = {
@@ -19,7 +18,7 @@ const REVIEW_TASKS = new Set(["test"]);
* *
* Replaces: issue-stage-transition.json * Replaces: issue-stage-transition.json
*/ */
export function issueLabel(config: Config, coder: CoderClient, gitea: GiteaClient) { export function issueLabel(config: Config, queue: TaskQueue) {
return async (c: Context) => { return async (c: Context) => {
const event = await c.req.json<IssueEvent>(); const event = await c.req.json<IssueEvent>();
@@ -55,18 +54,11 @@ export function issueLabel(config: Config, coder: CoderClient, gitea: GiteaClien
giteaToken, giteaToken,
}; };
console.log(`[issue-label] #${task.issueNumber} → ${taskType}`); const queued = queue.enqueue(task, {
useReviewAccount: false,
body: `🤖 SDLC stage transition: **${taskType}** — workspace queued.`,
});
const workspace = await coder.createWorkspace(task); return c.json({ ok: true, queued });
console.log(`[issue-label] workspace created: ${workspace.name}`);
await gitea.commentOnIssue(
org,
repo,
event.issue.number,
`🤖 SDLC stage transition: **${taskType}** — workspace created.`,
);
return c.json({ ok: true, workspace: workspace.name });
}; };
} }
+7 -15
View File
@@ -1,7 +1,6 @@
import type { Context } from "hono"; import type { Context } from "hono";
import type { IssueEvent, TaskRequest } from "../types.js"; import type { IssueEvent, TaskRequest } from "../types.js";
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { GiteaClient } from "../services/gitea.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
/** /**
@@ -9,7 +8,7 @@ import type { Config } from "../config.js";
* *
* Replaces: gitea-issue-triage.json * Replaces: gitea-issue-triage.json
*/ */
export function issueTriage(config: Config, coder: CoderClient, gitea: GiteaClient) { export function issueTriage(config: Config, queue: TaskQueue) {
return async (c: Context) => { return async (c: Context) => {
const event = await c.req.json<IssueEvent>(); const event = await c.req.json<IssueEvent>();
@@ -32,18 +31,11 @@ export function issueTriage(config: Config, coder: CoderClient, gitea: GiteaClie
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
console.log(`[issue-triage] #${task.issueNumber} → ${taskType}`); const queued = queue.enqueue(task, {
useReviewAccount: false,
body: `🤖 Issue received. Starting **${taskType}** stage.`,
});
const workspace = await coder.createWorkspace(task); return c.json({ ok: true, queued });
console.log(`[issue-triage] workspace created: ${workspace.name}`);
await gitea.commentOnIssue(
org,
repo,
event.issue.number,
`🤖 Issue received. Starting **${taskType}** stage.`,
);
return c.json({ ok: true, workspace: workspace.name });
}; };
} }
+4 -18
View File
@@ -1,4 +1,4 @@
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
import type { TaskRequest } from "../types.js"; import type { TaskRequest } from "../types.js";
@@ -9,15 +9,11 @@ export interface MaintenanceRepo {
} }
/** /**
* Weekly maintenance cron → create a maintenance workspace. * Weekly maintenance cron → enqueue a maintenance workspace per repo.
* *
* Replaces: scheduled-maintenance.json * Replaces: scheduled-maintenance.json
*/ */
export function runMaintenance( export function runMaintenance(config: Config, queue: TaskQueue, repos: MaintenanceRepo[]) {
config: Config,
coder: CoderClient,
repos: MaintenanceRepo[],
) {
return async () => { return async () => {
for (const entry of repos) { for (const entry of repos) {
const task: TaskRequest = { const task: TaskRequest = {
@@ -29,17 +25,7 @@ export function runMaintenance(
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
try { queue.enqueue(task);
const workspace = await coder.createWorkspace(task);
console.log(
`[maintenance] ${entry.org}/${entry.repo} → workspace: ${workspace.name}`,
);
} catch (err) {
console.error(
`[maintenance] failed for ${entry.org}/${entry.repo}:`,
err,
);
}
} }
}; };
} }
+4 -7
View File
@@ -1,6 +1,6 @@
import type { Context } from "hono"; import type { Context } from "hono";
import type { PullRequestEvent, TaskRequest } from "../types.js"; import type { PullRequestEvent, TaskRequest } from "../types.js";
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { GiteaClient } from "../services/gitea.js"; import type { GiteaClient } from "../services/gitea.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
@@ -10,7 +10,7 @@ import type { Config } from "../config.js";
* *
* Replaces: pr-review-trigger.json * Replaces: pr-review-trigger.json
*/ */
export function prReview(config: Config, coder: CoderClient, gitea: GiteaClient) { export function prReview(config: Config, queue: TaskQueue, gitea: GiteaClient) {
return async (c: Context) => { return async (c: Context) => {
const event = await c.req.json<PullRequestEvent>(); const event = await c.req.json<PullRequestEvent>();
@@ -40,11 +40,8 @@ export function prReview(config: Config, coder: CoderClient, gitea: GiteaClient)
giteaToken: config.giteaReviewToken, giteaToken: config.giteaReviewToken,
}; };
console.log(`[pr-review] PR #${prNumber} → review-code`); const queued = queue.enqueue(task);
const workspace = await coder.createWorkspace(task); return c.json({ ok: true, queued });
console.log(`[pr-review] workspace created: ${workspace.name}`);
return c.json({ ok: true, workspace: workspace.name });
}; };
} }
+7 -15
View File
@@ -1,7 +1,6 @@
import type { Context } from "hono"; import type { Context } from "hono";
import type { PullRequestReviewEvent, TaskRequest } from "../types.js"; import type { PullRequestReviewEvent, TaskRequest } from "../types.js";
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { GiteaClient } from "../services/gitea.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
/** /**
@@ -9,7 +8,7 @@ import type { Config } from "../config.js";
* *
* Replaces: pr-review-rework.json * Replaces: pr-review-rework.json
*/ */
export function prRework(config: Config, coder: CoderClient, gitea: GiteaClient) { export function prRework(config: Config, queue: TaskQueue) {
return async (c: Context) => { return async (c: Context) => {
const event = await c.req.json<PullRequestReviewEvent>(); const event = await c.req.json<PullRequestReviewEvent>();
@@ -43,18 +42,11 @@ export function prRework(config: Config, coder: CoderClient, gitea: GiteaClient)
giteaToken, giteaToken,
}; };
console.log(`[pr-rework] PR #${prNumber} ${reviewState} → ${taskType}`); const queued = queue.enqueue(task, {
useReviewAccount: false,
body: `🤖 Review outcome: **${taskType}** — workspace queued.`,
});
const workspace = await coder.createWorkspace(task); return c.json({ ok: true, queued });
console.log(`[pr-rework] workspace created: ${workspace.name}`);
await gitea.commentOnIssue(
org,
repo,
prNumber,
`🤖 Review outcome: starting **${taskType}** stage.`,
);
return c.json({ ok: true, workspace: workspace.name });
}; };
} }
+7 -15
View File
@@ -1,7 +1,6 @@
import type { Context } from "hono"; import type { Context } from "hono";
import type { IssueEvent, TaskRequest } from "../types.js"; import type { IssueEvent, TaskRequest } from "../types.js";
import type { CoderClient } from "../services/coder.js"; import type { TaskQueue } from "../queue.js";
import type { GiteaClient } from "../services/gitea.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
/** /**
@@ -10,7 +9,7 @@ import type { Config } from "../config.js";
* *
* Replaces: release-process.json * Replaces: release-process.json
*/ */
export function release(config: Config, coder: CoderClient, gitea: GiteaClient) { export function release(config: Config, queue: TaskQueue) {
return async (c: Context) => { return async (c: Context) => {
const event = await c.req.json<IssueEvent>(); const event = await c.req.json<IssueEvent>();
@@ -52,18 +51,11 @@ export function release(config: Config, coder: CoderClient, gitea: GiteaClient)
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
console.log(`[release] #${task.issueNumber} → release`); const queued = queue.enqueue(task, {
useReviewAccount: false,
body: `🤖 Release process started.`,
});
const workspace = await coder.createWorkspace(task); return c.json({ ok: true, queued });
console.log(`[release] workspace created: ${workspace.name}`);
await gitea.commentOnIssue(
org,
repo,
event.issue.number,
`🤖 Release process started.`,
);
return c.json({ ok: true, workspace: workspace.name });
}; };
} }
+25 -9
View File
@@ -6,6 +6,7 @@ import cron from "node-cron";
import { loadConfig } from "./config.js"; import { loadConfig } from "./config.js";
import { CoderClient } from "./services/coder.js"; import { CoderClient } from "./services/coder.js";
import { GiteaClient } from "./services/gitea.js"; import { GiteaClient } from "./services/gitea.js";
import { TaskQueue } from "./queue.js";
import { issueTriage } from "./handlers/issue-triage.js"; import { issueTriage } from "./handlers/issue-triage.js";
import { issueLabel } from "./handlers/issue-label.js"; import { issueLabel } from "./handlers/issue-label.js";
@@ -23,26 +24,41 @@ const config = loadConfig();
const coderClient = new CoderClient(config); const coderClient = new CoderClient(config);
const giteaClient = new GiteaClient(config); const giteaClient = new GiteaClient(config);
const concurrency = parseInt(process.env.QUEUE_CONCURRENCY || "2", 10);
const queue = new TaskQueue(coderClient, giteaClient, concurrency);
const app = new Hono(); const app = new Hono();
app.use("*", logger()); app.use("*", logger());
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Health check // Health check + queue status
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
app.get("/health", (c) => c.json({ status: "ok" })); app.get("/health", (c) => c.json({ status: "ok" }));
app.get("/queue", (c) => c.json(queue.status()));
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Webhook endpoints — same paths as the n8n webhooks so Gitea config // Webhook endpoints — same paths as the n8n webhooks so Gitea config
// can stay unchanged (just point to new host). // can stay unchanged (just point to new host).
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
app.post("/webhook/gitea-issue-triage", issueTriage(config, coderClient, giteaClient)); app.post("/webhook/gitea-issue-triage", issueTriage(config, queue));
app.post("/webhook/gitea-issue-label", issueLabel(config, coderClient, giteaClient)); app.post("/webhook/gitea-issue-label", issueLabel(config, queue));
app.post("/webhook/gitea-issue-comment", issueComment(config, coderClient)); app.post("/webhook/gitea-issue-comment", issueComment(config, queue));
app.post("/webhook/gitea-pr-review", prReview(config, coderClient, giteaClient)); app.post("/webhook/gitea-pr-review", prReview(config, queue, giteaClient));
app.post("/webhook/gitea-pr-review-rework", prRework(config, coderClient, giteaClient)); app.post("/webhook/gitea-pr-review-rework", prRework(config, queue));
app.post("/webhook/gitea-release", release(config, coderClient, giteaClient)); 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 });
});
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Scheduled maintenance cron (Monday 9:00 AM) // Scheduled maintenance cron (Monday 9:00 AM)
@@ -50,7 +66,7 @@ app.post("/webhook/gitea-release", release(config, coderClient, giteaClient));
const maintenanceRepos = parseMaintenanceRepos(); const maintenanceRepos = parseMaintenanceRepos();
if (maintenanceRepos.length > 0) { if (maintenanceRepos.length > 0) {
const maintenanceFn = runMaintenance(config, coderClient, maintenanceRepos); const maintenanceFn = runMaintenance(config, queue, maintenanceRepos);
cron.schedule("0 9 * * 1", () => { cron.schedule("0 9 * * 1", () => {
console.log("[cron] running weekly maintenance"); console.log("[cron] running weekly maintenance");
@@ -69,7 +85,7 @@ if (maintenanceRepos.length > 0) {
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
serve({ fetch: app.fetch, port: config.port }, (info) => { serve({ fetch: app.fetch, port: config.port }, (info) => {
console.log(`SDLC Orchestrator listening on :${info.port}`); console.log(`SDLC Orchestrator listening on :${info.port} (concurrency: ${concurrency})`);
}); });
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
+179
View File
@@ -0,0 +1,179 @@
import type { TaskRequest } from "./types.js";
import type { CoderClient } from "./services/coder.js";
import type { GiteaClient } from "./services/gitea.js";
interface QueuedTask {
task: TaskRequest;
comment?: { useReviewAccount: boolean; body: string };
}
interface ActiveWorkspace {
workspaceId: string;
workspaceName: string;
startedAt: Date;
}
/**
* Task queue with deduplication, concurrency control, and workspace lifecycle.
*
* - Tasks are keyed by `{taskType}-{repo}-{issue}` for dedup
* - If a task with the same key is pending, running, or has an active workspace, new requests are dropped
* - Concurrency is configurable (defaults to 2)
* - Tracks active workspaces and stops them on task-complete callback
*/
export class TaskQueue {
private pending = new Map<string, QueuedTask>();
private running = new Set<string>();
private active = new Map<string, ActiveWorkspace>();
private concurrency: number;
private coder: CoderClient;
private gitea: GiteaClient;
constructor(coder: CoderClient, gitea: GiteaClient, concurrency = 2) {
this.coder = coder;
this.gitea = gitea;
this.concurrency = concurrency;
}
/** Build a dedup key for a task. */
static key(task: TaskRequest): string {
return `${task.taskType}-${task.giteaRepo}-${task.issueNumber}`;
}
/** Build the workspace name (matches the dedup key). */
static workspaceName(task: TaskRequest): string {
return `${task.taskType}-${task.giteaRepo}-${task.issueNumber}`;
}
/**
* Enqueue a task. Returns true if queued, false if deduplicated (dropped).
*/
enqueue(
task: TaskRequest,
comment?: { useReviewAccount: boolean; body: string },
): boolean {
const key = TaskQueue.key(task);
if (this.active.has(key)) {
console.log(`[queue] dropped (active workspace): ${key}`);
return false;
}
if (this.running.has(key)) {
console.log(`[queue] dropped (running): ${key}`);
return false;
}
if (this.pending.has(key)) {
console.log(`[queue] dropped (pending): ${key}`);
return false;
}
this.pending.set(key, { task, comment });
console.log(
`[queue] enqueued: ${key} (pending: ${this.pending.size}, running: ${this.running.size}/${this.concurrency}, active: ${this.active.size})`,
);
this.drain();
return true;
}
/**
* Handle task-complete callback from a workspace.
* Stops the workspace via Coder API and removes it from active tracking.
*/
async taskComplete(workspaceName: string): Promise<boolean> {
const entry = this.active.get(workspaceName);
if (!entry) {
console.log(`[queue] task-complete for unknown workspace: ${workspaceName}`);
return false;
}
const elapsed = Math.round(
(Date.now() - entry.startedAt.getTime()) / 1000,
);
console.log(
`[queue] task complete: ${workspaceName} (ran for ${elapsed}s) — stopping workspace`,
);
this.active.delete(workspaceName);
try {
await this.coder.deleteWorkspace(entry.workspaceId);
console.log(`[queue] workspace deleted: ${workspaceName}`);
} catch (err) {
console.error(`[queue] failed to delete workspace ${workspaceName}:`, err);
}
// Drain in case pending tasks were waiting for capacity
this.drain();
return true;
}
/** Process queued tasks up to the concurrency limit. */
private drain(): void {
const totalInFlight = this.running.size + this.active.size;
while (totalInFlight + this.pending.size > 0 && this.running.size + this.active.size < this.concurrency && this.pending.size > 0) {
const [key, entry] = this.pending.entries().next().value!;
this.pending.delete(key);
this.running.add(key);
this.process(key, entry);
}
}
private async process(key: string, entry: QueuedTask): Promise<void> {
const { task, comment } = entry;
try {
console.log(`[queue] processing: ${key}`);
const workspace = await this.coder.createWorkspace(task);
console.log(`[queue] workspace created: ${workspace.name} (id: ${workspace.id})`);
// Move from running → active (workspace is now alive, waiting for callback)
this.active.set(key, {
workspaceId: workspace.id,
workspaceName: workspace.name,
startedAt: new Date(),
});
if (comment) {
await this.gitea.commentOnIssue(
task.giteaOrg,
task.giteaRepo,
task.issueNumber,
comment.body,
comment.useReviewAccount,
);
}
} catch (err) {
console.error(`[queue] failed: ${key}`, err);
} finally {
this.running.delete(key);
// Don't drain here — active workspaces count toward concurrency
}
}
/** Current queue status for health/debug endpoints. */
status(): {
pending: string[];
running: string[];
active: Record<string, { workspaceId: string; elapsed: number }>;
concurrency: number;
} {
const activeMap: Record<string, { workspaceId: string; elapsed: number }> = {};
for (const [key, entry] of this.active) {
activeMap[key] = {
workspaceId: entry.workspaceId,
elapsed: Math.round((Date.now() - entry.startedAt.getTime()) / 1000),
};
}
return {
pending: [...this.pending.keys()],
running: [...this.running.keys()],
active: activeMap,
concurrency: this.concurrency,
};
}
}
+108 -5
View File
@@ -1,25 +1,24 @@
import type { Config } from "../config.js"; import type { Config } from "../config.js";
import type { TaskRequest } from "../types.js"; import type { TaskRequest } from "../types.js";
import { TaskQueue } from "../queue.js";
export class CoderClient { export class CoderClient {
private baseUrl: string; private baseUrl: string;
private token: string; private token: string;
private templateId: string; private templateId: string;
private anthropicApiKey: string; private anthropicApiKey: string;
private callbackUrl: string;
constructor(config: Config) { constructor(config: Config) {
this.baseUrl = config.coderUrl.replace(/\/$/, ""); this.baseUrl = config.coderUrl.replace(/\/$/, "");
this.token = config.coderToken; this.token = config.coderToken;
this.templateId = config.coderTemplateId; this.templateId = config.coderTemplateId;
this.anthropicApiKey = config.anthropicApiKey; this.anthropicApiKey = config.anthropicApiKey;
this.callbackUrl = config.callbackUrl;
} }
async createWorkspace(task: TaskRequest): Promise<{ id: string; name: string }> { async createWorkspace(task: TaskRequest): Promise<{ id: string; name: string }> {
const timestamp = new Date() const name = TaskQueue.workspaceName(task);
.toISOString()
.replace(/[-:T]/g, "")
.slice(8, 14); // HHmmss
const name = `${task.taskType}-${task.issueNumber}-${timestamp}`;
const body = { const body = {
name, name,
@@ -32,6 +31,7 @@ export class CoderClient {
{ name: "gitea_org", value: task.giteaOrg }, { name: "gitea_org", value: task.giteaOrg },
{ name: "gitea_repo", value: task.giteaRepo }, { name: "gitea_repo", value: task.giteaRepo },
{ name: "anthropic_api_key", value: this.anthropicApiKey }, { name: "anthropic_api_key", value: this.anthropicApiKey },
{ name: "callback_url", value: `${this.callbackUrl}/webhook/task-complete/${name}` },
...(task.deployEnv ...(task.deployEnv
? [{ name: "deploy_env", value: task.deployEnv }] ? [{ name: "deploy_env", value: task.deployEnv }]
: []), : []),
@@ -50,6 +50,38 @@ export class CoderClient {
}, },
); );
// Handle 409 conflict — workspace with this name already exists.
// This happens when the orchestrator restarts and loses in-memory state.
// Only delete if the workspace is old (>30 min) or stopped/failed.
// If it's young and running, it may be mid-task — skip to avoid data loss.
if (res.status === 409) {
const existing = await this.findWorkspaceByName(name);
if (!existing) {
throw new Error(`Coder 409 but workspace "${name}" not found — possible race condition`);
}
const ageMs = Date.now() - new Date(existing.createdAt).getTime();
const ageMin = Math.round(ageMs / 60_000);
const stoppedStatuses = ["stopped", "failed", "canceled", "deleted"];
const isStopped = stoppedStatuses.includes(existing.latestBuildStatus);
const isStale = ageMin > 30;
if (isStopped || isStale) {
console.log(
`[coder] workspace "${name}" already exists (age: ${ageMin}m, status: ${existing.latestBuildStatus}) — deleting and retrying`,
);
await this.deleteWorkspace(existing.id);
await new Promise((resolve) => setTimeout(resolve, 5000));
return this.createWorkspace(task);
}
// Workspace is young and still running — don't kill it
console.log(
`[coder] workspace "${name}" is active (age: ${ageMin}m, status: ${existing.latestBuildStatus}) — skipping creation`,
);
return { id: existing.id, name: existing.name };
}
if (!res.ok) { if (!res.ok) {
const text = await res.text(); const text = await res.text();
throw new Error(`Coder API error ${res.status}: ${text}`); throw new Error(`Coder API error ${res.status}: ${text}`);
@@ -58,4 +90,75 @@ export class CoderClient {
const data = (await res.json()) as { id: string; name: string }; const data = (await res.json()) as { id: string; name: string };
return { id: data.id, name: data.name }; return { id: data.id, name: data.name };
} }
async stopWorkspace(workspaceId: string): Promise<void> {
const res = await fetch(
`${this.baseUrl}/api/v2/workspaces/${workspaceId}/builds`,
{
method: "POST",
headers: {
"Content-Type": "application/json",
"Coder-Session-Token": this.token,
},
body: JSON.stringify({ transition: "stop" }),
},
);
if (!res.ok) {
const text = await res.text();
throw new Error(`Coder stop error ${res.status}: ${text}`);
}
}
async deleteWorkspace(workspaceId: string): Promise<void> {
const res = await fetch(
`${this.baseUrl}/api/v2/workspaces/${workspaceId}`,
{
method: "DELETE",
headers: {
"Coder-Session-Token": this.token,
},
},
);
if (!res.ok) {
const text = await res.text();
throw new Error(`Coder delete error ${res.status}: ${text}`);
}
}
async findWorkspaceByName(name: string): Promise<{
id: string;
name: string;
createdAt: string;
latestBuildStatus: string;
} | null> {
const res = await fetch(
`${this.baseUrl}/api/v2/workspaces?q=name:${encodeURIComponent(name)}`,
{
headers: {
"Coder-Session-Token": this.token,
},
},
);
if (!res.ok) return null;
const data = (await res.json()) as {
workspaces: Array<{
id: string;
name: string;
created_at: string;
latest_build: { status: string };
}>;
};
const match = data.workspaces?.find((w) => w.name === name);
if (!match) return null;
return {
id: match.id,
name: match.name,
createdAt: match.created_at,
latestBuildStatus: match.latest_build?.status ?? "unknown",
};
}
} }