diff --git a/lib/fire-status.js b/lib/fire-status.js new file mode 100644 index 0000000..8796be2 --- /dev/null +++ b/lib/fire-status.js @@ -0,0 +1,5 @@ +/** Shared run status constants — keep free of fire/monitor imports. */ + +export const ACTIVE_RUN_STATUSES = new Set(['queued', 'running']) + +export const TERMINAL_RUN_STATUSES = new Set(['succeeded', 'failed', 'skipped']) diff --git a/lib/fire.js b/lib/fire.js index 21a6fe9..c3b269f 100644 --- a/lib/fire.js +++ b/lib/fire.js @@ -7,6 +7,14 @@ import { randomUUID } from 'node:crypto' import { wrapScheduledPrompt, scheduledPromptSource } from './prompt.js' import { decideDispatch } from './scheduler.js' import { appendRun, getJob, patchRun, upsertJob } from './store.js' +import { ACTIVE_RUN_STATUSES } from './fire-status.js' +import { + findLabelOverlapRun, + normalizeLabels, + projectJobState, +} from './monitor.js' + +export { ACTIVE_RUN_STATUSES } from './fire-status.js' export function createRunRecord(job, decision, now, trigger) { return { @@ -15,11 +23,14 @@ export function createRunRecord(job, decision, now, trigger) { scheduledAt: decision.scheduledAt ?? now, actualAt: now, status: 'queued', + stateEnteredAt: now, + exitedAt: null, trigger: trigger || 'schedule', sessionId: null, error: null, summary: '', reason: decision.reason || null, + outputRef: null, } } @@ -43,8 +54,48 @@ export function extractAssistantText(messages, maxChars = 4000) { return '' } -export function publicJob(job) { +export function publicJob(job, runs = null) { if (!job) return null + const jobRuns = Array.isArray(runs) + ? runs.filter((run) => run?.jobId === job.id) + : null + const projection = jobRuns + ? projectJobState(job, jobRuns) + : projectJobState(job, job.lastStatus + ? [{ status: job.lastStatus, stateEnteredAt: job.lastRunAt, actualAt: job.lastRunAt, exitedAt: job.lastRunAt }] + : []) + // When runs not provided, prefer lastStatus-based active hint only for running/queued tip. + let stateInfo = projection + if (!jobRuns && (job.lastStatus === 'running' || job.lastStatus === 'queued')) { + stateInfo = { + state: job.lastStatus === 'queued' ? 'pending' : 'running', + stateEnteredAt: job.lastRunAt || job.updatedAt || null, + activeRunId: null, + stuck: false, + } + } else if (!jobRuns && job.enabled === false) { + stateInfo = { + state: 'paused', + stateEnteredAt: job.updatedAt || job.createdAt || null, + activeRunId: null, + stuck: false, + } + } else if (!jobRuns && job.schedule?.kind === 'at' && job.nextRunAt == null + && (job.lastStatus === 'succeeded' || job.lastStatus === 'failed' || job.lastStatus === 'skipped')) { + stateInfo = { + state: job.lastStatus === 'succeeded' ? 'succeeded' : 'failed', + stateEnteredAt: job.lastRunAt || null, + activeRunId: null, + stuck: false, + } + } else if (!jobRuns) { + stateInfo = { + state: 'idle', + stateEnteredAt: job.lastRunAt || job.updatedAt || job.createdAt || null, + activeRunId: null, + stuck: false, + } + } return { id: job.id, name: job.name, @@ -68,6 +119,14 @@ export function publicJob(job) { } : { kind: 'dsh' }, mirrorToSession: job.mirrorToSession === true, + labels: normalizeLabels(job.labels), + ...job.watch ? { watch: job.watch } : {}, + ...job.progress ? { progress: job.progress } : {}, + ...job.report ? { report: job.report } : {}, + persistHistory: job.persistHistory || { kind: 'retain', limit: 200 }, + state: stateInfo.state, + stateEnteredAt: stateInfo.stateEnteredAt, + stuck: stateInfo.stuck === true, ...job.origin ? { origin: job.origin } : {}, schedule: job.schedule, createdAt: job.createdAt, @@ -80,6 +139,7 @@ export function publicJob(job) { /** * Claim one occurrence: overlap → skipped history row; otherwise queued run. + * Also skips when another job with the same role+task labels is already active. */ export function claimOccurrence(state, jobId, now, trigger, policies) { const job = getJob(state, jobId) @@ -89,7 +149,7 @@ export function claimOccurrence(state, jobId, now, trigger, policies) { throw error } const runs = (state.runs || []).filter((run) => run.jobId === jobId) - const decision = trigger === 'run-now' + let decision = trigger === 'run-now' ? decideRunNow(job, runs, now, policies) : decideDispatch({ job, @@ -100,17 +160,33 @@ export function claimOccurrence(state, jobId, now, trigger, policies) { graceMs: policies?.graceMs, }) + if (decision.action === 'fire') { + const overlap = findLabelOverlapRun(state, job) + if (overlap) { + decision = { + action: 'skip', + reason: 'label_overlap', + scheduledAt: decision.scheduledAt ?? now, + nextRunAt: Object.hasOwn(decision, 'nextRunAt') ? decision.nextRunAt : job.nextRunAt, + } + } + } + if (decision.action === 'wait') { - return { state, job: publicJob(job), decision, run: null } + return { state, job: publicJob(job, state.runs), decision, run: null } } const run = createRunRecord(job, decision, now, trigger) if (decision.action === 'skip') { run.status = 'skipped' + run.stateEnteredAt = now + run.exitedAt = now run.reason = decision.reason const oneshot = job.schedule?.kind === 'at' if (decision.reason === 'overlap') { run.summary = 'Skipped because a run is already queued or running' + } else if (decision.reason === 'label_overlap') { + run.summary = 'Skipped because another job with the same role+task labels is already running' } else if (decision.reason === 'misfire' && oneshot) { run.summary = 'Skipped: one-shot time missed (host downtime or delayed tick; no backlog). nextRunAt cleared.' run.error = 'misfire:one-shot' @@ -135,10 +211,11 @@ export function claimOccurrence(state, jobId, now, trigger, policies) { } let next = upsertJob(state, nextJob) if (!alreadySkipped) next = appendRun(next, run) - return { state: next, job: publicJob(nextJob), decision, run: alreadySkipped ? null : run } + return { state: next, job: publicJob(nextJob, next.runs), decision, run: alreadySkipped ? null : run } } run.status = 'queued' + run.stateEnteredAt = now const nextJob = { ...job, updatedAt: now, @@ -148,7 +225,7 @@ export function claimOccurrence(state, jobId, now, trigger, policies) { } let next = upsertJob(state, nextJob) next = appendRun(next, run) - return { state: next, job: publicJob(nextJob), decision, run } + return { state: next, job: publicJob(nextJob, next.runs), decision, run } } function decideRunNow(job, runs, now, policies) { @@ -159,8 +236,6 @@ function decideRunNow(job, runs, now, policies) { return { action: 'fire', scheduledAt: now, nextRunAt: job.nextRunAt } } -export const ACTIVE_RUN_STATUSES = new Set(['queued', 'running']) - function failRun(state, runId, job, now, message) { return settleRun(state, runId, { status: 'failed', @@ -177,12 +252,16 @@ export function settleRun(state, runId, terminal, now) { const status = terminal?.status && !ACTIVE_RUN_STATUSES.has(terminal.status) ? terminal.status : 'succeeded' + const prevStatus = run.status let next = patchRun(state, runId, { status, summary: terminal?.summary || '', error: terminal?.error || null, actualAt: now, + stateEnteredAt: prevStatus === status ? (run.stateEnteredAt || now) : now, + exitedAt: now, sessionId: run.sessionId, + ...terminal?.outputRef ? { outputRef: terminal.outputRef } : {}, }) if (job) { next = upsertJob(next, { @@ -222,7 +301,12 @@ export async function executeClaimedRun(state, runId, deps) { } const clock = deps.now || (() => Date.now()) const mark = clock() - let next = patchRun(state, runId, { status: 'running', actualAt: mark }) + let next = patchRun(state, runId, { + status: 'running', + actualAt: mark, + stateEnteredAt: mark, + exitedAt: null, + }) next = upsertJob(next, { ...job, lastRunAt: mark, lastStatus: 'running', updatedAt: mark }) const text = wrapScheduledPrompt(job.prompt, { @@ -235,7 +319,11 @@ export async function executeClaimedRun(state, runId, deps) { if (typeof deps.createAndPrompt !== 'function') { next = failRun(next, runId, job, clock(), 'session port unavailable') - return { state: next, run: (next.runs || []).find((row) => row.id === runId), job: publicJob(getJob(next, job.id)) } + return { + state: next, + run: (next.runs || []).find((row) => row.id === runId), + job: publicJob(getJob(next, job.id), next.runs), + } } try { @@ -256,12 +344,20 @@ export async function executeClaimedRun(state, runId, deps) { const waited = await deps.waitForTurn(sessionId) const ended = clock() next = settleRun(next, runId, waited || { status: 'succeeded', summary: created.summary }, ended) - return { state: next, run: (next.runs || []).find((row) => row.id === runId), job: publicJob(getJob(next, job.id)) } + return { + state: next, + run: (next.runs || []).find((row) => row.id === runId), + job: publicJob(getJob(next, job.id), next.runs), + } } if (ACTIVE_RUN_STATUSES.has(reported)) { // Leave running so overlap skip still works while a turn is in flight. // Caller (dispatchRun / live waitForTurn) must settle to a terminal status. - return { state: next, run: (next.runs || []).find((row) => row.id === runId), job: publicJob(getJob(next, job.id)) } + return { + state: next, + run: (next.runs || []).find((row) => row.id === runId), + job: publicJob(getJob(next, job.id), next.runs), + } } const ended = clock() next = settleRun(next, runId, { @@ -269,11 +365,19 @@ export async function executeClaimedRun(state, runId, deps) { summary: created.summary || 'Dispatched to a new session', error: created.error || null, }, ended) - return { state: next, run: (next.runs || []).find((row) => row.id === runId), job: publicJob(getJob(next, job.id)) } + return { + state: next, + run: (next.runs || []).find((row) => row.id === runId), + job: publicJob(getJob(next, job.id), next.runs), + } } catch (error) { const message = error instanceof Error ? error.message : String(error) next = failRun(next, runId, job, clock(), message) - return { state: next, run: (next.runs || []).find((row) => row.id === runId), job: publicJob(getJob(next, job.id)) } + return { + state: next, + run: (next.runs || []).find((row) => row.id === runId), + job: publicJob(getJob(next, job.id), next.runs), + } } } @@ -286,6 +390,8 @@ export function interruptActiveRuns(state, now, reason = 'host_interrupted') { error: reason, summary: 'Host stopped before this run finished', actualAt: now, + stateEnteredAt: now, + exitedAt: now, }) const job = getJob(next, run.jobId) if (job) { diff --git a/lib/host.js b/lib/host.js index 3219782..3204280 100644 --- a/lib/host.js +++ b/lib/host.js @@ -11,6 +11,13 @@ import { claimOccurrence, executeClaimedRun, extractAssistantText, interruptActi import { assertDeliveryAllowedForIdentity, deliverRunToIm, mergeDeliveryMention, mirrorRunToSession, normalizeDelivery, normalizeOrigin } from './delivery.js' import { apiError, resolveLocale } from './i18n.js' import { workspaceVisibleIds } from './isolation.js' +import { + formatWatchReport, + labelsMatch, + normalizeLabels, + readProgressSnapshot, + selectJobsByWatch, +} from './monitor.js' import { decideDispatch, nextFire, validateSchedule } from './scheduler.js' import { createJobRecord, @@ -18,6 +25,7 @@ import { DEFAULT_SETTINGS, getJob, listHistory, + listHistoryPage, listJobs, normalizeSettings, removeJob, @@ -337,23 +345,30 @@ async function readJsonBody(req) { return JSON.parse(text) } -function jobView(job) { - return publicJob(job) +function jobView(job, runs = null) { + return publicJob(job, runs) } function runView(run) { if (!run) return null return { id: run.id, + run_id: run.id, jobId: run.jobId, scheduledAt: run.scheduledAt, actualAt: run.actualAt, status: run.status, + state: run.status === 'queued' ? 'pending' + : run.status === 'skipped' ? 'failed' + : run.status, + stateEnteredAt: run.stateEnteredAt ?? run.actualAt ?? run.scheduledAt ?? null, + exitedAt: run.exitedAt ?? null, trigger: run.trigger, sessionId: run.sessionId, error: run.error, summary: run.summary, reason: run.reason, + output_ref: run.outputRef || (run.sessionId ? `session:${run.sessionId}` : null), } } @@ -420,13 +435,53 @@ export function createHostService(options = {}) { } } + /** + * Notify watcher jobs that watch this completed job (report.mode=on_complete). + */ + async function maybeNotifyWatchers(targetJob, targetRun) { + if (!targetJob || !targetRun) return + if (targetRun.status === 'queued' || targetRun.status === 'running') return + const state = await snapshot() + const watchers = listJobs(state).filter((job) => { + if (!job?.watch || !job.report) return false + if (job.report.mode !== 'on_complete') return false + if (job.report.delivery === 'none') return false + if (job.watch.taskId && job.watch.taskId === targetJob.id) return true + if (job.watch.labels && Object.keys(job.watch.labels).length) { + return labelsMatch(targetJob.labels, job.watch.labels, job.watch.match || 'all') + } + return false + }) + for (const watcher of watchers) { + try { + const text = formatWatchReport(watcher, [{ + id: targetJob.id, + name: targetJob.name, + state: targetRun.status === 'succeeded' ? 'succeeded' : 'failed', + stuck: false, + }]) + if (watcher.report.delivery === 'im' || watcher.delivery?.kind === 'im') { + const deliveryJob = watcher.report.delivery === 'im' || watcher.delivery?.kind === 'im' + ? { ...watcher, delivery: watcher.delivery?.kind === 'im' ? watcher.delivery : watcher.delivery } + : watcher + if (deliveryJob.delivery?.kind === 'im') { + await deliverRunToIm(deliveryJob, text, { dshIm: getDshIm() }) + } + } + } catch (error) { + logger.warn?.(`[dsh-ops-cron] watcher report failed for ${watcher.id}: ${error instanceof Error ? error.message : error}`) + } + } + } + async function settleAndNotify(jobId, runId, claimedDecision) { const state = await snapshot() const run = (state.runs || []).find((row) => row.id === runId) const job = getJob(state, jobId) await maybeDeliverIm(job, run) await maybeMirrorResult(job, run) - return { job: jobView(job), run: runView(run), decision: claimedDecision } + await maybeNotifyWatchers(job, run) + return { job: jobView(job, state.runs), run: runView(run), decision: claimedDecision } } async function createJob(input, identity = null, opts = {}) { @@ -491,7 +546,7 @@ export function createHostService(options = {}) { created = createJobRecord(safeInput, current, t) return upsertJob(current, created) }) - return jobView(created) + return jobView(created, []) } async function updateJob(jobId, patch, identity = null) { @@ -527,6 +582,13 @@ export function createHostService(options = {}) { : (patch.mirror_to_session !== undefined ? patch.mirror_to_session === true : job.mirrorToSession === true), + labels: patch.labels !== undefined ? patch.labels : job.labels, + watch: patch.watch !== undefined ? patch.watch : job.watch, + progress: patch.progress !== undefined ? patch.progress : job.progress, + report: patch.report !== undefined ? patch.report : job.report, + persistHistory: patch.persistHistory !== undefined + ? patch.persistHistory + : (patch.persist_history !== undefined ? patch.persist_history : job.persistHistory), ownerEmpNo: job.ownerEmpNo || UNASSIGNED_OWNER, ownerDisplayName: job.ownerDisplayName || '', } @@ -567,9 +629,13 @@ export function createHostService(options = {}) { if (patch.schedule === undefined && patch.enabled === undefined) { record.nextRunAt = job.nextRunAt } + // Clear optional monitor fields when explicitly set to null. + if (patch.watch === null) delete record.watch + if (patch.progress === null) delete record.progress + if (patch.report === null) delete record.report return upsertJob(current, record) }) - return jobView(getJob(state, jobId)) + return jobView(getJob(state, jobId), state.runs) } async function pauseJob(jobId, enabled, identity = null) { @@ -654,6 +720,7 @@ export function createHostService(options = {}) { const result = await dispatchRun(job.id, 'schedule') if (result.run) fired.push(result) } + await tickWatcherReports() return fired } finally { ticking = false @@ -853,7 +920,8 @@ export function createHostService(options = {}) { }) return changed ? next : current }) - const jobs = filterJobsForIdentity(listJobs(state), identity).map(jobView) + const jobs = filterJobsForIdentity(listJobs(state), identity) + .map((job) => jobView(job, state.runs)) write(200, { ok: true, jobs, @@ -881,7 +949,7 @@ export function createHostService(options = {}) { const existing = getJob(state, jobId) assertCanAccessJob(existing, identity) if (method === 'GET' && !rest) { - write(200, { ok: true, job: jobView(existing) }) + write(200, { ok: true, job: jobView(existing, state.runs) }) return } if (method === 'PATCH' && !rest) { @@ -965,6 +1033,62 @@ export function createHostService(options = {}) { return } + if (path === `${API_PREFIX}/query` && method === 'GET') { + const identity = await requireIdentity(req, write, getUdsAuth) + if (!identity) return + const taskId = url.searchParams.get('taskId') || url.searchParams.get('task_id') || '' + const stateFilter = url.searchParams.get('state') || '' + const match = url.searchParams.get('match') === 'any' ? 'any' : 'all' + const includeRuns = url.searchParams.get('include_runs') === '1' + || url.searchParams.get('include_runs') === 'true' + const decodeProgress = url.searchParams.get('decode_progress') === '1' + || url.searchParams.get('decode_progress') === 'true' + const labels = {} + for (const [key, value] of url.searchParams.entries()) { + if (key.startsWith('label.')) labels[key.slice(6)] = value + } + const result = await queryJobs({ + taskId: taskId || undefined, + labels: Object.keys(labels).length ? labels : undefined, + match, + state: stateFilter || undefined, + includeRuns, + decodeProgress, + }, identity) + write(200, { ok: true, ...result, viewer: viewerPayload(identity) }) + return + } + + if (path === `${API_PREFIX}/progress` && method === 'GET') { + const identity = await requireIdentity(req, write, getUdsAuth) + if (!identity) return + const taskId = url.searchParams.get('taskId') || url.searchParams.get('task_id') || '' + const labels = {} + for (const [key, value] of url.searchParams.entries()) { + if (key.startsWith('label.')) labels[key.slice(6)] = value + } + const match = url.searchParams.get('match') === 'any' ? 'any' : 'all' + const result = await getProgress({ + taskId: taskId || undefined, + labels: Object.keys(labels).length ? labels : undefined, + match, + }, identity) + write(200, { ok: true, ...result, viewer: viewerPayload(identity) }) + return + } + + if (path === `${API_PREFIX}/archive` && method === 'POST') { + const identity = await requireIdentity(req, write, getUdsAuth) + if (!identity) return + // Reserved: persist_history=archive:endpoint will push here later. + write(200, { + ok: true, + archived: false, + message: 'archive endpoint reserved; set persist_history=archive: on jobs for future use', + }) + return + } + if (path === `${API_PREFIX}/preview` && method === 'POST') { const identity = await requireIdentity(req, write, getUdsAuth) if (!identity) return @@ -994,7 +1118,7 @@ export function createHostService(options = {}) { const locale = resolveLocale(req) const code = error && error.code if (code === 'NOT_FOUND') return write(404, { ...apiError('not_found', locale), error: error.message || 'not_found' }) - if (code === 'INVALID_CRON' || code === 'INVALID_AT' || code === 'INVALID_SCHEDULE' || code === 'INVALID_JOB' || code === 'INVALID_TIMEZONE' || code === 'INVALID_CWD' || code === 'INVALID_DELIVERY' || code === 'IM_DELIVERY_FORBIDDEN' || code === 'IM_TARGET_FORBIDDEN' || code === 'IM_JOB_MISSING_CWD') { + if (code === 'INVALID_CRON' || code === 'INVALID_AT' || code === 'INVALID_SCHEDULE' || code === 'INVALID_JOB' || code === 'INVALID_TIMEZONE' || code === 'INVALID_CWD' || code === 'INVALID_DELIVERY' || code === 'IM_DELIVERY_FORBIDDEN' || code === 'IM_TARGET_FORBIDDEN' || code === 'IM_JOB_MISSING_CWD' || code === 'INVALID_WATCH' || code === 'INVALID_PROGRESS' || code === 'INVALID_REPORT' || code === 'INVALID_PERSIST_HISTORY' || code === 'INVALID_PROGRESS_PATH') { return write(400, { ok: false, error: error.message, code, message: error.message }) } if (code === 'PAYLOAD_TOO_LARGE') return write(413, apiError('payload_too_large', locale)) @@ -1005,6 +1129,161 @@ export function createHostService(options = {}) { } } + function filterJobsByQuery(jobs, query = {}) { + let out = [...jobs] + const labels = normalizeLabels(query.labels) + if (Object.keys(labels).length) { + const match = query.match === 'any' ? 'any' : 'all' + out = out.filter((job) => labelsMatch(job.labels, labels, match)) + } + if (query.taskId) { + out = out.filter((job) => job.id === query.taskId) + } + return out + } + + async function queryJobs(query = {}, identity = null) { + const state = await snapshot() + let jobs = listJobs(state) + if (identity) jobs = filterJobsForIdentity(jobs, identity) + const by = {} + if (query.taskId) by.taskId = String(query.taskId).trim() + if (query.labels) by.labels = normalizeLabels(query.labels) + if (query.match) by.match = query.match === 'any' ? 'any' : 'all' + if (by.taskId || (by.labels && Object.keys(by.labels).length)) { + jobs = selectJobsByWatch(jobs, by) + } else { + jobs = filterJobsByQuery(jobs, query) + } + const stateFilter = typeof query.state === 'string' ? query.state.trim() : '' + const views = jobs.map((job) => jobView(job, state.runs)) + const filtered = stateFilter + ? views.filter((job) => job.state === stateFilter) + : views + const includeRuns = query.includeRuns === true || query.include_runs === true + const decodeProgress = query.decodeProgress === true || query.decode_progress === true + const items = [] + for (const job of filtered) { + const item = { job } + if (includeRuns) { + item.runs = listHistory(state, job.id).map(runView) + } + if (decodeProgress) { + const raw = getJob(state, job.id) + item.progress = await readProgressSnapshot(raw) + } + items.push(item) + } + return { jobs: filtered, items, count: filtered.length } + } + + async function listRuns(taskId, opts = {}, identity = null) { + const id = String(taskId || '').trim() + if (!id) { + const error = new Error('task_id is required') + error.code = 'INVALID_JOB' + throw error + } + const state = await snapshot() + const job = getJob(state, id) + if (identity) assertCanAccessJob(job, identity) + else if (!job) { + const error = new Error('job not found') + error.code = 'NOT_FOUND' + throw error + } + const page = listHistoryPage(state, id, opts) + return { + task_id: id, + runs: page.runs.map(runView), + next_cursor: page.nextCursor, + total: page.total, + } + } + + async function getProgress(query = {}, identity = null) { + const state = await snapshot() + let jobs = listJobs(state) + if (identity) jobs = filterJobsForIdentity(jobs, identity) + const by = {} + if (query.taskId) by.taskId = String(query.taskId).trim() + if (query.labels) by.labels = normalizeLabels(query.labels) + if (query.match) by.match = query.match === 'any' ? 'any' : 'all' + if (query.by_watch && typeof query.by_watch === 'object') { + if (query.by_watch.taskId) by.taskId = String(query.by_watch.taskId).trim() + if (query.by_watch.labels) by.labels = normalizeLabels(query.by_watch.labels) + if (query.by_watch.match) by.match = query.by_watch.match === 'any' ? 'any' : 'all' + } + const selected = (by.taskId || (by.labels && Object.keys(by.labels).length)) + ? selectJobsByWatch(jobs, by) + : filterJobsByQuery(jobs, query) + const snapshots = [] + for (const job of selected) { + snapshots.push({ + job: jobView(job, state.runs), + progress: await readProgressSnapshot(job), + }) + } + return { snapshots, count: snapshots.length } + } + + /** + * Periodic / on_change report tick for watcher jobs. + */ + async function tickWatcherReports() { + const state = await snapshot() + const t = now() + for (const watcher of listJobs(state)) { + if (!watcher?.watch || !watcher.report) continue + if (watcher.report.delivery === 'none') continue + if (watcher.report.mode !== 'periodic' && watcher.report.mode !== 'on_change') continue + if (watcher.report.mode === 'periodic') { + const intervalMs = Math.max(1, Number(watcher.report.intervalMin) || 15) * 60_000 + if (watcher.lastReportAt && (t - watcher.lastReportAt) < intervalMs) continue + } + const targets = selectJobsByWatch(listJobs(state), watcher.watch) + .filter((job) => job.id !== watcher.id) + .map((job) => jobView(job, state.runs)) + if (!targets.length) continue + const progressById = {} + let changed = false + for (const target of targets) { + const raw = getJob(state, target.id) + const snap = await readProgressSnapshot(raw) + progressById[target.id] = snap + const hash = JSON.stringify(snap?.metrics || {}) + if (watcher.report.mode === 'on_change') { + const prev = watcher.lastProgressHash?.[target.id] + if (prev !== hash) changed = true + } + } + if (watcher.report.mode === 'on_change' && !changed) continue + const text = formatWatchReport(watcher, targets, progressById) + try { + if ((watcher.report.delivery === 'im' || watcher.delivery?.kind === 'im') + && watcher.delivery?.kind === 'im') { + await deliverRunToIm(watcher, text, { dshIm: getDshIm() }) + } + const hashes = { ...(watcher.lastProgressHash || {}) } + for (const [id, snap] of Object.entries(progressById)) { + hashes[id] = JSON.stringify(snap?.metrics || {}) + } + await withState((current) => { + const job = getJob(current, watcher.id) + if (!job) return current + return upsertJob(current, { + ...job, + lastReportAt: t, + lastProgressHash: hashes, + updatedAt: job.updatedAt, + }) + }) + } catch (error) { + logger.warn?.(`[dsh-ops-cron] periodic report failed for ${watcher.id}: ${error instanceof Error ? error.message : error}`) + } + } + } + return { store, now, @@ -1015,22 +1294,34 @@ export function createHostService(options = {}) { updateSettings, dispatchRun, tick, + tickWatcherReports, recover, startTimer, stopTimer, handleRequest, snapshot, getUdsAuth, - async listJobs(identity = null) { - const jobs = listJobs(await snapshot()) - const filtered = identity ? filterJobsForIdentity(jobs, identity) : jobs - return filtered.map(jobView) + queryJobs, + listRuns, + getProgress, + async listJobs(identity = null, query = null) { + const state = await snapshot() + let jobs = listJobs(state) + if (identity) jobs = filterJobsForIdentity(jobs, identity) + if (query) { + jobs = filterJobsByQuery(jobs, query) + const views = jobs.map((job) => jobView(job, state.runs)) + if (query.state) return views.filter((job) => job.state === query.state) + return views + } + return jobs.map((job) => jobView(job, state.runs)) }, async getJob(jobId, identity = null) { - const job = getJob(await snapshot(), jobId) + const state = await snapshot() + const job = getJob(state, jobId) if (!job) return null if (identity) assertCanAccessJob(job, identity) - return jobView(job) + return jobView(job, state.runs) }, async listHistory(jobId, identity = null) { const state = await snapshot() diff --git a/lib/index.js b/lib/index.js index db178ab..9bf7197 100644 --- a/lib/index.js +++ b/lib/index.js @@ -184,7 +184,7 @@ export function apply(ctx, config = {}) { ctx.inject(['tools'], (tctx) => { registerCronTools(tctx, service, { getDshIm, getAgentPresets }) - tctx.logger?.info?.('[dsh-ops-cron] tools cron_create/list/pause/resume/delete registered') + tctx.logger?.info?.('[dsh-ops-cron] tools cron_create/list/query/runs/progress/pause/resume/delete registered') }) ctx.inject(['skills'], (sctx) => { diff --git a/lib/monitor.js b/lib/monitor.js new file mode 100644 index 0000000..a93633e --- /dev/null +++ b/lib/monitor.js @@ -0,0 +1,389 @@ +/** + * Monitor extensions: labels, watch, progress, report, job state projection. + * Pure helpers — no Cordis / fs side effects except readProgressSnapshot. + */ + +import { readFile } from 'node:fs/promises' +import { isAbsolute, normalize, resolve, sep } from 'node:path' +import { ACTIVE_RUN_STATUSES } from './fire-status.js' + +export const LIFECYCLE_STATES = Object.freeze([ + 'idle', + 'pending', + 'running', + 'succeeded', + 'failed', + 'paused', +]) + +/** Map internal run.status → public lifecycle state. */ +export function runStatusToState(status) { + if (status === 'queued') return 'pending' + if (status === 'running') return 'running' + if (status === 'succeeded') return 'succeeded' + if (status === 'failed' || status === 'skipped') return 'failed' + return 'idle' +} + +/** + * Normalize string→string label map. Empty keys dropped; values coerced to string. + * @param {unknown} input + * @returns {Record} + */ +export function normalizeLabels(input) { + if (!input || typeof input !== 'object' || Array.isArray(input)) return {} + const out = {} + for (const [rawKey, rawVal] of Object.entries(input)) { + const key = String(rawKey || '').trim() + if (!key) continue + if (rawVal === undefined || rawVal === null) continue + out[key] = String(rawVal).trim() + } + return out +} + +/** + * @param {Record} jobLabels + * @param {Record} query + * @param {'any'|'all'} match + */ +export function labelsMatch(jobLabels, query, match = 'all') { + const labels = normalizeLabels(jobLabels) + const want = normalizeLabels(query) + const keys = Object.keys(want) + if (!keys.length) return true + if (match === 'any') { + return keys.some((key) => labels[key] === want[key]) + } + return keys.every((key) => labels[key] === want[key]) +} + +/** + * Watch declaration: taskId and/or labels (at least one). + * @param {unknown} input + */ +export function normalizeWatch(input) { + if (input == null || input === false) return null + if (typeof input !== 'object' || Array.isArray(input)) { + const error = new Error('watch must be an object with taskId and/or labels') + error.code = 'INVALID_WATCH' + throw error + } + const taskId = typeof input.taskId === 'string' + ? input.taskId.trim() + : (typeof input.task_id === 'string' ? input.task_id.trim() : '') + const labels = normalizeLabels(input.labels) + const match = input.match === 'any' ? 'any' : 'all' + const timeoutMin = Number(input.timeout_min ?? input.timeoutMin) + if (!taskId && !Object.keys(labels).length) { + const error = new Error('watch requires taskId or labels') + error.code = 'INVALID_WATCH' + throw error + } + const watch = { match } + if (taskId) watch.taskId = taskId + if (Object.keys(labels).length) watch.labels = labels + if (Number.isFinite(timeoutMin) && timeoutMin > 0) { + watch.timeoutMin = Math.min(24 * 60, Math.round(timeoutMin)) + } + return watch +} + +/** + * @param {unknown} input + */ +export function normalizeProgress(input) { + if (input == null || input === false) return null + if (typeof input !== 'object' || Array.isArray(input)) { + const error = new Error('progress must be an object') + error.code = 'INVALID_PROGRESS' + throw error + } + const channelSrc = input.channel && typeof input.channel === 'object' ? input.channel : {} + const file = typeof channelSrc.file === 'string' ? channelSrc.file.trim() : '' + const kind = typeof channelSrc.kind === 'string' + ? channelSrc.kind.trim().toLowerCase() + : (file ? 'file' : '') + if (kind && kind !== 'file' && kind !== 'dsh' && kind !== 'im' && kind !== 'grpc') { + const error = new Error('progress.channel.kind must be file|dsh|im|grpc') + error.code = 'INVALID_PROGRESS' + throw error + } + if (kind === 'file' && !file) { + const error = new Error('progress.channel.file is required when kind=file') + error.code = 'INVALID_PROGRESS' + throw error + } + const metricsIn = Array.isArray(input.metrics) ? input.metrics : [] + const metrics = [] + for (const row of metricsIn) { + if (!row || typeof row !== 'object') continue + const key = String(row.key || '').trim() + if (!key) continue + metrics.push({ + key, + label: typeof row.label === 'string' ? row.label.trim() : key, + unit: typeof row.unit === 'string' ? row.unit.trim() : '', + ...typeof row.derived === 'string' && row.derived.trim() + ? { derived: row.derived.trim() } + : {}, + }) + } + const progress = { metrics } + if (kind || file) { + progress.channel = { + kind: kind || 'file', + ...file ? { file } : {}, + } + } + return progress +} + +/** + * @param {unknown} input + * @param {object|null} [jobDelivery] fallback delivery when report.delivery omitted + */ +export function normalizeReport(input, jobDelivery = null) { + if (input == null || input === false) return null + if (typeof input !== 'object' || Array.isArray(input)) { + const error = new Error('report must be an object') + error.code = 'INVALID_REPORT' + throw error + } + const modeRaw = String(input.mode || 'on_complete').trim().toLowerCase() + const mode = modeRaw === 'on_change' || modeRaw === 'periodic' || modeRaw === 'on_complete' + ? modeRaw + : 'on_complete' + let deliveryKind = 'none' + if (input.delivery === 'none' || input.delivery === false) { + deliveryKind = 'none' + } else if (input.delivery === 'dsh' || input.delivery === 'im') { + deliveryKind = input.delivery + } else if (input.delivery && typeof input.delivery === 'object') { + deliveryKind = String(input.delivery.kind || 'none').toLowerCase() === 'im' ? 'im' : 'dsh' + } else if (jobDelivery?.kind === 'im') { + deliveryKind = 'im' + } else if (jobDelivery?.kind === 'dsh') { + deliveryKind = 'dsh' + } + const intervalMin = Number(input.interval_min ?? input.intervalMin) + const report = { mode, delivery: deliveryKind } + if (mode === 'periodic' && Number.isFinite(intervalMin) && intervalMin > 0) { + report.intervalMin = Math.min(24 * 60, Math.round(intervalMin)) + } + return report +} + +/** + * Persist history policy. + * - forever: keep all terminal runs for this job (no prune of its rows) + * - retain:N: keep last N terminals (default from settings) + * - archive:endpoint: treat as forever locally + stash endpoint for future archive job + * @param {unknown} input + * @param {number} [defaultRetain] + */ +export function normalizePersistHistory(input, defaultRetain = 200) { + if (input == null || input === '') { + return { kind: 'retain', limit: defaultRetain } + } + if (typeof input === 'object' && !Array.isArray(input)) { + const kind = String(input.kind || '').trim().toLowerCase() + if (kind === 'forever') return { kind: 'forever' } + if (kind === 'archive') { + const endpoint = typeof input.endpoint === 'string' ? input.endpoint.trim() : '' + return { kind: 'archive', endpoint: endpoint || '' } + } + if (kind === 'retain') { + const limit = Number(input.limit ?? input.n) + return { + kind: 'retain', + limit: Number.isInteger(limit) && limit >= 10 ? Math.min(50_000, limit) : defaultRetain, + } + } + } + const raw = String(input).trim().toLowerCase() + if (raw === 'forever') return { kind: 'forever' } + const retain = raw.match(/^retain:(\d+)$/) + if (retain) { + const limit = Number(retain[1]) + return { + kind: 'retain', + limit: Number.isInteger(limit) && limit >= 10 ? Math.min(50_000, limit) : defaultRetain, + } + } + const archive = raw.match(/^archive:(.+)$/) + if (archive) { + return { kind: 'archive', endpoint: archive[1].trim() } + } + const error = new Error('persist_history must be forever | retain:N | archive:endpoint') + error.code = 'INVALID_PERSIST_HISTORY' + throw error +} + +export function persistHistoryLimit(policy, settingsLimit = 200) { + if (!policy || policy.kind === 'forever' || policy.kind === 'archive') return Number.POSITIVE_INFINITY + if (policy.kind === 'retain' && Number.isInteger(policy.limit)) return policy.limit + return settingsLimit +} + +/** + * Project job + runs into a lifecycle state for listeners. + * @param {object} job + * @param {object[]} [runs] runs for this job (newest-first or any order) + * @param {number} [now] + */ +export function projectJobState(job, runs = [], now = Date.now()) { + if (!job) { + return { state: 'idle', stateEnteredAt: null, activeRunId: null, stuck: false } + } + if (job.enabled === false) { + return { + state: 'paused', + stateEnteredAt: job.updatedAt || job.createdAt || null, + activeRunId: null, + stuck: false, + } + } + const list = Array.isArray(runs) ? runs : [] + const active = list.find((run) => run && ACTIVE_RUN_STATUSES.has(run.status)) + if (active) { + const state = runStatusToState(active.status) + const entered = active.stateEnteredAt || active.actualAt || active.scheduledAt || null + const timeoutMs = Math.max(1, Number(job.timeoutMinutes) || 10) * 60_000 + const stuck = state === 'running' + && Number.isFinite(entered) + && (now - entered) > timeoutMs + return { state, stateEnteredAt: entered, activeRunId: active.id, stuck } + } + const oneshotDone = job.schedule?.kind === 'at' && job.nextRunAt == null + if (oneshotDone && (job.lastStatus === 'succeeded' || job.lastStatus === 'failed' || job.lastStatus === 'skipped')) { + const last = list.find((run) => run && !ACTIVE_RUN_STATUSES.has(run.status)) + || null + const state = job.lastStatus === 'succeeded' ? 'succeeded' : 'failed' + return { + state, + stateEnteredAt: last?.exitedAt || last?.stateEnteredAt || last?.actualAt || job.lastRunAt || null, + activeRunId: null, + stuck: false, + } + } + // Recurring (or one-shot waiting): idle between fires. + return { + state: 'idle', + stateEnteredAt: job.lastRunAt || job.updatedAt || job.createdAt || null, + activeRunId: null, + stuck: false, + } +} + +/** + * Jobs matching watch.taskId and/or watch.labels. + * @param {object[]} jobs + * @param {{ taskId?: string, labels?: Record, match?: 'any'|'all' }} by + */ +export function selectJobsByWatch(jobs, by = {}) { + const list = Array.isArray(jobs) ? jobs : [] + const taskId = typeof by.taskId === 'string' ? by.taskId.trim() : '' + const labels = normalizeLabels(by.labels) + const match = by.match === 'any' ? 'any' : 'all' + let out = list + if (taskId) out = out.filter((job) => job?.id === taskId) + if (Object.keys(labels).length) { + out = out.filter((job) => labelsMatch(job?.labels, labels, match)) + } + return out +} + +/** + * Label-scoped single-runner: another job sharing ALL of `guardLabels` has active run. + * Empty guardLabels → no cross-job guard. + * @param {object} state + * @param {object} job + * @param {Record} [guardLabels] defaults to job.labels when role+task present + */ +export function findLabelOverlapRun(state, job, guardLabels = null) { + const labels = normalizeLabels(guardLabels || job?.labels) + // Only enforce when both role and task are set (worker class), to avoid over-blocking. + if (!labels.role || !labels.task) return null + const jobs = (state?.jobs || []).filter((row) => row && row.id !== job.id && labelsMatch(row.labels, { + role: labels.role, + task: labels.task, + }, 'all')) + if (!jobs.length) return null + const jobIds = new Set(jobs.map((row) => row.id)) + return (state?.runs || []).find((run) => ( + run + && jobIds.has(run.jobId) + && ACTIVE_RUN_STATUSES.has(run.status) + )) || null +} + +/** + * Resolve progress file path under job cwd when relative. + * Rejects path escape outside cwd for relative paths; absolute paths allowed when under cwd or allowAbsolute. + */ +export function resolveProgressFilePath(job, filePath) { + const raw = String(filePath || '').trim() + if (!raw) return null + const cwd = String(job?.cwd || '').trim() + const resolved = isAbsolute(raw) ? normalize(raw) : resolve(cwd || process.cwd(), raw) + if (cwd) { + const root = normalize(resolve(cwd)) + const target = normalize(resolved) + const prefix = root.endsWith(sep) ? root : `${root}${sep}` + if (target !== root && !target.startsWith(prefix)) { + // Absolute path outside cwd: still allow if job declared it explicitly (channel.file absolute). + if (!isAbsolute(raw)) { + const error = new Error('progress file must stay under job cwd') + error.code = 'INVALID_PROGRESS_PATH' + throw error + } + } + } + return resolved +} + +/** + * Read and lightly validate a progress snapshot JSON file. + * @param {object} job + * @returns {Promise} + */ +export async function readProgressSnapshot(job) { + const file = job?.progress?.channel?.file + if (!file) return null + const path = resolveProgressFilePath(job, file) + if (!path) return null + try { + const raw = await readFile(path, 'utf8') + const parsed = JSON.parse(raw) + if (!parsed || typeof parsed !== 'object') return null + return { + taskId: parsed.taskId || job.id, + state: parsed.state || null, + updatedAt: parsed.updatedAt || null, + metrics: parsed.metrics && typeof parsed.metrics === 'object' ? parsed.metrics : {}, + path, + } + } catch (error) { + if (error && (error.code === 'ENOENT' || error instanceof SyntaxError)) { + return { taskId: job.id, state: null, updatedAt: null, metrics: {}, path, missing: true, error: error.message } + } + throw error + } +} + +/** + * Format a short progress / state report line for IM or mirror. + */ +export function formatWatchReport(job, targets, progressById = {}) { + const lines = [`[cron watch] ${job?.name || job?.id || 'watcher'}`] + for (const target of targets || []) { + const prog = progressById[target.id] + const stuck = target.stuck ? ' STUCK' : '' + const metrics = prog?.metrics && Object.keys(prog.metrics).length + ? ` metrics=${JSON.stringify(prog.metrics)}` + : '' + lines.push(`- ${target.name || target.id}: ${target.state}${stuck}${metrics}`) + } + return lines.join('\n') +} diff --git a/lib/store.js b/lib/store.js index 88f7ef4..d12b552 100644 --- a/lib/store.js +++ b/lib/store.js @@ -11,6 +11,15 @@ import { nextFire, validateSchedule } from './scheduler.js' import { normalizeDelivery, normalizeOrigin } from './delivery.js' import { normalizeAgentPresetId } from './preset.js' import { normalizeOwnerEmpNo, UNASSIGNED_OWNER } from './ownership.js' +import { + normalizeLabels, + normalizePersistHistory, + normalizeProgress, + normalizeReport, + normalizeWatch, + persistHistoryLimit, +} from './monitor.js' +import { ACTIVE_RUN_STATUSES } from './fire-status.js' export { UNASSIGNED_OWNER } @@ -116,6 +125,16 @@ export function createJobRecord(input, state, now) { const ownerDisplayName = typeof input.ownerDisplayName === 'string' ? input.ownerDisplayName.trim().slice(0, 80) : '' + const labels = normalizeLabels(input.labels) + const watch = input.watch !== undefined ? normalizeWatch(input.watch) : null + const progress = input.progress !== undefined ? normalizeProgress(input.progress) : null + const report = input.report !== undefined + ? normalizeReport(input.report, delivery) + : null + const persistHistory = normalizePersistHistory( + input.persistHistory ?? input.persist_history, + settings.historyLimit, + ) const job = { id: String(input.id || newId()), name, @@ -133,6 +152,11 @@ export function createJobRecord(input, state, now) { mirrorToSession, ownerEmpNo, ownerDisplayName, + labels, + ...watch ? { watch } : {}, + ...progress ? { progress } : {}, + ...report ? { report } : {}, + persistHistory, ...origin ? { origin } : {}, schedule: schedule.kind === 'cron' ? { kind: 'cron', expr: schedule.expr, timezone: schedule.timezone } @@ -157,17 +181,47 @@ export function createJobRecord(input, state, now) { return job } -export function pruneRuns(runs, historyLimit) { +export function pruneRuns(runs, historyLimit, jobs = []) { const list = Array.isArray(runs) ? runs : [] - const limit = Number.isInteger(historyLimit) ? historyLimit : DEFAULT_SETTINGS.historyLimit - const active = [] - const terminal = [] - for (const run of list) { - if (run && (run.status === 'queued' || run.status === 'running')) active.push(run) - else terminal.push(run) + const defaultLimit = Number.isInteger(historyLimit) ? historyLimit : DEFAULT_SETTINGS.historyLimit + const byJob = new Map() + for (const job of jobs || []) { + if (job?.id) byJob.set(job.id, job) } - terminal.sort((a, b) => (b.actualAt || b.scheduledAt || 0) - (a.actualAt || a.scheduledAt || 0)) - return [...active, ...terminal.slice(0, Math.max(10, limit))] + const active = [] + /** @type {Map} */ + const terminalsByJob = new Map() + const orphanTerminal = [] + for (const run of list) { + if (!run) continue + if (ACTIVE_RUN_STATUSES.has(run.status)) { + active.push(run) + continue + } + const jobId = run.jobId || '' + if (!jobId || !byJob.has(jobId)) { + orphanTerminal.push(run) + continue + } + if (!terminalsByJob.has(jobId)) terminalsByJob.set(jobId, []) + terminalsByJob.get(jobId).push(run) + } + const keptTerminal = [] + for (const [jobId, rows] of terminalsByJob) { + const job = byJob.get(jobId) + const policy = job?.persistHistory + || normalizePersistHistory(undefined, defaultLimit) + const limit = persistHistoryLimit(policy, defaultLimit) + rows.sort((a, b) => (b.actualAt || b.scheduledAt || 0) - (a.actualAt || a.scheduledAt || 0)) + if (!Number.isFinite(limit)) { + keptTerminal.push(...rows) + } else { + keptTerminal.push(...rows.slice(0, Math.max(10, limit))) + } + } + orphanTerminal.sort((a, b) => (b.actualAt || b.scheduledAt || 0) - (a.actualAt || a.scheduledAt || 0)) + keptTerminal.push(...orphanTerminal.slice(0, Math.max(10, defaultLimit))) + return [...active, ...keptTerminal] } export function listJobs(state) { @@ -178,10 +232,49 @@ export function getJob(state, jobId) { return (state?.jobs || []).find((job) => job.id === jobId) || null } -export function listHistory(state, jobId) { +/** + * Filter history with optional state / time / cursor pagination. + * @returns {{ runs: object[], nextCursor: string|null, total: number }} + */ +export function listHistoryPage(state, jobId, opts = {}) { const runs = state?.runs || [] - const filtered = jobId ? runs.filter((run) => run.jobId === jobId) : runs - return [...filtered].sort((a, b) => (b.actualAt || b.scheduledAt || 0) - (a.actualAt || a.scheduledAt || 0)) + let filtered = jobId ? runs.filter((run) => run.jobId === jobId) : [...runs] + const stateFilter = typeof opts.state === 'string' ? opts.state.trim() : '' + if (stateFilter) { + filtered = filtered.filter((run) => { + if (stateFilter === 'pending') return run.status === 'queued' + if (stateFilter === 'running') return run.status === 'running' + if (stateFilter === 'succeeded') return run.status === 'succeeded' + if (stateFilter === 'failed') return run.status === 'failed' || run.status === 'skipped' + return run.status === stateFilter + }) + } + const from = Number(opts.from) + const to = Number(opts.to) + if (Number.isFinite(from)) { + filtered = filtered.filter((run) => (run.actualAt || run.scheduledAt || 0) >= from) + } + if (Number.isFinite(to)) { + filtered = filtered.filter((run) => (run.actualAt || run.scheduledAt || 0) <= to) + } + filtered.sort((a, b) => (b.actualAt || b.scheduledAt || 0) - (a.actualAt || a.scheduledAt || 0)) + const cursor = typeof opts.cursor === 'string' ? opts.cursor.trim() : '' + if (cursor) { + const idx = filtered.findIndex((run) => run.id === cursor) + if (idx >= 0) filtered = filtered.slice(idx + 1) + } + const limit = Number(opts.limit) + const take = Number.isInteger(limit) && limit > 0 ? Math.min(500, limit) : filtered.length + const page = filtered.slice(0, take) + const nextCursor = page.length && filtered.length > page.length + ? page[page.length - 1].id + : null + return { runs: page, nextCursor, total: filtered.length } +} + +/** @deprecated-compatible: returns sorted run array (no pagination). */ +export function listHistory(state, jobId) { + return listHistoryPage(state, jobId).runs } export function upsertJob(state, job) { @@ -201,7 +294,7 @@ export function removeJob(state, jobId) { export function appendRun(state, run) { const isolated = applyRunIsolation(state, run) - isolated.runs = pruneRuns(isolated.runs, state.settings?.historyLimit) + isolated.runs = pruneRuns(isolated.runs, state.settings?.historyLimit, isolated.jobs) return isolated } @@ -212,7 +305,7 @@ export function patchRun(state, runId, patch) { let next = { ...state, runs } const updated = runs.find((run) => run.id === runId) if (updated?.sessionId) next = recordHiddenSession(next, updated.sessionId) - next.runs = pruneRuns(next.runs, next.settings?.historyLimit) + next.runs = pruneRuns(next.runs, next.settings?.historyLimit, next.jobs) return next } @@ -309,14 +402,77 @@ export function createStore(options = {}) { } } +function hydrateJob(row, settings) { + if (!row || !row.id) return null + const labels = normalizeLabels(row.labels) + let watch = null + let progress = null + let report = null + try { + watch = row.watch ? normalizeWatch(row.watch) : null + } catch { + watch = null + } + try { + progress = row.progress ? normalizeProgress(row.progress) : null + } catch { + progress = null + } + try { + report = row.report ? normalizeReport(row.report, row.delivery) : null + } catch { + report = null + } + let persistHistory + try { + persistHistory = normalizePersistHistory( + row.persistHistory ?? row.persist_history, + settings?.historyLimit, + ) + } catch { + persistHistory = normalizePersistHistory(undefined, settings?.historyLimit) + } + const hydrated = { + ...row, + labels, + persistHistory, + } + if (watch) hydrated.watch = watch + else delete hydrated.watch + if (progress) hydrated.progress = progress + else delete hydrated.progress + if (report) hydrated.report = report + else delete hydrated.report + return hydrated +} + +function hydrateRun(row) { + if (!row || !row.id) return null + return { + ...row, + stateEnteredAt: row.stateEnteredAt ?? row.actualAt ?? row.scheduledAt ?? null, + exitedAt: row.exitedAt ?? ( + row.status && !ACTIVE_RUN_STATUSES.has(row.status) ? (row.actualAt || null) : null + ), + outputRef: row.outputRef || null, + } +} + function hydrate(parsed) { const base = emptyState() if (!parsed || typeof parsed !== 'object') return base + const settings = normalizeSettings(parsed.settings) + const jobs = Array.isArray(parsed.jobs) + ? parsed.jobs.map((row) => hydrateJob(row, settings)).filter(Boolean) + : [] + const runs = Array.isArray(parsed.runs) + ? parsed.runs.map(hydrateRun).filter(Boolean) + : [] return { version: STORE_VERSION, - settings: normalizeSettings(parsed.settings), - jobs: Array.isArray(parsed.jobs) ? parsed.jobs.filter((row) => row && row.id) : [], - runs: Array.isArray(parsed.runs) ? parsed.runs.filter((row) => row && row.id) : [], + settings, + jobs, + runs, hiddenSessionIds: Array.isArray(parsed.hiddenSessionIds) ? parsed.hiddenSessionIds.map(String) : [], diff --git a/lib/tools.js b/lib/tools.js index 113b587..aca0a90 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -38,6 +38,14 @@ const JOB_SCHEMA = { delivery: { type: 'object', additionalProperties: true }, origin: { type: 'object', additionalProperties: true }, mirrorToSession: { type: 'boolean' }, + labels: { type: 'object', additionalProperties: { type: 'string' } }, + watch: { type: 'object', additionalProperties: true }, + progress: { type: 'object', additionalProperties: true }, + report: { type: 'object', additionalProperties: true }, + persistHistory: { type: 'object', additionalProperties: true }, + state: { type: 'string' }, + stateEnteredAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + stuck: { type: 'boolean' }, schedule: { type: 'object', additionalProperties: true }, createdAt: { type: 'number' }, updatedAt: { type: 'number' }, @@ -48,6 +56,90 @@ const JOB_SCHEMA = { }, } +const RUN_SCHEMA = { + type: 'object', + additionalProperties: true, + properties: { + id: { type: 'string' }, + run_id: { type: 'string' }, + jobId: { type: 'string' }, + status: { type: 'string' }, + state: { type: 'string' }, + stateEnteredAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + exitedAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + }, +} + +function parseLabelsArg(raw) { + if (raw == null || raw === '') return undefined + if (typeof raw === 'object' && !Array.isArray(raw)) { + const out = {} + for (const [k, v] of Object.entries(raw)) { + if (k) out[String(k)] = String(v ?? '') + } + return out + } + if (typeof raw === 'string') { + const trimmed = raw.trim() + if (!trimmed) return undefined + try { + const parsed = JSON.parse(trimmed) + if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) { + return parseLabelsArg(parsed) + } + } catch { + // key=value,key2=value2 + const out = {} + for (const part of trimmed.split(',')) { + const idx = part.indexOf('=') + if (idx <= 0) continue + out[part.slice(0, idx).trim()] = part.slice(idx + 1).trim() + } + return Object.keys(out).length ? out : undefined + } + } + return undefined +} + +function parseJsonObjectArg(raw, fieldName) { + if (raw == null || raw === '' || raw === false) return undefined + if (typeof raw === 'object' && !Array.isArray(raw)) return raw + if (typeof raw === 'string') { + try { + const parsed = JSON.parse(raw.trim()) + if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) return parsed + } catch { + const error = new Error(`${fieldName} must be a JSON object`) + error.code = 'INVALID_JOB' + throw error + } + } + const error = new Error(`${fieldName} must be an object`) + error.code = 'INVALID_JOB' + throw error +} + +function monitorFieldsFromArgs(args = {}) { + const out = {} + const labels = parseLabelsArg(args.labels) + if (labels) out.labels = labels + if (args.watch !== undefined) { + const watch = parseJsonObjectArg(args.watch, 'watch') + out.watch = watch === undefined ? null : watch + } + if (args.progress !== undefined) { + const progress = parseJsonObjectArg(args.progress, 'progress') + out.progress = progress === undefined ? null : progress + } + if (args.report !== undefined) { + const report = parseJsonObjectArg(args.report, 'report') + out.report = report === undefined ? null : report + } + if (args.persist_history !== undefined) out.persist_history = args.persist_history + if (args.persistHistory !== undefined) out.persistHistory = args.persistHistory + return out +} + function text(value) { return [{ type: 'text', text: String(value || '') }] } @@ -109,13 +201,17 @@ function jobLine(job) { if (!job) return '' const tz = job.schedule?.timezone || 'Asia/Shanghai' const when = job.nextRunAt ? formatInZone(job.nextRunAt, tz) : 'n/a' - const state = job.enabled === false ? 'paused' : 'enabled' + const life = job.state || (job.enabled === false ? 'paused' : 'enabled') const sched = job.schedule?.kind === 'at' ? `at ${job.schedule.at}` : (job.schedule?.expr || 'cron') const model = job.provider && job.model ? `${job.provider}/${job.model}` : 'default-model' const delivery = deliveryLine(job.delivery) - return `${job.name} [${state}] ${sched} tz=${tz} cwd=${job.cwd || '(recent workspace)'} model=${model} ${delivery} next=${when} id=${job.id}` + const labels = job.labels && Object.keys(job.labels).length + ? ` labels=${JSON.stringify(job.labels)}` + : '' + const stuck = job.stuck ? ' STUCK' : '' + return `${job.name} [${life}${stuck}] ${sched} tz=${tz} cwd=${job.cwd || '(recent workspace)'} model=${model} ${delivery}${labels} next=${when} id=${job.id}` } export function callerWorkingDirectory(exec) { @@ -292,6 +388,11 @@ export function cronToolDefinitions(service, deps = {}) { im_target_id: { type: 'string', description: 'Opaque targetId from IM 投递设置 when delivery=im (Web/sidebar only; IM chats cannot retarget).' }, agent_preset: { type: 'string', description: 'Agent preset id for scheduled runs. Omit to inherit from the current WhatsApp/IM chat or Host default.' }, mirror_to_session: { type: 'boolean', description: 'If true, after each fire append the run summary into the origin WhatsApp/Web session as a system-reminder (no new model turn). Default false — results stay in run history (and IM delivery when configured).' }, + labels: { type: 'object', additionalProperties: { type: 'string' }, description: 'String labels for discovery/monitoring, e.g. {"role":"worker","task":"theory","owner":"alice"}.' }, + watch: { type: 'object', additionalProperties: true, description: 'Listener binding: {"taskId":"..."} and/or {"labels":{...},"match":"any|all","timeout_min":30}.' }, + progress: { type: 'object', additionalProperties: true, description: 'Progress declaration: {"channel":{"file":"progress.json"},"metrics":[{"key":"done","label":"已完成"}]}.' }, + report: { type: 'object', additionalProperties: true, description: 'Watcher report: {"mode":"on_change|on_complete|periodic","delivery":"dsh|im|none","interval_min":15}.' }, + persist_history: { type: 'string', description: 'History policy: forever | retain:N | archive:endpoint. Default retain from settings.' }, }, required: ['name', 'prompt'], }, @@ -346,6 +447,7 @@ export function cronToolDefinitions(service, deps = {}) { mirrorToSession: args.mirror_to_session === true, ...origin ? { origin } : {}, agentPreset, + ...monitorFieldsFromArgs(args), }, identity, { fromImPeer: !!peer?.botId }) return { job } } catch (error) { @@ -362,12 +464,15 @@ export function cronToolDefinitions(service, deps = {}) { }, { name: 'cron_list', - description: 'List DSH 定时任务 jobs (sidebar scheduled jobs). Use this whenever the user asks what scheduled tasks exist. NEVER use crontab -l or /etc/cron* — those are OS crontabs, not this plugin.', + description: 'List DSH 定时任务 jobs (sidebar scheduled jobs). Filter by enabled_only, state (idle|pending|running|succeeded|failed|paused), labels, or task_id. NEVER use crontab -l.', parameters: { type: 'object', additionalProperties: false, properties: { enabled_only: { type: 'boolean', description: 'If true, omit paused jobs.' }, + state: { type: 'string', description: 'Lifecycle filter: idle|pending|running|succeeded|failed|paused.' }, + task_id: { type: 'string', description: 'Exact job id.' }, + labels: { type: 'object', additionalProperties: { type: 'string' }, description: 'Label filter (all keys must match unless used with cron_query match=any).' }, }, }, output: { @@ -390,12 +495,170 @@ export function cronToolDefinitions(service, deps = {}) { aborted(exec) const peer = await resolveCallerPeer(exec, getDshIm()) const identity = resolveToolIdentity(exec, service) - let jobs = await service.listJobs(identity) + const query = { + ...parseLabelsArg(args?.labels) ? { labels: parseLabelsArg(args.labels) } : {}, + ...args?.task_id ? { taskId: String(args.task_id).trim() } : {}, + ...args?.state ? { state: String(args.state).trim() } : {}, + } + let jobs = await service.listJobs(identity, Object.keys(query).length ? query : null) if (peer?.botId) jobs = jobs.filter((job) => jobVisibleToPeer(job, peer)) if (args?.enabled_only === true) jobs = jobs.filter((job) => job.enabled !== false) return { jobs, count: jobs.length } }, }, + { + name: 'cron_query', + description: 'Listener/monitor query: resolve jobs by taskId or labels (match any|all), optionally include run history and decode progress snapshots from progress.channel.file.', + parameters: { + type: 'object', + additionalProperties: false, + properties: { + task_id: { type: 'string', description: 'Watch a single job id.' }, + labels: { type: 'object', additionalProperties: { type: 'string' }, description: 'Label selector.' }, + match: { type: 'string', description: 'Label match mode: all (default) or any.' }, + state: { type: 'string', description: 'Filter projected lifecycle state.' }, + include_runs: { type: 'boolean', description: 'Include run instances per job.' }, + decode_progress: { type: 'boolean', description: 'Read progress.channel.file snapshots.' }, + }, + }, + output: { + schema: { + type: 'object', + additionalProperties: true, + properties: { + jobs: { type: 'array', items: JOB_SCHEMA }, + items: { type: 'array', items: { type: 'object', additionalProperties: true } }, + count: { type: 'integer' }, + }, + }, + render: (_args, value) => { + const rows = (value.jobs || []).map(jobLine) + return text(rows.length + ? `cron_query (${value.count}):\n${rows.join('\n')}` + : 'cron_query: no matching jobs.') + }, + }, + presentCall: () => ({ card: 'generic', title: 'cron_query' }), + async execute(args, exec) { + aborted(exec) + const peer = await resolveCallerPeer(exec, getDshIm()) + const identity = resolveToolIdentity(exec, service) + const result = await service.queryJobs({ + taskId: args?.task_id, + labels: parseLabelsArg(args?.labels), + match: args?.match === 'any' ? 'any' : 'all', + state: args?.state, + includeRuns: args?.include_runs === true, + decodeProgress: args?.decode_progress === true, + }, identity) + if (peer?.botId) { + result.jobs = result.jobs.filter((job) => jobVisibleToPeer(job, peer)) + result.items = (result.items || []).filter((item) => jobVisibleToPeer(item.job, peer)) + result.count = result.jobs.length + } + return result + }, + }, + { + name: 'cron_runs', + description: 'List run instances for a task (history). Supports state/from/to/limit/cursor. History defaults to retain:N; forever jobs keep all rows.', + parameters: { + type: 'object', + additionalProperties: false, + properties: { + task_id: { type: 'string', description: 'Job id.' }, + state: { type: 'string', description: 'Filter: pending|running|succeeded|failed.' }, + from: { type: 'number', description: 'Epoch ms lower bound.' }, + to: { type: 'number', description: 'Epoch ms upper bound.' }, + limit: { type: 'integer', description: 'Page size (max 500).' }, + cursor: { type: 'string', description: 'Pagination cursor (previous run id).' }, + }, + required: ['task_id'], + }, + output: { + schema: { + type: 'object', + additionalProperties: true, + properties: { + task_id: { type: 'string' }, + runs: { type: 'array', items: RUN_SCHEMA }, + next_cursor: { oneOf: [{ type: 'string' }, { type: 'null' }] }, + total: { type: 'integer' }, + }, + }, + render: (_args, value) => { + const rows = (value.runs || []).map((run) => ( + `${run.run_id || run.id} ${run.state || run.status} entered=${run.stateEnteredAt || '-'} exited=${run.exitedAt || '-'}` + )) + return text(rows.length + ? `Runs for ${value.task_id} (${value.total}):\n${rows.join('\n')}` + : `No runs for ${value.task_id}.`) + }, + }, + presentCall: (args) => ({ card: 'generic', title: 'cron_runs', content: String(args?.task_id || '') }), + async execute(args, exec) { + aborted(exec) + const peer = await resolveCallerPeer(exec, getDshIm()) + const identity = resolveToolIdentity(exec, service) + const taskId = String(args?.task_id || '').trim() + await requireOwnedJob(service, taskId, peer, identity) + return service.listRuns(taskId, { + state: args?.state, + from: args?.from, + to: args?.to, + limit: args?.limit, + cursor: args?.cursor, + }, identity) + }, + }, + { + name: 'cron_progress', + description: 'Read current progress snapshots for a task_id or label watch selector (progress.channel.file).', + parameters: { + type: 'object', + additionalProperties: false, + properties: { + task_id: { type: 'string', description: 'Job id.' }, + labels: { type: 'object', additionalProperties: { type: 'string' }, description: 'Label selector (by_watch).' }, + match: { type: 'string', description: 'any|all for labels.' }, + }, + }, + output: { + schema: { + type: 'object', + additionalProperties: true, + properties: { + snapshots: { type: 'array', items: { type: 'object', additionalProperties: true } }, + count: { type: 'integer' }, + }, + }, + render: (_args, value) => { + const rows = (value.snapshots || []).map((row) => { + const metrics = row.progress?.metrics + ? JSON.stringify(row.progress.metrics) + : (row.progress?.missing ? '(no file)' : '{}') + return `${row.job?.name || row.job?.id}: ${metrics}` + }) + return text(rows.length ? `Progress (${value.count}):\n${rows.join('\n')}` : 'No progress snapshots.') + }, + }, + presentCall: () => ({ card: 'generic', title: 'cron_progress' }), + async execute(args, exec) { + aborted(exec) + const peer = await resolveCallerPeer(exec, getDshIm()) + const identity = resolveToolIdentity(exec, service) + const result = await service.getProgress({ + taskId: args?.task_id, + labels: parseLabelsArg(args?.labels), + match: args?.match === 'any' ? 'any' : 'all', + }, identity) + if (peer?.botId) { + result.snapshots = (result.snapshots || []).filter((row) => jobVisibleToPeer(row.job, peer)) + result.count = result.snapshots.length + } + return result + }, + }, { name: 'cron_pause', description: 'Pause a scheduled task by id from cron_list / cron_create. It stays in 定时任务 but will not fire until resumed.', @@ -497,8 +760,9 @@ export function cronGuidanceText(nowMs = Date.now(), timeZone = 'Asia/Shanghai') 'Delivery: when chatting on WhatsApp/IM, omit delivery so the job defaults to im for the current chat (group→same group, DM→same DM). A 投递目标 is reused or auto-created. On Web/DSH, default is dsh. Or pass delivery=im with im_bot_id+im_target_id.', 'Session mirror: by default the run summary is NOT injected into the origin WhatsApp/Web chat. Pass mirror_to_session=true only when the user wants follow-up context in that chat. Full tool traces always stay in run history; IM delivery (when configured) is independent.', 'Agent preset: omit agent_preset to inherit (WhatsApp chat/group preset → creating session → Host default). Pass agent_preset to pin a preset for every scheduled run.', + 'Labels/monitor: pass labels={"role":"worker","task":"theory"} on workers; listeners pass watch={"taskId":"..."} or watch={"labels":{...},"match":"all"}. Use cron_query / cron_progress / cron_runs to poll state and progress.', 'When the user asks to look at, create, pause, resume, or delete 定时任务 / scheduled tasks / cron jobs:', - '1. If cron_list / cron_create / cron_pause / cron_resume / cron_delete are in your tool list, call them.', + '1. If cron_list / cron_create / cron_pause / cron_resume / cron_delete / cron_query / cron_runs / cron_progress are in your tool list, call them.', '2. If they are not listed, load skill "scheduled-tasks" via the skill tool (exact name), then retry the cron_* tools. Those tools are registered by the dsh-ops-cron plugin — they are not unlocked by skill_search/dev_tool_search (those names are obsolete).', '3. If cron_* still return unknown tool after loading the skill, the plugin failed to register (check Host logs for unsupported JSON schema). Tell the user; do not invent crontab workarounds.', '4. Never run crontab, never read /etc/cron*, and never say there are no tasks until cron_list has returned.', @@ -517,8 +781,11 @@ export function makeCronSkill() { content: `${cronGuidanceText()} Tools: -- cron_list — list jobs -- cron_create — "一分钟后" → after_minutes=1. Clock time → hour+minute only. Do not send at/expr at the same time. Pass cwd as the current workspace path when creating from a workspace chat. Pass provider+model or inherit the current session model. Delivery and agent preset auto from session, or pass delivery=im / agent_preset explicitly. Session mirror is off by default; pass mirror_to_session=true only if the user wants the summary injected into the origin chat. +- cron_list — list jobs (optional state/labels/task_id filters) +- cron_query — listener query by taskId or labels; optional include_runs / decode_progress +- cron_runs — run history for a task_id +- cron_progress — read progress.channel.file snapshots +- cron_create — "一分钟后" → after_minutes=1. Clock time → hour+minute only. Do not send at/expr at the same time. Pass cwd as the current workspace path when creating from a workspace chat. Pass provider+model or inherit the current session model. Delivery and agent preset auto from session, or pass delivery=im / agent_preset explicitly. Session mirror is off by default; pass mirror_to_session=true only if the user wants the summary injected into the origin chat. Optional labels/watch/progress/report/persist_history for monitor workers and listeners. - cron_pause / cron_resume / cron_delete — by id from cron_list `, } diff --git a/test/monitor.test.js b/test/monitor.test.js new file mode 100644 index 0000000..d7db75c --- /dev/null +++ b/test/monitor.test.js @@ -0,0 +1,219 @@ +import assert from 'node:assert/strict' +import { mkdtemp, writeFile, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { test } from 'node:test' +import { + labelsMatch, + normalizeLabels, + normalizePersistHistory, + normalizeProgress, + normalizeReport, + normalizeWatch, + projectJobState, + readProgressSnapshot, + selectJobsByWatch, +} from '../lib/monitor.js' +import { createJobRecord, emptyState, listHistoryPage, pruneRuns } from '../lib/store.js' +import { claimOccurrence, publicJob } from '../lib/fire.js' +import { createHostService } from '../lib/host.js' +import { cronToolDefinitions } from '../lib/tools.js' + +test('labelsMatch any/all', () => { + const job = { role: 'worker', task: 'theory', owner: 'alice' } + assert.equal(labelsMatch(job, { role: 'worker', task: 'theory' }, 'all'), true) + assert.equal(labelsMatch(job, { role: 'worker', owner: 'bob' }, 'all'), false) + assert.equal(labelsMatch(job, { role: 'worker', owner: 'bob' }, 'any'), true) +}) + +test('normalize watch / progress / report / persist_history', () => { + assert.deepEqual(normalizeWatch({ taskId: 'abc', timeout_min: 30 }).taskId, 'abc') + assert.equal(normalizeWatch({ labels: { task: 'theory' }, match: 'any' }).match, 'any') + assert.throws(() => normalizeWatch({}), { code: 'INVALID_WATCH' }) + const progress = normalizeProgress({ + channel: { file: '进度.json' }, + metrics: [{ key: 'done', label: '已完成', unit: '个' }], + }) + assert.equal(progress.channel.kind, 'file') + assert.equal(progress.metrics[0].key, 'done') + assert.equal(normalizeReport({ mode: 'periodic', delivery: 'im', interval_min: 5 }).intervalMin, 5) + assert.equal(normalizePersistHistory('forever').kind, 'forever') + assert.equal(normalizePersistHistory('retain:50').limit, 50) + assert.equal(normalizePersistHistory('archive:https://example/archive').endpoint, 'https://example/archive') +}) + +test('projectJobState maps run statuses', () => { + const job = { + id: 'j1', + enabled: true, + timeoutMinutes: 10, + schedule: { kind: 'cron', expr: '0 * * * *', timezone: 'Asia/Shanghai' }, + createdAt: 1, + updatedAt: 1, + } + assert.equal(projectJobState(job, []).state, 'idle') + assert.equal(projectJobState(job, [{ id: 'r1', status: 'queued', stateEnteredAt: 10 }]).state, 'pending') + assert.equal(projectJobState(job, [{ id: 'r1', status: 'running', stateEnteredAt: 10 }]).state, 'running') + assert.equal(projectJobState({ ...job, enabled: false }, []).state, 'paused') + const stuck = projectJobState( + job, + [{ id: 'r1', status: 'running', stateEnteredAt: Date.now() - 2 * 60 * 60_000 }], + Date.now(), + ) + assert.equal(stuck.stuck, true) +}) + +test('createJobRecord stores monitor fields; publicJob exposes state', () => { + const state = emptyState() + const now = Date.parse('2026-09-10T07:00:00.000Z') + const job = createJobRecord({ + name: 'worker', + prompt: 'work', + schedule: { kind: 'at', at: new Date(now + 60_000).toISOString(), timezone: 'Asia/Shanghai' }, + labels: { role: 'worker', task: 'theory' }, + progress: { channel: { file: 'p.json' }, metrics: [{ key: 'done' }] }, + persist_history: 'forever', + }, state, now) + assert.deepEqual(job.labels, { role: 'worker', task: 'theory' }) + assert.equal(job.persistHistory.kind, 'forever') + assert.equal(job.progress.channel.file, 'p.json') + const view = publicJob(job, []) + assert.equal(view.state, 'idle') + assert.deepEqual(view.labels, job.labels) +}) + +test('label overlap skips second worker', () => { + const now = Date.parse('2026-09-10T07:00:00.000Z') + let state = emptyState() + const a = createJobRecord({ + name: 'w1', + prompt: 'a', + schedule: { kind: 'at', at: new Date(now + 60_000).toISOString(), timezone: 'Asia/Shanghai' }, + labels: { role: 'worker', task: 'theory' }, + }, state, now) + const b = createJobRecord({ + name: 'w2', + prompt: 'b', + schedule: { kind: 'at', at: new Date(now + 60_000).toISOString(), timezone: 'Asia/Shanghai' }, + labels: { role: 'worker', task: 'theory' }, + }, state, now) + state = { ...state, jobs: [a, b], runs: [{ + id: 'run-a', + jobId: a.id, + status: 'running', + scheduledAt: now, + actualAt: now, + stateEnteredAt: now, + }] } + const claimed = claimOccurrence(state, b.id, now, 'run-now', state.settings) + assert.equal(claimed.decision.action, 'skip') + assert.equal(claimed.decision.reason, 'label_overlap') +}) + +test('pruneRuns forever keeps all terminals for that job', () => { + const jobs = [{ + id: 'j1', + persistHistory: { kind: 'forever' }, + }, { + id: 'j2', + persistHistory: { kind: 'retain', limit: 10 }, + }] + const runs = [] + for (let i = 0; i < 30; i++) { + runs.push({ id: `a${i}`, jobId: 'j1', status: 'succeeded', actualAt: i }) + runs.push({ id: `b${i}`, jobId: 'j2', status: 'succeeded', actualAt: i }) + } + const pruned = pruneRuns(runs, 10, jobs) + assert.equal(pruned.filter((r) => r.jobId === 'j1').length, 30) + assert.equal(pruned.filter((r) => r.jobId === 'j2').length, 10) +}) + +test('listHistoryPage paginates', () => { + const state = { + runs: [ + { id: 'r3', jobId: 'j', status: 'succeeded', actualAt: 30 }, + { id: 'r2', jobId: 'j', status: 'succeeded', actualAt: 20 }, + { id: 'r1', jobId: 'j', status: 'succeeded', actualAt: 10 }, + ], + } + const page1 = listHistoryPage(state, 'j', { limit: 2 }) + assert.equal(page1.runs.length, 2) + assert.equal(page1.nextCursor, 'r2') + const page2 = listHistoryPage(state, 'j', { limit: 2, cursor: page1.nextCursor }) + assert.equal(page2.runs.length, 1) + assert.equal(page2.runs[0].id, 'r1') +}) + +test('cron_query / cron_progress / cron_runs tools', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-cron-mon-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const progressPath = join(dir, 'progress.json') + await writeFile(progressPath, JSON.stringify({ + taskId: 'x', + state: 'running', + updatedAt: '2026-09-10T15:00:00+08:00', + metrics: { total: 10, done: 3, percent: 30 }, + }), 'utf8') + + const service = createHostService({ + filePath: join(dir, 'store.json'), + now: () => Date.parse('2026-09-10T07:00:00.000Z'), + sessionPort: { + async createAndPrompt() { + return { sessionId: 's1', status: 'succeeded', summary: 'ok' } + }, + }, + }) + const tools = Object.fromEntries(cronToolDefinitions(service).map((row) => [row.name, row])) + const created = await tools.cron_create.execute({ + name: 'theory-worker', + prompt: 'do theory', + after_minutes: 5, + labels: { role: 'worker', task: 'theory' }, + progress: { channel: { file: progressPath }, metrics: [{ key: 'done' }, { key: 'total' }] }, + persist_history: 'forever', + cwd: dir, + }, { agent: { session: { header: { cwd: dir } } } }) + assert.equal(created.job.labels.role, 'worker') + assert.equal(created.job.state, 'idle') + + const watcher = await tools.cron_create.execute({ + name: 'theory-watch', + prompt: 'watch workers via cron_query', + expr: '*/5 * * * *', + watch: { labels: { task: 'theory' }, match: 'all' }, + report: { mode: 'on_complete', delivery: 'none' }, + cwd: dir, + }, { agent: { session: { header: { cwd: dir } } } }) + assert.equal(watcher.job.watch.labels.task, 'theory') + + const queried = await tools.cron_query.execute({ + labels: { task: 'theory' }, + match: 'all', + decode_progress: true, + }, {}) + assert.ok(queried.count >= 1) + assert.ok(queried.items.some((item) => item.progress?.metrics?.done === 3)) + + const listed = await tools.cron_list.execute({ labels: { role: 'worker' }, state: 'idle' }, {}) + assert.ok(listed.jobs.some((job) => job.id === created.job.id)) + + await service.dispatchRun(created.job.id, 'run-now') + const runs = await tools.cron_runs.execute({ task_id: created.job.id, limit: 10 }, {}) + assert.ok(runs.runs.length >= 1) + assert.ok(runs.runs[0].stateEnteredAt) + + const progress = await tools.cron_progress.execute({ task_id: created.job.id }, {}) + assert.equal(progress.snapshots[0].progress.metrics.done, 3) + + assert.deepEqual( + selectJobsByWatch([created.job, watcher.job], { labels: { task: 'theory' } }).map((j) => j.name), + ['theory-worker'], + ) + assert.deepEqual( + selectJobsByWatch([created.job, watcher.job], watcher.job.watch).map((j) => j.name), + ['theory-worker'], + ) + assert.deepEqual(normalizeLabels({ a: 1 }), { a: '1' }) + assert.ok(await readProgressSnapshot(created.job)) +}) diff --git a/test/package.test.js b/test/package.test.js index 15d2fdd..8e97378 100644 --- a/test/package.test.js +++ b/test/package.test.js @@ -24,6 +24,9 @@ test('installable bundle declares host apply, client half, unique id, and no @de const tools = await readFile(join(root, 'lib/tools.js'), 'utf8') assert.match(tools, /name: 'cron_create'/) assert.match(tools, /name: 'cron_list'/) + assert.match(tools, /name: 'cron_query'/) + assert.match(tools, /name: 'cron_runs'/) + assert.match(tools, /name: 'cron_progress'/) assert.match(tools, /crontab/) assert.match(tools, /scheduled-tasks/) assert.match(tools, /hour/) diff --git a/test/tools.test.js b/test/tools.test.js index bb6bbad..ea77129 100644 --- a/test/tools.test.js +++ b/test/tools.test.js @@ -431,7 +431,16 @@ test('registerCronTools registers each definition and disposer unregisters', () }, } const off = registerCronTools(ctx, {}) - assert.deepEqual(registered, ['cron_create', 'cron_list', 'cron_pause', 'cron_resume', 'cron_delete']) + assert.deepEqual(registered, [ + 'cron_create', + 'cron_list', + 'cron_query', + 'cron_runs', + 'cron_progress', + 'cron_pause', + 'cron_resume', + 'cron_delete', + ]) off() assert.deepEqual(registered, []) }) @@ -440,6 +449,9 @@ test('cron tool output schemas never use type arrays (Host rejects them)', () => const defs = cronToolDefinitions({ async createJob() { return {} }, async listJobs() { return [] }, + async queryJobs() { return { jobs: [], items: [], count: 0 } }, + async listRuns() { return { task_id: '', runs: [], next_cursor: null, total: 0 } }, + async getProgress() { return { snapshots: [], count: 0 } }, async pauseJob() { return {} }, async resumeJob() { return {} }, async deleteJob() { return {} },