diff --git a/bun.lock b/bun.lock index d397692..d3cca5c 100644 --- a/bun.lock +++ b/bun.lock @@ -1,5 +1,6 @@ { "lockfileVersion": 1, + "configVersion": 0, "workspaces": { "": { "name": "devintern", @@ -43,6 +44,7 @@ "dependencies": { "@types/turndown": "^5.0.5", "commander": "^11.0.0", + "cron-parser": "^5.4.0", "dotenv": "^16.3.1", "p-queue": "^9.0.1", "turndown": "^7.2.0", @@ -1060,7 +1062,7 @@ "builder-util-runtime": ["builder-util-runtime@9.7.0", "", { "dependencies": { "debug": "^4.3.4", "sax": "^1.2.4" } }, "sha512-g/kR520giAFYkSXTzcmF3kqQq7wi8F6N6SzeDgZrqTBN+VHdmgWOyTdD1yD7AATDId/yXLvuP34CxW46/BwCdw=="], - "bun-types": ["bun-types@1.3.14", "", { "dependencies": { "@types/node": "*" } }, "sha512-4N0ig0fEomHt5R0KCFWjovxow98rIoRwKolrYdCcknNwMekCXRnWEUvgu5soYV8QXtVsrUD8B95MBOZGPvr6KQ=="], + "bun-types": ["bun-types@1.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-iIKw23BspnQQYd3prITOBxeUsxBHnwzX6YJfGMuNOZzeNcMmVqzIIVGRm1l69ogaPQmb4wB6BN8mA5bE9YuC5Q=="], "bundle-name": ["bundle-name@4.1.0", "", { "dependencies": { "run-applescript": "^7.0.0" } }, "sha512-tjwM5exMg6BGRI+kNmTntNsvdZS1X8BFYS6tnJ2hdH0kVxM6/eVZ2xy+FqStSWvYmtfFMDLIxurorHwDKfDz5Q=="], @@ -1166,6 +1168,8 @@ "crelt": ["crelt@1.0.7", "", {}, "sha512-aK6BbWfhf4U/wCcLHKPJl/xa6VkVstRaPywWtMKGwuOLc/wZTyQYuoxgvZnNsBvv7Kg3YTBQYYBCggcviQczuA=="], + "cron-parser": ["cron-parser@5.10.0", "", { "dependencies": { "luxon": "^3.7.2" } }, "sha512-izNAxJyRWUP8ljBoDSub5WyrVOUlT4SLGShswE7eoRBpp6QUsSycYxLBMJlbshgPBMcPT/nrfgjNY2918ayv2A=="], + "cross-dirname": ["cross-dirname@0.1.0", "", {}, "sha512-+R08/oI0nl3vfPcqftZRpytksBXDzOUveBq/NBVx0sUp1axwzPQrKinNx5yd5sxPu8j1wIy8AfnVQ+5eFdha6Q=="], "cross-spawn": ["cross-spawn@7.0.6", "", { "dependencies": { "path-key": "^3.1.0", "shebang-command": "^2.0.0", "which": "^2.0.1" } }, "sha512-uV2QOWP2nWzsy2aMp8aRibhi9dlzF5Hgh5SHaB9OiTGEyDTiJJyx0uy51QXdyWbtAHNua4XJzUKca3OzKUd3vA=="], @@ -1606,6 +1610,8 @@ "lucide-react": ["lucide-react@1.28.0", "", { "peerDependencies": { "react": "^16.5.1 || ^17.0.0 || ^18.0.0 || ^19.0.0" } }, "sha512-fARAFJULsGuDDydjp6+6blekG/sBIM29TerzLjc9bQUKAcEfrSc4ZQKb25KRz4OMKd87cZTb5dgq0w/T6KufVg=="], + "luxon": ["luxon@3.7.2", "", {}, "sha512-vtEhXh/gNjI9Yg1u4jX/0YVPMvxzHuGgCm6tC5kZyb08yjGWGnqAjGJvcXbqQR2P3MyMEFnRbpcdFS6PBcLqew=="], + "magic-string": ["magic-string@0.30.21", "", { "dependencies": { "@jridgewell/sourcemap-codec": "^1.5.5" } }, "sha512-vd2F4YUyEXKGcLHoq+TEyCjxueSeHnFxyyjNp80yg0XV4vUhnDer/lvvlqM/arB5bXQN5K2/3oinyCRyx8T2CQ=="], "markdown-table": ["markdown-table@3.0.4", "", {}, "sha512-wiYz4+JrLyb/DqW2hkFJxP7Vd7JuTDm77fvbM8VfEQdmSMqcImWeeRbHwZjBjIFki/VaMK2BhFi7oUUZeM5bqw=="], @@ -2232,6 +2238,20 @@ "@babel/helper-create-class-features-plugin/semver": ["semver@6.3.1", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-BR7VvDCVHO+q2xBEWskxS6DJE1qRnb7DxzUrogb71CWoSficBxYsiAGd+Kl0mmq/MprG9yArRkyrQxTO6XjMzA=="], + "@devintern/agent-harness/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@devintern/auth/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@devintern/dashboard-ui/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@devintern/license-check/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@devintern/task-trackers/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@devintern/text-formatter/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@devintern/utils/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + "@dotenvx/dotenvx/dotenv": ["dotenv@17.4.2", "", {}, "sha512-nI4U3TottKAcAD9LLud4Cb7b2QztQMUEfHbvhTH09bqXTxnSie8WnjPALV/WMCrJZ6UV/qHJ6L03OqO3LcdYZw=="], "@dotenvx/dotenvx/env-paths": ["env-paths@2.2.1", "", {}, "sha512-+h1lkLKhZMTYjog1VEpJNG7NZJWcuc2DDk/qsqSTRRCOXiLjeQ1d1/udrUGhqMxUgAlwKNZ0cf2uqan5GLuS2A=="], @@ -2258,6 +2278,8 @@ "@electron/windows-sign/fs-extra": ["fs-extra@11.4.0", "", { "dependencies": { "graceful-fs": "^4.2.0", "jsonfile": "^6.0.1", "universalify": "^2.0.0" } }, "sha512-EQsFzMUJkCKGr1ePqlYADkIUmHW1s3ZXr5Yqy6wbGrfUCphpl2maM/kyOIRA2HpP3AaFQTZXD4ldjek+nccddA=="], + "@getdevintern/license-policy/@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + "@malept/flatpak-bundler/fs-extra": ["fs-extra@9.1.0", "", { "dependencies": { "at-least-node": "^1.0.0", "graceful-fs": "^4.2.0", "jsonfile": "^6.0.1", "universalify": "^2.0.0" } }, "sha512-hcg3ZmepS30/7BSFqRvoo3DOMQu7IjqxO5nCDt+zM9XWjb33Wg7ziNT+Qvqbuc3+gWpzO02JubVyk2G4Zvo1OQ=="], "@radix-ui/react-accordion/@radix-ui/react-compose-refs": ["@radix-ui/react-compose-refs@1.1.5", "", { "peerDependencies": { "@types/react": "*", "react": "^16.8 || ^17.0 || ^18.0 || ^19.0 || ^19.0.0-rc" }, "optionalPeers": ["@types/react"] }, "sha512-+48PbAAbq3didjJxa+OaWY2ZwgAKsNiRGyeHKszblZMQ+kcpd9pAaT11cMkGEie0vsOi3QdeTE6d5Fe3Gn61kA=="], @@ -2350,6 +2372,8 @@ "@tailwindcss/oxide-wasm32-wasi/tslib": ["tslib@2.8.1", "", { "bundled": true }, "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w=="], + "@types/bun/bun-types": ["bun-types@1.3.14", "", { "dependencies": { "@types/node": "*" } }, "sha512-4N0ig0fEomHt5R0KCFWjovxow98rIoRwKolrYdCcknNwMekCXRnWEUvgu5soYV8QXtVsrUD8B95MBOZGPvr6KQ=="], + "@types/cacheable-request/@types/node": ["@types/node@22.20.1", "", { "dependencies": { "undici-types": "~6.21.0" } }, "sha512-EANqOCF9QFyra+4pfxUcX9STKJpCLjMbObVzljIJomAWSnuSIEAvyzEU53GaajbXJEgdh0iEcPL+DGvpUd4k1Q=="], "@types/fs-extra/@types/node": ["@types/node@24.13.3", "", { "dependencies": { "undici-types": "~7.18.0" } }, "sha512-Dh8vAsV36ig5wa9OX4pXvMc9D3Veibfw2wix0CUwYODLD8nkj9UsLjASr49nPg+2eKzxhBV+v7L8pXvT4e639Q=="], @@ -2372,6 +2396,8 @@ "buffer-image-size/@types/node": ["@types/node@22.20.1", "", { "dependencies": { "undici-types": "~6.21.0" } }, "sha512-EANqOCF9QFyra+4pfxUcX9STKJpCLjMbObVzljIJomAWSnuSIEAvyzEU53GaajbXJEgdh0iEcPL+DGvpUd4k1Q=="], + "bun-types/@types/node": ["@types/node@24.13.3", "", { "dependencies": { "undici-types": "~7.18.0" } }, "sha512-Dh8vAsV36ig5wa9OX4pXvMc9D3Veibfw2wix0CUwYODLD8nkj9UsLjASr49nPg+2eKzxhBV+v7L8pXvT4e639Q=="], + "cacheable-request/get-stream": ["get-stream@5.2.0", "", { "dependencies": { "pump": "^3.0.0" } }, "sha512-nBF+F1rAZVCu/p7rjzgA+Yb4lfYXrpl7a6VmJrU8wF9I1CKvP/QwPNZHnOlwbTkY6dvtFIzFMSyQXbLoTQPRpA=="], "chalk/ansi-styles": ["ansi-styles@4.3.0", "", { "dependencies": { "color-convert": "^2.0.1" } }, "sha512-zbB9rCJAT1rbjiVDb2hqKFHNYLxgtk8NURxZ3IZwD3F6NtxbXZQCnnSi1Lkx+IDohdPlFp222wVALIheZJQSEg=="], @@ -2546,6 +2572,8 @@ "axios/https-proxy-agent/agent-base": ["agent-base@6.0.2", "", { "dependencies": { "debug": "4" } }, "sha512-RZNwNclF7+MS/8bDg70amg32dyeZGZxiDuQmZxKLAlQjr3jGyLx+4Kkk58UO7D2QdgFIQCovuSuZESne6RG6XQ=="], + "bun-types/@types/node/undici-types": ["undici-types@7.18.2", "", {}, "sha512-AsuCzffGHJybSaRrmr5eHr81mwJU3kjw6M+uprWvCXiNeN9SOGwQ3Jn8jb8m3Z6izVgknn1R0FTCEAP2QrLY/w=="], + "cliui/string-width/is-fullwidth-code-point": ["is-fullwidth-code-point@3.0.0", "", {}, "sha512-zymm5+u+sCsSWyD9qNaejV3DFvhCKclKdizYaJUuHA83RLjb7nSuGnddCHGv0hk+KY7BMAlsWeK4Ueg6EV6XQg=="], "cliui/strip-ansi/ansi-regex": ["ansi-regex@5.0.1", "", {}, "sha512-quJQXlTSUGL2LH9SUXo8VwsY4soanhgo6LNSm84E1LBcE8s3O0wpdiRzyR9z/ZZJMlMWv37qOOb9pdJlMUEKFQ=="], diff --git a/docs/code/dashboard.md b/docs/code/dashboard.md index 8a0f74e..8410a3b 100644 --- a/docs/code/dashboard.md +++ b/docs/code/dashboard.md @@ -3,12 +3,12 @@ title: "Observability Dashboard" description: "A local web dashboard for worker run history: per-task timelines, stage-by-stage outcomes, and aggregate stats" section: "Server Automation" order: 2 -dateModified: 2026-07-03 +dateModified: 2026-08-24 --- # Observability Dashboard -`devintern dashboard` serves a local web dashboard over the worker's run history: every task and PR mention the worker handled, the stages each run went through (feasibility, implementation, self-review, change requests, outcome), and aggregate stats like success rate and runs per week. +`devintern dashboard` serves a local web dashboard over the worker's run history: every task, PR mention, and scheduled automation the worker handled, the stages each run went through (feasibility, implementation, self-review, change requests, outcome), and aggregate stats like success rate and runs per week. All data is read from the worker's local database (`.devintern-code/queue.db`). Nothing is uploaded anywhere: the dashboard runs on your machine and binds to localhost by default. @@ -28,7 +28,7 @@ The standalone command reads the database in read-only mode, so it is safe to ru ## What it shows -- **Run list**: every run with its status, task key, origin (tracker task or PR mention), agent harness, PR link, and duration. Filter by status or origin. +- **Run list**: every run with its status, task key or automation id, origin (tracker task, PR mention, or scheduled), agent harness, PR link, and duration. Filter by status or origin (`origin=scheduled` isolates automation runs). - **Run detail**: a stage-by-stage timeline for one run: the feasibility verdict, the implementation summary, each self-review iteration, each human change request and how it was handled, and the final outcome. - **Stats**: runs per week, success and escalation rates, median run duration, and a per-harness breakdown over a selectable window (7, 30, or 90 days, or all time). - **Worker status**: whether the daemon is running, queued and failed events, open agent PRs, and per-source poll cursors. @@ -52,7 +52,7 @@ The dashboard is backed by a small read-only JSON API you can use directly, for | Endpoint | Returns | | --------------------------- | --------------------------------------------------------------------- | -| `GET /api/runs` | Paginated run list (`limit`, `offset`, `status`, `origin`, `taskKey`) | +| `GET /api/runs` | Paginated run list (`limit`, `offset`, `status`, `origin`, `taskKey`); `origin=scheduled` is supported | | `GET /api/runs/:id` | One run with its stage timeline | | `GET /api/stats?window=30d` | Aggregate stats (`7d`, `30d`, `90d`, or `all`) | | `GET /api/worker` | Worker liveness, queue counts, agent PRs, poll cursors | diff --git a/docs/code/worker.md b/docs/code/worker.md index 21b12c9..eb62f34 100644 --- a/docs/code/worker.md +++ b/docs/code/worker.md @@ -3,7 +3,7 @@ title: "Worker Daemon" description: "Run devintern as a single long-running worker that reacts to PR reviews and tracker changes" section: "Server Automation" order: 0 -dateModified: 2026-08-17 +dateModified: 2026-08-24 --- # Worker Daemon @@ -34,6 +34,76 @@ devintern worker --query "status=todo" --listen `devintern serve` still works as a deprecated alias for `devintern worker --listen`. +## Recurring automations + +For a single repository, put recurring work in `.devintern-code/automations.toml`: + +```toml +[[automations]] +id = "dependency-health" +enabled = true +interval = "6h" +prompt = """Pick one outdated dependency and upgrade it within the same major version. +Run the test suite; if anything breaks, revert the upgrade instead of fixing forward.""" + +[[automations]] +id = "flaky-test-triage" +enabled = true +cron = "0 9 * * 1" +prompt = """Re-run the test suite twice and look for flaky tests. +For each flaky test, add a short comment explaining the suspected race condition. +Do not change production code.""" +``` + +Every entry needs a stable unique `id`, boolean `enabled`, non-empty `prompt`, and exactly one schedule. Intervals use positive minutes, hours, or days (`15m`, `6h`, `1d`). Cron expressions have five fields and use the worker host's timezone in v1; persisted occurrence times are UTC. + +Configuration is validated as a group at worker startup and changes require a restart. An automation file is itself a valid event source, so `devintern worker` stays running without `--query` or `--listen` when at least one automation entry is configured (disabled entries are validated but not scheduled). + +### What an automation is + +Automations are independent of your task tracker: **the prompt is the task**. Each occurrence writes the prompt to a local markdown task file and feeds it through exactly the same pipeline as any other task — clarity check, planning, implementation, commit, PR creation, auto-review, run records. Nothing is created in your tracker, so no tracker credentials are needed for automation-only workers. + +Concretely, each occurrence: + +1. Writes `.devintern-code/automations//.md` (resolved like every other devintern state: the nearest `.devintern-code` walking up from the working directory). +2. Spawns the normal CLI on that file as a subprocess, so the run gets its own branch, commits, and — by default — a pull request. +3. Records the attempt with the `scheduled` origin and the automation id, so you can filter scheduled runs in the [dashboard](./dashboard.md). + +Because the occurrence is just a markdown task, you can reproduce or rerun any occurrence by hand: + +```bash +devintern .devintern-code/automations/dependency-health/2026-08-24T09-00-00-000Z.md +``` + +### Writing good prompts + +The prompt replaces the ticket description the agent would normally read, so treat it like you would write a task for a new teammate: + +- **Scope it to one change per run.** "Apply one safe improvement" produces reviewable PRs; "clean up the repository" produces sprawling ones. +- **State the guardrails.** What not to touch, when to stop, what must pass (`Run the test suite before committing`). +- **Say what done means.** The pipeline's incomplete-detection reads the agent output; concrete success criteria make escalations rare. +- Prefer recurring maintenance work (dependency bumps within a major, flaky-test triage, changelog refreshes, TODO sweeps) over open-ended feature work. + +### Tuning how occurrences run + +Occurrences use the same flag defaults as polled tasks: `WORKER_TASK_ARGS` overrides them (default `--create-pr`). For example, set `WORKER_TASK_ARGS=--auto-review` to have every automated PR go through the review loop too, or clear it to keep runs local without PRs. Note this variable applies to polled tracker tasks as well. + +### Schedule semantics + +Schedule cursors and claims live in `queue.db`. Missed occurrences coalesce to at most one immediate run after startup. The occurrence cursor advances atomically when claimed, so a crash does not replay a possibly completed run. Active claims receive heartbeats; after two minutes without a heartbeat a later due occurrence may recover the stale claim. If the same automation is still active at its next occurrence, that occurrence is logged and skipped without creating a run record. If the repository lock is held by another task, the occurrence is also skipped. This is an at-most-once policy: skipped occurrences are not replayed. + +On shutdown the scheduler stops its timer, terminates active automation subprocess groups, waits for them to exit, and leaves their claims recoverable in SQLite. + +### Troubleshooting + +| Symptom | Likely cause | +| ---------------------------- | ------------------------------------------------------------ | +| No occurrences fire after editing the TOML | Config is loaded at startup — restart the worker. Startup validation errors name the offending entry. | +| `occurrence skipped: previous run is active` | The previous occurrence still runs (or its lease is stale). Long prompts may simply need a longer schedule. | +| `occurrence skipped: repository is busy` | Another task holds the repo run lock; the next occurrence will retry. | +| Scheduled runs missing from the dashboard | Filter the run list by origin `scheduled`; check the worker has an automation license (startup log). | +| Task files pile up under `.devintern-code/automations/` | They are small and safe to delete — they are only run inputs; the durable record is the run history in `queue.db`. | + ## Polling mode With `--query` (or `WORKER_TASK_QUERY`), the worker polls your tracker on an interval (default 60 seconds) and runs every task that matches the query. The query uses the same language as batch `--query` runs for your tracker, so "ready" means whatever your query says, for example a status or label. @@ -113,7 +183,7 @@ Mention matching requires a resolvable bot identity, so this team/automation fea - Events are persisted to a local SQLite queue (`.devintern-code/queue.db`) before processing, so a crash or restart never loses accepted work. - Duplicate webhook deliveries are detected by GitHub's delivery id and skipped. - Review feedback is processed before new task pickup: a human waiting on feedback beats a ticket that can wait a minute. -- One task runs at a time per repository. +- One task or scheduled automation runs at a time per repository. ## Instant events with the relay diff --git a/docs/code/workspaces.md b/docs/code/workspaces.md index 1619a4c..27a683c 100644 --- a/docs/code/workspaces.md +++ b/docs/code/workspaces.md @@ -3,7 +3,7 @@ title: "Workspaces (Multi-Repo Fleet)" description: "Drive many repositories with one devintern worker: a single workspace.toml, routing rules, and per-task worktrees" section: "Server Automation" order: 1 -dateModified: 2026-08-17 +dateModified: 2026-08-24 --- # Workspaces (Multi-Repo Fleet) @@ -52,11 +52,34 @@ project = "BACK" repo = "frontend" project = "WEB" labels = ["frontend"] + +[[automations]] +id = "backend-maintenance" +enabled = true +interval = "6h" +repo = "backend" +prompt = "Inspect the backend and implement one safe maintenance improvement." + +[[automations]] +id = "weekly-frontend-cleanup" +enabled = true +cron = "0 9 * * 1" +repo = "frontend" +prompt = "Review the frontend and clean up one source of recurring noise." ``` - `[defaults].tracker` picks the tracker for the fleet query; any tracker with polling support works (Jira, Linear, GitHub Issues, Azure DevOps, Asana, Trello, Markdown). - Repo names must be unique and filesystem-safe; they become directory names under `repos/` and `worktrees/`. - Rule criteria combine with AND; list values (`components`, `labels`) match when the task carries any of them. Comparisons are case-insensitive. `project` matches the task key prefix for `PROJ-123` style keys (Jira, Linear); trackers with numeric or opaque ids route via labels or components. +- `[[automations]]` uses the same schema as single-repo `.devintern-code/automations.toml`. An entry must name `repo` when the workspace has more than one repository. See [Worker Daemon → Recurring automations](./worker.md#recurring-automations) for prompt-writing guidance and schedule semantics. + +### How workspace automations differ from single-repo ones + +The scheduling is identical; only where the work runs changes: + +- Each occurrence runs in the repo's persistent base worktree (`~/.devintern/worktrees//base`) with the same layered environment as review work: shared `.env` → repo `env_file` → `[repos.env]`. +- It takes the normal per-repo run lock, so it never mutates a checkout concurrently with a task or PR run. +- Occurrence task files land under the workspace home (`~/.devintern/automations//`), next to `repos/`, `worktrees/`, and the central database — not inside the repo worktrees. ## Creating a workspace @@ -94,7 +117,7 @@ devintern worker --workspace /path/to/workspace.toml devintern worker --no-workspace # force single-repo mode in the current repo ``` -The fleet query comes from `[defaults].task_query`, or `--query` / `WORKER_TASK_QUERY` to override. `--listen` (direct webhooks) is single-repo and cannot be combined with workspace mode. +The fleet query comes from `[defaults].task_query`, or `--query` / `WORKER_TASK_QUERY` to override. A workspace with automations can omit the query and run as an automation-only worker. `--listen` (direct webhooks) is single-repo and cannot be combined with workspace mode. Workspace and automation configuration is loaded at startup; restart the worker after editing it. Schedule state and leases for automations live in the central workspace database. One systemd unit runs the whole fleet: diff --git a/packages/code/package.json b/packages/code/package.json index 565e225..4052cd1 100644 --- a/packages/code/package.json +++ b/packages/code/package.json @@ -61,6 +61,7 @@ "dependencies": { "@types/turndown": "^5.0.5", "commander": "^11.0.0", + "cron-parser": "^5.4.0", "dotenv": "^16.3.1", "p-queue": "^9.0.1", "turndown": "^7.2.0" diff --git a/packages/code/src/index.ts b/packages/code/src/index.ts index a9c0f6e..6b912fa 100755 --- a/packages/code/src/index.ts +++ b/packages/code/src/index.ts @@ -626,6 +626,27 @@ if (process.argv[2] === "init") { const dbPath = resolveQueueDbPath(); const acquirers = []; + const { loadSingleRepoAutomations } = await import("./lib/automation-config"); + const automations = loadSingleRepoAutomations(); + if (automations.length > 0) { + const { AutomationAcquirer } = await import("./lib/automation-acquirer"); + acquirers.push( + new AutomationAcquirer({ + automations, + dbPath, + resolveContext: async () => { + const runLock = new LockManager(process.cwd()); + if (!runLock.acquire().success) return null; + return { + cwd: process.cwd(), + env: { ...process.env }, + release: () => runLock.release(), + }; + }, + }), + ); + } + if (workerQuery) { const trackerType = process.env.TASK_TRACKER || "jira"; const { supportsPolling, trackersSupportingPolling } = @@ -1598,10 +1619,14 @@ async function processSingleTask(taskKey: string, taskIndex = 0, totalTasks = 1) } // Structured run record for this attempt (skips above are not attempts). + // Scheduled automations run through this same pipeline with their prompt + // materialized as a markdown task; env markers attribute those runs. + const scheduledAutomationId = process.env.DEVINTERN_AUTOMATION_ID; beginRun({ - origin: "task", + origin: scheduledAutomationId ? "scheduled" : "task", taskKey: workflowKey, tracker: process.env.TASK_TRACKER || "jira", + ...(scheduledAutomationId ? { automationId: scheduledAutomationId } : {}), }); if (!isMarkdownTaskTracker(tracker)) { diff --git a/packages/code/src/lib/automation-acquirer.ts b/packages/code/src/lib/automation-acquirer.ts new file mode 100644 index 0000000..71b3622 --- /dev/null +++ b/packages/code/src/lib/automation-acquirer.ts @@ -0,0 +1,394 @@ +import { randomUUID } from "crypto"; +import { mkdirSync, writeFileSync } from "fs"; +import { spawn } from "child_process"; +import type { ChildProcess } from "child_process"; +import { join } from "path"; + +import { CronExpressionParser } from "cron-parser"; +import { resolveConfigDir } from "@devintern/utils"; + +import type { Acquirer } from "../worker"; +import type { AutomationConfig } from "./automation-config"; +import { AutomationStateStore } from "./automation-state"; +import { workerTaskArgs } from "./task-polling-acquirer"; + +/** Environment markers the task pipeline reads to attribute scheduled runs. */ +export const AUTOMATION_ORIGIN_ENV = "DEVINTERN_RUN_ORIGIN"; +export const AUTOMATION_ID_ENV = "DEVINTERN_AUTOMATION_ID"; + +const LEASE_MS = 2 * 60_000; +const HEARTBEAT_MS = 30_000; +const TERMINATION_GRACE_MS = 5_000; +export const MAX_TIMER_DELAY_MS = 2_147_483_647; + +export interface AutomationRunContext { + cwd: string; + env: Record; + repo?: string; + /** + * Explicit directory for occurrence task files. When omitted, files go + * under the nearest `.devintern-code` found above {@linkcode cwd}. + */ + taskFileDir?: string; + release(): void | Promise; +} + +export interface SpawnedAutomationRun { + completion: Promise; + terminate(): void; +} + +export interface AutomationAcquirerOptions { + automations: AutomationConfig[]; + dbPath: string; + resolveContext: (automation: AutomationConfig) => Promise; + now?: () => number; + spawnRun?: (automation: AutomationConfig, context: AutomationRunContext) => SpawnedAutomationRun; + leaseMs?: number; + heartbeatMs?: number; + terminationGraceMs?: number; + setTimer?: (callback: () => void, delay: number) => ReturnType; + clearTimer?: (timer: ReturnType) => void; +} + +interface ActiveAutomationRun { + run: SpawnedAutomationRun; + lifecycle: Promise; +} + +/** Calculate the first future occurrence after `afterMs` (cron uses host timezone). */ +export function nextAutomationDue(automation: AutomationConfig, afterMs: number): number { + if (automation.intervalMs) return afterMs + automation.intervalMs; + if (!automation.cron) throw new Error(`Automation "${automation.id}" has no schedule`); + return CronExpressionParser.parse(automation.cron, { currentDate: new Date(afterMs) }) + .next() + .getTime(); +} + +/** One-timer scheduler with durable UTC cursors and per-automation leases. */ +export class AutomationAcquirer implements Acquirer { + readonly name = "scheduled-automations"; + private options: AutomationAcquirerOptions; + private store: AutomationStateStore; + private owner = `${process.pid}:${randomUUID()}`; + private timer: ReturnType | null = null; + private active = new Map(); + private tickPromise: Promise | null = null; + private stopped = true; + + constructor(options: AutomationAcquirerOptions) { + this.options = options; + this.store = new AutomationStateStore(options.dbPath); + } + + async start(): Promise { + this.stopped = false; + const now = this.now(); + for (const automation of this.options.automations.filter((item) => item.enabled)) { + this.store.register(automation, nextAutomationDue(automation, now)); + } + console.log( + `⏰ Scheduling ${this.options.automations.filter((item) => item.enabled).length} enabled automation(s)`, + ); + await this.tick(); + } + + async stop(): Promise { + this.stopped = true; + if (this.timer) (this.options.clearTimer ?? clearTimeout)(this.timer); + this.timer = null; + await this.tickPromise; + for (const active of this.active.values()) active.run.terminate(); + await Promise.allSettled([...this.active.values()].map((active) => active.lifecycle)); + this.store.close(); + } + + /** Public for deterministic tests; production calls it through one setTimeout. */ + async tick(): Promise { + if (this.stopped) return; + if (this.tickPromise) return this.tickPromise; + const tickPromise = this.runTick(); + this.tickPromise = tickPromise; + try { + await tickPromise; + } finally { + if (this.tickPromise === tickPromise) this.tickPromise = null; + } + } + + private async runTick(): Promise { + const leaseMs = this.options.leaseMs ?? LEASE_MS; + for (const automation of this.options.automations.filter((item) => item.enabled)) { + let state = this.store.get(automation.id); + if (!state) continue; + + if (this.active.has(automation.id)) { + const active = this.active.get(automation.id) as ActiveAutomationRun; + if (!this.store.heartbeat(automation.id, this.owner, this.now(), leaseMs)) { + console.warn(`⏭️ [automation:${automation.id}] terminating: lease was lost`); + active.run.terminate(); + await active.lifecycle; + } + state = this.store.get(automation.id); + if (!state) continue; + } + if (state.nextDueAt > this.now()) continue; + + const overlapNow = this.now(); + if (state.leaseOwner && (state.leaseExpiresAt ?? 0) > overlapNow) { + const skipNow = this.now(); + const nextDue = nextAutomationDue(automation, skipNow); + if (this.store.skipOverlap(automation.id, skipNow, nextDue)) { + console.warn( + `⏭️ [automation:${automation.id}] occurrence skipped: previous run is active`, + ); + } + continue; + } + const claimNow = this.now(); + const nextDue = nextAutomationDue(automation, claimNow); + if (!this.store.claim(automation.id, this.owner, claimNow, nextDue, leaseMs)) continue; + + let context: AutomationRunContext | null = null; + try { + let ownsClaim = true; + const heartbeatMs = Math.min(this.options.heartbeatMs ?? HEARTBEAT_MS, leaseMs / 2); + const preparationHeartbeat = setInterval( + () => { + if (!this.store.heartbeat(automation.id, this.owner, this.now(), leaseMs)) { + ownsClaim = false; + } + }, + Math.max(1, heartbeatMs), + ); + preparationHeartbeat.unref(); + try { + context = await this.options.resolveContext(automation); + } finally { + clearInterval(preparationHeartbeat); + } + if (!context) { + console.warn(`⏭️ [automation:${automation.id}] occurrence skipped: repository is busy`); + this.store.release(automation.id, this.owner); + continue; + } + if (this.stopped) { + this.store.release(automation.id, this.owner); + await context.release(); + continue; + } + ownsClaim &&= this.store.heartbeat(automation.id, this.owner, this.now(), leaseMs); + if (!ownsClaim) { + console.warn(`⏭️ [automation:${automation.id}] occurrence skipped: lease was lost`); + await context.release(); + continue; + } + console.log(`\n⏰ [automation:${automation.id}] starting scheduled run`); + const run = this.options.spawnRun + ? this.options.spawnRun(automation, context) + : defaultSpawnRun(automation, context, this.options.terminationGraceMs); + const active: ActiveAutomationRun = { + run, + lifecycle: Promise.resolve(), + }; + active.lifecycle = run.completion + .then((ok) => + console.log( + ok + ? `✅ [automation:${automation.id}] completed` + : `⚠️ [automation:${automation.id}] did not complete cleanly`, + ), + ) + .catch((error) => + console.error(`❌ [automation:${automation.id}] ${(error as Error).message}`), + ) + .finally(async () => { + try { + this.store.release(automation.id, this.owner); + await context?.release(); + } finally { + if (this.active.get(automation.id) === active) this.active.delete(automation.id); + this.scheduleNext(); + } + }); + this.active.set(automation.id, active); + } catch (error) { + this.store.release(automation.id, this.owner); + await context?.release(); + console.error(`❌ [automation:${automation.id}] ${(error as Error).message}`); + } + } + this.scheduleNext(); + } + + private now(): number { + return (this.options.now ?? Date.now)(); + } + + private scheduleNext(): void { + if (this.stopped) return; + if (this.timer) (this.options.clearTimer ?? clearTimeout)(this.timer); + const now = this.now(); + const dueTimes = this.options.automations + .filter((item) => item.enabled) + .map((item) => this.store.get(item.id)?.nextDueAt) + .filter((value): value is number => value !== undefined); + const heartbeatAt = + this.active.size > 0 ? now + (this.options.heartbeatMs ?? HEARTBEAT_MS) : Infinity; + const wakeAt = Math.min(heartbeatAt, ...dueTimes); + if (!Number.isFinite(wakeAt)) return; + this.timer = (this.options.setTimer ?? setTimeout)( + () => void this.tick(), + Math.min(MAX_TIMER_DELAY_MS, Math.max(0, wakeAt - now)), + ); + } +} + +/** + * Directory that receives one markdown task file per occurrence. + * + * When the context does not pin an explicit directory, resolved like every + * other durable-state location: the nearest existing `.devintern-code` found + * by walking up from the run cwd, so a worker launched from a subfolder + * reuses the project's config directory instead of creating a stray one + * beside the cwd. Falls back to `/.devintern-code` (e.g. inside + * disposable worktrees where no parent config exists). + */ +export function automationTaskDir(context: AutomationRunContext): string { + return ( + context.taskFileDir ?? + join( + resolveConfigDir({ configDirName: ".devintern-code", startDir: context.cwd }), + "automations", + ) + ); +} + +/** + * Materialize the automation prompt as a local markdown task file so the + * regular task pipeline can process it like any other tracker-less task. + */ +export function writeAutomationTaskFile( + automation: AutomationConfig, + context: AutomationRunContext, +): string { + const dir = join(automationTaskDir(context), automation.id); + mkdirSync(dir, { recursive: true }); + const stamp = new Date().toISOString().replace(/[:.]/g, "-"); + const filePath = join(dir, `${stamp}.md`); + const body = [ + "---", + "type: Task", + "---", + "", + `# ${automation.id}`, + "", + automation.prompt.trim(), + "", + ].join("\n"); + writeFileSync(filePath, body); + return filePath; +} + +function defaultSpawnRun( + automation: AutomationConfig, + context: AutomationRunContext, + terminationGraceMs?: number, +): SpawnedAutomationRun { + const taskFile = writeAutomationTaskFile(automation, context); + const env: Record = { + ...context.env, + [AUTOMATION_ORIGIN_ENV]: "scheduled", + [AUTOMATION_ID_ENV]: automation.id, + }; + return spawnAutomationProcess( + process.execPath, + [process.argv[1] as string, taskFile, ...workerTaskArgs()], + { + cwd: context.cwd, + env, + terminationGraceMs, + }, + ); +} + +/** + * Spawn one isolated automation subprocess (the normal CLI pipeline for the + * materialized task file) and terminate its process tree within a bound. + */ +export function spawnAutomationProcess( + executable: string, + args: string[], + options: { + cwd: string; + env: Record; + terminationGraceMs?: number; + }, +): SpawnedAutomationRun { + const detached = process.platform !== "win32"; + const child: ChildProcess = spawn(executable, args, { + cwd: options.cwd, + env: options.env, + stdio: "inherit", + detached, + }); + let terminating = false; + let exitResult: boolean | undefined; + let terminationTimer: ReturnType | undefined; + let settle!: (ok: boolean) => void; + const completion = new Promise((resolve) => { + settle = resolve; + }); + const kill = (signal: NodeJS.Signals) => { + if (child.pid && detached) { + try { + process.kill(-child.pid, signal); + return; + } catch { + // Fall through to the direct child. + } + } + try { + child.kill(signal); + } catch { + // The child may already have exited. + } + }; + const processGroupAlive = () => { + if (!child.pid || !detached) return false; + try { + process.kill(-child.pid, 0); + return true; + } catch { + return false; + } + }; + const beginTermination = () => { + if (terminating || exitResult !== undefined) return; + terminating = true; + kill("SIGTERM"); + terminationTimer = setTimeout(() => { + kill("SIGKILL"); + settle(false); + }, options.terminationGraceMs ?? TERMINATION_GRACE_MS); + }; + + child.once("close", (code) => { + exitResult = code === 0; + if (!terminating) settle(exitResult); + else if (!processGroupAlive()) { + if (terminationTimer) clearTimeout(terminationTimer); + settle(false); + } + }); + child.once("error", () => { + exitResult = false; + if (terminationTimer) clearTimeout(terminationTimer); + settle(false); + }); + + return { + completion, + terminate: beginTermination, + }; +} diff --git a/packages/code/src/lib/automation-config.ts b/packages/code/src/lib/automation-config.ts new file mode 100644 index 0000000..26ff287 --- /dev/null +++ b/packages/code/src/lib/automation-config.ts @@ -0,0 +1,146 @@ +import { existsSync, readFileSync } from "fs"; +import { join } from "path"; + +import { CronExpressionParser } from "cron-parser"; + +import { parseToml } from "./workspace/toml"; + +export interface AutomationConfig { + id: string; + enabled: boolean; + prompt: string; + cron?: string; + interval?: string; + intervalMs?: number; + repo?: string; +} + +export const SINGLE_REPO_AUTOMATIONS_PATH = ".devintern-code/automations.toml"; +const AUTOMATION_ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._-]*$/; +const DURATION_PATTERN = /^(\d+)([mhd])$/; +const MAX_DATE_MS = 8_640_000_000_000_000; + +/** Parse a documented interval (`15m`, `6h`, or `1d`) into milliseconds. */ +export function parseAutomationInterval(value: string, nowMs = Date.now()): number | null { + const match = value.match(DURATION_PATTERN); + if (!match) return null; + const amount = Number(match[1]); + if (!Number.isSafeInteger(amount) || amount < 1) return null; + const unitMs = match[2] === "m" ? 60_000 : match[2] === "h" ? 3_600_000 : 86_400_000; + const intervalMs = amount * unitMs; + const dueAt = nowMs + intervalMs; + if (!Number.isSafeInteger(intervalMs) || !Number.isSafeInteger(dueAt) || dueAt > MAX_DATE_MS) { + return null; + } + return intervalMs; +} + +/** Validate and normalize `[[automations]]` tables, collecting every error. */ +export function parseAutomationEntries( + value: unknown, + options: { sourceLabel: string; repoNames?: Set }, +): { automations: AutomationConfig[]; errors: string[] } { + const errors: string[] = []; + const automations: AutomationConfig[] = []; + const ids = new Set(); + if (value === undefined || value === null) return { automations, errors }; + if (!Array.isArray(value)) { + return { automations, errors: ["[[automations]] must be an array of tables."] }; + } + + for (const [index, raw] of value.entries()) { + const label = `[[automations]][${index}]`; + if (!raw || typeof raw !== "object" || Array.isArray(raw)) { + errors.push(`${label} must be a table.`); + continue; + } + const table = raw as Record; + const stringValue = (key: string): string | undefined => { + const item = table[key]; + if (item === undefined || item === null) return undefined; + if (typeof item !== "string" || !item.trim()) { + errors.push(`${label}.${key} must be a non-empty string.`); + return undefined; + } + return key === "prompt" ? item : item.trim(); + }; + + const id = stringValue("id"); + if (!id) errors.push(`${label}.id is required.`); + else if (!AUTOMATION_ID_PATTERN.test(id)) + errors.push(`${label}.id must contain only letters, digits, ".", "_", or "-".`); + else if (ids.has(id)) errors.push(`Duplicate automation id "${id}".`); + else ids.add(id); + + const enabledValue = table.enabled; + if (typeof enabledValue !== "boolean") errors.push(`${label}.enabled must be a boolean.`); + const prompt = stringValue("prompt"); + if (!prompt) errors.push(`${label}.prompt is required.`); + + const cron = stringValue("cron"); + const interval = stringValue("interval"); + if (Boolean(cron) === Boolean(interval)) { + errors.push(`${label} must set exactly one of cron or interval.`); + } + if (cron) { + if (cron.split(/\s+/).length !== 5) { + errors.push(`${label}.cron must be a five-field cron expression.`); + } else { + try { + CronExpressionParser.parse(cron); + } catch (error) { + errors.push(`${label}.cron is invalid: ${(error as Error).message}`); + } + } + } + const intervalMs = interval ? parseAutomationInterval(interval) : undefined; + if (interval && intervalMs === null) { + errors.push(`${label}.interval must use a positive duration such as 15m, 6h, or 1d.`); + } + + const repo = stringValue("repo"); + if (repo && options.repoNames && !options.repoNames.has(repo)) { + errors.push(`${label}.repo "${repo}" does not match any [[repos]] name.`); + } + + if ( + id && + typeof enabledValue === "boolean" && + prompt && + Boolean(cron) !== Boolean(interval) && + (!interval || intervalMs !== null) + ) { + automations.push({ + id, + enabled: enabledValue, + prompt, + cron, + interval, + intervalMs: intervalMs ?? undefined, + repo, + }); + } + } + return { automations, errors }; +} + +/** Parse the single-repository automation TOML file. */ +export function parseAutomationConfig(text: string, sourceLabel = SINGLE_REPO_AUTOMATIONS_PATH) { + let document: Record; + try { + document = parseToml(text); + } catch (error) { + throw new Error(`Failed to parse ${sourceLabel}: ${(error as Error).message}`); + } + const result = parseAutomationEntries(document.automations, { sourceLabel }); + if (result.errors.length > 0) { + throw new Error(`Invalid ${sourceLabel}:\n- ${result.errors.join("\n- ")}`); + } + return result.automations; +} + +/** Load single-repository automations, returning an empty list when absent. */ +export function loadSingleRepoAutomations(baseDir = process.cwd()): AutomationConfig[] { + const path = join(baseDir, SINGLE_REPO_AUTOMATIONS_PATH); + return existsSync(path) ? parseAutomationConfig(readFileSync(path, "utf8"), path) : []; +} diff --git a/packages/code/src/lib/automation-state.ts b/packages/code/src/lib/automation-state.ts new file mode 100644 index 0000000..c2385ed --- /dev/null +++ b/packages/code/src/lib/automation-state.ts @@ -0,0 +1,123 @@ +import { Database } from "bun:sqlite"; + +import type { AutomationConfig } from "./automation-config"; +import { prepareQueueDbDirectory } from "./webhook-queue"; + +export interface AutomationScheduleState { + automationId: string; + lastScheduledAt?: number; + nextDueAt: number; + leaseOwner?: string; + leaseExpiresAt?: number; + heartbeatAt?: number; +} + +/** Durable schedule cursors and overlap leases stored beside the worker queue. */ +export class AutomationStateStore { + private db: Database; + + constructor(dbPath: string) { + prepareQueueDbDirectory(dbPath); + this.db = new Database(dbPath); + this.db.run("PRAGMA busy_timeout = 5000"); + this.db.run("PRAGMA journal_mode = WAL"); + this.db.run(` + CREATE TABLE IF NOT EXISTS automation_schedules ( + automation_id TEXT PRIMARY KEY, + schedule_spec TEXT NOT NULL, + last_scheduled_at INTEGER, + next_due_at INTEGER NOT NULL, + lease_owner TEXT, + lease_expires_at INTEGER, + heartbeat_at INTEGER + ) + `); + } + + /** Create schedule state once; restarts retain the prior interval anchor. */ + register(automation: AutomationConfig, nextDueAt: number): void { + const spec = automation.cron ? `cron:${automation.cron}` : `interval:${automation.interval}`; + const existing = this.db + .query("SELECT schedule_spec FROM automation_schedules WHERE automation_id = ?") + .get(automation.id) as { schedule_spec: string } | null; + if (!existing) { + this.db.run( + `INSERT INTO automation_schedules (automation_id, schedule_spec, next_due_at) + VALUES (?, ?, ?)`, + [automation.id, spec, nextDueAt], + ); + } else if (existing.schedule_spec !== spec) { + this.db.run( + `UPDATE automation_schedules SET schedule_spec = ?, next_due_at = ?, + last_scheduled_at = NULL WHERE automation_id = ?`, + [spec, nextDueAt, automation.id], + ); + } + } + + get(automationId: string): AutomationScheduleState | null { + const row = this.db + .query("SELECT * FROM automation_schedules WHERE automation_id = ?") + .get(automationId) as Record | null; + if (!row) return null; + return { + automationId: row.automation_id as string, + lastScheduledAt: (row.last_scheduled_at as number | null) ?? undefined, + nextDueAt: row.next_due_at as number, + leaseOwner: (row.lease_owner as string | null) ?? undefined, + leaseExpiresAt: (row.lease_expires_at as number | null) ?? undefined, + heartbeatAt: (row.heartbeat_at as number | null) ?? undefined, + }; + } + + /** Atomically claim one due occurrence and advance its cursor before execution. */ + claim( + automationId: string, + owner: string, + now: number, + nextDueAt: number, + leaseMs: number, + ): boolean { + const result = this.db.run( + `UPDATE automation_schedules + SET last_scheduled_at = next_due_at, next_due_at = ?, lease_owner = ?, + lease_expires_at = ?, heartbeat_at = ? + WHERE automation_id = ? AND next_due_at <= ? + AND (lease_owner IS NULL OR lease_expires_at <= ?)`, + [nextDueAt, owner, now + leaseMs, now, automationId, now, now], + ); + return result.changes === 1; + } + + /** Coalesce a due occurrence skipped because the previous run is still active. */ + skipOverlap(automationId: string, now: number, nextDueAt: number): boolean { + const result = this.db.run( + `UPDATE automation_schedules SET last_scheduled_at = next_due_at, next_due_at = ? + WHERE automation_id = ? AND next_due_at <= ? + AND lease_owner IS NOT NULL AND lease_expires_at > ?`, + [nextDueAt, automationId, now, now], + ); + return result.changes === 1; + } + + heartbeat(automationId: string, owner: string, now: number, leaseMs: number): boolean { + const result = this.db.run( + `UPDATE automation_schedules SET heartbeat_at = ?, lease_expires_at = ? + WHERE automation_id = ? AND lease_owner = ?`, + [now, now + leaseMs, automationId, owner], + ); + return result.changes === 1; + } + + release(automationId: string, owner: string): void { + this.db.run( + `UPDATE automation_schedules SET lease_owner = NULL, lease_expires_at = NULL, + heartbeat_at = NULL WHERE automation_id = ? AND lease_owner = ?`, + [automationId, owner], + ); + } + + close(): void { + this.db.close(); + } +} diff --git a/packages/code/src/lib/dashboard-api.ts b/packages/code/src/lib/dashboard-api.ts index 627f937..ab1c774 100644 --- a/packages/code/src/lib/dashboard-api.ts +++ b/packages/code/src/lib/dashboard-api.ts @@ -26,7 +26,7 @@ const RUN_STATUSES: RunStatus[] = [ "escalated", "abandoned", ]; -const RUN_ORIGINS: RunOrigin[] = ["task", "pr_mention"]; +const RUN_ORIGINS: RunOrigin[] = ["task", "pr_mention", "scheduled"]; const STATS_WINDOWS: Record = { "7d": 7 * 24 * 60 * 60 * 1000, diff --git a/packages/code/src/lib/init-scaffold.ts b/packages/code/src/lib/init-scaffold.ts index 5c341a0..be51c56 100644 --- a/packages/code/src/lib/init-scaffold.ts +++ b/packages/code/src/lib/init-scaffold.ts @@ -445,6 +445,7 @@ export function scaffoldProject(options: ScaffoldOptions = {}): boolean { ".devintern-code/*", "!.devintern-code/settings.json", "!.devintern-code/.env.example", + "!.devintern-code/automations.toml", ]; try { diff --git a/packages/code/src/lib/run-recorder.ts b/packages/code/src/lib/run-recorder.ts index fabb0e3..da333bb 100644 --- a/packages/code/src/lib/run-recorder.ts +++ b/packages/code/src/lib/run-recorder.ts @@ -15,7 +15,7 @@ import { Database } from "bun:sqlite"; import { prepareQueueDbDirectory, resolveQueueDbPath } from "./webhook-queue"; -export type RunOrigin = "task" | "pr_mention"; +export type RunOrigin = "task" | "pr_mention" | "scheduled"; export type RunStatus = | "in_progress" @@ -50,6 +50,7 @@ export interface RunMeta { branch?: string; repo?: string; prNumber?: number; + automationId?: string; } export interface RunRecord extends RunMeta { @@ -57,6 +58,8 @@ export interface RunRecord extends RunMeta { /** 1-based attempt number for the task (null-ish for pr_mention runs). */ attempt?: number; prUrl?: string; + ticketKey?: string; + ticketUrl?: string; status: RunStatus; outcomeReason?: string; startedAt: number; @@ -181,7 +184,10 @@ export class RunStore { outcome_reason TEXT, started_at INTEGER NOT NULL, finished_at INTEGER, - attempt INTEGER + attempt INTEGER, + automation_id TEXT, + ticket_key TEXT, + ticket_url TEXT ) `); @@ -190,6 +196,16 @@ export class RunStore { if (!columns.some((c) => c.name === "attempt")) { this.db.run("ALTER TABLE runs ADD COLUMN attempt INTEGER"); } + if (!columns.some((c) => c.name === "automation_id")) { + this.db.run("ALTER TABLE runs ADD COLUMN automation_id TEXT"); + } + if (!columns.some((c) => c.name === "ticket_key")) { + this.db.run("ALTER TABLE runs ADD COLUMN ticket_key TEXT"); + } + if (!columns.some((c) => c.name === "ticket_url")) { + this.db.run("ALTER TABLE runs ADD COLUMN ticket_url TEXT"); + } + this.db.run("CREATE INDEX IF NOT EXISTS idx_runs_automation_id ON runs(automation_id)"); this.db.run(` CREATE INDEX IF NOT EXISTS idx_runs_task_key ON runs(task_key) @@ -221,8 +237,9 @@ export class RunStore { createRun(meta: RunMeta): number { const attempt = meta.taskKey ? this.countRuns(meta.taskKey) + 1 : null; const result = this.db.run( - `INSERT INTO runs (origin, task_key, tracker, harness, branch, repo, pr_number, status, started_at, attempt) - VALUES (?, ?, ?, ?, ?, ?, ?, 'in_progress', ?, ?)`, + `INSERT INTO runs (origin, task_key, tracker, harness, branch, repo, pr_number, + automation_id, status, started_at, attempt) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'in_progress', ?, ?)`, [ meta.origin, meta.taskKey ?? null, @@ -231,6 +248,7 @@ export class RunStore { meta.branch ?? null, meta.repo ?? null, meta.prNumber ?? null, + meta.automationId ?? null, Date.now(), attempt, ], @@ -408,7 +426,7 @@ export class RunStore { escalated: 0, abandoned: 0, }; - const byOrigin: Record = { task: 0, pr_mention: 0 }; + const byOrigin: Record = { task: 0, pr_mention: 0, scheduled: 0 }; const weekCounts = new Map(); const harnesses = new Map< string, @@ -418,7 +436,7 @@ export class RunStore { for (const row of rows) { byStatus[row.status] += 1; - byOrigin[row.origin] += 1; + byOrigin[row.origin] = (byOrigin[row.origin] ?? 0) + 1; const week = weekStartIso(row.started_at); weekCounts.set(week, (weekCounts.get(week) ?? 0) + 1); @@ -502,6 +520,9 @@ export class RunStore { startedAt: row.started_at as number, finishedAt: (row.finished_at as number | null) ?? undefined, attempt: (row.attempt as number | null) ?? undefined, + automationId: (row.automation_id as string | null) ?? undefined, + ticketKey: (row.ticket_key as string | null) ?? undefined, + ticketUrl: (row.ticket_url as string | null) ?? undefined, }; } diff --git a/packages/code/src/lib/workspace/config.ts b/packages/code/src/lib/workspace/config.ts index 7605e33..e3e3736 100644 --- a/packages/code/src/lib/workspace/config.ts +++ b/packages/code/src/lib/workspace/config.ts @@ -1,6 +1,8 @@ import { readFileSync } from "fs"; import { supportsPolling, trackersSupportingPolling } from "../tracker-capabilities"; +import { parseAutomationEntries } from "../automation-config"; +import type { AutomationConfig } from "../automation-config"; import { parseToml } from "./toml"; /** Workspace-wide settings from the `[workspace]` table. */ @@ -53,6 +55,7 @@ export interface WorkspaceConfig { defaults: WorkspaceDefaults; repos: RepoConfig[]; routing: RoutingRule[]; + automations: AutomationConfig[]; } export const DEFAULT_WORKTREES_TTL_DAYS = 7; @@ -255,11 +258,23 @@ export function parseWorkspaceConfig( routing.push(rule); } + const automationResult = parseAutomationEntries(document.automations, { + sourceLabel, + repoNames, + }); + errors.push(...automationResult.errors); + if (errors.length > 0) { throw new Error(`Invalid ${sourceLabel}:\n- ${errors.join("\n- ")}`); } - return { workspace: { worktreesTtlDays }, defaults, repos, routing }; + return { + workspace: { worktreesTtlDays }, + defaults, + repos, + routing, + automations: automationResult.automations, + }; } /** diff --git a/packages/code/src/lib/workspace/init.ts b/packages/code/src/lib/workspace/init.ts index 39fbc8a..6df771a 100644 --- a/packages/code/src/lib/workspace/init.ts +++ b/packages/code/src/lib/workspace/init.ts @@ -46,6 +46,17 @@ default_branch = "main" # repo = "backend" # project = "BACK" # task key prefix (BACK-123) # labels = ["backend"] # any-of; AND-ed with the other criteria + +# Recurring work is loaded once when the worker starts. Each occurrence runs +# the prompt through the normal task pipeline as a local markdown task. +# Cron uses the worker host timezone; interval values support m, h, and d. +# +# [[automations]] +# id = "weekday-maintenance" +# enabled = true +# cron = "0 9 * * 1-5" # exactly one of cron / interval +# prompt = "Inspect dependency health and fix one safe issue." +# repo = "backend" # required in a multi-repo workspace `; const ENV_TEMPLATE = `# Shared workspace environment: tracker credentials, GITHUB_TOKEN and/or diff --git a/packages/code/src/lib/workspace/workspace-worker.ts b/packages/code/src/lib/workspace/workspace-worker.ts index 03d58cb..f99540b 100644 --- a/packages/code/src/lib/workspace/workspace-worker.ts +++ b/packages/code/src/lib/workspace/workspace-worker.ts @@ -8,6 +8,8 @@ * in the central workspace DB. */ +import { join } from "path"; + import { LockManager } from "../lock-manager"; import { TaskPollingAcquirer, runTaskViaCli, workerTaskArgs } from "../task-polling-acquirer"; import type { ChangeDetector } from "../change-detector"; @@ -27,6 +29,8 @@ import type { RoutableTask } from "./router"; import { createRepoRunLock, createWorkspaceLock, openWorkspaceState } from "./state"; import type { RoutingSkipStore } from "./state"; import { RepoManager } from "./repo-manager"; +import { AutomationAcquirer } from "../automation-acquirer"; +import type { AutomationConfig } from "../automation-config"; /** Task shape the fleet acquirer needs (structural subset of `Task`). */ export interface FleetTask { @@ -68,6 +72,48 @@ export interface WorkspaceTaskAcquirerDeps { repoLock?: (repoName: string) => LockManager; } +interface RepoRunLockLike { + acquire(): { success: boolean; message: string; pid?: number }; + release(): void; +} + +/** Resolve a scheduled run context while holding the repo lock during preparation. */ +export async function resolveWorkspaceAutomationContext( + automation: AutomationConfig, + config: WorkspaceConfig, + workspaceDir: string, + repoManager: RepoManagerLike, + repoLock: (repoName: string) => RepoRunLockLike = (name) => createRepoRunLock(name, workspaceDir), +) { + // Keep occurrence task files in the workspace home (next to repos/, + // worktrees/, and the central DB) instead of inside a repo worktree. + const taskFileDir = join(workspaceDir, "automations"); + const repo = automation.repo + ? findRepo(config, automation.repo) + : config.repos.length === 1 + ? config.repos[0] + : undefined; + if (!repo) return { cwd: workspaceDir, env: { ...process.env }, taskFileDir, release() {} }; + + const lock = repoLock(repo.name); + if (!lock.acquire().success) return null; + try { + await repoManager.ensureBareClone(repo); + await repoManager.fetch(repo.name); + const cwd = await repoManager.ensureBaseWorktree(repo); + return { + cwd, + env: buildRepoEnv(repo, workspaceDir), + repo: repo.name, + taskFileDir, + release: () => lock.release(), + }; + } catch (error) { + lock.release(); + throw error; + } +} + /** Per-task CLI args: workspace defaults win, then the usual env/default. */ export function fleetTaskArgs(config: WorkspaceConfig): string[] { const raw = config.defaults.workerTaskArgs; @@ -258,7 +304,7 @@ export async function runWorkspaceWorker(options: RunWorkspaceWorkerOptions): Pr process.env.WEBHOOK_QUEUE_DB = workspaceDbPath(workspaceDir); const query = options.query ?? config.defaults.taskQuery; - if (!query) { + if (!query && config.automations.length === 0) { console.error( "❌ Workspace mode needs a task query: set [defaults].task_query in workspace.toml " + "or pass --query.", @@ -266,18 +312,6 @@ export async function runWorkspaceWorker(options: RunWorkspaceWorkerOptions): Pr process.exit(1); } - const { TaskTrackerManager } = await import("../task-tracker-manager"); - const { createChangeDetector } = await import("../change-detector"); - const tracker = new TaskTrackerManager().getClient(); - const detector = createChangeDetector(config.defaults.tracker, (q) => tracker.searchTasks(q)); - if (!detector) { - console.error( - `❌ Could not initialize the ${config.defaults.tracker} change detector. ` + - "Check the tracker's required variables in the workspace .env.", - ); - process.exit(1); - } - const state = openWorkspaceState(workspaceDir); const repoManager = new RepoManager(workspaceDir); @@ -291,34 +325,71 @@ export async function runWorkspaceWorker(options: RunWorkspaceWorkerOptions): Pr } } - const acquirers: import("../../worker").Acquirer[] = [ - createWorkspaceTaskAcquirer({ - config, - workspaceDir, - workerState: state.workerState, - queue: state.queue, - skips: state.skips, - repoManager, - detector, - searchTasks: (q) => tracker.searchTasks(q), - query, - intervalSeconds: options.intervalSeconds, - verbose: options.verbose, - }), - ]; + const acquirers: import("../../worker").Acquirer[] = []; - acquirers.push( - ...(await buildFleetEventAcquirers({ - config, - workspaceDir, - state, - repoManager, - searchTasks: (q) => tracker.searchTasks(q), - query, - intervalSeconds: options.intervalSeconds, - verbose: options.verbose, - })), - ); + if (config.automations.length > 0) { + const semanticErrors: string[] = []; + for (const automation of config.automations) { + if (!automation.repo && config.repos.length !== 1) { + semanticErrors.push( + `Automation "${automation.id}" must set repo when the workspace has multiple repositories.`, + ); + } + } + if (semanticErrors.length > 0) { + throw new Error(`Invalid ${configPath}:\n- ${semanticErrors.join("\n- ")}`); + } + + acquirers.push( + new AutomationAcquirer({ + automations: config.automations, + dbPath: state.dbPath, + resolveContext: (automation) => + resolveWorkspaceAutomationContext(automation, config, workspaceDir, repoManager), + }), + ); + } + + if (query) { + const { TaskTrackerManager } = await import("../task-tracker-manager"); + const { createChangeDetector } = await import("../change-detector"); + const tracker = new TaskTrackerManager().getClient(); + const detector = createChangeDetector(config.defaults.tracker, (q) => tracker.searchTasks(q)); + if (!detector) { + console.error( + `❌ Could not initialize the ${config.defaults.tracker} change detector. ` + + "Check the tracker's required variables in the workspace .env.", + ); + process.exit(1); + } + acquirers.push( + createWorkspaceTaskAcquirer({ + config, + workspaceDir, + workerState: state.workerState, + queue: state.queue, + skips: state.skips, + repoManager, + detector, + searchTasks: (q) => tracker.searchTasks(q), + query, + intervalSeconds: options.intervalSeconds, + verbose: options.verbose, + }), + ); + acquirers.push( + ...(await buildFleetEventAcquirers({ + config, + workspaceDir, + state, + repoManager, + searchTasks: (q) => tracker.searchTasks(q), + query, + intervalSeconds: options.intervalSeconds, + verbose: options.verbose, + })), + ); + } if (options.ui) { const { startDashboardServer } = await import("../../dashboard-server"); diff --git a/packages/code/tests/automation-acquirer.test.ts b/packages/code/tests/automation-acquirer.test.ts new file mode 100644 index 0000000..1bd3606 --- /dev/null +++ b/packages/code/tests/automation-acquirer.test.ts @@ -0,0 +1,465 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { existsSync, mkdirSync, readFileSync, rmSync } from "fs"; +import { join } from "path"; +import { tmpdir } from "os"; + +import { + AUTOMATION_ID_ENV, + AUTOMATION_ORIGIN_ENV, + AutomationAcquirer, + MAX_TIMER_DELAY_MS, + nextAutomationDue, + spawnAutomationProcess, + writeAutomationTaskFile, +} from "../src/lib/automation-acquirer"; +import { AutomationStateStore } from "../src/lib/automation-state"; +import type { AutomationConfig } from "../src/lib/automation-config"; + +describe("AutomationAcquirer", () => { + const dbPaths: string[] = []; + afterEach(() => { + for (const path of dbPaths.splice(0)) { + for (const suffix of ["", "-wal", "-shm"]) rmSync(`${path}${suffix}`, { force: true }); + } + }); + + test("calculates future cron and interval due times", () => { + const after = new Date("2026-08-21T12:34:00Z").getTime(); + const interval: AutomationConfig = { + id: "i", + enabled: true, + prompt: "p", + interval: "15m", + intervalMs: 900_000, + }; + const cron: AutomationConfig = { + id: "c", + enabled: true, + prompt: "p", + cron: "*/5 * * * *", + }; + expect(nextAutomationDue(interval, after)).toBe(after + 900_000); + expect(nextAutomationDue(cron, after)).toBeGreaterThan(after); + }); + + test.each([ + { + label: "interval", + now: 0, + schedule: { interval: "30d", intervalMs: 30 * 86_400_000 }, + }, + { + label: "cron", + now: new Date(2026, 0, 2).getTime(), + schedule: { cron: "0 0 1 1 *" }, + }, + ])("caps long $label timer delays to the runtime maximum", async ({ now, schedule }) => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + const delays: number[] = []; + const setTimer = (_callback: () => void, delay: number) => { + delays.push(delay); + return 1 as unknown as ReturnType; + }; + const acquirer = new AutomationAcquirer({ + automations: [ + { + id: `long-${schedule.cron ? "cron" : "interval"}`, + enabled: true, + prompt: "p", + ...schedule, + }, + ], + dbPath, + resolveContext: async () => null, + now: () => now, + setTimer, + clearTimer: () => {}, + }); + + await acquirer.start(); + expect(delays.at(-1)).toBe(MAX_TIMER_DELAY_MS); + await acquirer.stop(); + }); + + test("materializes the prompt as a markdown task file with attribution env", () => { + const baseDir = join(tmpdir(), `acquirer-task-${Date.now()}-${Math.random()}`); + const context = { cwd: baseDir, env: {}, repo: "api", release() {} }; + const filePath = writeAutomationTaskFile( + { + id: "dependency-health", + enabled: true, + prompt: "Inspect dependency health.\nApply one improvement.", + interval: "1d", + intervalMs: 86_400_000, + }, + context, + ); + + expect(filePath).toContain(join(".devintern-code", "automations", "dependency-health")); + expect(filePath.endsWith(".md")).toBe(true); + const content = readFileSync(filePath, "utf8"); + expect(content).toContain("# dependency-health"); + expect(content).toContain("Apply one improvement."); + rmSync(baseDir, { recursive: true, force: true }); + }); + + test("writes task files into the project config dir when launched from a subfolder", () => { + const project = join(tmpdir(), `acquirer-sub-${Date.now()}-${Math.random()}`); + const subfolder = join(project, "packages", "app"); + mkdirSync(join(project, ".devintern-code"), { recursive: true }); + mkdirSync(subfolder, { recursive: true }); + + const filePath = writeAutomationTaskFile( + { + id: "sub-folder", + enabled: true, + prompt: "work", + interval: "1d", + intervalMs: 86_400_000, + }, + { cwd: subfolder, env: {}, release() {} }, + ); + + expect(filePath.startsWith(join(project, ".devintern-code", "automations"))).toBe(true); + rmSync(project, { recursive: true, force: true }); + }); + + test("falls back to the run cwd's config dir when no parent config exists", () => { + const worktree = join(tmpdir(), `acquirer-fallback-${Date.now()}-${Math.random()}`); + mkdirSync(worktree, { recursive: true }); + + const filePath = writeAutomationTaskFile( + { + id: "fallback", + enabled: true, + prompt: "work", + interval: "1d", + intervalMs: 86_400_000, + }, + { cwd: worktree, env: {}, release() {} }, + ); + + expect(filePath.startsWith(join(worktree, ".devintern-code", "automations"))).toBe(true); + rmSync(worktree, { recursive: true, force: true }); + }); + + test("honors an explicit task file directory from the run context", () => { + const worktree = join(tmpdir(), `acquirer-explicit-${Date.now()}-${Math.random()}`); + const workspaceHome = join(tmpdir(), `acquirer-home-${Date.now()}-${Math.random()}`); + mkdirSync(worktree, { recursive: true }); + + const filePath = writeAutomationTaskFile( + { + id: "explicit", + enabled: true, + prompt: "work", + interval: "1d", + intervalMs: 86_400_000, + }, + { cwd: worktree, env: {}, taskFileDir: join(workspaceHome, "automations"), release() {} }, + ); + + expect(filePath.startsWith(join(workspaceHome, "automations", "explicit"))).toBe(true); + expect(existsSync(join(worktree, ".devintern-code"))).toBe(false); + rmSync(worktree, { recursive: true, force: true }); + rmSync(workspaceHome, { recursive: true, force: true }); + }); + + test("escalates to SIGKILL when an automation ignores SIGTERM", async () => { + const run = spawnAutomationProcess( + process.execPath, + ["-e", 'process.on("SIGTERM", () => {}); setInterval(() => {}, 1000)'], + { cwd: process.cwd(), env: process.env, terminationGraceMs: 50 }, + ); + await new Promise((resolve) => setTimeout(resolve, 100)); + const startedAt = Date.now(); + run.terminate(); + + expect(await run.completion).toBe(false); + expect(Date.now() - startedAt).toBeLessThan(1_000); + }); + + test("completes cleanly when the spawned pipeline exits zero", async () => { + const run = spawnAutomationProcess(process.execPath, ["-e", "process.exit(0)"], { + cwd: process.cwd(), + env: process.env, + }); + + expect(await run.completion).toBe(true); + }); + + test("disabled entries never receive schedule state", async () => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + const automation: AutomationConfig = { + id: "off", + enabled: false, + prompt: "p", + interval: "1h", + intervalMs: 3_600_000, + }; + const acquirer = new AutomationAcquirer({ + automations: [automation], + dbPath, + resolveContext: async () => null, + now: () => 100, + }); + await acquirer.start(); + await acquirer.stop(); + const store = new AutomationStateStore(dbPath); + expect(store.get("off")).toBeNull(); + store.close(); + }); + + test("coalesces a missed occurrence to one run through the task pipeline", async () => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + let now = 0; + let runs = 0; + let resolveRun!: (ok: boolean) => void; + const automation: AutomationConfig = { + id: "due", + enabled: true, + prompt: "p", + interval: "15m", + intervalMs: 900_000, + }; + const acquirer = new AutomationAcquirer({ + automations: [automation], + dbPath, + now: () => now, + resolveContext: async () => ({ cwd: "/tmp", env: {}, release() {} }), + spawnRun: () => { + runs += 1; + return { completion: new Promise((resolve) => (resolveRun = resolve)), terminate() {} }; + }, + }); + await acquirer.start(); + now = 10_000_000; + await acquirer.tick(); + expect(runs).toBe(1); + await acquirer.tick(); + expect(runs).toBe(1); + resolveRun(true); + await new Promise((resolve) => setTimeout(resolve, 0)); + await acquirer.stop(); + }); + + test("heartbeats a claim while context resolution exceeds the lease", async () => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + let preparationStarted!: () => void; + const started = new Promise((resolve) => (preparationStarted = resolve)); + let firstRuns = 0; + let secondContexts = 0; + const automation: AutomationConfig = { + id: "slow-context", + enabled: true, + prompt: "p", + interval: "10ms", + intervalMs: 10, + }; + const first = new AutomationAcquirer({ + automations: [automation], + dbPath, + leaseMs: 40, + heartbeatMs: 10, + resolveContext: async () => { + preparationStarted(); + await new Promise((resolve) => setTimeout(resolve, 120)); + return { cwd: "/tmp", env: {}, release() {} }; + }, + spawnRun: () => { + firstRuns += 1; + return { completion: Promise.resolve(true), terminate() {} }; + }, + }); + const second = new AutomationAcquirer({ + automations: [automation], + dbPath, + leaseMs: 40, + heartbeatMs: 10, + resolveContext: async () => { + secondContexts += 1; + return { cwd: "/tmp", env: {}, release() {} }; + }, + spawnRun: () => ({ completion: Promise.resolve(true), terminate() {} }), + }); + + await first.start(); + await started; + await new Promise((resolve) => setTimeout(resolve, 70)); + await second.start(); + expect(secondContexts).toBe(0); + await second.stop(); + + await new Promise((resolve) => setTimeout(resolve, 70)); + expect(firstRuns).toBe(1); + await first.stop(); + }); + + test("uses the actual claim time for later automations after slow preparation", async () => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + let now = 0; + const finishRuns: Array<(ok: boolean) => void> = []; + const automations: AutomationConfig[] = ["first", "second"].map((id) => ({ + id, + enabled: true, + prompt: "p", + interval: "10ms", + intervalMs: 10, + })); + const acquirer = new AutomationAcquirer({ + automations, + dbPath, + now: () => now, + leaseMs: 40, + setTimer: () => 1 as unknown as ReturnType, + clearTimer: () => {}, + resolveContext: async (automation) => { + if (automation.id === "first") now = 100; + return { cwd: "/tmp", env: {}, release() {} }; + }, + spawnRun: () => { + let finish!: (ok: boolean) => void; + const completion = new Promise((resolve) => (finish = resolve)); + finishRuns.push(finish); + return { completion, terminate: () => finish(false) }; + }, + }); + + await acquirer.start(); + now = 10; + await acquirer.tick(); + + const store = new AutomationStateStore(dbPath); + expect(store.get("second")?.heartbeatAt).toBe(100); + expect(store.get("second")?.leaseExpiresAt).toBe(140); + store.close(); + for (const finish of finishRuns) finish(true); + await acquirer.stop(); + }); + + test("terminates and cleans up an active run after lease ownership changes", async () => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + let now = 0; + let finishFirst!: (ok: boolean) => void; + let finishSecond!: (ok: boolean) => void; + let terminated = false; + let released = false; + const automation: AutomationConfig = { + id: "lost-lease", + enabled: true, + prompt: "p", + interval: "10ms", + intervalMs: 10, + }; + const timerOptions = { + setTimer: () => 1 as unknown as ReturnType, + clearTimer: () => {}, + }; + const first = new AutomationAcquirer({ + automations: [automation], + dbPath, + now: () => now, + leaseMs: 40, + ...timerOptions, + resolveContext: async () => ({ + cwd: "/tmp", + env: {}, + release: () => { + released = true; + }, + }), + spawnRun: () => ({ + completion: new Promise((resolve) => (finishFirst = resolve)), + terminate: () => { + terminated = true; + finishFirst(false); + }, + }), + }); + const second = new AutomationAcquirer({ + automations: [automation], + dbPath, + now: () => now, + leaseMs: 40, + ...timerOptions, + resolveContext: async () => ({ cwd: "/tmp", env: {}, release() {} }), + spawnRun: () => ({ + completion: new Promise((resolve) => (finishSecond = resolve)), + terminate: () => finishSecond(false), + }), + }); + + await first.start(); + now = 10; + await first.tick(); + now = 100; + await second.start(); + await first.tick(); + + expect(terminated).toBe(true); + expect(released).toBe(true); + await first.stop(); + await second.stop(); + }); + + test("stop waits for active run cleanup before closing state", async () => { + const dbPath = join(tmpdir(), `acquirer-${Date.now()}-${Math.random()}.db`); + dbPaths.push(dbPath); + let now = 0; + let resolveRun!: (ok: boolean) => void; + let finishRelease!: () => void; + let releaseStarted = false; + let stopFinished = false; + const acquirer = new AutomationAcquirer({ + automations: [ + { + id: "shutdown", + enabled: true, + prompt: "p", + interval: "1m", + intervalMs: 60_000, + }, + ], + dbPath, + now: () => now, + setTimer: () => 1 as unknown as ReturnType, + clearTimer: () => {}, + resolveContext: async () => ({ + cwd: "/tmp", + env: {}, + release: async () => { + releaseStarted = true; + await new Promise((resolve) => (finishRelease = resolve)); + }, + }), + spawnRun: () => ({ + completion: new Promise((resolve) => (resolveRun = resolve)), + terminate: () => resolveRun(false), + }), + }); + await acquirer.start(); + now = 60_000; + await acquirer.tick(); + + const stopping = acquirer.stop().then(() => { + stopFinished = true; + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(releaseStarted).toBe(true); + expect(stopFinished).toBe(false); + finishRelease(); + await stopping; + expect(stopFinished).toBe(true); + }); + + test("exports the scheduled-run attribution env marker names", () => { + expect(AUTOMATION_ORIGIN_ENV).toBe("DEVINTERN_RUN_ORIGIN"); + expect(AUTOMATION_ID_ENV).toBe("DEVINTERN_AUTOMATION_ID"); + }); +}); diff --git a/packages/code/tests/automation-config.test.ts b/packages/code/tests/automation-config.test.ts new file mode 100644 index 0000000..dc3de79 --- /dev/null +++ b/packages/code/tests/automation-config.test.ts @@ -0,0 +1,88 @@ +import { describe, expect, test } from "bun:test"; + +import { parseAutomationConfig, parseAutomationInterval } from "../src/lib/automation-config"; +import { parseWorkspaceConfig } from "../src/lib/workspace/config"; + +describe("automation configuration", () => { + test("parses cron and interval entries including multiline prompts", () => { + const entries = parseAutomationConfig(` +[[automations]] +id = "daily-review" +enabled = true +cron = "0 9 * * 1-5" +prompt = """Review the repository. +Fix one issue.""" + +[[automations]] +id = "maintenance" +enabled = false +interval = "6h" +repo = "api" +prompt = "Apply one safe maintenance improvement" +`); + expect(entries).toHaveLength(2); + expect(entries[0]?.prompt).toContain("\n"); + expect(entries[1]?.intervalMs).toBe(6 * 60 * 60 * 1000); + expect(entries[1]?.enabled).toBe(false); + expect(entries[1]?.repo).toBe("api"); + }); + + test("accepts only documented positive duration units", () => { + expect(parseAutomationInterval("15m")).toBe(900_000); + expect(parseAutomationInterval("1d")).toBe(86_400_000); + expect(parseAutomationInterval("30s")).toBeNull(); + expect(parseAutomationInterval("0h")).toBeNull(); + }); + + test("rejects unsafe intervals and dates outside the runtime range", () => { + expect(parseAutomationInterval("9007199254740991d", 0)).toBeNull(); + expect(parseAutomationInterval("100000000d", 0)).toBe(8_640_000_000_000_000); + expect(parseAutomationInterval("100000001d", 0)).toBeNull(); + expect(parseAutomationInterval("1d", 8_640_000_000_000_000)).toBeNull(); + }); + + test("collects duplicate, prompt, and schedule errors", () => { + let message = ""; + try { + parseAutomationConfig(` +[[automations]] +id = "same" +enabled = "yes" +cron = "bad" +interval = "2w" +prompt = "" + +[[automations]] +id = "same" +enabled = true +prompt = "ok" +`); + } catch (error) { + message = (error as Error).message; + } + expect(message).toContain('Duplicate automation id "same"'); + expect(message).toContain("enabled must be a boolean"); + expect(message).toContain("prompt is required"); + expect(message).toContain("exactly one of cron or interval"); + }); + + test("workspace validation rejects unknown repositories", () => { + expect(() => + parseWorkspaceConfig(` +[defaults] +tracker = "jira" + +[[repos]] +name = "api" +remote = "https://github.com/acme/api.git" + +[[automations]] +id = "other" +enabled = true +interval = "1h" +prompt = "work" +repo = "web" +`), + ).toThrow(/does not match any \[\[repos\]\] name/); + }); +}); diff --git a/packages/code/tests/automation-state.test.ts b/packages/code/tests/automation-state.test.ts new file mode 100644 index 0000000..5ab8ed3 --- /dev/null +++ b/packages/code/tests/automation-state.test.ts @@ -0,0 +1,53 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { rmSync } from "fs"; +import { join } from "path"; +import { tmpdir } from "os"; + +import { AutomationStateStore } from "../src/lib/automation-state"; +import type { AutomationConfig } from "../src/lib/automation-config"; + +const AUTOMATION: AutomationConfig = { + id: "cleanup", + enabled: true, + prompt: "clean up", + interval: "15m", + intervalMs: 900_000, +}; + +describe("AutomationStateStore", () => { + const paths: string[] = []; + afterEach(() => { + for (const path of paths.splice(0)) { + for (const suffix of ["", "-wal", "-shm"]) rmSync(`${path}${suffix}`, { force: true }); + } + }); + + test("persists interval anchors and claims across restart", () => { + const path = join(tmpdir(), `automation-${Date.now()}-${Math.random()}.db`); + paths.push(path); + const first = new AutomationStateStore(path); + first.register(AUTOMATION, 1_000); + expect(first.claim(AUTOMATION.id, "worker-a", 1_000, 901_000, 60_000)).toBe(true); + first.close(); + + const reopened = new AutomationStateStore(path); + const state = reopened.get(AUTOMATION.id); + expect(state?.lastScheduledAt).toBe(1_000); + expect(state?.nextDueAt).toBe(901_000); + expect(state?.leaseOwner).toBe("worker-a"); + reopened.close(); + }); + + test("prevents overlap, advances skipped occurrence, and recovers stale leases", () => { + const path = join(tmpdir(), `automation-${Date.now()}-${Math.random()}.db`); + paths.push(path); + const store = new AutomationStateStore(path); + store.register(AUTOMATION, 100); + expect(store.claim(AUTOMATION.id, "worker-a", 100, 200, 500)).toBe(true); + expect(store.claim(AUTOMATION.id, "worker-b", 200, 300, 50)).toBe(false); + expect(store.skipOverlap(AUTOMATION.id, 200, 300)).toBe(true); + expect(store.claim(AUTOMATION.id, "worker-b", 600, 700, 50)).toBe(true); + expect(store.get(AUTOMATION.id)?.leaseOwner).toBe("worker-b"); + store.close(); + }); +}); diff --git a/packages/code/tests/dashboard-api.test.ts b/packages/code/tests/dashboard-api.test.ts index 2044326..a17825d 100644 --- a/packages/code/tests/dashboard-api.test.ts +++ b/packages/code/tests/dashboard-api.test.ts @@ -40,7 +40,7 @@ describe("dashboard API", () => { store: RunStore, options: { taskKey?: string; - origin?: "task" | "pr_mention"; + origin?: "task" | "pr_mention" | "scheduled"; harness?: string; status?: "succeeded" | "failed" | "escalated" | "abandoned" | "deferred"; prUrl?: string; @@ -66,15 +66,16 @@ describe("dashboard API", () => { seedRun(store, { taskKey: "PROJ-1", status: "succeeded" }); seedRun(store, { taskKey: "PROJ-2", status: "failed" }); seedRun(store, { origin: "pr_mention", status: "succeeded" }); + seedRun(store, { origin: "scheduled", status: "succeeded" }); store.close(); const all = handleRuns(data, new URLSearchParams()); expect(all.status).toBe(200); - expect((all.body as { total: number }).total).toBe(3); + expect((all.body as { total: number }).total).toBe(4); const paged = handleRuns(data, new URLSearchParams("limit=2&offset=2")); - expect((paged.body as { runs: unknown[] }).runs.length).toBe(1); - expect((paged.body as { total: number }).total).toBe(3); + expect((paged.body as { runs: unknown[] }).runs.length).toBe(2); + expect((paged.body as { total: number }).total).toBe(4); const failed = handleRuns(data, new URLSearchParams("status=failed")); const failedBody = failed.body as { runs: { taskKey?: string }[]; total: number }; @@ -84,6 +85,9 @@ describe("dashboard API", () => { const mentions = handleRuns(data, new URLSearchParams("origin=pr_mention")); expect((mentions.body as { total: number }).total).toBe(1); + const scheduled = handleRuns(data, new URLSearchParams("origin=scheduled")); + expect((scheduled.body as { total: number }).total).toBe(1); + const byKey = handleRuns(data, new URLSearchParams("taskKey=PROJ-1")); expect((byKey.body as { total: number }).total).toBe(1); }); diff --git a/packages/code/tests/run-recorder.test.ts b/packages/code/tests/run-recorder.test.ts index 7b91038..86f6159 100644 --- a/packages/code/tests/run-recorder.test.ts +++ b/packages/code/tests/run-recorder.test.ts @@ -47,6 +47,22 @@ describe("RunStore", () => { expect(run?.prNumber).toBe(42); }); + test("scheduled runs persist automation metadata", () => { + const id = store.createRun({ + origin: "scheduled", + automationId: "weekly-plan", + repo: "backend", + harness: "codex", + }); + store.finishRun(id, "succeeded"); + + expect(store.getRun(id)).toMatchObject({ + origin: "scheduled", + automationId: "weekly-plan", + }); + expect(store.getStats(null).byOrigin.scheduled).toBe(1); + }); + test("stages accumulate in order with structured detail", () => { const id = store.createRun({ origin: "task", taskKey: "PROJ-2" }); store.addStage(id, "feasibility", "succeeded", "clear enough", '{"clarityScore":8}'); @@ -159,6 +175,8 @@ describe("RunStore", () => { const migrated = new RunStore(legacyPath); const id = migrated.createRun({ origin: "task", taskKey: "OLD-1" }); expect(migrated.getRun(id)?.attempt).toBe(2); + const scheduled = migrated.createRun({ origin: "scheduled", automationId: "new-schedule" }); + expect(migrated.getRun(scheduled)?.automationId).toBe("new-schedule"); migrated.close(); for (const suffix of ["", "-wal", "-shm"]) { rmSync(`${legacyPath}${suffix}`, { force: true }); diff --git a/packages/code/tests/workspace-worker.test.ts b/packages/code/tests/workspace-worker.test.ts index d5da803..e2e1843 100644 --- a/packages/code/tests/workspace-worker.test.ts +++ b/packages/code/tests/workspace-worker.test.ts @@ -5,7 +5,11 @@ import { tmpdir } from "os"; import { parseWorkspaceConfig } from "../src/lib/workspace/config"; import type { RepoConfig } from "../src/lib/workspace/config"; -import { createWorkspaceTaskAcquirer, fleetTaskArgs } from "../src/lib/workspace/workspace-worker"; +import { + createWorkspaceTaskAcquirer, + fleetTaskArgs, + resolveWorkspaceAutomationContext, +} from "../src/lib/workspace/workspace-worker"; import type { FleetTask, RepoManagerLike } from "../src/lib/workspace/workspace-worker"; import { createRepoRunLock, openWorkspaceState } from "../src/lib/workspace/state"; import type { WorkspaceState } from "../src/lib/workspace/state"; @@ -224,3 +228,59 @@ describe("fleetTaskArgs", () => { expect(fleetTaskArgs(CONFIG)).toEqual(["--create-pr", "--auto-review"]); }); }); + +describe("resolveWorkspaceAutomationContext", () => { + test("does not prepare the repository when its run lock is unavailable", async () => { + const workspaceDir = join( + tmpdir(), + `ws-automation-${Date.now()}-${Math.random().toString(36).slice(2)}`, + ); + const repoManager = new FakeRepoManager(workspaceDir); + const context = await resolveWorkspaceAutomationContext( + { + id: "scheduled", + enabled: true, + prompt: "work", + interval: "1h", + intervalMs: 3_600_000, + repo: "backend", + }, + CONFIG, + workspaceDir, + repoManager, + () => ({ + acquire: () => ({ success: false, message: "busy" }), + release() {}, + }), + ); + + expect(context).toBeNull(); + expect(repoManager.calls).toEqual([]); + rmSync(workspaceDir, { recursive: true, force: true }); + }); + + test("pins occurrence task files to the workspace home", async () => { + const workspaceDir = join( + tmpdir(), + `ws-automation-dir-${Date.now()}-${Math.random().toString(36).slice(2)}`, + ); + const repoManager = new FakeRepoManager(workspaceDir); + const context = await resolveWorkspaceAutomationContext( + { + id: "scheduled", + enabled: true, + prompt: "work", + interval: "1h", + intervalMs: 3_600_000, + repo: "backend", + }, + CONFIG, + workspaceDir, + repoManager, + ); + + expect(context?.taskFileDir).toBe(join(workspaceDir, "automations")); + expect(context?.cwd).toContain(join("worktrees", "backend", "base")); + rmSync(workspaceDir, { recursive: true, force: true }); + }); +}); diff --git a/packages/dashboard-ui/src/components/RunResult.test.tsx b/packages/dashboard-ui/src/components/RunResult.test.tsx new file mode 100644 index 0000000..5f6d9af --- /dev/null +++ b/packages/dashboard-ui/src/components/RunResult.test.tsx @@ -0,0 +1,23 @@ +import { expect, test } from "bun:test"; +import { createElement } from "react"; +import { renderToStaticMarkup } from "react-dom/server"; + +import { RunResult } from "@/components/RunResult"; + +test("renders a ticket key as text when its tracker provides no URL", () => { + const html = renderToStaticMarkup(createElement(RunResult, { run: { ticketKey: "ENG-42" } })); + + expect(html).toContain("ENG-42"); + expect(html).not.toContain(" { + const html = renderToStaticMarkup( + createElement(RunResult, { + run: { ticketKey: "ENG-42", ticketUrl: "https://tracker.test/ENG-42" }, + }), + ); + + expect(html).toContain("ENG-42"); + expect(html).toContain('href="https://tracker.test/ENG-42"'); +}); diff --git a/packages/dashboard-ui/src/components/RunResult.tsx b/packages/dashboard-ui/src/components/RunResult.tsx new file mode 100644 index 0000000..eafcf4d --- /dev/null +++ b/packages/dashboard-ui/src/components/RunResult.tsx @@ -0,0 +1,27 @@ +import { ExternalLink } from "lucide-react"; + +import type { RunRecord } from "@/lib/api"; + +/** Link to a run result when possible, retaining ticket keys without tracker URLs. */ +export function RunResult({ + run, +}: { + run: Pick; +}) { + const url = run.prUrl ?? run.ticketUrl; + const label = run.ticketKey ?? (run.prNumber ? `#${run.prNumber}` : url ? "#PR" : undefined); + if (!label) return ; + if (!url) return {label}; + return ( + event.stopPropagation()} + className="inline-flex items-center gap-1 text-primary hover:underline" + > + {label} + + + ); +} diff --git a/packages/dashboard-ui/src/lib/api.ts b/packages/dashboard-ui/src/lib/api.ts index 2fe276c..26341ac 100644 --- a/packages/dashboard-ui/src/lib/api.ts +++ b/packages/dashboard-ui/src/lib/api.ts @@ -5,7 +5,7 @@ import { useCallback, useEffect, useRef, useState } from "react"; -export type RunOrigin = "task" | "pr_mention"; +export type RunOrigin = "task" | "pr_mention" | "scheduled"; export type RunStatus = | "in_progress" @@ -18,6 +18,9 @@ export type RunStatus = export interface RunRecord { id: number; origin: RunOrigin; + automationId?: string; + ticketKey?: string; + ticketUrl?: string; taskKey?: string; tracker?: string; harness?: string; diff --git a/packages/dashboard-ui/src/views/RunDetailView.tsx b/packages/dashboard-ui/src/views/RunDetailView.tsx index 5c7e3ad..0b587f7 100644 --- a/packages/dashboard-ui/src/views/RunDetailView.tsx +++ b/packages/dashboard-ui/src/views/RunDetailView.tsx @@ -1,6 +1,7 @@ import { useState } from "react"; -import { ArrowLeft, ChevronDown, ChevronRight, ExternalLink } from "lucide-react"; +import { ArrowLeft, ChevronDown, ChevronRight } from "lucide-react"; +import { RunResult } from "@/components/RunResult"; import { EmptyState, StageBadge, StatusBadge } from "@/components/shared"; import { StageDetailFields } from "@/components/StageDetailFields"; import { Button } from "@/components/ui/button"; @@ -118,6 +119,7 @@ export function RunDetailView({ runId, onBack }: { runId: number; onBack: () =>

{data.run.taskKey ?? + data.run.automationId ?? (data.run.prNumber ? `PR #${data.run.prNumber}` : `Run ${data.run.id}`)}

@@ -129,7 +131,13 @@ export function RunDetailView({ runId, onBack }: { runId: number; onBack: () => @@ -143,24 +151,7 @@ export function RunDetailView({ runId, onBack }: { runId: number; onBack: () => ) } /> - - #{data.run.prNumber ?? "PR"} - - - ) : ( - "–" - ) - } - /> + } /> void }) { @@ -74,7 +75,7 @@ export function RunsView({ onOpenRun }: { onOpenRun: (id: number) => void }) { {data && data.runs.length === 0 ? ( ) : null} @@ -84,10 +85,10 @@ export function RunsView({ onOpenRun }: { onOpenRun: (id: number) => void }) { Status - Task + Work Origin Harness - PR + Result Duration Started @@ -99,29 +100,22 @@ export function RunsView({ onOpenRun }: { onOpenRun: (id: number) => void }) { - {run.taskKey ?? (run.prNumber ? `PR #${run.prNumber}` : `run ${run.id}`)} + {run.taskKey ?? + run.automationId ?? + (run.prNumber ? `PR #${run.prNumber}` : `run ${run.id}`)} - {run.origin === "task" ? "task" : "PR mention"} + {run.origin === "task" + ? "task" + : run.origin === "scheduled" + ? "scheduled" + : "PR mention"} {run.harness ?? "–"} - {run.prUrl ? ( - event.stopPropagation()} - className="inline-flex items-center gap-1 text-primary hover:underline" - > - #{run.prNumber ?? "PR"} - - - ) : ( - - )} + {run.finishedAt ? formatDuration(run.finishedAt - run.startedAt) : "…"} diff --git a/packages/dashboard-ui/src/views/StatsView.tsx b/packages/dashboard-ui/src/views/StatsView.tsx index 5beca04..0606f98 100644 --- a/packages/dashboard-ui/src/views/StatsView.tsx +++ b/packages/dashboard-ui/src/views/StatsView.tsx @@ -81,7 +81,7 @@ export function StatsView() {