diff --git a/SCHEDULER_DISCOVERY.md b/SCHEDULER_DISCOVERY.md new file mode 100644 index 00000000..b1f6c74e --- /dev/null +++ b/SCHEDULER_DISCOVERY.md @@ -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` (`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` and + `scheduledJobRuns?: Record` 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. diff --git a/src/config/server-timing.ts b/src/config/server-timing.ts index adbc6efe..0c47e8ec 100644 --- a/src/config/server-timing.ts +++ b/src/config/server-timing.ts @@ -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; diff --git a/src/scheduler/scheduler-input.ts b/src/scheduler/scheduler-input.ts new file mode 100644 index 00000000..c0e73dcf --- /dev/null +++ b/src/scheduler/scheduler-input.ts @@ -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; +} diff --git a/src/scheduler/scheduler-service.ts b/src/scheduler/scheduler-service.ts new file mode 100644 index 00000000..6cbc2456 --- /dev/null +++ b/src/scheduler/scheduler-service.ts @@ -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 => 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): 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 { + 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 { + 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 { + 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 { + 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 => { + 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() }); + } +} diff --git a/src/scheduler/scheduler-time.ts b/src/scheduler/scheduler-time.ts new file mode 100644 index 00000000..3d006f00 --- /dev/null +++ b/src/scheduler/scheduler-time.ts @@ -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}`; +} diff --git a/src/state-store.ts b/src/state-store.ts index ffab99f4..f5a24c44 100644 --- a/src/state-store.ts +++ b/src/state-store.ts @@ -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 { + 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 { + 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; diff --git a/src/types/app-state.ts b/src/types/app-state.ts index 56809a6b..d6797714 100644 --- a/src/types/app-state.ts +++ b/src/types/app-state.ts @@ -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; + /** Scheduled job run history, keyed by run ID. */ + scheduledJobRuns?: Record; } // ========== Default Configuration ========== diff --git a/src/types/scheduler.ts b/src/types/scheduler.ts new file mode 100644 index 00000000..d8b2552b --- /dev/null +++ b/src/types/scheduler.ts @@ -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; +} diff --git a/src/web/ports/index.ts b/src/web/ports/index.ts index 2373015f..3c952dc5 100644 --- a/src/web/ports/index.ts +++ b/src/web/ports/index.ts @@ -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'; diff --git a/src/web/ports/scheduler-port.ts b/src/web/ports/scheduler-port.ts new file mode 100644 index 00000000..d8632bb9 --- /dev/null +++ b/src/web/ports/scheduler-port.ts @@ -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; +} diff --git a/src/web/routes/index.ts b/src/web/routes/index.ts index 9adaec30..05cbb5ef 100644 --- a/src/web/routes/index.ts +++ b/src/web/routes/index.ts @@ -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'; diff --git a/src/web/routes/scheduler-routes.ts b/src/web/routes/scheduler-routes.ts new file mode 100644 index 00000000..3330c70c --- /dev/null +++ b/src/web/routes/scheduler-routes.ts @@ -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(); + }); +} diff --git a/src/web/schemas.ts b/src/web/schemas.ts index a6dc1416..b2eca9a8 100644 --- a/src/web/schemas.ts +++ b/src/web/schemas.ts @@ -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, 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'), diff --git a/src/web/server.ts b/src/web/server.ts index 61324c2b..4022b3b7 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -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 = new Map(); private scheduledRuns: Map = 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( () => { diff --git a/src/web/sse-events.ts b/src/web/sse-events.ts index 5d57e781..f5c638e2 100644 --- a/src/web/sse-events.ts +++ b/src/web/sse-events.ts @@ -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, diff --git a/test/scheduler-time.test.ts b/test/scheduler-time.test.ts new file mode 100644 index 00000000..b1e60383 --- /dev/null +++ b/test/scheduler-time.test.ts @@ -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 { + 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'); + }); +});