Compare commits

...
10 Commits
Author SHA1 Message Date
Daniel SamsonandClaude Opus 4.6 546d2615c7 feat: re-trigger tester when human comments on a PR
Publish Image / publish (push) Successful in 22s
PR comments fire as issue_comment webhooks in Gitea. Previously the
handler only re-triggered the analyst for regular issues. Now it
detects PR comments via the pull_request field and queues a test
workspace using the review bot account.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 20:57:46 +01:00
Daniel SamsonandClaude Opus 4.6 f693d56762 Handle Gitea review webhook type field alongside state
Publish Image / publish (push) Successful in 23s
Gitea webhook payloads use "type" (e.g. "pull_request_review_rejected")
instead of "state" (e.g. "REQUEST_CHANGES"). Map both formats so
the pr-rework handler works with all Gitea webhook variations.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 18:13:51 +01:00
Daniel SamsonandClaude Opus 4.6 a72ea9ad1d Fix BullMQ job dedup — remove completed/failed jobs
Publish Image / publish (push) Successful in 21s
Completed jobs with the same jobId block new jobs from being added.
Add removeOnComplete and removeOnFail to all job options so the same
task can be re-triggered after completion.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 16:37:52 +01:00
Daniel SamsonandClaude Opus 4.6 2cdf68ddff Fix pr-rework crash — guard against missing review.state
Publish Image / publish (push) Successful in 23s
Gitea sends review objects without a state field in some cases.
Use optional chaining and log the review payload for debugging.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 16:26:58 +01:00
Daniel SamsonandClaude Opus 4.6 81aee2f7fc Route spec PR review feedback to rework-spec instead of rework-pr
Publish Image / publish (push) Successful in 23s
When a review requests changes on a docs/ branch (architect PR), trigger
the rework-spec task type instead of rework-pr. This re-runs the architect
persona to fix the design docs based on review feedback.

- Add rework-spec to Coder template task_type options and lightweight stages
- Detect architect PRs via docs/ branch prefix in pr-rework handler

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 16:18:00 +01:00
Daniel SamsonandClaude Opus 4.6 f495ea391f Fix issue-label and pr-rework webhook handlers
Publish Image / publish (push) Successful in 22s
- issue-label: Accept action "created" from Gitea (in addition to
  "labeled" and "label_updated") — Gitea sends "created" for new labels
- pr-rework: Guard against missing review object to prevent TypeError
  crash when Gitea sends review webhook without review data

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 16:12:04 +01:00
Daniel SamsonandClaude Opus 4.6 d19798dfea Add debug logging to issue-label handler
Publish Image / publish (push) Successful in 23s
Log the webhook action and labels to diagnose why label events are
being ignored.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 15:56:19 +01:00
Daniel SamsonandClaude Opus 4.6 0812a9be37 Fix workspace deletion — use build transition API
Publish Image / publish (push) Successful in 35s
The plain DELETE endpoint returns 405 for stopped workspaces. Use
POST /workspaces/:id/builds with {"transition":"delete"} instead,
which is the correct Coder API for workspace deletion.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 15:51:24 +01:00
Daniel SamsonandClaude Opus 4.6 66649022ca Fix BullMQ jobId format — colons are not allowed
Publish Image / publish (push) Successful in 26s
BullMQ throws "Custom Id cannot contain :" when jobId includes colons.
Replace colon separators with dashes in all jobId values.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-12 15:48:47 +01:00
Daniel SamsonandClaude Opus 4.6 13e96e9356 Migrate task queue from in-memory to BullMQ + Redis
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>
2026-04-12 15:40:12 +01:00
18 changed files with 785 additions and 190 deletions
+5 -1
View File
@@ -128,6 +128,10 @@ data "coder_parameter" "task_type" {
name = "Rework PR" name = "Rework PR"
value = "rework-pr" value = "rework-pr"
} }
option {
name = "Rework Spec"
value = "rework-spec"
}
option { option {
name = "Release" name = "Release"
value = "release" value = "release"
@@ -470,7 +474,7 @@ resource "coder_agent" "main" {
git fetch origin develop 2>/dev/null && git checkout develop 2>/dev/null || true git fetch origin develop 2>/dev/null && git checkout develop 2>/dev/null || true
# Lightweight stages only need the code, not a full build # Lightweight stages only need the code, not a full build
LIGHT_STAGES="analyse architect release maintenance" LIGHT_STAGES="analyse architect rework-spec release maintenance"
if echo "$LIGHT_STAGES" | grep -qw "$TASK_TYPE"; then if echo "$LIGHT_STAGES" | grep -qw "$TASK_TYPE"; then
echo "Lightweight stage ($TASK_TYPE) — skipping build" echo "Lightweight stage ($TASK_TYPE) — skipping build"
else else
+69
View File
@@ -0,0 +1,69 @@
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: redis-data
namespace: sdlc-orchestrator
spec:
accessModes:
- ReadWriteOnce
resources:
requests:
storage: 1Gi
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: redis
namespace: sdlc-orchestrator
spec:
replicas: 1
selector:
matchLabels:
app: redis
template:
metadata:
labels:
app: redis
spec:
containers:
- name: redis
image: redis:7-alpine
args: ["--appendonly", "yes"]
ports:
- containerPort: 6379
volumeMounts:
- name: data
mountPath: /data
resources:
requests:
cpu: 50m
memory: 64Mi
limits:
cpu: 200m
memory: 256Mi
livenessProbe:
exec:
command: ["redis-cli", "ping"]
initialDelaySeconds: 5
periodSeconds: 30
readinessProbe:
exec:
command: ["redis-cli", "ping"]
initialDelaySeconds: 3
periodSeconds: 10
volumes:
- name: data
persistentVolumeClaim:
claimName: redis-data
---
apiVersion: v1
kind: Service
metadata:
name: redis
namespace: sdlc-orchestrator
spec:
selector:
app: redis
ports:
- port: 6379
targetPort: 6379
+3
View File
@@ -15,8 +15,11 @@ stringData:
GITEA_DEV_TOKEN: "<gitea-claude-dev-token>" GITEA_DEV_TOKEN: "<gitea-claude-dev-token>"
GITEA_REVIEW_TOKEN: "<gitea-claude-review-token>" GITEA_REVIEW_TOKEN: "<gitea-claude-review-token>"
ANTHROPIC_API_KEY: "<anthropic-api-key>" ANTHROPIC_API_KEY: "<anthropic-api-key>"
CLAUDE_OAUTH_TOKEN: "<claude-code-max-oauth-token>"
GITEA_URL: "https://gitea.samson.media" GITEA_URL: "https://gitea.samson.media"
BOT_DEV_USERNAME: "claude-dev" BOT_DEV_USERNAME: "claude-dev"
BOT_REVIEW_USERNAME: "claude-review" BOT_REVIEW_USERNAME: "claude-review"
CALLBACK_URL: "https://sdlc.samson.media"
REDIS_URL: "redis://redis.sdlc-orchestrator.svc.cluster.local:6379"
# Comma-separated: org/repo=clone_url # Comma-separated: org/repo=clone_url
MAINTENANCE_REPOS: "claude/babble=https://gitea.samson.media/claude/babble.git" MAINTENANCE_REPOS: "claude/babble=https://gitea.samson.media/claude/babble.git"
+293 -29
View File
@@ -9,12 +9,12 @@
"version": "1.0.0", "version": "1.0.0",
"dependencies": { "dependencies": {
"@hono/node-server": "^1.13.8", "@hono/node-server": "^1.13.8",
"bullmq": "^5.34.8",
"hono": "^4.7.6", "hono": "^4.7.6",
"node-cron": "^3.0.3" "ioredis": "^5.6.1"
}, },
"devDependencies": { "devDependencies": {
"@types/node": "^22.15.3", "@types/node": "^22.15.3",
"@types/node-cron": "^3.0.11",
"eslint": "^9.25.1", "eslint": "^9.25.1",
"prettier": "^3.5.3", "prettier": "^3.5.3",
"tsx": "^4.19.4", "tsx": "^4.19.4",
@@ -671,6 +671,90 @@
"url": "https://github.com/sponsors/nzakas" "url": "https://github.com/sponsors/nzakas"
} }
}, },
"node_modules/@ioredis/commands": {
"version": "1.5.1",
"resolved": "https://registry.npmjs.org/@ioredis/commands/-/commands-1.5.1.tgz",
"integrity": "sha512-JH8ZL/ywcJyR9MmJ5BNqZllXNZQqQbnVZOqpPQqE1vHiFgAw4NHbvE0FOduNU8IX9babitBT46571OnPTT0Zcw==",
"license": "MIT"
},
"node_modules/@msgpackr-extract/msgpackr-extract-darwin-arm64": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/@msgpackr-extract/msgpackr-extract-darwin-arm64/-/msgpackr-extract-darwin-arm64-3.0.3.tgz",
"integrity": "sha512-QZHtlVgbAdy2zAqNA9Gu1UpIuI8Xvsd1v8ic6B2pZmeFnFcMWiPLfWXh7TVw4eGEZ/C9TH281KwhVoeQUKbyjw==",
"cpu": [
"arm64"
],
"license": "MIT",
"optional": true,
"os": [
"darwin"
]
},
"node_modules/@msgpackr-extract/msgpackr-extract-darwin-x64": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/@msgpackr-extract/msgpackr-extract-darwin-x64/-/msgpackr-extract-darwin-x64-3.0.3.tgz",
"integrity": "sha512-mdzd3AVzYKuUmiWOQ8GNhl64/IoFGol569zNRdkLReh6LRLHOXxU4U8eq0JwaD8iFHdVGqSy4IjFL4reoWCDFw==",
"cpu": [
"x64"
],
"license": "MIT",
"optional": true,
"os": [
"darwin"
]
},
"node_modules/@msgpackr-extract/msgpackr-extract-linux-arm": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/@msgpackr-extract/msgpackr-extract-linux-arm/-/msgpackr-extract-linux-arm-3.0.3.tgz",
"integrity": "sha512-fg0uy/dG/nZEXfYilKoRe7yALaNmHoYeIoJuJ7KJ+YyU2bvY8vPv27f7UKhGRpY6euFYqEVhxCFZgAUNQBM3nw==",
"cpu": [
"arm"
],
"license": "MIT",
"optional": true,
"os": [
"linux"
]
},
"node_modules/@msgpackr-extract/msgpackr-extract-linux-arm64": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/@msgpackr-extract/msgpackr-extract-linux-arm64/-/msgpackr-extract-linux-arm64-3.0.3.tgz",
"integrity": "sha512-YxQL+ax0XqBJDZiKimS2XQaf+2wDGVa1enVRGzEvLLVFeqa5kx2bWbtcSXgsxjQB7nRqqIGFIcLteF/sHeVtQg==",
"cpu": [
"arm64"
],
"license": "MIT",
"optional": true,
"os": [
"linux"
]
},
"node_modules/@msgpackr-extract/msgpackr-extract-linux-x64": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/@msgpackr-extract/msgpackr-extract-linux-x64/-/msgpackr-extract-linux-x64-3.0.3.tgz",
"integrity": "sha512-cvwNfbP07pKUfq1uH+S6KJ7dT9K8WOE4ZiAcsrSes+UY55E/0jLYc+vq+DO7jlmqRb5zAggExKm0H7O/CBaesg==",
"cpu": [
"x64"
],
"license": "MIT",
"optional": true,
"os": [
"linux"
]
},
"node_modules/@msgpackr-extract/msgpackr-extract-win32-x64": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/@msgpackr-extract/msgpackr-extract-win32-x64/-/msgpackr-extract-win32-x64-3.0.3.tgz",
"integrity": "sha512-x0fWaQtYp4E6sktbsdAqnehxDgEc/VwM7uLsRCYWaiGu0ykYdZPiS8zCWdnjHwyiumousxfBm4SO31eXqwEZhQ==",
"cpu": [
"x64"
],
"license": "MIT",
"optional": true,
"os": [
"win32"
]
},
"node_modules/@types/estree": { "node_modules/@types/estree": {
"version": "1.0.8", "version": "1.0.8",
"resolved": "https://registry.npmjs.org/@types/estree/-/estree-1.0.8.tgz", "resolved": "https://registry.npmjs.org/@types/estree/-/estree-1.0.8.tgz",
@@ -695,13 +779,6 @@
"undici-types": "~6.21.0" "undici-types": "~6.21.0"
} }
}, },
"node_modules/@types/node-cron": {
"version": "3.0.11",
"resolved": "https://registry.npmjs.org/@types/node-cron/-/node-cron-3.0.11.tgz",
"integrity": "sha512-0ikrnug3/IyneSHqCBeslAhlK2aBfYek1fGo4bP4QnZPmiqSGRK+Oy7ZMisLWkesffJvQ1cqAcBnJC+8+nxIAg==",
"dev": true,
"license": "MIT"
},
"node_modules/acorn": { "node_modules/acorn": {
"version": "8.16.0", "version": "8.16.0",
"resolved": "https://registry.npmjs.org/acorn/-/acorn-8.16.0.tgz", "resolved": "https://registry.npmjs.org/acorn/-/acorn-8.16.0.tgz",
@@ -783,6 +860,34 @@
"concat-map": "0.0.1" "concat-map": "0.0.1"
} }
}, },
"node_modules/bullmq": {
"version": "5.73.4",
"resolved": "https://registry.npmjs.org/bullmq/-/bullmq-5.73.4.tgz",
"integrity": "sha512-Q+NeFLtdKSD3GDPYSX4pH+Mc9E4OZVKimXwrnZ5WmndNy31COMy4vQV9zfhgfHGSUFrlpsBicfKYbSjx9FbO+A==",
"license": "MIT",
"dependencies": {
"cron-parser": "4.9.0",
"ioredis": "5.10.1",
"msgpackr": "1.11.5",
"node-abort-controller": "3.1.1",
"semver": "7.7.4",
"tslib": "2.8.1",
"uuid": "11.1.0"
}
},
"node_modules/bullmq/node_modules/uuid": {
"version": "11.1.0",
"resolved": "https://registry.npmjs.org/uuid/-/uuid-11.1.0.tgz",
"integrity": "sha512-0/A9rDy9P7cJ+8w1c9WD9V//9Wj15Ce2MPz8Ri6032usz+NfePxx5AcN3bN+r6ZL6jEo066/yNYB3tn4pQEx+A==",
"funding": [
"https://github.com/sponsors/broofa",
"https://github.com/sponsors/ctavan"
],
"license": "MIT",
"bin": {
"uuid": "dist/esm/bin/uuid"
}
},
"node_modules/callsites": { "node_modules/callsites": {
"version": "3.1.0", "version": "3.1.0",
"resolved": "https://registry.npmjs.org/callsites/-/callsites-3.1.0.tgz", "resolved": "https://registry.npmjs.org/callsites/-/callsites-3.1.0.tgz",
@@ -810,6 +915,15 @@
"url": "https://github.com/chalk/chalk?sponsor=1" "url": "https://github.com/chalk/chalk?sponsor=1"
} }
}, },
"node_modules/cluster-key-slot": {
"version": "1.1.2",
"resolved": "https://registry.npmjs.org/cluster-key-slot/-/cluster-key-slot-1.1.2.tgz",
"integrity": "sha512-RMr0FhtfXemyinomL4hrWcYJxmX6deFdCxpJzhDttxgO1+bcCnkk+9drydLVDmAMG7NE6aN/fl4F7ucU/90gAA==",
"license": "Apache-2.0",
"engines": {
"node": ">=0.10.0"
}
},
"node_modules/color-convert": { "node_modules/color-convert": {
"version": "2.0.1", "version": "2.0.1",
"resolved": "https://registry.npmjs.org/color-convert/-/color-convert-2.0.1.tgz", "resolved": "https://registry.npmjs.org/color-convert/-/color-convert-2.0.1.tgz",
@@ -837,6 +951,18 @@
"dev": true, "dev": true,
"license": "MIT" "license": "MIT"
}, },
"node_modules/cron-parser": {
"version": "4.9.0",
"resolved": "https://registry.npmjs.org/cron-parser/-/cron-parser-4.9.0.tgz",
"integrity": "sha512-p0SaNjrHOnQeR8/VnfGbmg9te2kfyYSQ7Sc/j/6DtPL3JQvKxmjO9TSjNFpujqV3vEYYBvNNvXSxzyksBWAx1Q==",
"license": "MIT",
"dependencies": {
"luxon": "^3.2.1"
},
"engines": {
"node": ">=12.0.0"
}
},
"node_modules/cross-spawn": { "node_modules/cross-spawn": {
"version": "7.0.6", "version": "7.0.6",
"resolved": "https://registry.npmjs.org/cross-spawn/-/cross-spawn-7.0.6.tgz", "resolved": "https://registry.npmjs.org/cross-spawn/-/cross-spawn-7.0.6.tgz",
@@ -856,7 +982,6 @@
"version": "4.4.3", "version": "4.4.3",
"resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz", "resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz",
"integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==", "integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==",
"dev": true,
"license": "MIT", "license": "MIT",
"dependencies": { "dependencies": {
"ms": "^2.1.3" "ms": "^2.1.3"
@@ -877,6 +1002,25 @@
"dev": true, "dev": true,
"license": "MIT" "license": "MIT"
}, },
"node_modules/denque": {
"version": "2.1.0",
"resolved": "https://registry.npmjs.org/denque/-/denque-2.1.0.tgz",
"integrity": "sha512-HVQE3AAb/pxF8fQAoiqpvg9i3evqug3hoiwakOyZAwJm+6vZehbkYXZ0l4JxS+I3QxM97v5aaRNhj8v5oBhekw==",
"license": "Apache-2.0",
"engines": {
"node": ">=0.10"
}
},
"node_modules/detect-libc": {
"version": "2.1.2",
"resolved": "https://registry.npmjs.org/detect-libc/-/detect-libc-2.1.2.tgz",
"integrity": "sha512-Btj2BOOO83o3WyH59e8MgXsxEQVcarkUOpEYrubB0urwnN10yQ364rsiByU11nZlqWYZm05i/of7io4mzihBtQ==",
"license": "Apache-2.0",
"optional": true,
"engines": {
"node": ">=8"
}
},
"node_modules/esbuild": { "node_modules/esbuild": {
"version": "0.27.7", "version": "0.27.7",
"resolved": "https://registry.npmjs.org/esbuild/-/esbuild-0.27.7.tgz", "resolved": "https://registry.npmjs.org/esbuild/-/esbuild-0.27.7.tgz",
@@ -1268,6 +1412,30 @@
"node": ">=0.8.19" "node": ">=0.8.19"
} }
}, },
"node_modules/ioredis": {
"version": "5.10.1",
"resolved": "https://registry.npmjs.org/ioredis/-/ioredis-5.10.1.tgz",
"integrity": "sha512-HuEDBTI70aYdx1v6U97SbNx9F1+svQKBDo30o0b9fw055LMepzpOOd0Ccg9Q6tbqmBSJaMuY0fB7yw9/vjBYCA==",
"license": "MIT",
"dependencies": {
"@ioredis/commands": "1.5.1",
"cluster-key-slot": "^1.1.0",
"debug": "^4.3.4",
"denque": "^2.1.0",
"lodash.defaults": "^4.2.0",
"lodash.isarguments": "^3.1.0",
"redis-errors": "^1.2.0",
"redis-parser": "^3.0.0",
"standard-as-callback": "^2.1.0"
},
"engines": {
"node": ">=12.22.0"
},
"funding": {
"type": "opencollective",
"url": "https://opencollective.com/ioredis"
}
},
"node_modules/is-extglob": { "node_modules/is-extglob": {
"version": "2.1.1", "version": "2.1.1",
"resolved": "https://registry.npmjs.org/is-extglob/-/is-extglob-2.1.1.tgz", "resolved": "https://registry.npmjs.org/is-extglob/-/is-extglob-2.1.1.tgz",
@@ -1372,6 +1540,18 @@
"url": "https://github.com/sponsors/sindresorhus" "url": "https://github.com/sponsors/sindresorhus"
} }
}, },
"node_modules/lodash.defaults": {
"version": "4.2.0",
"resolved": "https://registry.npmjs.org/lodash.defaults/-/lodash.defaults-4.2.0.tgz",
"integrity": "sha512-qjxPLHd3r5DnsdGacqOMU6pb/avJzdh9tFX2ymgoZE27BmjXrNy/y4LoaiTeAb+O3gL8AfpJGtqfX/ae2leYYQ==",
"license": "MIT"
},
"node_modules/lodash.isarguments": {
"version": "3.1.0",
"resolved": "https://registry.npmjs.org/lodash.isarguments/-/lodash.isarguments-3.1.0.tgz",
"integrity": "sha512-chi4NHZlZqZD18a0imDHnZPrDeBbTtVN7GXMwuGdRH9qotxAjYs3aVLKc7zNOG9eddR5Ksd8rvFEBc9SsggPpg==",
"license": "MIT"
},
"node_modules/lodash.merge": { "node_modules/lodash.merge": {
"version": "4.6.2", "version": "4.6.2",
"resolved": "https://registry.npmjs.org/lodash.merge/-/lodash.merge-4.6.2.tgz", "resolved": "https://registry.npmjs.org/lodash.merge/-/lodash.merge-4.6.2.tgz",
@@ -1379,6 +1559,15 @@
"dev": true, "dev": true,
"license": "MIT" "license": "MIT"
}, },
"node_modules/luxon": {
"version": "3.7.2",
"resolved": "https://registry.npmjs.org/luxon/-/luxon-3.7.2.tgz",
"integrity": "sha512-vtEhXh/gNjI9Yg1u4jX/0YVPMvxzHuGgCm6tC5kZyb08yjGWGnqAjGJvcXbqQR2P3MyMEFnRbpcdFS6PBcLqew==",
"license": "MIT",
"engines": {
"node": ">=12"
}
},
"node_modules/minimatch": { "node_modules/minimatch": {
"version": "3.1.5", "version": "3.1.5",
"resolved": "https://registry.npmjs.org/minimatch/-/minimatch-3.1.5.tgz", "resolved": "https://registry.npmjs.org/minimatch/-/minimatch-3.1.5.tgz",
@@ -1396,9 +1585,39 @@
"version": "2.1.3", "version": "2.1.3",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz",
"integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==",
"dev": true,
"license": "MIT" "license": "MIT"
}, },
"node_modules/msgpackr": {
"version": "1.11.5",
"resolved": "https://registry.npmjs.org/msgpackr/-/msgpackr-1.11.5.tgz",
"integrity": "sha512-UjkUHN0yqp9RWKy0Lplhh+wlpdt9oQBYgULZOiFhV3VclSF1JnSQWZ5r9gORQlNYaUKQoR8itv7g7z1xDDuACA==",
"license": "MIT",
"optionalDependencies": {
"msgpackr-extract": "^3.0.2"
}
},
"node_modules/msgpackr-extract": {
"version": "3.0.3",
"resolved": "https://registry.npmjs.org/msgpackr-extract/-/msgpackr-extract-3.0.3.tgz",
"integrity": "sha512-P0efT1C9jIdVRefqjzOQ9Xml57zpOXnIuS+csaB4MdZbTdmGDLo8XhzBG1N7aO11gKDDkJvBLULeFTo46wwreA==",
"hasInstallScript": true,
"license": "MIT",
"optional": true,
"dependencies": {
"node-gyp-build-optional-packages": "5.2.2"
},
"bin": {
"download-msgpackr-prebuilds": "bin/download-prebuilds.js"
},
"optionalDependencies": {
"@msgpackr-extract/msgpackr-extract-darwin-arm64": "3.0.3",
"@msgpackr-extract/msgpackr-extract-darwin-x64": "3.0.3",
"@msgpackr-extract/msgpackr-extract-linux-arm": "3.0.3",
"@msgpackr-extract/msgpackr-extract-linux-arm64": "3.0.3",
"@msgpackr-extract/msgpackr-extract-linux-x64": "3.0.3",
"@msgpackr-extract/msgpackr-extract-win32-x64": "3.0.3"
}
},
"node_modules/natural-compare": { "node_modules/natural-compare": {
"version": "1.4.0", "version": "1.4.0",
"resolved": "https://registry.npmjs.org/natural-compare/-/natural-compare-1.4.0.tgz", "resolved": "https://registry.npmjs.org/natural-compare/-/natural-compare-1.4.0.tgz",
@@ -1406,16 +1625,25 @@
"dev": true, "dev": true,
"license": "MIT" "license": "MIT"
}, },
"node_modules/node-cron": { "node_modules/node-abort-controller": {
"version": "3.0.3", "version": "3.1.1",
"resolved": "https://registry.npmjs.org/node-cron/-/node-cron-3.0.3.tgz", "resolved": "https://registry.npmjs.org/node-abort-controller/-/node-abort-controller-3.1.1.tgz",
"integrity": "sha512-dOal67//nohNgYWb+nWmg5dkFdIwDm8EpeGYMekPMrngV3637lqnX0lbUcCtgibHTz6SEz7DAIjKvKDFYCnO1A==", "integrity": "sha512-AGK2yQKIjRuqnc6VkX2Xj5d+QW8xZ87pa1UK6yA6ouUyuxfHuMP6umE5QK7UmTeOAymo+Zx1Fxiuw9rVx8taHQ==",
"license": "ISC", "license": "MIT"
"dependencies": {
"uuid": "8.3.2"
}, },
"engines": { "node_modules/node-gyp-build-optional-packages": {
"node": ">=6.0.0" "version": "5.2.2",
"resolved": "https://registry.npmjs.org/node-gyp-build-optional-packages/-/node-gyp-build-optional-packages-5.2.2.tgz",
"integrity": "sha512-s+w+rBWnpTMwSFbaE0UXsRlg7hU4FjekKU4eyAih5T8nJuNZT1nNsskXpxmeqSK9UzkBl6UgRlnKc8hz8IEqOw==",
"license": "MIT",
"optional": true,
"dependencies": {
"detect-libc": "^2.0.1"
},
"bin": {
"node-gyp-build-optional-packages": "bin.js",
"node-gyp-build-optional-packages-optional": "optional.js",
"node-gyp-build-optional-packages-test": "build-test.js"
} }
}, },
"node_modules/optionator": { "node_modules/optionator": {
@@ -1537,6 +1765,27 @@
"node": ">=6" "node": ">=6"
} }
}, },
"node_modules/redis-errors": {
"version": "1.2.0",
"resolved": "https://registry.npmjs.org/redis-errors/-/redis-errors-1.2.0.tgz",
"integrity": "sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w==",
"license": "MIT",
"engines": {
"node": ">=4"
}
},
"node_modules/redis-parser": {
"version": "3.0.0",
"resolved": "https://registry.npmjs.org/redis-parser/-/redis-parser-3.0.0.tgz",
"integrity": "sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==",
"license": "MIT",
"dependencies": {
"redis-errors": "^1.0.0"
},
"engines": {
"node": ">=4"
}
},
"node_modules/resolve-from": { "node_modules/resolve-from": {
"version": "4.0.0", "version": "4.0.0",
"resolved": "https://registry.npmjs.org/resolve-from/-/resolve-from-4.0.0.tgz", "resolved": "https://registry.npmjs.org/resolve-from/-/resolve-from-4.0.0.tgz",
@@ -1557,6 +1806,18 @@
"url": "https://github.com/privatenumber/resolve-pkg-maps?sponsor=1" "url": "https://github.com/privatenumber/resolve-pkg-maps?sponsor=1"
} }
}, },
"node_modules/semver": {
"version": "7.7.4",
"resolved": "https://registry.npmjs.org/semver/-/semver-7.7.4.tgz",
"integrity": "sha512-vFKC2IEtQnVhpT78h1Yp8wzwrf8CM+MzKMHGJZfBtzhZNycRFnXsHk6E5TxIkkMsgNS7mdX3AGB7x2QM2di4lA==",
"license": "ISC",
"bin": {
"semver": "bin/semver.js"
},
"engines": {
"node": ">=10"
}
},
"node_modules/shebang-command": { "node_modules/shebang-command": {
"version": "2.0.0", "version": "2.0.0",
"resolved": "https://registry.npmjs.org/shebang-command/-/shebang-command-2.0.0.tgz", "resolved": "https://registry.npmjs.org/shebang-command/-/shebang-command-2.0.0.tgz",
@@ -1580,6 +1841,12 @@
"node": ">=8" "node": ">=8"
} }
}, },
"node_modules/standard-as-callback": {
"version": "2.1.0",
"resolved": "https://registry.npmjs.org/standard-as-callback/-/standard-as-callback-2.1.0.tgz",
"integrity": "sha512-qoRRSyROncaz1z0mvYqIE4lCd9p2R90i6GxW3uZv5ucSu8tU7B5HXUP1gG8pVZsYNVaXjk8ClXHPttLyxAL48A==",
"license": "MIT"
},
"node_modules/strip-json-comments": { "node_modules/strip-json-comments": {
"version": "3.1.1", "version": "3.1.1",
"resolved": "https://registry.npmjs.org/strip-json-comments/-/strip-json-comments-3.1.1.tgz", "resolved": "https://registry.npmjs.org/strip-json-comments/-/strip-json-comments-3.1.1.tgz",
@@ -1606,6 +1873,12 @@
"node": ">=8" "node": ">=8"
} }
}, },
"node_modules/tslib": {
"version": "2.8.1",
"resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz",
"integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==",
"license": "0BSD"
},
"node_modules/tsx": { "node_modules/tsx": {
"version": "4.21.0", "version": "4.21.0",
"resolved": "https://registry.npmjs.org/tsx/-/tsx-4.21.0.tgz", "resolved": "https://registry.npmjs.org/tsx/-/tsx-4.21.0.tgz",
@@ -1670,15 +1943,6 @@
"punycode": "^2.1.0" "punycode": "^2.1.0"
} }
}, },
"node_modules/uuid": {
"version": "8.3.2",
"resolved": "https://registry.npmjs.org/uuid/-/uuid-8.3.2.tgz",
"integrity": "sha512-+NYs2QeMWy+GWFOEm9xnn6HCDp0l7QBD7ml8zLUmJ+93Q5NF0NocErnwkTkXVFNiX3/fpC6afS8Dhb/gz7R7eg==",
"license": "MIT",
"bin": {
"uuid": "dist/bin/uuid"
}
},
"node_modules/which": { "node_modules/which": {
"version": "2.0.2", "version": "2.0.2",
"resolved": "https://registry.npmjs.org/which/-/which-2.0.2.tgz", "resolved": "https://registry.npmjs.org/which/-/which-2.0.2.tgz",
+2 -2
View File
@@ -13,12 +13,12 @@
}, },
"dependencies": { "dependencies": {
"@hono/node-server": "^1.13.8", "@hono/node-server": "^1.13.8",
"bullmq": "^5.34.8",
"hono": "^4.7.6", "hono": "^4.7.6",
"node-cron": "^3.0.3" "ioredis": "^5.6.1"
}, },
"devDependencies": { "devDependencies": {
"@types/node": "^22.15.3", "@types/node": "^22.15.3",
"@types/node-cron": "^3.0.11",
"tsx": "^4.19.4", "tsx": "^4.19.4",
"typescript": "^5.8.3", "typescript": "^5.8.3",
"eslint": "^9.25.1", "eslint": "^9.25.1",
+12
View File
@@ -25,6 +25,14 @@ export interface Config {
callbackUrl: string; callbackUrl: string;
/** Optional webhook secret for verifying Gitea signatures */ /** Optional webhook secret for verifying Gitea signatures */
webhookSecret?: string; webhookSecret?: string;
/** Redis connection URL (e.g. redis://redis:6379) */
redisUrl: string;
/** Max concurrent workspace creations (default: 2) */
queueConcurrency: number;
/** Max concurrent cleanup jobs (default: 5) */
cleanupConcurrency: number;
/** Minutes before a workspace is considered stale (default: 120) */
staleWorkspaceMinutes: number;
} }
function required(name: string): string { function required(name: string): string {
@@ -50,5 +58,9 @@ export function loadConfig(): Config {
giteaUrl: process.env.GITEA_URL || "https://gitea.samson.media", giteaUrl: process.env.GITEA_URL || "https://gitea.samson.media",
callbackUrl: required("CALLBACK_URL"), callbackUrl: required("CALLBACK_URL"),
webhookSecret: process.env.WEBHOOK_SECRET, webhookSecret: process.env.WEBHOOK_SECRET,
redisUrl: required("REDIS_URL"),
queueConcurrency: parseInt(process.env.QUEUE_CONCURRENCY || "2", 10),
cleanupConcurrency: parseInt(process.env.CLEANUP_CONCURRENCY || "5", 10),
staleWorkspaceMinutes: parseInt(process.env.STALE_WORKSPACE_MINUTES || "120", 10),
}; };
} }
+27 -4
View File
@@ -4,7 +4,9 @@ import type { TaskQueue } from "../queue.js";
import type { Config } from "../config.js"; import type { Config } from "../config.js";
/** /**
* Issue comment created or edited → re-trigger analyst if still in analysis. * Issue comment created or edited:
* - On a PR → re-trigger tester
* - On an issue still in analysis → re-trigger analyst
* *
* Replaces: issue-comment-reply.json * Replaces: issue-comment-reply.json
*/ */
@@ -23,6 +25,29 @@ export function issueComment(config: Config, queue: TaskQueue) {
return c.text("ignored: bot comment", 200); return c.text("ignored: bot comment", 200);
} }
const [org, repo] = event.repository.full_name.split("/");
const isPullRequest = !!event.issue.pull_request;
if (isPullRequest) {
// Comment on a PR → re-trigger tester
const task: TaskRequest = {
taskType: "test",
issueNumber: event.issue.number,
giteaOrg: org,
giteaRepo: repo,
repoCloneUrl: event.repository.clone_url,
giteaToken: config.giteaReviewToken,
};
const queued = await queue.enqueue(task, {
useReviewAccount: true,
body: "🤖 Re-triggering **test** — workspace queued.",
});
return c.json({ ok: true, queued });
}
// Regular issue comment → re-trigger analyst if still in analysis
const labels = event.issue.labels.map((l) => l.name); const labels = event.issue.labels.map((l) => l.name);
const pastAnalysis = const pastAnalysis =
labels.includes("ready-for-architecture") || labels.includes("ready-for-architecture") ||
@@ -32,8 +57,6 @@ export function issueComment(config: Config, queue: TaskQueue) {
return c.text("ignored: issue past analysis stage", 200); return c.text("ignored: issue past analysis stage", 200);
} }
const [org, repo] = event.repository.full_name.split("/");
const task: TaskRequest = { const task: TaskRequest = {
taskType: "analyse", taskType: "analyse",
issueNumber: event.issue.number, issueNumber: event.issue.number,
@@ -43,7 +66,7 @@ export function issueComment(config: Config, queue: TaskQueue) {
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
const queued = queue.enqueue(task); const queued = await queue.enqueue(task);
return c.json({ ok: true, queued }); return c.json({ ok: true, queued });
}; };
+4 -2
View File
@@ -22,7 +22,9 @@ 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>();
if (event.action !== "label_updated" && event.action !== "labeled") { console.log(`[issue-label] action="${event.action}" issue=#${event.issue.number} labels=[${event.issue.labels.map((l) => l.name).join(", ")}]`);
if (!["labeled", "label_updated", "created"].includes(event.action)) {
return c.text("ignored: not a label event", 200); return c.text("ignored: not a label event", 200);
} }
@@ -54,7 +56,7 @@ export function issueLabel(config: Config, queue: TaskQueue) {
giteaToken, giteaToken,
}; };
const queued = queue.enqueue(task, { const queued = await queue.enqueue(task, {
useReviewAccount: false, useReviewAccount: false,
body: `🤖 SDLC stage transition: **${taskType}** — workspace queued.`, body: `🤖 SDLC stage transition: **${taskType}** — workspace queued.`,
}); });
+1 -1
View File
@@ -31,7 +31,7 @@ export function issueTriage(config: Config, queue: TaskQueue) {
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
const queued = queue.enqueue(task, { const queued = await queue.enqueue(task, {
useReviewAccount: false, useReviewAccount: false,
body: `🤖 Issue received. Starting **${taskType}** stage.`, body: `🤖 Issue received. Starting **${taskType}** stage.`,
}); });
+1 -1
View File
@@ -25,7 +25,7 @@ export function runMaintenance(config: Config, queue: TaskQueue, repos: Maintena
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
queue.enqueue(task); await queue.enqueue(task);
} }
}; };
} }
+1 -1
View File
@@ -40,7 +40,7 @@ export function prReview(config: Config, queue: TaskQueue, gitea: GiteaClient) {
giteaToken: config.giteaReviewToken, giteaToken: config.giteaReviewToken,
}; };
const queued = queue.enqueue(task); const queued = await queue.enqueue(task);
return c.json({ ok: true, queued }); return c.json({ ok: true, queued });
}; };
+29 -3
View File
@@ -16,15 +16,41 @@ export function prRework(config: Config, queue: TaskQueue) {
return c.text("ignored: not a reviewed event", 200); return c.text("ignored: not a reviewed event", 200);
} }
if (!event.review) {
console.log(`[pr-rework] ignored: no review object`);
return c.text("ignored: no review data", 200);
}
const [org, repo] = event.repository.full_name.split("/"); const [org, repo] = event.repository.full_name.split("/");
const prNumber = event.pull_request.number; const prNumber = event.pull_request.number;
const reviewState = event.review.state.toLowerCase();
// Gitea webhook payloads use "state" in the API but "type" in webhooks
// state: "REQUEST_CHANGES" | "APPROVED" | "COMMENT"
// type: "pull_request_review_rejected" | "pull_request_review_approved" | ...
let reviewState: string;
if (event.review.state) {
reviewState = event.review.state.toLowerCase();
} else if (event.review.type) {
const typeMap: Record<string, string> = {
pull_request_review_rejected: "request_changes",
pull_request_review_approved: "approved",
pull_request_review_comment: "comment",
};
reviewState = typeMap[event.review.type] || event.review.type;
} else {
console.log(`[pr-rework] ignored: no review state or type (review: ${JSON.stringify(event.review)})`);
return c.text("ignored: no review state", 200);
}
console.log(`[pr-rework] PR #${prNumber} review: ${reviewState}`);
let taskType: string | null = null; let taskType: string | null = null;
let giteaToken: string; let giteaToken: string;
const branch = event.pull_request.head.ref;
const isArchitectPR = branch.startsWith("docs/");
if (reviewState === "request_changes") { if (reviewState === "request_changes") {
taskType = "rework-pr"; taskType = isArchitectPR ? "rework-spec" : "rework-pr";
giteaToken = config.giteaDevToken; giteaToken = config.giteaDevToken;
} else if (reviewState === "approved") { } else if (reviewState === "approved") {
taskType = "test"; taskType = "test";
@@ -42,7 +68,7 @@ export function prRework(config: Config, queue: TaskQueue) {
giteaToken, giteaToken,
}; };
const queued = queue.enqueue(task, { const queued = await queue.enqueue(task, {
useReviewAccount: false, useReviewAccount: false,
body: `🤖 Review outcome: **${taskType}** — workspace queued.`, body: `🤖 Review outcome: **${taskType}** — workspace queued.`,
}); });
+1 -1
View File
@@ -51,7 +51,7 @@ export function release(config: Config, queue: TaskQueue) {
giteaToken: config.giteaDevToken, giteaToken: config.giteaDevToken,
}; };
const queued = queue.enqueue(task, { const queued = await queue.enqueue(task, {
useReviewAccount: false, useReviewAccount: false,
body: `🤖 Release process started.`, body: `🤖 Release process started.`,
}); });
+31 -20
View File
@@ -1,7 +1,6 @@
import { Hono } from "hono"; import { Hono } from "hono";
import { logger } from "hono/logger"; import { logger } from "hono/logger";
import { serve } from "@hono/node-server"; import { serve } from "@hono/node-server";
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";
@@ -23,9 +22,7 @@ import { runMaintenance, type MaintenanceRepo } from "./handlers/maintenance.js"
const config = loadConfig(); const config = loadConfig();
const coderClient = new CoderClient(config); const coderClient = new CoderClient(config);
const giteaClient = new GiteaClient(config); const giteaClient = new GiteaClient(config);
const queue = new TaskQueue(config, coderClient, giteaClient);
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());
@@ -35,7 +32,7 @@ app.use("*", logger());
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
app.get("/health", (c) => c.json({ status: "ok" })); app.get("/health", (c) => c.json({ status: "ok" }));
app.get("/queue", (c) => c.json(queue.status())); app.get("/queue", async (c) => c.json(await queue.status()));
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Webhook endpoints — same paths as the n8n webhooks so Gitea config // Webhook endpoints — same paths as the n8n webhooks so Gitea config
@@ -61,37 +58,51 @@ app.post("/webhook/task-complete/:name", async (c) => {
}); });
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Scheduled maintenance cron (Monday 9:00 AM) // Maintenance (scheduled via BullMQ repeatable job)
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
const maintenanceRepos = parseMaintenanceRepos(); const maintenanceRepos = parseMaintenanceRepos();
if (maintenanceRepos.length > 0) { if (maintenanceRepos.length > 0) {
const maintenanceFn = runMaintenance(config, queue, maintenanceRepos); const maintenanceFn = runMaintenance(config, queue, maintenanceRepos);
cron.schedule("0 9 * * 1", () => { // The maintenance function is called by the queue's repeatable job system.
console.log("[cron] running weekly maintenance"); // We store the function reference for the maintenance handler to use.
maintenanceFn().catch((err) => console.log(`Maintenance configured for ${maintenanceRepos.length} repo(s) (Monday 9:00 AM via BullMQ)`);
console.error("[cron] maintenance failed:", err),
);
});
console.log(
`Maintenance cron scheduled for ${maintenanceRepos.length} repo(s)`,
);
} }
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Start // Start
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
async function start() {
// Initialize repeatable jobs (stale sweep)
await queue.init();
// Reconcile queue state with Coder before accepting traffic // Reconcile queue state with Coder before accepting traffic
queue.reconcile().catch((err) => await queue.reconcile();
console.error("[startup] reconcile failed:", err),
);
serve({ fetch: app.fetch, port: config.port }, (info) => { serve({ fetch: app.fetch, port: config.port }, (info) => {
console.log(`SDLC Orchestrator listening on :${info.port} (concurrency: ${concurrency})`); 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 // Helpers
+272 -106
View File
@@ -1,173 +1,236 @@
import { Queue, Worker, type Job } from "bullmq";
import Redis from "ioredis";
import type { TaskRequest } from "./types.js"; import type { TaskRequest } from "./types.js";
import type { CoderClient } from "./services/coder.js"; import type { CoderClient } from "./services/coder.js";
import type { GiteaClient } from "./services/gitea.js"; import type { GiteaClient } from "./services/gitea.js";
import type { Config } from "./config.js";
import { dedupKey, workspaceName } from "./workspace-name.js";
interface QueuedTask { const ACTIVE_HASH = "sdlc:active-workspaces";
interface CreateJobData {
task: TaskRequest; task: TaskRequest;
comment?: { useReviewAccount: boolean; body: string }; comment?: { useReviewAccount: boolean; body: string };
} }
interface ActiveWorkspace { interface CleanupJobData {
workspaceId: string; workspaceId: string;
workspaceName: string; workspaceName: string;
startedAt: Date; }
interface ActiveEntry {
workspaceId: string;
workspaceName: string;
startedAt: number;
} }
/** /**
* Task queue with deduplication, concurrency control, and workspace lifecycle. * BullMQ-backed task queue with Redis persistence, retry, and async cleanup.
* *
* - Tasks are keyed by `{taskType}-{repo}-{issue}` for dedup * Two queues:
* - If a task with the same key is pending, running, or has an active workspace, new requests are dropped * - `workspace-create`: creates Coder workspaces (concurrency-limited)
* - Concurrency is configurable (defaults to 2) * - `workspace-cleanup`: stops and deletes workspaces (higher concurrency, with retry)
* - Tracks active workspaces and stops them on task-complete callback *
* Active workspaces (created, awaiting callback) are tracked in a Redis hash
* for dedup across restarts.
*/ */
export class TaskQueue { export class TaskQueue {
private pending = new Map<string, QueuedTask>(); private redis: Redis;
private running = new Set<string>(); private createQueue: Queue<CreateJobData>;
private active = new Map<string, ActiveWorkspace>(); private cleanupQueue: Queue<CleanupJobData | { sweep: true }>;
private concurrency: number; private createWorker: Worker<CreateJobData>;
private cleanupWorker: Worker<CleanupJobData | { sweep: true }>;
private coder: CoderClient; private coder: CoderClient;
private gitea: GiteaClient; private gitea: GiteaClient;
private config: Config;
constructor(coder: CoderClient, gitea: GiteaClient, concurrency = 2) { constructor(config: Config, coder: CoderClient, gitea: GiteaClient) {
this.config = config;
this.coder = coder; this.coder = coder;
this.gitea = gitea; this.gitea = gitea;
this.concurrency = concurrency;
this.redis = new Redis(config.redisUrl, { maxRetriesPerRequest: null });
const connection = { connection: this.redis };
this.createQueue = new Queue("workspace-create", connection);
this.cleanupQueue = new Queue("workspace-cleanup", connection);
this.createWorker = new Worker(
"workspace-create",
(job) => this.processCreate(job),
{ ...connection, concurrency: config.queueConcurrency },
);
this.cleanupWorker = new Worker(
"workspace-cleanup",
(job) => this.processCleanup(job),
{ ...connection, concurrency: config.cleanupConcurrency },
);
this.createWorker.on("failed", (job, err) => {
console.error(`[queue] create job failed: ${job?.id}`, err.message);
});
this.cleanupWorker.on("failed", (job, err) => {
console.error(`[queue] cleanup job failed: ${job?.id}`, err.message);
});
}
/**
* Initialize repeatable jobs (stale sweep).
* Call after construction — separated because it's async.
*/
async init(): Promise<void> {
await this.cleanupQueue.add(
"stale-sweep",
{ sweep: true } as { sweep: true },
{
repeat: { every: 10 * 60 * 1000 },
jobId: "stale-sweep",
},
);
console.log("[queue] stale sweep scheduled every 10 minutes");
} }
/** /**
* Reconcile in-memory state with Coder on startup. * Reconcile in-memory state with Coder on startup.
* Re-adopts any running workspaces so callbacks and dedup work correctly. * Re-adopts any running workspaces so dedup works correctly.
*/ */
async reconcile(): Promise<void> { async reconcile(): Promise<void> {
const workspaces = await this.coder.listWorkspaces(); const workspaces = await this.coder.listWorkspaces();
const runningStatuses = ["starting", "running", "started"]; const runningStatuses = ["starting", "running", "started"];
let adopted = 0;
for (const ws of workspaces) { for (const ws of workspaces) {
if (!runningStatuses.includes(ws.latestBuildStatus)) continue; if (!runningStatuses.includes(ws.latestBuildStatus)) continue;
// Only adopt workspaces that match our naming pattern: {taskType}-{repo}-{issue} // Only adopt workspaces that match our naming pattern
const parts = ws.name.match(/^(.+?)-(.+?)-(\d+)$/); const parts = ws.name.match(/^(.+?)-(.+?)-(\d+)$/);
if (!parts) continue; if (!parts) continue;
this.active.set(ws.name, { const existing = await this.redis.hget(ACTIVE_HASH, ws.name);
if (existing) continue;
const entry: ActiveEntry = {
workspaceId: ws.id, workspaceId: ws.id,
workspaceName: ws.name, workspaceName: ws.name,
startedAt: new Date(), // approximate — we don't know the real start time startedAt: Date.now(),
}); };
await this.redis.hset(ACTIVE_HASH, ws.name, JSON.stringify(entry));
adopted++;
console.log(`[queue] reconciled: adopted workspace "${ws.name}" (status: ${ws.latestBuildStatus})`); console.log(`[queue] reconciled: adopted workspace "${ws.name}" (status: ${ws.latestBuildStatus})`);
} }
if (this.active.size > 0) { // Clean stale entries from the hash that no longer exist in Coder
console.log(`[queue] reconcile complete: ${this.active.size} active workspace(s) adopted`); const activeEntries = await this.redis.hgetall(ACTIVE_HASH);
const coderNames = new Set(workspaces.map((w) => w.name));
for (const key of Object.keys(activeEntries)) {
if (!coderNames.has(key)) {
await this.redis.hdel(ACTIVE_HASH, key);
console.log(`[queue] reconciled: removed stale entry "${key}" (no matching workspace)`);
} }
} }
/** Build a dedup key for a task. */ if (adopted > 0) {
static key(task: TaskRequest): string { console.log(`[queue] reconcile complete: ${adopted} workspace(s) adopted`);
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 a task. Returns true if queued, false if deduplicated (dropped).
*/ */
enqueue( async enqueue(
task: TaskRequest, task: TaskRequest,
comment?: { useReviewAccount: boolean; body: string }, comment?: { useReviewAccount: boolean; body: string },
): boolean { ): Promise<boolean> {
const key = TaskQueue.key(task); const key = dedupKey(task);
if (this.active.has(key)) { // Check if workspace is already active (created, awaiting callback)
const activeEntry = await this.redis.hget(ACTIVE_HASH, key);
if (activeEntry) {
console.log(`[queue] dropped (active workspace): ${key}`); console.log(`[queue] dropped (active workspace): ${key}`);
return false; return false;
} }
if (this.running.has(key)) { try {
console.log(`[queue] dropped (running): ${key}`); await this.createQueue.add(
return false; "create-workspace",
} { task, comment },
{
if (this.pending.has(key)) { jobId: `create-${key}`,
console.log(`[queue] dropped (pending): ${key}`); attempts: 3,
return false; backoff: { type: "exponential", delay: 5000 },
} removeOnComplete: true,
removeOnFail: true,
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})`,
); );
} catch (err: unknown) {
// BullMQ throws when a job with the same ID already exists
if (err instanceof Error && err.message.includes("duplicated")) {
console.log(`[queue] dropped (already queued): ${key}`);
return false;
}
throw err;
}
this.drain(); console.log(`[queue] enqueued: ${key}`);
return true; return true;
} }
/** /**
* Handle task-complete callback from a workspace. * Handle task-complete callback from a workspace.
* Stops the workspace via Coder API and removes it from active tracking. * Enqueues a cleanup job instead of blocking.
*/ */
async taskComplete(workspaceName: string): Promise<boolean> { async taskComplete(workspaceName: string): Promise<boolean> {
const entry = this.active.get(workspaceName); const raw = await this.redis.hget(ACTIVE_HASH, workspaceName);
if (!entry) { if (!raw) {
console.log(`[queue] task-complete for unknown workspace: ${workspaceName}`); console.log(`[queue] task-complete for unknown workspace: ${workspaceName}`);
return false; return false;
} }
const elapsed = Math.round( const entry: ActiveEntry = JSON.parse(raw);
(Date.now() - entry.startedAt.getTime()) / 1000, const elapsed = Math.round((Date.now() - entry.startedAt) / 1000);
console.log(`[queue] task complete: ${workspaceName} (ran for ${elapsed}s) — queueing cleanup`);
// Remove from active immediately so new tasks for this key can be queued
await this.redis.hdel(ACTIVE_HASH, workspaceName);
await this.cleanupQueue.add(
"cleanup-workspace",
{ workspaceId: entry.workspaceId, workspaceName: entry.workspaceName },
{
jobId: `cleanup-${workspaceName}`,
attempts: 5,
backoff: { type: "exponential", delay: 10000 },
removeOnComplete: true,
removeOnFail: true,
},
); );
console.log(
`[queue] task complete: ${workspaceName} (ran for ${elapsed}s) — stopping workspace`,
);
this.active.delete(workspaceName);
try {
await this.coder.stopWorkspace(entry.workspaceId);
console.log(`[queue] workspace stopped: ${workspaceName} — waiting for shutdown`);
// Wait for the workspace to fully stop before deleting
await new Promise((resolve) => setTimeout(resolve, 10_000));
await this.coder.deleteWorkspace(entry.workspaceId);
console.log(`[queue] workspace deleted: ${workspaceName}`);
} catch (err) {
console.error(`[queue] failed to clean up workspace ${workspaceName}:`, err);
}
// Drain in case pending tasks were waiting for capacity
this.drain();
return true; return true;
} }
/** Process queued tasks up to the concurrency limit. */ // ---------------------------------------------------------------------------
private drain(): void { // Workers
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> { private async processCreate(job: Job<CreateJobData>): Promise<void> {
const { task, comment } = entry; const { task, comment } = job.data;
const key = dedupKey(task);
console.log(`[queue] processing create: ${key} (attempt ${job.attemptsMade + 1})`);
try {
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} (id: ${workspace.id})`); console.log(`[queue] workspace created: ${workspace.name} (id: ${workspace.id})`);
// Move from running → active (workspace is now alive, waiting for callback) // Track as active
this.active.set(key, { const entry: ActiveEntry = {
workspaceId: workspace.id, workspaceId: workspace.id,
workspaceName: workspace.name, workspaceName: workspace.name,
startedAt: new Date(), startedAt: Date.now(),
}); };
await this.redis.hset(ACTIVE_HASH, key, JSON.stringify(entry));
if (comment) { if (comment) {
await this.gitea.commentOnIssue( await this.gitea.commentOnIssue(
@@ -178,34 +241,137 @@ export class TaskQueue {
comment.useReviewAccount, comment.useReviewAccount,
); );
} }
}
private async processCleanup(
job: Job<CleanupJobData | { sweep: true }>,
): Promise<void> {
if ("sweep" in job.data) {
await this.staleSweep();
return;
}
const { workspaceId, workspaceName: wsName } = job.data as CleanupJobData;
console.log(`[queue] cleanup: stopping ${wsName} (attempt ${job.attemptsMade + 1})`);
await this.coder.stopWorkspace(workspaceId);
// Poll until stopped (max 5 minutes)
const maxPolls = 30;
const pollInterval = 10_000;
const doneStatuses = ["stopped", "failed", "canceled", "deleted"];
for (let i = 0; i < maxPolls; i++) {
await new Promise((r) => setTimeout(r, pollInterval));
const ws = await this.coder.findWorkspaceByName(wsName);
if (!ws || doneStatuses.includes(ws.latestBuildStatus)) {
break;
}
console.log(`[queue] cleanup: ${wsName} still ${ws.latestBuildStatus}, polling...`);
}
await this.coder.deleteWorkspace(workspaceId);
console.log(`[queue] workspace deleted: ${wsName}`);
}
private async staleSweep(): Promise<void> {
console.log("[queue] running stale sweep");
const workspaces = await this.coder.listWorkspaces();
const now = Date.now();
const staleMs = this.config.staleWorkspaceMinutes * 60 * 1000;
const doneStatuses = ["stopped", "failed", "canceled", "deleted"];
for (const ws of workspaces) {
// Only manage workspaces that match our naming pattern
const parts = ws.name.match(/^(.+?)-(.+?)-(\d+)$/);
if (!parts) continue;
if (doneStatuses.includes(ws.latestBuildStatus)) {
// Stopped/failed workspace — clean up
console.log(`[queue] stale sweep: deleting ${ws.latestBuildStatus} workspace "${ws.name}"`);
try {
await this.coder.deleteWorkspace(ws.id);
await this.redis.hdel(ACTIVE_HASH, ws.name);
} catch (err) { } catch (err) {
console.error(`[queue] failed: ${key}`, err); console.error(`[queue] stale sweep: failed to delete "${ws.name}":`, err);
} finally { }
this.running.delete(key); continue;
// Don't drain here — active workspaces count toward concurrency }
// Check if running workspace is stale based on our tracking
const raw = await this.redis.hget(ACTIVE_HASH, ws.name);
if (raw) {
const entry: ActiveEntry = JSON.parse(raw);
if (now - entry.startedAt > staleMs) {
console.log(`[queue] stale sweep: workspace "${ws.name}" exceeded ${this.config.staleWorkspaceMinutes}m — stopping`);
try {
await this.coder.stopWorkspace(ws.id);
await this.redis.hdel(ACTIVE_HASH, ws.name);
// Enqueue cleanup to handle the stop→delete lifecycle
await this.cleanupQueue.add(
"cleanup-workspace",
{ workspaceId: ws.id, workspaceName: ws.name },
{
jobId: `cleanup-stale-${ws.name}-${Date.now()}`,
attempts: 5,
backoff: { type: "exponential", delay: 10000 },
removeOnComplete: true,
removeOnFail: true,
},
);
} catch (err) {
console.error(`[queue] stale sweep: failed to stop "${ws.name}":`, err);
}
}
}
} }
} }
/** Current queue status for health/debug endpoints. */ // ---------------------------------------------------------------------------
status(): { // Status & lifecycle
pending: string[]; // ---------------------------------------------------------------------------
running: string[];
active: Record<string, { workspaceId: string; elapsed: number }>; async status(): Promise<{
pending: number;
active: number;
activeWorkspaces: Record<string, { workspaceId: string; elapsed: number }>;
failed: number;
concurrency: number; concurrency: number;
} { }> {
const activeMap: Record<string, { workspaceId: string; elapsed: number }> = {}; const [waiting, activeJobs, failed] = await Promise.all([
for (const [key, entry] of this.active) { this.createQueue.getWaitingCount(),
activeMap[key] = { this.createQueue.getActiveCount(),
this.createQueue.getFailedCount(),
]);
const activeEntries = await this.redis.hgetall(ACTIVE_HASH);
const activeWorkspaces: Record<string, { workspaceId: string; elapsed: number }> = {};
const now = Date.now();
for (const [key, raw] of Object.entries(activeEntries)) {
const entry: ActiveEntry = JSON.parse(raw);
activeWorkspaces[key] = {
workspaceId: entry.workspaceId, workspaceId: entry.workspaceId,
elapsed: Math.round((Date.now() - entry.startedAt.getTime()) / 1000), elapsed: Math.round((now - entry.startedAt) / 1000),
}; };
} }
return { return {
pending: [...this.pending.keys()], pending: waiting + activeJobs,
running: [...this.running.keys()], active: Object.keys(activeWorkspaces).length,
active: activeMap, activeWorkspaces,
concurrency: this.concurrency, failed,
concurrency: this.config.queueConcurrency,
}; };
} }
async gracefulShutdown(): Promise<void> {
console.log("[queue] shutting down workers...");
await Promise.all([
this.createWorker.close(),
this.cleanupWorker.close(),
]);
this.redis.disconnect();
console.log("[queue] shutdown complete");
}
} }
+6 -4
View File
@@ -1,6 +1,6 @@
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"; import { workspaceName as buildWorkspaceName } from "../workspace-name.js";
export class CoderClient { export class CoderClient {
private baseUrl: string; private baseUrl: string;
@@ -20,7 +20,7 @@ export class CoderClient {
} }
async createWorkspace(task: TaskRequest): Promise<{ id: string; name: string }> { async createWorkspace(task: TaskRequest): Promise<{ id: string; name: string }> {
const name = TaskQueue.workspaceName(task); const name = buildWorkspaceName(task);
const body = { const body = {
name, name,
@@ -109,12 +109,14 @@ export class CoderClient {
async deleteWorkspace(workspaceId: string): Promise<void> { async deleteWorkspace(workspaceId: string): Promise<void> {
const res = await fetch( const res = await fetch(
`${this.baseUrl}/api/v2/workspaces/${workspaceId}`, `${this.baseUrl}/api/v2/workspaces/${workspaceId}/builds`,
{ {
method: "DELETE", method: "POST",
headers: { headers: {
"Content-Type": "application/json",
"Coder-Session-Token": this.token, "Coder-Session-Token": this.token,
}, },
body: JSON.stringify({ transition: "delete" }),
}, },
); );
+3 -1
View File
@@ -30,6 +30,7 @@ export interface GiteaIssue {
state: string; state: string;
user: GiteaUser; user: GiteaUser;
labels: GiteaLabel[]; labels: GiteaLabel[];
pull_request?: { merged: boolean } | null;
} }
export interface GiteaPullRequest { export interface GiteaPullRequest {
@@ -46,7 +47,8 @@ export interface GiteaPullRequest {
export interface GiteaReview { export interface GiteaReview {
id: number; id: number;
body: string; body: string;
state: string; // "approved", "request_changes", "comment" state?: string; // "approved", "request_changes", "comment" (API response)
type?: string; // "pull_request_review_rejected" | "pull_request_review_approved" (webhook payload)
user: GiteaUser; user: GiteaUser;
} }
+11
View File
@@ -0,0 +1,11 @@
import type { TaskRequest } from "./types.js";
/** Build a dedup key for a task. */
export function dedupKey(task: TaskRequest): string {
return `${task.taskType}-${task.giteaRepo}-${task.issueNumber}`;
}
/** Build the workspace name (matches the dedup key). */
export function workspaceName(task: TaskRequest): string {
return `${task.taskType}-${task.giteaRepo}-${task.issueNumber}`;
}