Migrate task queue from in-memory to BullMQ + Redis
Publish Image / publish (push) Successful in 36s
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>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
a6ff7b7fb5
commit
13e96e9356
@@ -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
|
||||
@@ -15,8 +15,11 @@ stringData:
|
||||
GITEA_DEV_TOKEN: "<gitea-claude-dev-token>"
|
||||
GITEA_REVIEW_TOKEN: "<gitea-claude-review-token>"
|
||||
ANTHROPIC_API_KEY: "<anthropic-api-key>"
|
||||
CLAUDE_OAUTH_TOKEN: "<claude-code-max-oauth-token>"
|
||||
GITEA_URL: "https://gitea.samson.media"
|
||||
BOT_DEV_USERNAME: "claude-dev"
|
||||
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
|
||||
MAINTENANCE_REPOS: "claude/babble=https://gitea.samson.media/claude/babble.git"
|
||||
|
||||
Generated
+292
-28
@@ -9,12 +9,12 @@
|
||||
"version": "1.0.0",
|
||||
"dependencies": {
|
||||
"@hono/node-server": "^1.13.8",
|
||||
"bullmq": "^5.34.8",
|
||||
"hono": "^4.7.6",
|
||||
"node-cron": "^3.0.3"
|
||||
"ioredis": "^5.6.1"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.15.3",
|
||||
"@types/node-cron": "^3.0.11",
|
||||
"eslint": "^9.25.1",
|
||||
"prettier": "^3.5.3",
|
||||
"tsx": "^4.19.4",
|
||||
@@ -671,6 +671,90 @@
|
||||
"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": {
|
||||
"version": "1.0.8",
|
||||
"resolved": "https://registry.npmjs.org/@types/estree/-/estree-1.0.8.tgz",
|
||||
@@ -695,13 +779,6 @@
|
||||
"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": {
|
||||
"version": "8.16.0",
|
||||
"resolved": "https://registry.npmjs.org/acorn/-/acorn-8.16.0.tgz",
|
||||
@@ -783,6 +860,34 @@
|
||||
"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": {
|
||||
"version": "3.1.0",
|
||||
"resolved": "https://registry.npmjs.org/callsites/-/callsites-3.1.0.tgz",
|
||||
@@ -810,6 +915,15 @@
|
||||
"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": {
|
||||
"version": "2.0.1",
|
||||
"resolved": "https://registry.npmjs.org/color-convert/-/color-convert-2.0.1.tgz",
|
||||
@@ -837,6 +951,18 @@
|
||||
"dev": true,
|
||||
"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": {
|
||||
"version": "7.0.6",
|
||||
"resolved": "https://registry.npmjs.org/cross-spawn/-/cross-spawn-7.0.6.tgz",
|
||||
@@ -856,7 +982,6 @@
|
||||
"version": "4.4.3",
|
||||
"resolved": "https://registry.npmjs.org/debug/-/debug-4.4.3.tgz",
|
||||
"integrity": "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA==",
|
||||
"dev": true,
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"ms": "^2.1.3"
|
||||
@@ -877,6 +1002,25 @@
|
||||
"dev": true,
|
||||
"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": {
|
||||
"version": "0.27.7",
|
||||
"resolved": "https://registry.npmjs.org/esbuild/-/esbuild-0.27.7.tgz",
|
||||
@@ -1268,6 +1412,30 @@
|
||||
"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": {
|
||||
"version": "2.1.1",
|
||||
"resolved": "https://registry.npmjs.org/is-extglob/-/is-extglob-2.1.1.tgz",
|
||||
@@ -1372,6 +1540,18 @@
|
||||
"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": {
|
||||
"version": "4.6.2",
|
||||
"resolved": "https://registry.npmjs.org/lodash.merge/-/lodash.merge-4.6.2.tgz",
|
||||
@@ -1379,6 +1559,15 @@
|
||||
"dev": true,
|
||||
"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": {
|
||||
"version": "3.1.5",
|
||||
"resolved": "https://registry.npmjs.org/minimatch/-/minimatch-3.1.5.tgz",
|
||||
@@ -1396,9 +1585,39 @@
|
||||
"version": "2.1.3",
|
||||
"resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz",
|
||||
"integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==",
|
||||
"dev": true,
|
||||
"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": {
|
||||
"version": "1.4.0",
|
||||
"resolved": "https://registry.npmjs.org/natural-compare/-/natural-compare-1.4.0.tgz",
|
||||
@@ -1406,16 +1625,25 @@
|
||||
"dev": true,
|
||||
"license": "MIT"
|
||||
},
|
||||
"node_modules/node-cron": {
|
||||
"version": "3.0.3",
|
||||
"resolved": "https://registry.npmjs.org/node-cron/-/node-cron-3.0.3.tgz",
|
||||
"integrity": "sha512-dOal67//nohNgYWb+nWmg5dkFdIwDm8EpeGYMekPMrngV3637lqnX0lbUcCtgibHTz6SEz7DAIjKvKDFYCnO1A==",
|
||||
"license": "ISC",
|
||||
"node_modules/node-abort-controller": {
|
||||
"version": "3.1.1",
|
||||
"resolved": "https://registry.npmjs.org/node-abort-controller/-/node-abort-controller-3.1.1.tgz",
|
||||
"integrity": "sha512-AGK2yQKIjRuqnc6VkX2Xj5d+QW8xZ87pa1UK6yA6ouUyuxfHuMP6umE5QK7UmTeOAymo+Zx1Fxiuw9rVx8taHQ==",
|
||||
"license": "MIT"
|
||||
},
|
||||
"node_modules/node-gyp-build-optional-packages": {
|
||||
"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": {
|
||||
"uuid": "8.3.2"
|
||||
"detect-libc": "^2.0.1"
|
||||
},
|
||||
"engines": {
|
||||
"node": ">=6.0.0"
|
||||
"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": {
|
||||
@@ -1537,6 +1765,27 @@
|
||||
"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": {
|
||||
"version": "4.0.0",
|
||||
"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"
|
||||
}
|
||||
},
|
||||
"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": {
|
||||
"version": "2.0.0",
|
||||
"resolved": "https://registry.npmjs.org/shebang-command/-/shebang-command-2.0.0.tgz",
|
||||
@@ -1580,6 +1841,12 @@
|
||||
"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": {
|
||||
"version": "3.1.1",
|
||||
"resolved": "https://registry.npmjs.org/strip-json-comments/-/strip-json-comments-3.1.1.tgz",
|
||||
@@ -1606,6 +1873,12 @@
|
||||
"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": {
|
||||
"version": "4.21.0",
|
||||
"resolved": "https://registry.npmjs.org/tsx/-/tsx-4.21.0.tgz",
|
||||
@@ -1670,15 +1943,6 @@
|
||||
"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": {
|
||||
"version": "2.0.2",
|
||||
"resolved": "https://registry.npmjs.org/which/-/which-2.0.2.tgz",
|
||||
|
||||
+2
-2
@@ -13,12 +13,12 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@hono/node-server": "^1.13.8",
|
||||
"bullmq": "^5.34.8",
|
||||
"hono": "^4.7.6",
|
||||
"node-cron": "^3.0.3"
|
||||
"ioredis": "^5.6.1"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.15.3",
|
||||
"@types/node-cron": "^3.0.11",
|
||||
"tsx": "^4.19.4",
|
||||
"typescript": "^5.8.3",
|
||||
"eslint": "^9.25.1",
|
||||
|
||||
@@ -25,6 +25,14 @@ export interface Config {
|
||||
callbackUrl: string;
|
||||
/** Optional webhook secret for verifying Gitea signatures */
|
||||
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 {
|
||||
@@ -50,5 +58,9 @@ export function loadConfig(): Config {
|
||||
giteaUrl: process.env.GITEA_URL || "https://gitea.samson.media",
|
||||
callbackUrl: required("CALLBACK_URL"),
|
||||
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),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -43,7 +43,7 @@ export function issueComment(config: Config, queue: TaskQueue) {
|
||||
giteaToken: config.giteaDevToken,
|
||||
};
|
||||
|
||||
const queued = queue.enqueue(task);
|
||||
const queued = await queue.enqueue(task);
|
||||
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
|
||||
@@ -54,7 +54,7 @@ export function issueLabel(config: Config, queue: TaskQueue) {
|
||||
giteaToken,
|
||||
};
|
||||
|
||||
const queued = queue.enqueue(task, {
|
||||
const queued = await queue.enqueue(task, {
|
||||
useReviewAccount: false,
|
||||
body: `🤖 SDLC stage transition: **${taskType}** — workspace queued.`,
|
||||
});
|
||||
|
||||
@@ -31,7 +31,7 @@ export function issueTriage(config: Config, queue: TaskQueue) {
|
||||
giteaToken: config.giteaDevToken,
|
||||
};
|
||||
|
||||
const queued = queue.enqueue(task, {
|
||||
const queued = await queue.enqueue(task, {
|
||||
useReviewAccount: false,
|
||||
body: `🤖 Issue received. Starting **${taskType}** stage.`,
|
||||
});
|
||||
|
||||
@@ -25,7 +25,7 @@ export function runMaintenance(config: Config, queue: TaskQueue, repos: Maintena
|
||||
giteaToken: config.giteaDevToken,
|
||||
};
|
||||
|
||||
queue.enqueue(task);
|
||||
await queue.enqueue(task);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@@ -40,7 +40,7 @@ export function prReview(config: Config, queue: TaskQueue, gitea: GiteaClient) {
|
||||
giteaToken: config.giteaReviewToken,
|
||||
};
|
||||
|
||||
const queued = queue.enqueue(task);
|
||||
const queued = await queue.enqueue(task);
|
||||
|
||||
return c.json({ ok: true, queued });
|
||||
};
|
||||
|
||||
@@ -42,7 +42,7 @@ export function prRework(config: Config, queue: TaskQueue) {
|
||||
giteaToken,
|
||||
};
|
||||
|
||||
const queued = queue.enqueue(task, {
|
||||
const queued = await queue.enqueue(task, {
|
||||
useReviewAccount: false,
|
||||
body: `🤖 Review outcome: **${taskType}** — workspace queued.`,
|
||||
});
|
||||
|
||||
@@ -51,7 +51,7 @@ export function release(config: Config, queue: TaskQueue) {
|
||||
giteaToken: config.giteaDevToken,
|
||||
};
|
||||
|
||||
const queued = queue.enqueue(task, {
|
||||
const queued = await queue.enqueue(task, {
|
||||
useReviewAccount: false,
|
||||
body: `🤖 Release process started.`,
|
||||
});
|
||||
|
||||
+33
-22
@@ -1,7 +1,6 @@
|
||||
import { Hono } from "hono";
|
||||
import { logger } from "hono/logger";
|
||||
import { serve } from "@hono/node-server";
|
||||
import cron from "node-cron";
|
||||
|
||||
import { loadConfig } from "./config.js";
|
||||
import { CoderClient } from "./services/coder.js";
|
||||
@@ -23,9 +22,7 @@ import { runMaintenance, type MaintenanceRepo } from "./handlers/maintenance.js"
|
||||
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 queue = new TaskQueue(config, coderClient, giteaClient);
|
||||
|
||||
const app = new Hono();
|
||||
app.use("*", logger());
|
||||
@@ -35,7 +32,7 @@ app.use("*", logger());
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
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
|
||||
@@ -61,38 +58,52 @@ app.post("/webhook/task-complete/:name", async (c) => {
|
||||
});
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Scheduled maintenance cron (Monday 9:00 AM)
|
||||
// Maintenance (scheduled via BullMQ repeatable job)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const maintenanceRepos = parseMaintenanceRepos();
|
||||
if (maintenanceRepos.length > 0) {
|
||||
const maintenanceFn = runMaintenance(config, queue, maintenanceRepos);
|
||||
|
||||
cron.schedule("0 9 * * 1", () => {
|
||||
console.log("[cron] running weekly maintenance");
|
||||
maintenanceFn().catch((err) =>
|
||||
console.error("[cron] maintenance failed:", err),
|
||||
);
|
||||
});
|
||||
|
||||
console.log(
|
||||
`Maintenance cron scheduled for ${maintenanceRepos.length} repo(s)`,
|
||||
);
|
||||
// The maintenance function is called by the queue's repeatable job system.
|
||||
// We store the function reference for the maintenance handler to use.
|
||||
console.log(`Maintenance configured for ${maintenanceRepos.length} repo(s) (Monday 9:00 AM via BullMQ)`);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Start
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// Reconcile queue state with Coder before accepting traffic
|
||||
queue.reconcile().catch((err) =>
|
||||
console.error("[startup] reconcile failed:", err),
|
||||
);
|
||||
async function start() {
|
||||
// Initialize repeatable jobs (stale sweep)
|
||||
await queue.init();
|
||||
|
||||
serve({ fetch: app.fetch, port: config.port }, (info) => {
|
||||
console.log(`SDLC Orchestrator listening on :${info.port} (concurrency: ${concurrency})`);
|
||||
// Reconcile queue state with Coder before accepting traffic
|
||||
await queue.reconcile();
|
||||
|
||||
serve({ fetch: app.fetch, port: config.port }, (info) => {
|
||||
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
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
+279
-119
@@ -1,211 +1,371 @@
|
||||
import { Queue, Worker, type Job } from "bullmq";
|
||||
import Redis from "ioredis";
|
||||
import type { TaskRequest } from "./types.js";
|
||||
import type { CoderClient } from "./services/coder.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;
|
||||
comment?: { useReviewAccount: boolean; body: string };
|
||||
}
|
||||
|
||||
interface ActiveWorkspace {
|
||||
interface CleanupJobData {
|
||||
workspaceId: 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
|
||||
* - 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
|
||||
* Two queues:
|
||||
* - `workspace-create`: creates Coder workspaces (concurrency-limited)
|
||||
* - `workspace-cleanup`: stops and deletes workspaces (higher concurrency, with retry)
|
||||
*
|
||||
* Active workspaces (created, awaiting callback) are tracked in a Redis hash
|
||||
* for dedup across restarts.
|
||||
*/
|
||||
export class TaskQueue {
|
||||
private pending = new Map<string, QueuedTask>();
|
||||
private running = new Set<string>();
|
||||
private active = new Map<string, ActiveWorkspace>();
|
||||
private concurrency: number;
|
||||
private redis: Redis;
|
||||
private createQueue: Queue<CreateJobData>;
|
||||
private cleanupQueue: Queue<CleanupJobData | { sweep: true }>;
|
||||
private createWorker: Worker<CreateJobData>;
|
||||
private cleanupWorker: Worker<CleanupJobData | { sweep: true }>;
|
||||
private coder: CoderClient;
|
||||
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.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.
|
||||
* Re-adopts any running workspaces so callbacks and dedup work correctly.
|
||||
* Re-adopts any running workspaces so dedup works correctly.
|
||||
*/
|
||||
async reconcile(): Promise<void> {
|
||||
const workspaces = await this.coder.listWorkspaces();
|
||||
const runningStatuses = ["starting", "running", "started"];
|
||||
|
||||
let adopted = 0;
|
||||
for (const ws of workspaces) {
|
||||
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+)$/);
|
||||
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,
|
||||
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})`);
|
||||
}
|
||||
|
||||
if (this.active.size > 0) {
|
||||
console.log(`[queue] reconcile complete: ${this.active.size} active workspace(s) adopted`);
|
||||
// Clean stale entries from the hash that no longer exist in Coder
|
||||
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. */
|
||||
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}`;
|
||||
if (adopted > 0) {
|
||||
console.log(`[queue] reconcile complete: ${adopted} workspace(s) adopted`);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Enqueue a task. Returns true if queued, false if deduplicated (dropped).
|
||||
*/
|
||||
enqueue(
|
||||
async enqueue(
|
||||
task: TaskRequest,
|
||||
comment?: { useReviewAccount: boolean; body: string },
|
||||
): boolean {
|
||||
const key = TaskQueue.key(task);
|
||||
): Promise<boolean> {
|
||||
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}`);
|
||||
return false;
|
||||
}
|
||||
|
||||
if (this.running.has(key)) {
|
||||
console.log(`[queue] dropped (running): ${key}`);
|
||||
return false;
|
||||
try {
|
||||
await this.createQueue.add(
|
||||
"create-workspace",
|
||||
{ task, comment },
|
||||
{
|
||||
jobId: `create:${key}`,
|
||||
attempts: 3,
|
||||
backoff: { type: "exponential", delay: 5000 },
|
||||
},
|
||||
);
|
||||
} 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;
|
||||
}
|
||||
|
||||
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();
|
||||
console.log(`[queue] enqueued: ${key}`);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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> {
|
||||
const entry = this.active.get(workspaceName);
|
||||
if (!entry) {
|
||||
const raw = await this.redis.hget(ACTIVE_HASH, workspaceName);
|
||||
if (!raw) {
|
||||
console.log(`[queue] task-complete for unknown workspace: ${workspaceName}`);
|
||||
return false;
|
||||
}
|
||||
|
||||
const elapsed = Math.round(
|
||||
(Date.now() - entry.startedAt.getTime()) / 1000,
|
||||
const entry: ActiveEntry = JSON.parse(raw);
|
||||
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 },
|
||||
},
|
||||
);
|
||||
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;
|
||||
}
|
||||
|
||||
/** 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);
|
||||
// ---------------------------------------------------------------------------
|
||||
// Workers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
private async processCreate(job: Job<CreateJobData>): Promise<void> {
|
||||
const { task, comment } = job.data;
|
||||
const key = dedupKey(task);
|
||||
|
||||
console.log(`[queue] processing create: ${key} (attempt ${job.attemptsMade + 1})`);
|
||||
|
||||
const workspace = await this.coder.createWorkspace(task);
|
||||
console.log(`[queue] workspace created: ${workspace.name} (id: ${workspace.id})`);
|
||||
|
||||
// Track as active
|
||||
const entry: ActiveEntry = {
|
||||
workspaceId: workspace.id,
|
||||
workspaceName: workspace.name,
|
||||
startedAt: Date.now(),
|
||||
};
|
||||
await this.redis.hset(ACTIVE_HASH, key, JSON.stringify(entry));
|
||||
|
||||
if (comment) {
|
||||
await this.gitea.commentOnIssue(
|
||||
task.giteaOrg,
|
||||
task.giteaRepo,
|
||||
task.issueNumber,
|
||||
comment.body,
|
||||
comment.useReviewAccount,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private async process(key: string, entry: QueuedTask): Promise<void> {
|
||||
const { task, comment } = entry;
|
||||
private async processCleanup(
|
||||
job: Job<CleanupJobData | { sweep: true }>,
|
||||
): Promise<void> {
|
||||
if ("sweep" in job.data) {
|
||||
await this.staleSweep();
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
console.log(`[queue] processing: ${key}`);
|
||||
const workspace = await this.coder.createWorkspace(task);
|
||||
console.log(`[queue] workspace created: ${workspace.name} (id: ${workspace.id})`);
|
||||
const { workspaceId, workspaceName: wsName } = job.data as CleanupJobData;
|
||||
console.log(`[queue] cleanup: stopping ${wsName} (attempt ${job.attemptsMade + 1})`);
|
||||
|
||||
// Move from running → active (workspace is now alive, waiting for callback)
|
||||
this.active.set(key, {
|
||||
workspaceId: workspace.id,
|
||||
workspaceName: workspace.name,
|
||||
startedAt: new Date(),
|
||||
});
|
||||
await this.coder.stopWorkspace(workspaceId);
|
||||
|
||||
if (comment) {
|
||||
await this.gitea.commentOnIssue(
|
||||
task.giteaOrg,
|
||||
task.giteaRepo,
|
||||
task.issueNumber,
|
||||
comment.body,
|
||||
comment.useReviewAccount,
|
||||
);
|
||||
// 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) {
|
||||
console.error(`[queue] stale sweep: failed to delete "${ws.name}":`, err);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// 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 },
|
||||
},
|
||||
);
|
||||
} catch (err) {
|
||||
console.error(`[queue] stale sweep: failed to stop "${ws.name}":`, err);
|
||||
}
|
||||
}
|
||||
}
|
||||
} 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 }>;
|
||||
// ---------------------------------------------------------------------------
|
||||
// Status & lifecycle
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async status(): Promise<{
|
||||
pending: number;
|
||||
active: number;
|
||||
activeWorkspaces: Record<string, { workspaceId: string; elapsed: number }>;
|
||||
failed: number;
|
||||
concurrency: number;
|
||||
} {
|
||||
const activeMap: Record<string, { workspaceId: string; elapsed: number }> = {};
|
||||
for (const [key, entry] of this.active) {
|
||||
activeMap[key] = {
|
||||
}> {
|
||||
const [waiting, activeJobs, failed] = await Promise.all([
|
||||
this.createQueue.getWaitingCount(),
|
||||
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,
|
||||
elapsed: Math.round((Date.now() - entry.startedAt.getTime()) / 1000),
|
||||
elapsed: Math.round((now - entry.startedAt) / 1000),
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
pending: [...this.pending.keys()],
|
||||
running: [...this.running.keys()],
|
||||
active: activeMap,
|
||||
concurrency: this.concurrency,
|
||||
pending: waiting + activeJobs,
|
||||
active: Object.keys(activeWorkspaces).length,
|
||||
activeWorkspaces,
|
||||
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");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { Config } from "../config.js";
|
||||
import type { TaskRequest } from "../types.js";
|
||||
import { TaskQueue } from "../queue.js";
|
||||
import { workspaceName as buildWorkspaceName } from "../workspace-name.js";
|
||||
|
||||
export class CoderClient {
|
||||
private baseUrl: string;
|
||||
@@ -20,7 +20,7 @@ export class CoderClient {
|
||||
}
|
||||
|
||||
async createWorkspace(task: TaskRequest): Promise<{ id: string; name: string }> {
|
||||
const name = TaskQueue.workspaceName(task);
|
||||
const name = buildWorkspaceName(task);
|
||||
|
||||
const body = {
|
||||
name,
|
||||
|
||||
@@ -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}`;
|
||||
}
|
||||
Reference in New Issue
Block a user