mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
feat(scheduler): add cron-style scheduled jobs (backend)
Adds a saved/named scheduling layer on top of Codeman's existing session primitives. Distinct from the legacy run-now ScheduledRun concept. - types/scheduler.ts: ScheduledJob + ScheduledJobRun - state-store: persist scheduledJobs/scheduledJobRuns in ~/.codeman/state.json - scheduler/scheduler-time.ts: pure once/interval/daily/weekly next-run math - scheduler/scheduler-service.ts: CRUD, Run Now, due-checker tick, run history; reuses SessionPort (create -> start -> writeViaMux) for launches - web/routes/scheduler-routes.ts: /api/scheduler/jobs CRUD + run + history - web/schemas.ts: zod validation with schedule-type-aware refinements - web/sse-events.ts: scheduler:* events - server.ts: wire service into route context + 30s background tick loop - test/scheduler-time.test.ts: 14 unit tests for next-run calculations Phase 1 discovery recorded in SCHEDULER_DISCOVERY.md. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Rp7JhmQXcYJhmxFMdZuuah
This commit is contained in:
@@ -0,0 +1,138 @@
|
||||
# SCHEDULER_DISCOVERY.md
|
||||
|
||||
Phase 1 deliverable for the "Add Scheduling to Codeman" build brief.
|
||||
This documents the existing Codeman architecture and the smallest integration
|
||||
points for a cron-style scheduler. **No session/tmux logic will be rebuilt** —
|
||||
the new code is purely a trigger + persistence + history layer on top of the
|
||||
existing primitives.
|
||||
|
||||
Stack: `aicodeman` v1.2.1 — Fastify 5 backend, `node-pty` + tmux sessions,
|
||||
vanilla-JS SPA frontend served as static assets, JSON file state store, zod
|
||||
validation, ports-based dependency injection.
|
||||
|
||||
---
|
||||
|
||||
## 0. Critical finding: an existing `ScheduledRun` is NOT a scheduler
|
||||
|
||||
Codeman already has a `ScheduledRun` concept (`/api/scheduled`,
|
||||
`src/web/ports/infra-port.ts:14-26`, `src/web/server.ts:1480-1605`). It is a
|
||||
**run-now, duration-bounded autonomous loop**: given `{prompt, workingDir,
|
||||
durationMinutes}` it immediately spawns/kills throwaway sessions in a loop until
|
||||
the duration elapses. It has **no** time-based triggering, recurrence
|
||||
(once/interval/daily/weekly), enable/disable, next-run calculation, run history,
|
||||
or persistence across restarts.
|
||||
|
||||
Therefore the brief's core (the calendar/cron trigger layer) does **not** exist
|
||||
and must be built. The execution primitives it sits on top of **do** exist and
|
||||
will be reused. To honor brief §16 ("do not rename existing core concepts"), the
|
||||
new feature is named **`ScheduledJob`** (with **`ScheduledJobRun`** history
|
||||
records), kept distinct from the existing `ScheduledRun`.
|
||||
|
||||
---
|
||||
|
||||
## 1. Where session creation happens
|
||||
|
||||
- Canonical create flow: `POST /api/sessions`,
|
||||
`src/web/routes/session-routes.ts:262-438`.
|
||||
- `new Session({ workingDir, mode, ... })` (`src/session.ts:421-570`)
|
||||
- `ctx.addSession(session)` → `ctx.setupSessionListeners(session)` →
|
||||
`ctx.persistSessionState(session)` (all via `SessionPort`).
|
||||
- `SessionPort` interface: `src/web/ports/session-port.ts:8-16`.
|
||||
- **Integration point:** the scheduler service will mirror this exact sequence
|
||||
(create → addSession → setupSessionListeners → start) via `SessionPort`,
|
||||
not reimplement it.
|
||||
|
||||
## 2. Where agent/session types are defined
|
||||
|
||||
- `type SessionMode = 'claude' | 'shell' | 'opencode' | 'codex' | 'gemini'`
|
||||
(`src/types/session.ts:43-44`). `shell` covers the brief's "Terminal/custom".
|
||||
- CLI availability resolvers in `src/utils/{claude,codex,gemini,opencode}-cli-resolver.ts`.
|
||||
- **Integration point:** the job's `agentType` reuses `SessionMode` verbatim.
|
||||
|
||||
## 3. Where input is sent into a session
|
||||
|
||||
- Raw / paste: `session.write(data)` (`src/session.ts:2243-2247`) — direct PTY write.
|
||||
- Typed (recommended): `session.writeViaMux(data)` (`src/session.ts:2301-2311`)
|
||||
— tmux `send-keys`, falls back to PTY. Submit requires trailing `\r`.
|
||||
- **Integration point:** prompt delivery uses `writeViaMux` (typed) by default,
|
||||
`write` (paste) as the alternate `input_mode`.
|
||||
|
||||
## 4. Where active sessions are listed
|
||||
|
||||
- `ctx.sessions: ReadonlyMap<string, Session>` (`SessionPort`).
|
||||
- Filters: `Array.from(ctx.sessions.values()).filter(s => s.mode === X)` and
|
||||
`.isBusy()` / `.isIdle()` (`src/session-manager.ts:220-247`).
|
||||
- **Integration point:** the §8 multi-session warning queries this map.
|
||||
|
||||
## 5. Where session kill/delete is handled
|
||||
|
||||
- `ctx.cleanupSession(sessionId, killMux?, reason?)`
|
||||
(`SessionPort`; impl `src/web/server.ts:997-1152`). Underlying
|
||||
`session.stop(killMux)` at `src/session.ts:2498-2585`.
|
||||
- The scheduler does **not** kill sessions it launches (the brief wants them
|
||||
visible in the normal session UI); cleanup stays user-driven.
|
||||
|
||||
## 6. How session state is stored / 7. Existing persistence
|
||||
|
||||
- JSON file store: `~/.codeman/state.json` (+ `state-inner.json` for Ralph).
|
||||
`StateStore` class `src/state-store.ts:71`; `AppState` interface
|
||||
`src/types/app-state.ts:99-114`.
|
||||
- Pattern: declare a field on `AppState`, add typed get/set methods on
|
||||
`StateStore` that mutate in-memory state and call the debounced `save()`
|
||||
(500ms debounce, atomic temp-file+rename, `.bak` backup, circuit breaker).
|
||||
- **Integration point:** add `scheduledJobs?: Record<string, ScheduledJob>` and
|
||||
`scheduledJobRuns?: Record<string, ScheduledJobRun>` to `AppState`, with
|
||||
matching `StateStore` accessors. No new DB (brief §6 forbids Postgres/Redis).
|
||||
|
||||
## 8. Where backend routes live
|
||||
|
||||
- Route modules: `src/web/routes/*.ts`; barrel `src/web/routes/index.ts`;
|
||||
registered in `WebServer.setupRoutes()` `src/web/server.ts:858-876` with a
|
||||
single `ctx` object from `createRouteContext()` (`src/web/server.ts:553-613`)
|
||||
that satisfies all port interfaces.
|
||||
- Validation: zod schemas in `src/web/schemas.ts`, applied via
|
||||
`parseBody(Schema, req.body)` (`src/web/route-helpers.ts:101-111`).
|
||||
- Errors: `createErrorResponse(ApiErrorCode.X, msg)` / `ApiResponse`
|
||||
(`src/types/api.ts`), auto-mapped to HTTP status by a `preSerialization` hook
|
||||
(`src/web/server.ts:644-659`).
|
||||
- SSE: `ctx.broadcast(SseEvent.X, data)` (`EventPort`,
|
||||
`src/web/sse-events.ts`); frontend mirror in `src/web/public/constants.js`.
|
||||
- **Integration point:** new `scheduler-routes.ts` registered alongside the
|
||||
others; new zod schema; new `SseEvent` constants for job list/run changes.
|
||||
|
||||
## 9. Where frontend pages/components live
|
||||
|
||||
- Vanilla-JS SPA: single `src/web/public/index.html` + feature mixin files
|
||||
(`Object.assign(CodemanApp.prototype, {...})`). API via `api-client.js`
|
||||
(`_apiJson/_apiPost/_apiDelete`). Build = esbuild minify + content-hash, no
|
||||
bundler (`scripts/build.mjs`).
|
||||
- UI is panels/modals toggled by JS classes; forms use `.form-row` / `.modal`
|
||||
conventions (`styles.css`). SSE handler map in `app.js`.
|
||||
- **Integration point:** add a new `scheduler-ui.js` mixin + a panel/modal in
|
||||
`index.html` + nav entry, following the orchestrator/respawn panel pattern.
|
||||
|
||||
## 10. Background-loop pattern (for the due-checker)
|
||||
|
||||
- Established pattern: `this.cleanup.setInterval(fn, intervalMs, {description})`
|
||||
in `WebServer.start()` (`src/web/server.ts:~1942-1966`), auto-disposed in
|
||||
`WebServer.stop()` via `this.cleanup.dispose()` (`src/web/server.ts:2336`).
|
||||
RalphLoop (`src/ralph-loop.ts:268-286`) shows the self-rescheduling guard idiom.
|
||||
- **Integration point:** register a 30s scheduler tick via `cleanup.setInterval`;
|
||||
no manual shutdown wiring needed.
|
||||
|
||||
---
|
||||
|
||||
## Smallest integration points (summary)
|
||||
|
||||
| New piece | Reuses | Location |
|
||||
| --- | --- | --- |
|
||||
| `ScheduledJob` / `ScheduledJobRun` types | — (new) | `src/types/scheduler.ts` |
|
||||
| Persistence | `StateStore` / `AppState` | `src/types/app-state.ts`, `src/state-store.ts` |
|
||||
| Next-run time math | — (new, pure, unit-tested) | `src/scheduler/scheduler-time.ts` |
|
||||
| Launch + send prompt | `SessionPort` (`addSession`/listeners/`writeViaMux`) | `src/scheduler/scheduler-service.ts` |
|
||||
| Background due loop | `cleanup.setInterval` pattern | `src/scheduler/scheduler-loop.ts` |
|
||||
| Routes + schema | route/ports/zod/SSE patterns | `src/web/routes/scheduler-routes.ts`, `src/web/schemas.ts`, `src/web/sse-events.ts` |
|
||||
| UI | panel/modal/mixin conventions | `src/web/public/scheduler-ui.js`, `index.html` |
|
||||
|
||||
Nothing in the session, tmux, persistence, routing, or SSE subsystems is
|
||||
rewritten — the scheduler is additive and calls existing services.
|
||||
@@ -51,6 +51,19 @@ export const SCHEDULED_CLEANUP_INTERVAL = 5 * 60 * 1000;
|
||||
/** Completed scheduled run max age before cleanup (ms) */
|
||||
export const SCHEDULED_RUN_MAX_AGE = 60 * 60 * 1000;
|
||||
|
||||
// ============================================================================
|
||||
// Scheduled Jobs (cron-style scheduler)
|
||||
// ============================================================================
|
||||
|
||||
/** How often the scheduler loop wakes to check for due jobs (ms). */
|
||||
export const SCHEDULER_TICK_INTERVAL = 30 * 1000;
|
||||
|
||||
/** Max attempts (× 500ms) to poll a launched session for CLI readiness before sending the prompt. */
|
||||
export const SCHEDULER_READY_MAX_ATTEMPTS = 60;
|
||||
|
||||
/** Extra settle delay after CLI readiness is detected, before sending the prompt (ms). */
|
||||
export const SCHEDULER_READY_SETTLE_MS = 2000;
|
||||
|
||||
/** Session limit retry wait before retrying (ms) */
|
||||
export const SESSION_LIMIT_WAIT_MS = 5000;
|
||||
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
/**
|
||||
* @fileoverview Input shape for creating/updating a scheduled job. This is the
|
||||
* user-settable subset of `ScheduledJob` (server-maintained bookkeeping fields
|
||||
* such as nextRunAt / lastStatus are excluded). Produced by the zod schema.
|
||||
*/
|
||||
|
||||
import type { ConcurrencyPolicy, InputMode, PromptMode, ScheduleType } from '../types/scheduler.js';
|
||||
import type { SessionMode } from '../types/session.js';
|
||||
|
||||
export type { ScheduledJob, ScheduledJobRun, ScheduledJobRunStatus, TriggerType } from '../types/scheduler.js';
|
||||
|
||||
export interface ScheduledJobInput {
|
||||
name: string;
|
||||
agentType: SessionMode;
|
||||
workingDir: string;
|
||||
launchCommand?: string;
|
||||
promptMode: PromptMode;
|
||||
promptText?: string;
|
||||
promptFilePath?: string;
|
||||
inputMode: InputMode;
|
||||
scheduleType: ScheduleType;
|
||||
runAt?: number;
|
||||
intervalMinutes?: number;
|
||||
dailyTime?: string;
|
||||
weeklyDays?: number[];
|
||||
weeklyTime?: string;
|
||||
enabled: boolean;
|
||||
notes?: string;
|
||||
concurrencyPolicy: ConcurrencyPolicy;
|
||||
}
|
||||
@@ -0,0 +1,357 @@
|
||||
/**
|
||||
* @fileoverview Scheduler service: CRUD for scheduled jobs, manual Run Now,
|
||||
* the background due-job tick, and run-history recording.
|
||||
*
|
||||
* It does NOT own session/tmux logic — it reuses Codeman's existing session
|
||||
* layer (create → addSession → setupSessionListeners → startInteractive/Shell →
|
||||
* send prompt via writeViaMux/write), mirroring the "quick start" route flow.
|
||||
*/
|
||||
|
||||
import { v4 as uuidv4 } from 'uuid';
|
||||
import { readFile } from 'node:fs/promises';
|
||||
import { statSync } from 'node:fs';
|
||||
import { Session } from '../session.js';
|
||||
import { SseEvent } from '../web/sse-events.js';
|
||||
import { getErrorMessage } from '../types/api.js';
|
||||
import { MAX_CONCURRENT_SESSIONS } from '../config/map-limits.js';
|
||||
import { SCHEDULER_READY_MAX_ATTEMPTS, SCHEDULER_READY_SETTLE_MS } from '../config/server-timing.js';
|
||||
import { computeNextRunAt, dueKeyFor } from './scheduler-time.js';
|
||||
import type { SessionPort, EventPort, ConfigPort, InfraPort } from '../web/ports/index.js';
|
||||
import type { ScheduledJob, ScheduledJobRun, ScheduledJobRunStatus, TriggerType } from '../types/scheduler.js';
|
||||
import type { ScheduledJobInput } from './scheduler-input.js';
|
||||
|
||||
/** The subset of the route context the scheduler depends on. */
|
||||
export type SchedulerDeps = SessionPort & EventPort & ConfigPort & InfraPort;
|
||||
|
||||
const delay = (ms: number): Promise<void> => new Promise((r) => setTimeout(r, ms));
|
||||
|
||||
export class SchedulerService {
|
||||
constructor(private readonly deps: SchedulerDeps) {}
|
||||
|
||||
private get store() {
|
||||
return this.deps.store;
|
||||
}
|
||||
|
||||
// ───────────────────────────── Reads ─────────────────────────────
|
||||
|
||||
listJobs(): ScheduledJob[] {
|
||||
return Object.values(this.store.getScheduledJobs());
|
||||
}
|
||||
|
||||
getJob(id: string): ScheduledJob | null {
|
||||
return this.store.getScheduledJob(id);
|
||||
}
|
||||
|
||||
listRuns(jobId?: string): ScheduledJobRun[] {
|
||||
const all = Object.values(this.store.getScheduledJobRuns());
|
||||
const filtered = jobId ? all.filter((r) => r.scheduledJobId === jobId) : all;
|
||||
return filtered.sort((a, b) => b.startedAt - a.startedAt);
|
||||
}
|
||||
|
||||
/** Number of active sessions of a given agent type (for the multi-session warning). */
|
||||
countActiveAgents(agentType: string): number {
|
||||
let n = 0;
|
||||
for (const s of this.deps.sessions.values()) if (s.mode === agentType) n++;
|
||||
return n;
|
||||
}
|
||||
|
||||
// ──────────────────────────── Mutations ───────────────────────────
|
||||
|
||||
createJob(input: ScheduledJobInput): ScheduledJob {
|
||||
const now = Date.now();
|
||||
const job: ScheduledJob = {
|
||||
id: uuidv4(),
|
||||
name: input.name,
|
||||
agentType: input.agentType,
|
||||
workingDir: input.workingDir,
|
||||
launchCommand: input.launchCommand,
|
||||
promptMode: input.promptMode,
|
||||
promptText: input.promptText,
|
||||
promptFilePath: input.promptFilePath,
|
||||
inputMode: input.inputMode,
|
||||
scheduleType: input.scheduleType,
|
||||
runAt: input.runAt,
|
||||
intervalMinutes: input.intervalMinutes,
|
||||
dailyTime: input.dailyTime,
|
||||
weeklyDays: input.weeklyDays,
|
||||
weeklyTime: input.weeklyTime,
|
||||
enabled: input.enabled,
|
||||
notes: input.notes,
|
||||
concurrencyPolicy: input.concurrencyPolicy,
|
||||
createdAt: now,
|
||||
updatedAt: now,
|
||||
lastRunAt: null,
|
||||
nextRunAt: null,
|
||||
lastStatus: null,
|
||||
lastDueKey: null,
|
||||
};
|
||||
job.nextRunAt = job.enabled ? computeNextRunAt(job, now) : null;
|
||||
this.store.setScheduledJob(job.id, job);
|
||||
this.broadcastListChanged();
|
||||
return job;
|
||||
}
|
||||
|
||||
updateJob(id: string, patch: Partial<ScheduledJobInput>): ScheduledJob | null {
|
||||
const existing = this.getJob(id);
|
||||
if (!existing) return null;
|
||||
const now = Date.now();
|
||||
const updated: ScheduledJob = {
|
||||
...existing,
|
||||
...patch,
|
||||
id: existing.id,
|
||||
createdAt: existing.createdAt,
|
||||
updatedAt: now,
|
||||
// Editing a job re-arms it: clear the one-time completion + dup-guard so a
|
||||
// changed schedule can fire again.
|
||||
completedOnce: false,
|
||||
lastDueKey: null,
|
||||
};
|
||||
updated.nextRunAt = updated.enabled ? computeNextRunAt(updated, now) : null;
|
||||
this.store.setScheduledJob(updated.id, updated);
|
||||
this.broadcastListChanged();
|
||||
return updated;
|
||||
}
|
||||
|
||||
setEnabled(id: string, enabled: boolean): ScheduledJob | null {
|
||||
const existing = this.getJob(id);
|
||||
if (!existing) return null;
|
||||
const now = Date.now();
|
||||
existing.enabled = enabled;
|
||||
existing.updatedAt = now;
|
||||
existing.nextRunAt = enabled ? computeNextRunAt(existing, now) : null;
|
||||
this.store.setScheduledJob(existing.id, existing);
|
||||
this.broadcastListChanged();
|
||||
return existing;
|
||||
}
|
||||
|
||||
deleteJob(id: string): boolean {
|
||||
if (!this.getJob(id)) return false;
|
||||
this.store.removeScheduledJob(id);
|
||||
for (const run of this.listRuns(id)) this.store.removeScheduledJobRun(run.id);
|
||||
this.deps.broadcast(SseEvent.SchedulerJobDeleted, { id });
|
||||
this.broadcastListChanged();
|
||||
return true;
|
||||
}
|
||||
|
||||
// ──────────────────────────── Execution ───────────────────────────
|
||||
|
||||
/** Manual Run Now — always launches regardless of schedule/enabled state. */
|
||||
async runNow(id: string): Promise<ScheduledJobRun | null> {
|
||||
const job = this.getJob(id);
|
||||
if (!job) return null;
|
||||
return this.launch(job, 'manual_run_now');
|
||||
}
|
||||
|
||||
/**
|
||||
* Background tick: launch every enabled job whose next run is due. Advances
|
||||
* each job's schedule and guards against double-launching the same due time.
|
||||
*/
|
||||
async tickDueJobs(now: number = Date.now()): Promise<void> {
|
||||
for (const job of this.listJobs()) {
|
||||
if (!job.enabled || job.nextRunAt == null || job.nextRunAt > now) continue;
|
||||
|
||||
const key = dueKeyFor(job.id, job.nextRunAt);
|
||||
if (job.lastDueKey === key) {
|
||||
// This due time was already consumed (overlap/restart) — just advance.
|
||||
this.advanceAfterFire(job, now);
|
||||
continue;
|
||||
}
|
||||
|
||||
// Optional concurrency policy for AUTOMATIC runs.
|
||||
if (job.concurrencyPolicy === 'skip_if_same_agent_running' && this.countActiveAgents(job.agentType) > 0) {
|
||||
job.lastDueKey = key;
|
||||
this.advanceAfterFire(job, now);
|
||||
continue;
|
||||
}
|
||||
|
||||
job.lastDueKey = key;
|
||||
// Advance the schedule BEFORE launching so a slow launch can't be
|
||||
// re-triggered by the next tick.
|
||||
this.advanceAfterFire(job, now);
|
||||
this.launch(job, 'scheduled').catch((err) =>
|
||||
console.error(`[scheduler] launch failed for job ${job.id}:`, getErrorMessage(err))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/** Recompute nextRunAt for loaded jobs on boot (e.g. after a restart). */
|
||||
init(): void {
|
||||
const now = Date.now();
|
||||
for (const job of this.listJobs()) {
|
||||
const isDeadOnce = job.scheduleType === 'once' && job.completedOnce;
|
||||
if (job.enabled && job.nextRunAt == null && !isDeadOnce) {
|
||||
job.nextRunAt = computeNextRunAt(job, now);
|
||||
this.store.setScheduledJob(job.id, job);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ──────────────────────────── Internals ───────────────────────────
|
||||
|
||||
private advanceAfterFire(job: ScheduledJob, now: number): void {
|
||||
if (job.scheduleType === 'once') {
|
||||
job.completedOnce = true;
|
||||
job.enabled = false;
|
||||
job.nextRunAt = null;
|
||||
} else {
|
||||
job.nextRunAt = computeNextRunAt(job, now);
|
||||
}
|
||||
job.updatedAt = now;
|
||||
this.store.setScheduledJob(job.id, job);
|
||||
this.broadcastListChanged();
|
||||
}
|
||||
|
||||
private async launch(job: ScheduledJob, trigger: TriggerType): Promise<ScheduledJobRun> {
|
||||
const run: ScheduledJobRun = {
|
||||
id: uuidv4(),
|
||||
scheduledJobId: job.id,
|
||||
sessionId: null,
|
||||
sessionName: null,
|
||||
startedAt: Date.now(),
|
||||
finishedAt: null,
|
||||
status: 'created',
|
||||
triggerType: trigger,
|
||||
createdSessionUrl: null,
|
||||
};
|
||||
this.store.setScheduledJobRun(run.id, run);
|
||||
this.deps.broadcast(SseEvent.SchedulerRunCreated, run);
|
||||
|
||||
// Resolve the prompt.
|
||||
let prompt: string;
|
||||
try {
|
||||
prompt = await this.resolvePrompt(job);
|
||||
} catch (err) {
|
||||
return this.failRun(job, run, `Prompt error: ${getErrorMessage(err)}`);
|
||||
}
|
||||
|
||||
// Validate working directory.
|
||||
try {
|
||||
if (!statSync(job.workingDir).isDirectory()) {
|
||||
return this.failRun(job, run, 'workingDir is not a directory');
|
||||
}
|
||||
} catch {
|
||||
return this.failRun(job, run, 'workingDir does not exist');
|
||||
}
|
||||
|
||||
// Respect the global session cap.
|
||||
if (this.deps.sessions.size >= MAX_CONCURRENT_SESSIONS) {
|
||||
return this.failRun(job, run, `Maximum concurrent sessions (${MAX_CONCURRENT_SESSIONS}) reached`);
|
||||
}
|
||||
|
||||
// Create + start the session (mirrors the quick-start route flow).
|
||||
let session: Session;
|
||||
try {
|
||||
const mode = job.agentType;
|
||||
const globalNice = await this.deps.getGlobalNiceConfig();
|
||||
const modelConfig = await this.deps.getModelConfig();
|
||||
const claudeModeConfig = await this.deps.getClaudeModeConfig();
|
||||
const model = mode !== 'shell' ? modelConfig?.defaultModel || undefined : undefined;
|
||||
session = new Session({
|
||||
workingDir: job.workingDir,
|
||||
mode,
|
||||
name: job.name,
|
||||
mux: this.deps.mux,
|
||||
useMux: true,
|
||||
niceConfig: globalNice,
|
||||
model,
|
||||
claudeMode: claudeModeConfig.claudeMode,
|
||||
allowedTools: claudeModeConfig.allowedTools,
|
||||
});
|
||||
this.deps.addSession(session);
|
||||
this.store.incrementSessionsCreated();
|
||||
this.deps.persistSessionState(session);
|
||||
await this.deps.setupSessionListeners(session);
|
||||
this.deps.broadcast(SseEvent.SessionCreated, this.deps.getSessionStateWithRespawn(session));
|
||||
if (mode === 'shell') {
|
||||
await session.startShell();
|
||||
} else {
|
||||
await session.startInteractive();
|
||||
}
|
||||
this.deps.broadcast(SseEvent.SessionInteractive, { id: session.id, mode });
|
||||
} catch (err) {
|
||||
return this.failRun(job, run, `Session launch failed: ${getErrorMessage(err)}`);
|
||||
}
|
||||
|
||||
run.sessionId = session.id;
|
||||
run.sessionName = session.name;
|
||||
run.createdSessionUrl = `/?session=${session.id}`;
|
||||
run.status = 'session_started';
|
||||
this.store.setScheduledJobRun(run.id, run);
|
||||
this.deps.broadcast(SseEvent.SchedulerRunUpdated, run);
|
||||
this.updateJobLastStatus(job.id, 'session_started');
|
||||
|
||||
// Send the prompt once the CLI is ready (async; does not block the caller).
|
||||
this.sendPromptWhenReady(session.id, prompt, job, run);
|
||||
return run;
|
||||
}
|
||||
|
||||
private async resolvePrompt(job: ScheduledJob): Promise<string> {
|
||||
if (job.promptMode === 'prompt_file_path') {
|
||||
if (!job.promptFilePath) throw new Error('prompt file path is empty');
|
||||
return readFile(job.promptFilePath, 'utf-8');
|
||||
}
|
||||
return job.promptText ?? '';
|
||||
}
|
||||
|
||||
private sendPromptWhenReady(sessionId: string, prompt: string, job: ScheduledJob, run: ScheduledJobRun): void {
|
||||
setImmediate(() => {
|
||||
const poll = async (): Promise<void> => {
|
||||
if (job.agentType !== 'shell') {
|
||||
for (let attempt = 0; attempt < SCHEDULER_READY_MAX_ATTEMPTS; attempt++) {
|
||||
await delay(500);
|
||||
const s = this.deps.sessions.get(sessionId);
|
||||
if (!s) return; // session was removed
|
||||
const buf = s.getTerminalBuffer().slice(-2048);
|
||||
if (buf.includes('❯') || buf.includes('tokens')) break;
|
||||
}
|
||||
await delay(SCHEDULER_READY_SETTLE_MS);
|
||||
} else {
|
||||
await delay(1000);
|
||||
}
|
||||
const s = this.deps.sessions.get(sessionId);
|
||||
if (!s) return;
|
||||
try {
|
||||
const payload = prompt.endsWith('\r') ? prompt : `${prompt}\r`;
|
||||
if (job.inputMode === 'paste') {
|
||||
s.write(payload);
|
||||
} else {
|
||||
await s.writeViaMux(payload);
|
||||
}
|
||||
run.status = 'prompt_sent';
|
||||
run.finishedAt = Date.now();
|
||||
this.store.setScheduledJobRun(run.id, run);
|
||||
this.deps.broadcast(SseEvent.SchedulerRunUpdated, run);
|
||||
this.updateJobLastStatus(job.id, 'prompt_sent');
|
||||
} catch (err) {
|
||||
this.failRun(job, run, `Failed to send prompt: ${getErrorMessage(err)}`);
|
||||
}
|
||||
};
|
||||
poll().catch((err) => console.error('[scheduler] sendPromptWhenReady error:', getErrorMessage(err)));
|
||||
});
|
||||
}
|
||||
|
||||
private failRun(job: ScheduledJob, run: ScheduledJobRun, message: string): ScheduledJobRun {
|
||||
run.status = 'failed';
|
||||
run.errorMessage = message;
|
||||
run.finishedAt = Date.now();
|
||||
this.store.setScheduledJobRun(run.id, run);
|
||||
this.deps.broadcast(SseEvent.SchedulerRunUpdated, run);
|
||||
this.updateJobLastStatus(job.id, 'failed');
|
||||
return run;
|
||||
}
|
||||
|
||||
private updateJobLastStatus(jobId: string, status: ScheduledJobRunStatus): void {
|
||||
const fresh = this.store.getScheduledJob(jobId);
|
||||
if (!fresh) return;
|
||||
const now = Date.now();
|
||||
fresh.lastStatus = status;
|
||||
fresh.lastRunAt = now;
|
||||
fresh.updatedAt = now;
|
||||
this.store.setScheduledJob(fresh.id, fresh);
|
||||
this.broadcastListChanged();
|
||||
}
|
||||
|
||||
private broadcastListChanged(): void {
|
||||
this.deps.broadcast(SseEvent.SchedulerJobsChanged, { jobs: this.listJobs() });
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,82 @@
|
||||
/**
|
||||
* @fileoverview Pure next-run-time calculations for the scheduler.
|
||||
*
|
||||
* All functions are pure and take an explicit `after` timestamp (epoch ms) so
|
||||
* they are deterministic and unit-testable. Times use the SERVER'S LOCAL
|
||||
* timezone for v0.1 (per the build brief) — daily/weekly wall-clock times are
|
||||
* interpreted via the host's local time.
|
||||
*/
|
||||
|
||||
import type { ScheduledJob } from '../types/scheduler.js';
|
||||
|
||||
/** Parse an 'HH:MM' (24-hour) string into hours/minutes, or null if invalid. */
|
||||
export function parseHHMM(value: string | undefined): { hours: number; minutes: number } | null {
|
||||
if (!value) return null;
|
||||
const m = /^(\d{1,2}):(\d{2})$/.exec(value.trim());
|
||||
if (!m) return null;
|
||||
const hours = Number(m[1]);
|
||||
const minutes = Number(m[2]);
|
||||
if (hours < 0 || hours > 23 || minutes < 0 || minutes > 59) return null;
|
||||
return { hours, minutes };
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the epoch-ms timestamp for `hours:minutes` (local time) on the day of
|
||||
* `base`, shifted by `dayOffset` days.
|
||||
*/
|
||||
function atLocalTime(base: number, hours: number, minutes: number, dayOffset: number): number {
|
||||
const d = new Date(base);
|
||||
d.setHours(hours, minutes, 0, 0);
|
||||
d.setDate(d.getDate() + dayOffset);
|
||||
return d.getTime();
|
||||
}
|
||||
|
||||
/**
|
||||
* Compute the next fire time strictly relevant to `after`, or null if the job
|
||||
* has no future run (e.g. a completed one-time job, or invalid config).
|
||||
*
|
||||
* For `once`, returns the absolute `runAt` (even if already in the past, so a
|
||||
* missed one-time job still fires once) until it has `completedOnce`.
|
||||
*/
|
||||
export function computeNextRunAt(job: ScheduledJob, after: number): number | null {
|
||||
switch (job.scheduleType) {
|
||||
case 'once': {
|
||||
if (job.completedOnce) return null;
|
||||
return typeof job.runAt === 'number' ? job.runAt : null;
|
||||
}
|
||||
case 'interval': {
|
||||
const minutes = job.intervalMinutes;
|
||||
if (!minutes || minutes <= 0) return null;
|
||||
return after + minutes * 60_000;
|
||||
}
|
||||
case 'daily': {
|
||||
const t = parseHHMM(job.dailyTime);
|
||||
if (!t) return null;
|
||||
let next = atLocalTime(after, t.hours, t.minutes, 0);
|
||||
if (next <= after) next = atLocalTime(after, t.hours, t.minutes, 1);
|
||||
return next;
|
||||
}
|
||||
case 'weekly': {
|
||||
const t = parseHHMM(job.weeklyTime);
|
||||
if (!t) return null;
|
||||
const days = (job.weeklyDays ?? []).filter((d) => d >= 0 && d <= 6);
|
||||
if (days.length === 0) return null;
|
||||
for (let offset = 0; offset <= 7; offset++) {
|
||||
const cand = atLocalTime(after, t.hours, t.minutes, offset);
|
||||
if (cand > after && days.includes(new Date(cand).getDay())) return cand;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
default:
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Duplicate-launch guard key: identifies a specific due time for a job. The
|
||||
* scheduler records the key it last consumed so an overlapping or restarted
|
||||
* loop will not launch the same due time twice.
|
||||
*/
|
||||
export function dueKeyFor(jobId: string, fireTime: number): string {
|
||||
return `${jobId}:${fireTime}`;
|
||||
}
|
||||
@@ -272,6 +272,12 @@ export class StateStore {
|
||||
if (this.state.tokenStats) {
|
||||
parts.push(`"tokenStats":${JSON.stringify(this.state.tokenStats)}`);
|
||||
}
|
||||
if (this.state.scheduledJobs) {
|
||||
parts.push(`"scheduledJobs":${JSON.stringify(this.state.scheduledJobs)}`);
|
||||
}
|
||||
if (this.state.scheduledJobRuns) {
|
||||
parts.push(`"scheduledJobRuns":${JSON.stringify(this.state.scheduledJobRuns)}`);
|
||||
}
|
||||
|
||||
return `{${parts.join(',')}}`;
|
||||
}
|
||||
@@ -568,6 +574,51 @@ export class StateStore {
|
||||
this.save();
|
||||
}
|
||||
|
||||
// ========== Scheduled Job Methods (cron-style scheduler) ==========
|
||||
|
||||
/** Returns all scheduled jobs keyed by job ID. */
|
||||
getScheduledJobs(): Record<string, import('./types/scheduler.js').ScheduledJob> {
|
||||
if (!this.state.scheduledJobs) this.state.scheduledJobs = {};
|
||||
return this.state.scheduledJobs;
|
||||
}
|
||||
|
||||
/** Returns a scheduled job by ID, or null if not found. */
|
||||
getScheduledJob(id: string): import('./types/scheduler.js').ScheduledJob | null {
|
||||
return this.state.scheduledJobs?.[id] ?? null;
|
||||
}
|
||||
|
||||
/** Sets a scheduled job and triggers a debounced save. */
|
||||
setScheduledJob(id: string, job: import('./types/scheduler.js').ScheduledJob): void {
|
||||
if (!this.state.scheduledJobs) this.state.scheduledJobs = {};
|
||||
this.state.scheduledJobs[id] = job;
|
||||
this.save();
|
||||
}
|
||||
|
||||
/** Removes a scheduled job and triggers a debounced save. */
|
||||
removeScheduledJob(id: string): void {
|
||||
if (this.state.scheduledJobs) delete this.state.scheduledJobs[id];
|
||||
this.save();
|
||||
}
|
||||
|
||||
/** Returns all scheduled job runs keyed by run ID. */
|
||||
getScheduledJobRuns(): Record<string, import('./types/scheduler.js').ScheduledJobRun> {
|
||||
if (!this.state.scheduledJobRuns) this.state.scheduledJobRuns = {};
|
||||
return this.state.scheduledJobRuns;
|
||||
}
|
||||
|
||||
/** Sets a scheduled job run (history record) and triggers a debounced save. */
|
||||
setScheduledJobRun(id: string, run: import('./types/scheduler.js').ScheduledJobRun): void {
|
||||
if (!this.state.scheduledJobRuns) this.state.scheduledJobRuns = {};
|
||||
this.state.scheduledJobRuns[id] = run;
|
||||
this.save();
|
||||
}
|
||||
|
||||
/** Removes a scheduled job run and triggers a debounced save. */
|
||||
removeScheduledJobRun(id: string): void {
|
||||
if (this.state.scheduledJobRuns) delete this.state.scheduledJobRuns[id];
|
||||
this.save();
|
||||
}
|
||||
|
||||
/** Returns the application configuration. */
|
||||
getConfig() {
|
||||
return this.state.config;
|
||||
|
||||
@@ -23,6 +23,7 @@ import type { SessionState } from './session.js';
|
||||
import type { TaskState } from './task.js';
|
||||
import type { RalphLoopState } from './ralph.js';
|
||||
import type { RespawnConfig } from './respawn.js';
|
||||
import type { ScheduledJob, ScheduledJobRun } from './scheduler.js';
|
||||
|
||||
// ========== Global Stats Types ==========
|
||||
|
||||
@@ -111,6 +112,10 @@ export interface AppState {
|
||||
tokenStats?: TokenStats;
|
||||
/** Orchestrator Loop state (phased plan execution) */
|
||||
orchestrator?: import('./orchestrator.js').OrchestratorPersistState;
|
||||
/** Cron-style scheduled jobs, keyed by job ID. */
|
||||
scheduledJobs?: Record<string, ScheduledJob>;
|
||||
/** Scheduled job run history, keyed by run ID. */
|
||||
scheduledJobRuns?: Record<string, ScheduledJobRun>;
|
||||
}
|
||||
|
||||
// ========== Default Configuration ==========
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
/**
|
||||
* @fileoverview Scheduled Jobs (cron-style scheduler) type definitions.
|
||||
*
|
||||
* NOTE: This is intentionally distinct from the existing `ScheduledRun` concept
|
||||
* (see src/web/ports/infra-port.ts), which is a run-now, duration-bounded
|
||||
* autonomous loop. A `ScheduledJob` is a SAVED, NAMED job with a recurring
|
||||
* schedule (once/interval/daily/weekly), enable/disable, next-run calculation,
|
||||
* and a history of `ScheduledJobRun` records. The two do not interact.
|
||||
*
|
||||
* Persisted to `~/.codeman/state.json` via StateStore (see AppState).
|
||||
*/
|
||||
|
||||
import type { SessionMode } from './session.js';
|
||||
|
||||
/** How a job's fire times are computed. */
|
||||
export type ScheduleType = 'once' | 'interval' | 'daily' | 'weekly';
|
||||
|
||||
/** Where the prompt text comes from. */
|
||||
export type PromptMode = 'inline_text' | 'prompt_file_path';
|
||||
|
||||
/** How the prompt is delivered into the session. */
|
||||
export type InputMode = 'paste' | 'typed';
|
||||
|
||||
/** Lifecycle status of a single job execution. */
|
||||
export type ScheduledJobRunStatus = 'created' | 'session_started' | 'prompt_sent' | 'failed';
|
||||
|
||||
/** What triggered a run. */
|
||||
export type TriggerType = 'scheduled' | 'manual_run_now';
|
||||
|
||||
/** What to do for an AUTOMATIC run when sessions of the same agent already exist. */
|
||||
export type ConcurrencyPolicy = 'warn_only' | 'skip_if_same_agent_running';
|
||||
|
||||
/**
|
||||
* A saved, named scheduled job.
|
||||
*/
|
||||
export interface ScheduledJob {
|
||||
id: string;
|
||||
name: string;
|
||||
/** Reuses Codeman's existing session modes; 'shell' covers Terminal/custom. */
|
||||
agentType: SessionMode;
|
||||
workingDir: string;
|
||||
/** Optional custom launch command (only meaningful for 'shell' mode). */
|
||||
launchCommand?: string;
|
||||
|
||||
promptMode: PromptMode;
|
||||
promptText?: string;
|
||||
promptFilePath?: string;
|
||||
inputMode: InputMode;
|
||||
|
||||
scheduleType: ScheduleType;
|
||||
/** once: absolute epoch-ms fire time. */
|
||||
runAt?: number;
|
||||
/** interval: minutes between fires. */
|
||||
intervalMinutes?: number;
|
||||
/** daily: 'HH:MM' (24h, server-local time). */
|
||||
dailyTime?: string;
|
||||
/** weekly: weekdays 0–6 (0=Sunday). */
|
||||
weeklyDays?: number[];
|
||||
/** weekly: 'HH:MM' (24h, server-local time). */
|
||||
weeklyTime?: string;
|
||||
|
||||
enabled: boolean;
|
||||
notes?: string;
|
||||
/** Applies to automatic (scheduled) runs only. Manual Run Now always warns client-side. */
|
||||
concurrencyPolicy: ConcurrencyPolicy;
|
||||
|
||||
// ── Bookkeeping (server-maintained) ─────────────────────────────────────
|
||||
createdAt: number;
|
||||
updatedAt: number;
|
||||
lastRunAt: number | null;
|
||||
nextRunAt: number | null;
|
||||
lastStatus: ScheduledJobRunStatus | null;
|
||||
/** Duplicate-launch guard: identifies the most recent due-time consumed. */
|
||||
lastDueKey: string | null;
|
||||
/** True once a 'once' job has fired (it is also disabled). */
|
||||
completedOnce?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* A single execution of a scheduled job (history record).
|
||||
*/
|
||||
export interface ScheduledJobRun {
|
||||
id: string;
|
||||
scheduledJobId: string;
|
||||
sessionId: string | null;
|
||||
sessionName: string | null;
|
||||
startedAt: number;
|
||||
finishedAt: number | null;
|
||||
status: ScheduledJobRunStatus;
|
||||
errorMessage?: string;
|
||||
triggerType: TriggerType;
|
||||
/** Best-effort deep link to the created session in the web UI. */
|
||||
createdSessionUrl: string | null;
|
||||
}
|
||||
@@ -13,3 +13,4 @@ export type { ConfigPort } from './config-port.js';
|
||||
export type { InfraPort, ScheduledRun } from './infra-port.js';
|
||||
export type { AuthPort } from './auth-port.js';
|
||||
export type { OrchestratorPort } from './orchestrator-port.js';
|
||||
export type { SchedulerPort } from './scheduler-port.js';
|
||||
|
||||
@@ -0,0 +1,10 @@
|
||||
/**
|
||||
* @fileoverview Scheduler port — exposes the cron-style SchedulerService to
|
||||
* route handlers via the shared route context.
|
||||
*/
|
||||
|
||||
import type { SchedulerService } from '../../scheduler/scheduler-service.js';
|
||||
|
||||
export interface SchedulerPort {
|
||||
readonly scheduler: SchedulerService;
|
||||
}
|
||||
@@ -7,6 +7,7 @@ export { registerTeamRoutes } from './team-routes.js';
|
||||
export { registerMuxRoutes } from './mux-routes.js';
|
||||
export { registerFileRoutes } from './file-routes.js';
|
||||
export { registerScheduledRoutes } from './scheduled-routes.js';
|
||||
export { registerSchedulerRoutes } from './scheduler-routes.js';
|
||||
export { registerSystemRoutes } from './system-routes.js';
|
||||
export { registerHookEventRoutes } from './hook-event-routes.js';
|
||||
export { registerStatusTelemetryRoutes } from './status-telemetry-routes.js';
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
/**
|
||||
* @fileoverview Scheduled Jobs routes (cron-style scheduler).
|
||||
*
|
||||
* CRUD + enable/disable + Run Now + run history for `ScheduledJob`s. These are
|
||||
* separate from the legacy `/api/scheduled` (ScheduledRun) endpoints — see
|
||||
* SCHEDULER_DISCOVERY.md §0.
|
||||
*/
|
||||
|
||||
import { FastifyInstance } from 'fastify';
|
||||
import { ApiErrorCode, createErrorResponse } from '../../types.js';
|
||||
import { ScheduledJobSchema, ScheduledJobUpdateSchema, ScheduledJobEnabledSchema } from '../schemas.js';
|
||||
import { parseBody } from '../route-helpers.js';
|
||||
import type { SchedulerPort } from '../ports/index.js';
|
||||
|
||||
export function registerSchedulerRoutes(app: FastifyInstance, ctx: SchedulerPort): void {
|
||||
// ── Jobs ────────────────────────────────────────────────────────────────
|
||||
|
||||
app.get('/api/scheduler/jobs', async () => {
|
||||
return ctx.scheduler.listJobs();
|
||||
});
|
||||
|
||||
app.post('/api/scheduler/jobs', async (req) => {
|
||||
const body = parseBody(ScheduledJobSchema, req.body, 'Invalid scheduled job');
|
||||
return { job: ctx.scheduler.createJob(body) };
|
||||
});
|
||||
|
||||
app.get('/api/scheduler/jobs/:id', async (req) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const job = ctx.scheduler.getJob(id);
|
||||
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Scheduled job not found');
|
||||
return job;
|
||||
});
|
||||
|
||||
app.put('/api/scheduler/jobs/:id', async (req) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const body = parseBody(ScheduledJobUpdateSchema, req.body, 'Invalid scheduled job update');
|
||||
const job = ctx.scheduler.updateJob(id, body);
|
||||
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Scheduled job not found');
|
||||
return { job };
|
||||
});
|
||||
|
||||
app.delete('/api/scheduler/jobs/:id', async (req) => {
|
||||
const { id } = req.params as { id: string };
|
||||
if (!ctx.scheduler.deleteJob(id)) {
|
||||
return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Scheduled job not found');
|
||||
}
|
||||
return {};
|
||||
});
|
||||
|
||||
app.put('/api/scheduler/jobs/:id/enabled', async (req) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const { enabled } = parseBody(ScheduledJobEnabledSchema, req.body, 'Invalid request body');
|
||||
const job = ctx.scheduler.setEnabled(id, enabled);
|
||||
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Scheduled job not found');
|
||||
return { job };
|
||||
});
|
||||
|
||||
// ── Run Now ──────────────────────────────────────────────────────────────
|
||||
|
||||
app.post('/api/scheduler/jobs/:id/run', async (req) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const job = ctx.scheduler.getJob(id);
|
||||
if (!job) return createErrorResponse(ApiErrorCode.NOT_FOUND, 'Scheduled job not found');
|
||||
const run = await ctx.scheduler.runNow(id);
|
||||
return { run, activeAgents: ctx.scheduler.countActiveAgents(job.agentType) };
|
||||
});
|
||||
|
||||
// ── Run history ──────────────────────────────────────────────────────────
|
||||
|
||||
app.get('/api/scheduler/jobs/:id/runs', async (req) => {
|
||||
const { id } = req.params as { id: string };
|
||||
return ctx.scheduler.listRuns(id);
|
||||
});
|
||||
|
||||
app.get('/api/scheduler/runs', async () => {
|
||||
return ctx.scheduler.listRuns();
|
||||
});
|
||||
}
|
||||
@@ -574,6 +574,65 @@ export const ScheduledRunSchema = z.object({
|
||||
durationMinutes: z.number().int().min(1).max(14400).optional(),
|
||||
});
|
||||
|
||||
// ========== Scheduled Jobs (cron-style scheduler) ==========
|
||||
|
||||
/** 'HH:MM' 24-hour time. */
|
||||
const hhmmSchema = z.string().regex(/^([01]?\d|2[0-3]):[0-5]\d$/, 'Time must be HH:MM (24-hour)');
|
||||
|
||||
/** Shared field shape for creating/updating a scheduled job. */
|
||||
const ScheduledJobBaseSchema = z.object({
|
||||
name: z.string().min(1).max(200),
|
||||
agentType: z.enum(['claude', 'shell', 'opencode', 'codex', 'gemini']),
|
||||
workingDir: safePathSchema,
|
||||
launchCommand: z.string().max(2000).optional(),
|
||||
promptMode: z.enum(['inline_text', 'prompt_file_path']),
|
||||
promptText: z.string().max(100000).optional(),
|
||||
promptFilePath: safePathSchema.optional(),
|
||||
inputMode: z.enum(['paste', 'typed']),
|
||||
scheduleType: z.enum(['once', 'interval', 'daily', 'weekly']),
|
||||
runAt: z.number().int().positive().optional(),
|
||||
intervalMinutes: z.number().int().min(1).max(525600).optional(),
|
||||
dailyTime: hhmmSchema.optional(),
|
||||
weeklyDays: z.array(z.number().int().min(0).max(6)).min(1).max(7).optional(),
|
||||
weeklyTime: hhmmSchema.optional(),
|
||||
enabled: z.boolean(),
|
||||
notes: z.string().max(2000).optional(),
|
||||
concurrencyPolicy: z.enum(['warn_only', 'skip_if_same_agent_running']),
|
||||
});
|
||||
|
||||
/** Cross-field validation: required fields depend on promptMode + scheduleType. */
|
||||
function refineScheduledJob(val: z.infer<typeof ScheduledJobBaseSchema>, ctx: z.RefinementCtx): void {
|
||||
const add = (message: string, path: string) => ctx.addIssue({ code: 'custom', message, path: [path] });
|
||||
|
||||
if (val.promptMode === 'inline_text' && !val.promptText) {
|
||||
add('promptText is required when promptMode is inline_text', 'promptText');
|
||||
}
|
||||
if (val.promptMode === 'prompt_file_path' && !val.promptFilePath) {
|
||||
add('promptFilePath is required when promptMode is prompt_file_path', 'promptFilePath');
|
||||
}
|
||||
if (val.scheduleType === 'once' && val.runAt === undefined) {
|
||||
add('runAt is required for a one-time schedule', 'runAt');
|
||||
}
|
||||
if (val.scheduleType === 'interval' && val.intervalMinutes === undefined) {
|
||||
add('intervalMinutes is required for an interval schedule', 'intervalMinutes');
|
||||
}
|
||||
if (val.scheduleType === 'daily' && !val.dailyTime) {
|
||||
add('dailyTime is required for a daily schedule', 'dailyTime');
|
||||
}
|
||||
if (val.scheduleType === 'weekly' && (!val.weeklyTime || !val.weeklyDays?.length)) {
|
||||
add('weeklyDays and weeklyTime are required for a weekly schedule', 'weeklyTime');
|
||||
}
|
||||
}
|
||||
|
||||
/** POST /api/scheduler/jobs — full job definition. */
|
||||
export const ScheduledJobSchema = ScheduledJobBaseSchema.superRefine(refineScheduledJob);
|
||||
|
||||
/** PUT /api/scheduler/jobs/:id — partial update. */
|
||||
export const ScheduledJobUpdateSchema = ScheduledJobBaseSchema.partial();
|
||||
|
||||
/** PUT /api/scheduler/jobs/:id/enabled */
|
||||
export const ScheduledJobEnabledSchema = z.object({ enabled: z.boolean() });
|
||||
|
||||
/** POST /api/cases/link */
|
||||
export const LinkCaseSchema = z.object({
|
||||
name: z.string().regex(/^[a-zA-Z0-9_-]+$/, 'Invalid case name format'),
|
||||
|
||||
@@ -152,8 +152,10 @@ import {
|
||||
registerClipboardRoutes,
|
||||
registerSearchRoutes,
|
||||
registerOrchestratorRoutes,
|
||||
registerSchedulerRoutes,
|
||||
registerWsRoutes,
|
||||
} from './routes/index.js';
|
||||
import { SchedulerService } from '../scheduler/scheduler-service.js';
|
||||
|
||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||
|
||||
@@ -175,6 +177,7 @@ import {
|
||||
ITERATION_PAUSE_MS,
|
||||
STATS_COLLECTION_INTERVAL_MS,
|
||||
INACTIVITY_TIMEOUT_MS,
|
||||
SCHEDULER_TICK_INTERVAL,
|
||||
} from '../config/server-timing.js';
|
||||
|
||||
/**
|
||||
@@ -223,6 +226,8 @@ export class WebServer extends EventEmitter {
|
||||
// Store session listener references for explicit cleanup (prevents memory leaks)
|
||||
private sessionListenerRefs: Map<string, SessionListenerRefs> = new Map();
|
||||
private scheduledRuns: Map<string, ScheduledRun> = new Map();
|
||||
/** Cron-style scheduler service (assigned in setupRoutes). */
|
||||
private schedulerService!: SchedulerService;
|
||||
private sse: SseStreamManager;
|
||||
private store = getStore();
|
||||
private port: number;
|
||||
@@ -873,6 +878,13 @@ export class WebServer extends EventEmitter {
|
||||
registerClipboardRoutes(this.app, ctx);
|
||||
registerSearchRoutes(this.app, ctx);
|
||||
registerOrchestratorRoutes(this.app, ctx);
|
||||
|
||||
// Cron-style scheduler: build the service from the same context, recompute
|
||||
// due times for any persisted jobs, then expose it to its routes.
|
||||
this.schedulerService = new SchedulerService(ctx);
|
||||
this.schedulerService.init();
|
||||
registerSchedulerRoutes(this.app, { ...ctx, scheduler: this.schedulerService });
|
||||
|
||||
registerWsRoutes(this.app, ctx, () => this.getHostPolicy());
|
||||
}
|
||||
|
||||
@@ -1947,6 +1959,17 @@ export class WebServer extends EventEmitter {
|
||||
{ description: 'scheduled runs cleanup' }
|
||||
);
|
||||
|
||||
// Start the cron-style scheduler loop (fires due ScheduledJobs).
|
||||
this.cleanup.setInterval(
|
||||
() => {
|
||||
this.schedulerService.tickDueJobs().catch((err) => {
|
||||
console.error('[scheduler] tick failed:', getErrorMessage(err));
|
||||
});
|
||||
},
|
||||
SCHEDULER_TICK_INTERVAL,
|
||||
{ description: 'scheduled jobs due-checker' }
|
||||
);
|
||||
|
||||
// Start SSE client health check timer (prevents memory leaks from dead connections)
|
||||
this.cleanup.setInterval(
|
||||
() => {
|
||||
|
||||
@@ -240,6 +240,17 @@ export const ScheduledLog = 'scheduled:log' as const;
|
||||
/** Scheduled run deleted. */
|
||||
export const ScheduledDeleted = 'scheduled:deleted' as const;
|
||||
|
||||
// ─── Scheduled Jobs (cron-style scheduler) ───────────────────────────────────
|
||||
|
||||
/** The scheduled-jobs list changed (created/updated/enabled/run-status). Payload: { jobs }. */
|
||||
export const SchedulerJobsChanged = 'scheduler:jobsChanged' as const;
|
||||
/** A scheduled job was deleted. Payload: { id }. */
|
||||
export const SchedulerJobDeleted = 'scheduler:jobDeleted' as const;
|
||||
/** A scheduled-job run (history record) was created. Payload: ScheduledJobRun. */
|
||||
export const SchedulerRunCreated = 'scheduler:runCreated' as const;
|
||||
/** A scheduled-job run (history record) was updated. Payload: ScheduledJobRun. */
|
||||
export const SchedulerRunUpdated = 'scheduler:runUpdated' as const;
|
||||
|
||||
// ─── Teams ───────────────────────────────────────────────────────────────────
|
||||
|
||||
/** Agent team created. */
|
||||
@@ -469,6 +480,12 @@ export const SseEvent = {
|
||||
ScheduledLog,
|
||||
ScheduledDeleted,
|
||||
|
||||
// Scheduled jobs (cron-style scheduler)
|
||||
SchedulerJobsChanged,
|
||||
SchedulerJobDeleted,
|
||||
SchedulerRunCreated,
|
||||
SchedulerRunUpdated,
|
||||
|
||||
// Teams
|
||||
TeamCreated,
|
||||
TeamUpdated,
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
/**
|
||||
* Unit tests for the scheduler's pure next-run-time calculations.
|
||||
* Timezone-independent: daily/weekly expectations are asserted via local
|
||||
* Date getters rather than hardcoded epoch values.
|
||||
*/
|
||||
|
||||
import { describe, it, expect } from 'vitest';
|
||||
import { parseHHMM, computeNextRunAt, dueKeyFor } from '../src/scheduler/scheduler-time.js';
|
||||
import type { ScheduledJob } from '../src/types/scheduler.js';
|
||||
|
||||
function baseJob(partial: Partial<ScheduledJob>): ScheduledJob {
|
||||
return {
|
||||
id: 'j1',
|
||||
name: 'test',
|
||||
agentType: 'claude',
|
||||
workingDir: '/tmp',
|
||||
promptMode: 'inline_text',
|
||||
promptText: 'hi',
|
||||
inputMode: 'typed',
|
||||
scheduleType: 'once',
|
||||
enabled: true,
|
||||
concurrencyPolicy: 'warn_only',
|
||||
createdAt: 0,
|
||||
updatedAt: 0,
|
||||
lastRunAt: null,
|
||||
nextRunAt: null,
|
||||
lastStatus: null,
|
||||
lastDueKey: null,
|
||||
...partial,
|
||||
};
|
||||
}
|
||||
|
||||
describe('parseHHMM', () => {
|
||||
it('parses valid 24h times', () => {
|
||||
expect(parseHHMM('09:30')).toEqual({ hours: 9, minutes: 30 });
|
||||
expect(parseHHMM('23:59')).toEqual({ hours: 23, minutes: 59 });
|
||||
expect(parseHHMM('0:00')).toEqual({ hours: 0, minutes: 0 });
|
||||
});
|
||||
it('rejects invalid times', () => {
|
||||
expect(parseHHMM('24:00')).toBeNull();
|
||||
expect(parseHHMM('12:60')).toBeNull();
|
||||
expect(parseHHMM('9:5')).toBeNull(); // minutes must be 2 digits
|
||||
expect(parseHHMM('abc')).toBeNull();
|
||||
expect(parseHHMM(undefined)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('computeNextRunAt — once', () => {
|
||||
it('returns runAt even when already in the past (missed one-time job still fires)', () => {
|
||||
const job = baseJob({ scheduleType: 'once', runAt: 1000 });
|
||||
expect(computeNextRunAt(job, 500)).toBe(1000);
|
||||
expect(computeNextRunAt(job, 5000)).toBe(1000);
|
||||
});
|
||||
it('returns null once completed', () => {
|
||||
const job = baseJob({ scheduleType: 'once', runAt: 1000, completedOnce: true });
|
||||
expect(computeNextRunAt(job, 500)).toBeNull();
|
||||
});
|
||||
it('returns null with no runAt', () => {
|
||||
expect(computeNextRunAt(baseJob({ scheduleType: 'once' }), 0)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('computeNextRunAt — interval', () => {
|
||||
it('adds intervalMinutes to the after time', () => {
|
||||
const job = baseJob({ scheduleType: 'interval', intervalMinutes: 60 });
|
||||
expect(computeNextRunAt(job, 1000)).toBe(1000 + 60 * 60_000);
|
||||
});
|
||||
it('returns null with no/invalid interval', () => {
|
||||
expect(computeNextRunAt(baseJob({ scheduleType: 'interval' }), 0)).toBeNull();
|
||||
expect(computeNextRunAt(baseJob({ scheduleType: 'interval', intervalMinutes: 0 }), 0)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('computeNextRunAt — daily', () => {
|
||||
it('schedules today when the time is still ahead', () => {
|
||||
const after = new Date(2026, 0, 1, 10, 0, 0).getTime();
|
||||
const next = computeNextRunAt(baseJob({ scheduleType: 'daily', dailyTime: '14:30' }), after)!;
|
||||
const d = new Date(next);
|
||||
expect(d.getHours()).toBe(14);
|
||||
expect(d.getMinutes()).toBe(30);
|
||||
expect(d.getDate()).toBe(1);
|
||||
expect(next).toBeGreaterThan(after);
|
||||
});
|
||||
it('rolls to tomorrow when the time has passed', () => {
|
||||
const after = new Date(2026, 0, 1, 16, 0, 0).getTime();
|
||||
const next = computeNextRunAt(baseJob({ scheduleType: 'daily', dailyTime: '14:30' }), after)!;
|
||||
const d = new Date(next);
|
||||
expect(d.getHours()).toBe(14);
|
||||
expect(d.getDate()).toBe(2);
|
||||
expect(next).toBeGreaterThan(after);
|
||||
});
|
||||
it('returns null with no time', () => {
|
||||
expect(computeNextRunAt(baseJob({ scheduleType: 'daily' }), 0)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('computeNextRunAt — weekly', () => {
|
||||
it('finds the next selected weekday at the configured time', () => {
|
||||
const after = new Date(2026, 0, 1, 12, 0, 0).getTime();
|
||||
const targetDay = (new Date(after).getDay() + 2) % 7;
|
||||
const job = baseJob({ scheduleType: 'weekly', weeklyDays: [targetDay], weeklyTime: '08:00' });
|
||||
const next = computeNextRunAt(job, after)!;
|
||||
const d = new Date(next);
|
||||
expect(d.getDay()).toBe(targetDay);
|
||||
expect(d.getHours()).toBe(8);
|
||||
expect(next).toBeGreaterThan(after);
|
||||
// Within the coming week.
|
||||
expect(next - after).toBeLessThanOrEqual(7 * 24 * 60 * 60_000);
|
||||
});
|
||||
it('picks the soonest of multiple selected days', () => {
|
||||
const after = new Date(2026, 0, 1, 12, 0, 0).getTime();
|
||||
const soon = (new Date(after).getDay() + 1) % 7;
|
||||
const later = (new Date(after).getDay() + 3) % 7;
|
||||
const job = baseJob({ scheduleType: 'weekly', weeklyDays: [later, soon], weeklyTime: '09:00' });
|
||||
const next = computeNextRunAt(job, after)!;
|
||||
expect(new Date(next).getDay()).toBe(soon);
|
||||
});
|
||||
it('returns null with no days or no time', () => {
|
||||
expect(computeNextRunAt(baseJob({ scheduleType: 'weekly', weeklyTime: '09:00' }), 0)).toBeNull();
|
||||
expect(computeNextRunAt(baseJob({ scheduleType: 'weekly', weeklyDays: [1] }), 0)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('dueKeyFor', () => {
|
||||
it('combines job id and fire time', () => {
|
||||
expect(dueKeyFor('j1', 123)).toBe('j1:123');
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user