import { getSql } from "../connection"; import { errMsg } from "../../utils/errors"; import { CronExpressionParser } from "cron-parser"; import { parseDuration } from "../../utils/duration"; import { computeInitialNextRun } from "../../utils/schedule"; import { getConfig } from "../../utils/config"; import { validateJobName } from "../../utils/job-workspace"; import type { ScheduleType, JobLifecycle } from "../../types"; /** Validate that a schedule string matches its declared type. Throws on mismatch. */ function validateSchedule(schedule: string, scheduleType: ScheduleType): void { switch (scheduleType) { case "cron": try { CronExpressionParser.parse(schedule); } catch (err) { throw new Error(`Invalid cron expression "${schedule}": ${errMsg(err)}`); } break; case "interval": try { parseDuration(schedule); } catch (err) { throw new Error(`Invalid interval "${schedule}": ${errMsg(err)}`); } break; case "once": { const d = new Date(schedule); if (isNaN(d.getTime())) { throw new Error(`Invalid timestamp "${schedule}": expected ISO 8601 date`); } break; } } } export interface Job { name: string; schedule: string; prompt: string; status: JobLifecycle; always: boolean; scheduleType: ScheduleType; agent: string | null; employee: string | null; model: string | null; stateless: boolean; nextRunAt: string | null; lastRunAt: string | null; createdAt: string; updatedAt: string; } const COLS = "name, schedule, prompt, status, always, schedule_type, agent, employee, model, stateless, next_run_at, last_run_at, created_at, updated_at"; function toJob(r: Record): Job { // Support both old (enabled boolean) and new (status text) schema let status: JobLifecycle; if (typeof r.status === "string" && ["active", "disabled", "archived"].includes(r.status)) { status = r.status as JobLifecycle; } else { status = r.enabled ? "active" : "disabled"; } return { name: r.name, schedule: r.schedule, prompt: r.prompt, status, always: r.always ?? false, scheduleType: r.schedule_type || "cron", agent: r.agent || null, employee: r.employee || null, model: r.model || null, stateless: r.stateless ?? false, nextRunAt: r.next_run_at ? String(r.next_run_at) : null, lastRunAt: r.last_run_at ? String(r.last_run_at) : null, createdAt: String(r.created_at), updatedAt: String(r.updated_at), }; } async function notifyChange(): Promise { const sql = getSql(); await sql`SELECT pg_notify('nia_jobs', '')`; } export async function create( name: string, schedule: string, prompt: string, always = false, scheduleType: ScheduleType = "cron", nextRunAt?: Date, agent?: string, stateless = false, model?: string, employee?: string, ): Promise { validateJobName(name); validateSchedule(schedule, scheduleType); const existing = await get(name); if (existing) { throw new Error(`Job "${name}" already exists. Use \`nia job remove ${name}\` first, or choose a different name.`); } const sql = getSql(); await sql` INSERT INTO jobs (name, schedule, prompt, always, schedule_type, next_run_at, agent, stateless, model, employee, status) VALUES (${name}, ${schedule}, ${prompt}, ${always}, ${scheduleType}, ${nextRunAt ?? null}, ${agent ?? null}, ${stateless}, ${model ?? null}, ${employee ?? null}, 'active') `; await notifyChange(); } export async function list(): Promise { const sql = getSql(); const rows = await sql`SELECT ${sql.unsafe(COLS)} FROM jobs ORDER BY name`; return rows.map(toJob); } export async function get(name: string): Promise { const sql = getSql(); const rows = await sql`SELECT ${sql.unsafe(COLS)} FROM jobs WHERE name = ${name}`; return rows.length > 0 ? toJob(rows[0]) : null; } export async function update( name: string, fields: Partial<{ schedule: string; prompt: string; status: JobLifecycle; always: boolean; agent: string | null; employee: string | null; model: string | null; stateless: boolean; scheduleType: ScheduleType; }>, ): Promise { const sql = getSql(); const existing = await get(name); if (!existing) return false; const schedule = fields.schedule ?? existing.schedule; const scheduleType = fields.scheduleType ?? existing.scheduleType; const prompt = fields.prompt ?? existing.prompt; const status = fields.status ?? existing.status; const always = fields.always ?? existing.always; const agent = fields.agent !== undefined ? fields.agent : existing.agent; const employee = fields.employee !== undefined ? fields.employee : existing.employee; const model = fields.model !== undefined ? fields.model : existing.model; const stateless = fields.stateless ?? existing.stateless; const scheduleChanged = fields.schedule !== undefined || fields.scheduleType !== undefined; if (scheduleChanged) { validateSchedule(schedule, scheduleType); } if (scheduleChanged) { const nextRun = computeInitialNextRun(scheduleType, schedule, getConfig().timezone); await sql` UPDATE jobs SET schedule = ${schedule}, schedule_type = ${scheduleType}, prompt = ${prompt}, status = ${status}, always = ${always}, agent = ${agent}, employee = ${employee}, model = ${model}, stateless = ${stateless}, next_run_at = ${nextRun}, updated_at = NOW() WHERE name = ${name} `; } else { await sql` UPDATE jobs SET schedule = ${schedule}, schedule_type = ${scheduleType}, prompt = ${prompt}, status = ${status}, always = ${always}, agent = ${agent}, employee = ${employee}, model = ${model}, stateless = ${stateless}, updated_at = NOW() WHERE name = ${name} `; } await notifyChange(); return true; } export async function remove(name: string): Promise { const sql = getSql(); const result = await sql`DELETE FROM jobs WHERE name = ${name}`; if (result.count > 0) await notifyChange(); return result.count > 0; } export async function listEnabled(): Promise { const sql = getSql(); const rows = await sql`SELECT ${sql.unsafe(COLS)} FROM jobs WHERE status = 'active' ORDER BY name`; return rows.map(toJob); } export async function listDue(): Promise { const sql = getSql(); const rows = await sql` SELECT ${sql.unsafe(COLS)} FROM jobs WHERE status = 'active' AND next_run_at <= NOW() ORDER BY next_run_at `; return rows.map(toJob); } export async function markRun(name: string, nextRunAt: Date | null): Promise { const sql = getSql(); if (nextRunAt) { await sql`UPDATE jobs SET last_run_at = NOW(), next_run_at = ${nextRunAt}, updated_at = NOW() WHERE name = ${name}`; } else { await sql`UPDATE jobs SET last_run_at = NOW(), status = 'disabled', updated_at = NOW() WHERE name = ${name}`; } await notifyChange(); }