Compare commits
2
Commits
v1.0.1
..
7145fb536a
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7145fb536a | ||
|
|
f23fb4e2e0 |
@@ -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"
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1,14 +1,14 @@
|
||||
import type { Context } from "hono";
|
||||
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";
|
||||
|
||||
/**
|
||||
* 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
|
||||
*/
|
||||
export function issueComment(config: Config, coder: CoderClient) {
|
||||
export function issueComment(config: Config, queue: TaskQueue) {
|
||||
return async (c: Context) => {
|
||||
const event = await c.req.json<IssueCommentEvent>();
|
||||
|
||||
@@ -43,11 +43,8 @@ export function issueComment(config: Config, coder: CoderClient) {
|
||||
giteaToken: config.giteaDevToken,
|
||||
};
|
||||
|
||||
console.log(`[issue-comment] #${task.issueNumber} → re-trigger analyst`);
|
||||
const queued = queue.enqueue(task);
|
||||
|
||||
const workspace = await coder.createWorkspace(task);
|
||||
console.log(`[issue-comment] workspace created: ${workspace.name}`);
|
||||
|
||||
return c.json({ ok: true, workspace: workspace.name });
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import type { Context } from "hono";
|
||||
import type { IssueEvent, TaskRequest } from "../types.js";
|
||||
import type { CoderClient } from "../services/coder.js";
|
||||
import type { GiteaClient } from "../services/gitea.js";
|
||||
import type { TaskQueue } from "../queue.js";
|
||||
import type { Config } from "../config.js";
|
||||
|
||||
const LABEL_TO_TASK: Record<string, string> = {
|
||||
@@ -19,7 +18,7 @@ const REVIEW_TASKS = new Set(["test"]);
|
||||
*
|
||||
* 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) => {
|
||||
const event = await c.req.json<IssueEvent>();
|
||||
|
||||
@@ -55,18 +54,11 @@ export function issueLabel(config: Config, coder: CoderClient, gitea: GiteaClien
|
||||
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);
|
||||
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 });
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import type { Context } from "hono";
|
||||
import type { IssueEvent, TaskRequest } from "../types.js";
|
||||
import type { CoderClient } from "../services/coder.js";
|
||||
import type { GiteaClient } from "../services/gitea.js";
|
||||
import type { TaskQueue } from "../queue.js";
|
||||
import type { Config } from "../config.js";
|
||||
|
||||
/**
|
||||
@@ -9,7 +8,7 @@ import type { Config } from "../config.js";
|
||||
*
|
||||
* 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) => {
|
||||
const event = await c.req.json<IssueEvent>();
|
||||
|
||||
@@ -32,18 +31,11 @@ export function issueTriage(config: Config, coder: CoderClient, gitea: GiteaClie
|
||||
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);
|
||||
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 });
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
}
|
||||
|
||||
@@ -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 { 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
|
||||
*/
|
||||
export function runMaintenance(
|
||||
config: Config,
|
||||
coder: CoderClient,
|
||||
repos: MaintenanceRepo[],
|
||||
) {
|
||||
export function runMaintenance(config: Config, queue: TaskQueue, repos: MaintenanceRepo[]) {
|
||||
return async () => {
|
||||
for (const entry of repos) {
|
||||
const task: TaskRequest = {
|
||||
@@ -29,17 +25,7 @@ export function runMaintenance(
|
||||
giteaToken: config.giteaDevToken,
|
||||
};
|
||||
|
||||
try {
|
||||
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,
|
||||
);
|
||||
}
|
||||
queue.enqueue(task);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { Context } from "hono";
|
||||
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 { Config } from "../config.js";
|
||||
|
||||
@@ -10,7 +10,7 @@ import type { Config } from "../config.js";
|
||||
*
|
||||
* 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) => {
|
||||
const event = await c.req.json<PullRequestEvent>();
|
||||
|
||||
@@ -40,11 +40,8 @@ export function prReview(config: Config, coder: CoderClient, gitea: GiteaClient)
|
||||
giteaToken: config.giteaReviewToken,
|
||||
};
|
||||
|
||||
console.log(`[pr-review] PR #${prNumber} → review-code`);
|
||||
const queued = queue.enqueue(task);
|
||||
|
||||
const workspace = await coder.createWorkspace(task);
|
||||
console.log(`[pr-review] workspace created: ${workspace.name}`);
|
||||
|
||||
return c.json({ ok: true, workspace: workspace.name });
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
}
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import type { Context } from "hono";
|
||||
import type { PullRequestReviewEvent, TaskRequest } from "../types.js";
|
||||
import type { CoderClient } from "../services/coder.js";
|
||||
import type { GiteaClient } from "../services/gitea.js";
|
||||
import type { TaskQueue } from "../queue.js";
|
||||
import type { Config } from "../config.js";
|
||||
|
||||
/**
|
||||
@@ -9,7 +8,7 @@ import type { Config } from "../config.js";
|
||||
*
|
||||
* 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) => {
|
||||
const event = await c.req.json<PullRequestReviewEvent>();
|
||||
|
||||
@@ -43,18 +42,11 @@ export function prRework(config: Config, coder: CoderClient, gitea: GiteaClient)
|
||||
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);
|
||||
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 });
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
}
|
||||
|
||||
+7
-15
@@ -1,7 +1,6 @@
|
||||
import type { Context } from "hono";
|
||||
import type { IssueEvent, TaskRequest } from "../types.js";
|
||||
import type { CoderClient } from "../services/coder.js";
|
||||
import type { GiteaClient } from "../services/gitea.js";
|
||||
import type { TaskQueue } from "../queue.js";
|
||||
import type { Config } from "../config.js";
|
||||
|
||||
/**
|
||||
@@ -10,7 +9,7 @@ import type { Config } from "../config.js";
|
||||
*
|
||||
* Replaces: release-process.json
|
||||
*/
|
||||
export function release(config: Config, coder: CoderClient, gitea: GiteaClient) {
|
||||
export function release(config: Config, queue: TaskQueue) {
|
||||
return async (c: Context) => {
|
||||
const event = await c.req.json<IssueEvent>();
|
||||
|
||||
@@ -52,18 +51,11 @@ export function release(config: Config, coder: CoderClient, gitea: GiteaClient)
|
||||
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);
|
||||
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 });
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
}
|
||||
|
||||
+25
-9
@@ -6,6 +6,7 @@ import cron from "node-cron";
|
||||
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";
|
||||
@@ -23,26 +24,41 @@ const config = loadConfig();
|
||||
const coderClient = new CoderClient(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();
|
||||
app.use("*", logger());
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Health check
|
||||
// Health check + queue status
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
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
|
||||
// can stay unchanged (just point to new host).
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
app.post("/webhook/gitea-issue-triage", issueTriage(config, coderClient, giteaClient));
|
||||
app.post("/webhook/gitea-issue-label", issueLabel(config, coderClient, giteaClient));
|
||||
app.post("/webhook/gitea-issue-comment", issueComment(config, coderClient));
|
||||
app.post("/webhook/gitea-pr-review", prReview(config, coderClient, giteaClient));
|
||||
app.post("/webhook/gitea-pr-review-rework", prRework(config, coderClient, giteaClient));
|
||||
app.post("/webhook/gitea-release", release(config, coderClient, giteaClient));
|
||||
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 });
|
||||
});
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Scheduled maintenance cron (Monday 9:00 AM)
|
||||
@@ -50,7 +66,7 @@ app.post("/webhook/gitea-release", release(config, coderClient, giteaClient));
|
||||
|
||||
const maintenanceRepos = parseMaintenanceRepos();
|
||||
if (maintenanceRepos.length > 0) {
|
||||
const maintenanceFn = runMaintenance(config, coderClient, maintenanceRepos);
|
||||
const maintenanceFn = runMaintenance(config, queue, maintenanceRepos);
|
||||
|
||||
cron.schedule("0 9 * * 1", () => {
|
||||
console.log("[cron] running weekly maintenance");
|
||||
@@ -69,7 +85,7 @@ if (maintenanceRepos.length > 0) {
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
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
@@ -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,
|
||||
};
|
||||
}
|
||||
}
|
||||
+58
-5
@@ -1,25 +1,24 @@
|
||||
import type { Config } from "../config.js";
|
||||
import type { TaskRequest } from "../types.js";
|
||||
import { TaskQueue } from "../queue.js";
|
||||
|
||||
export class CoderClient {
|
||||
private baseUrl: string;
|
||||
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 }> {
|
||||
const timestamp = new Date()
|
||||
.toISOString()
|
||||
.replace(/[-:T]/g, "")
|
||||
.slice(8, 14); // HHmmss
|
||||
const name = `${task.taskType}-${task.issueNumber}-${timestamp}`;
|
||||
const name = TaskQueue.workspaceName(task);
|
||||
|
||||
const body = {
|
||||
name,
|
||||
@@ -32,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 }]
|
||||
: []),
|
||||
@@ -58,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<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