Add workspace lifecycle management via task-complete callback
Publish Image / publish (push) Successful in 20s
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>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
f23fb4e2e0
commit
7145fb536a
@@ -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"
|
||||||
@@ -443,6 +452,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
|
||||||
@@ -519,6 +529,13 @@ resource "coder_agent" "main" {
|
|||||||
claude --print --dangerously-skip-permissions "/project:$TASK_TYPE $ARGS" 2>&1 | tee ~/task-output.log
|
claude --print --dangerously-skip-permissions "/project:$TASK_TYPE $ARGS" 2>&1 | tee ~/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
|
||||||
}
|
}
|
||||||
@@ -657,6 +674,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"
|
||||||
|
|||||||
@@ -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,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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-pr-review-rework", prRework(config, queue));
|
||||||
app.post("/webhook/gitea-release", release(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)
|
// Scheduled maintenance cron (Monday 9:00 AM)
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
+75
-7
@@ -7,16 +7,24 @@ interface QueuedTask {
|
|||||||
comment?: { useReviewAccount: boolean; body: string };
|
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
|
* - 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)
|
* - Concurrency is configurable (defaults to 2)
|
||||||
|
* - Tracks active workspaces and stops them on task-complete callback
|
||||||
*/
|
*/
|
||||||
export class TaskQueue {
|
export class TaskQueue {
|
||||||
private pending = new Map<string, QueuedTask>();
|
private pending = new Map<string, QueuedTask>();
|
||||||
private running = new Set<string>();
|
private running = new Set<string>();
|
||||||
|
private active = new Map<string, ActiveWorkspace>();
|
||||||
private concurrency: number;
|
private concurrency: number;
|
||||||
private coder: CoderClient;
|
private coder: CoderClient;
|
||||||
private gitea: GiteaClient;
|
private gitea: GiteaClient;
|
||||||
@@ -46,6 +54,11 @@ export class TaskQueue {
|
|||||||
): boolean {
|
): boolean {
|
||||||
const key = TaskQueue.key(task);
|
const key = TaskQueue.key(task);
|
||||||
|
|
||||||
|
if (this.active.has(key)) {
|
||||||
|
console.log(`[queue] dropped (active workspace): ${key}`);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
if (this.running.has(key)) {
|
if (this.running.has(key)) {
|
||||||
console.log(`[queue] dropped (running): ${key}`);
|
console.log(`[queue] dropped (running): ${key}`);
|
||||||
return false;
|
return false;
|
||||||
@@ -58,16 +71,50 @@ export class TaskQueue {
|
|||||||
|
|
||||||
this.pending.set(key, { task, comment });
|
this.pending.set(key, { task, comment });
|
||||||
console.log(
|
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();
|
this.drain();
|
||||||
return true;
|
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. */
|
/** Process queued tasks up to the concurrency limit. */
|
||||||
private drain(): void {
|
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!;
|
const [key, entry] = this.pending.entries().next().value!;
|
||||||
this.pending.delete(key);
|
this.pending.delete(key);
|
||||||
this.running.add(key);
|
this.running.add(key);
|
||||||
@@ -81,7 +128,14 @@ export class TaskQueue {
|
|||||||
try {
|
try {
|
||||||
console.log(`[queue] processing: ${key}`);
|
console.log(`[queue] processing: ${key}`);
|
||||||
const workspace = await this.coder.createWorkspace(task);
|
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) {
|
if (comment) {
|
||||||
await this.gitea.commentOnIssue(
|
await this.gitea.commentOnIssue(
|
||||||
@@ -96,15 +150,29 @@ export class TaskQueue {
|
|||||||
console.error(`[queue] failed: ${key}`, err);
|
console.error(`[queue] failed: ${key}`, err);
|
||||||
} finally {
|
} finally {
|
||||||
this.running.delete(key);
|
this.running.delete(key);
|
||||||
this.drain();
|
// Don't drain here — active workspaces count toward concurrency
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Current queue status for health/debug endpoints. */
|
/** Current queue status for health/debug endpoints. */
|
||||||
status(): { pending: string[]; running: string[]; concurrency: number } {
|
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 {
|
return {
|
||||||
pending: [...this.pending.keys()],
|
pending: [...this.pending.keys()],
|
||||||
running: [...this.running.keys()],
|
running: [...this.running.keys()],
|
||||||
|
active: activeMap,
|
||||||
concurrency: this.concurrency,
|
concurrency: this.concurrency,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,12 +7,14 @@ export class CoderClient {
|
|||||||
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 }> {
|
||||||
@@ -29,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 }]
|
||||||
: []),
|
: []),
|
||||||
@@ -55,4 +58,57 @@ 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 } | 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;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user