From 7145fb536af69c7265ead921307d23e6c8c72126 Mon Sep 17 00:00:00 2001 From: Daniel Samson <12231216+daniel-samson@users.noreply.github.com> Date: Sun, 12 Apr 2026 13:37:08 +0100 Subject: [PATCH] Add workspace lifecycle management via task-complete callback - 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 --- coder/cloudflare-worker/main.tf | 22 +++++++++ src/config.ts | 3 ++ src/index.ts | 11 +++++ src/queue.ts | 82 ++++++++++++++++++++++++++++++--- src/services/coder.ts | 56 ++++++++++++++++++++++ 5 files changed, 167 insertions(+), 7 deletions(-) diff --git a/coder/cloudflare-worker/main.tf b/coder/cloudflare-worker/main.tf index 4720f29..3241103 100644 --- a/coder/cloudflare-worker/main.tf +++ b/coder/cloudflare-worker/main.tf @@ -208,6 +208,15 @@ data "coder_parameter" "anthropic_api_key" { 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" { name = "deploy_env" display_name = "Deploy Environment" @@ -443,6 +452,7 @@ resource "coder_agent" "main" { ISSUE_NUMBER = data.coder_parameter.issue_number.value REPO_CLONE_URL = data.coder_parameter.repo_clone_url.value DEPLOY_ENV = data.coder_parameter.deploy_env.value + CALLBACK_URL = data.coder_parameter.callback_url.value } startup_script = <<-EOT @@ -519,6 +529,13 @@ resource "coder_agent" "main" { claude --print --dangerously-skip-permissions "/project:$TASK_TYPE $ARGS" 2>&1 | tee ~/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 EOT } @@ -657,6 +674,11 @@ resource "kubernetes_deployment_v1" "workspace" { value = data.coder_parameter.deploy_env.value } + env { + name = "CALLBACK_URL" + value = data.coder_parameter.callback_url.value + } + resources { requests = { cpu = "${tonumber(data.coder_parameter.cpu.value) * 500}m" diff --git a/src/config.ts b/src/config.ts index 00b77e2..d8388f0 100644 --- a/src/config.ts +++ b/src/config.ts @@ -19,6 +19,8 @@ export interface Config { botReviewUsername: string; /** Gitea base URL (e.g. https://gitea.samson.media) */ giteaUrl: string; + /** Base URL for task-complete callbacks (e.g. https://sdlc.samson.media) */ + callbackUrl: string; /** Optional webhook secret for verifying Gitea signatures */ webhookSecret?: string; } @@ -43,6 +45,7 @@ export function loadConfig(): Config { botDevUsername: process.env.BOT_DEV_USERNAME || "claude-dev", botReviewUsername: process.env.BOT_REVIEW_USERNAME || "claude-review", giteaUrl: process.env.GITEA_URL || "https://gitea.samson.media", + callbackUrl: required("CALLBACK_URL"), webhookSecret: process.env.WEBHOOK_SECRET, }; } diff --git a/src/index.ts b/src/index.ts index 6310191..d047ce7 100644 --- a/src/index.ts +++ b/src/index.ts @@ -49,6 +49,17 @@ 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 }); +}); + // --------------------------------------------------------------------------- // Scheduled maintenance cron (Monday 9:00 AM) // --------------------------------------------------------------------------- diff --git a/src/queue.ts b/src/queue.ts index b9f1d20..7ac0f3e 100644 --- a/src/queue.ts +++ b/src/queue.ts @@ -7,16 +7,24 @@ interface QueuedTask { comment?: { useReviewAccount: boolean; body: string }; } +interface ActiveWorkspace { + workspaceId: string; + workspaceName: string; + startedAt: Date; +} + /** - * Task queue with deduplication and concurrency control. + * 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 already pending or running, new requests are dropped + * - 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(); private running = new Set(); + private active = new Map(); private concurrency: number; private coder: CoderClient; private gitea: GiteaClient; @@ -46,6 +54,11 @@ export class TaskQueue { ): 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; @@ -58,16 +71,50 @@ export class TaskQueue { this.pending.set(key, { task, comment }); console.log( - `[queue] enqueued: ${key} (pending: ${this.pending.size}, running: ${this.running.size}/${this.concurrency})`, + `[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 { + 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 { - while (this.running.size < this.concurrency && this.pending.size > 0) { + 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); @@ -81,7 +128,14 @@ export class TaskQueue { try { console.log(`[queue] processing: ${key}`); const workspace = await this.coder.createWorkspace(task); - console.log(`[queue] workspace created: ${workspace.name}`); + 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( @@ -96,15 +150,29 @@ export class TaskQueue { console.error(`[queue] failed: ${key}`, err); } finally { this.running.delete(key); - this.drain(); + // Don't drain here — active workspaces count toward concurrency } } /** Current queue status for health/debug endpoints. */ - status(): { pending: string[]; running: string[]; concurrency: number } { + status(): { + pending: string[]; + running: string[]; + active: Record; + concurrency: number; + } { + const activeMap: Record = {}; + 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, }; } diff --git a/src/services/coder.ts b/src/services/coder.ts index 0a89e1d..655dabf 100644 --- a/src/services/coder.ts +++ b/src/services/coder.ts @@ -7,12 +7,14 @@ export class CoderClient { private token: string; private templateId: string; private anthropicApiKey: string; + private callbackUrl: string; constructor(config: Config) { this.baseUrl = config.coderUrl.replace(/\/$/, ""); this.token = config.coderToken; this.templateId = config.coderTemplateId; this.anthropicApiKey = config.anthropicApiKey; + this.callbackUrl = config.callbackUrl; } async createWorkspace(task: TaskRequest): Promise<{ id: string; name: string }> { @@ -29,6 +31,7 @@ export class CoderClient { { name: "gitea_org", value: task.giteaOrg }, { name: "gitea_repo", value: task.giteaRepo }, { name: "anthropic_api_key", value: this.anthropicApiKey }, + { name: "callback_url", value: `${this.callbackUrl}/webhook/task-complete/${name}` }, ...(task.deployEnv ? [{ name: "deploy_env", value: task.deployEnv }] : []), @@ -55,4 +58,57 @@ export class CoderClient { const data = (await res.json()) as { id: string; name: string }; return { id: data.id, name: data.name }; } + + async stopWorkspace(workspaceId: string): Promise { + 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 { + 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 } | 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 }> }; + const match = data.workspaces?.find((w) => w.name === name); + return match ? { id: match.id } : null; + } }