Add cron monitor APIs: labels, watch, progress, and query tools.

Extend jobs with lifecycle projection and history policies so listeners can poll state at scale without breaking existing cron_* flows.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-10 15:23:42 +08:00
parent 43facd608a
commit b6934b20a3
10 changed files with 1501 additions and 53 deletions

View file

@ -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:<url> 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()