diff --git a/lib/fire.js b/lib/fire.js index 7cbe2ba..7a1e751 100644 --- a/lib/fire.js +++ b/lib/fire.js @@ -10,6 +10,7 @@ import { appendRun, getJob, patchRun, upsertJob } from './store.js' import { ACTIVE_RUN_STATUSES } from './fire-status.js' import { findLabelOverlapRun, + isRetriggerable, normalizeLabels, projectJobState, } from './monitor.js' @@ -127,6 +128,7 @@ export function publicJob(job, runs = null) { state: stateInfo.state, stateEnteredAt: stateInfo.stateEnteredAt, stuck: stateInfo.stuck === true, + retriggerable: isRetriggerable(job, jobRuns || []), ...job.origin ? { origin: job.origin } : {}, schedule: job.schedule, createdAt: job.createdAt, diff --git a/lib/host.js b/lib/host.js index 2381e2e..3554d5e 100644 --- a/lib/host.js +++ b/lib/host.js @@ -11,6 +11,7 @@ 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 { ACTIVE_RUN_STATUSES } from './fire-status.js' import { formatWatchReport, labelsMatch, @@ -699,6 +700,101 @@ export function createHostService(options = {}) { return updateJob(jobId, { enabled }, identity) } + /** + * Re-arm a consumed one-shot job. + * afterMinutes omitted/0 → immediate run-now; >0 → reschedule next and wait for tick. + * Auto-enables paused jobs. Rejects cron jobs and in-flight/pending-next oneshots. + */ + async function retriggerJob(jobId, opts = {}, identity = null) { + const afterRaw = opts.afterMinutes ?? opts.after_minutes + const afterMinutes = afterRaw === undefined || afterRaw === null || afterRaw === '' + ? 0 + : Number(afterRaw) + if (!Number.isFinite(afterMinutes) || afterMinutes < 0) { + const error = new Error('after_minutes must be >= 0') + error.code = 'INVALID_RETRIGGER' + throw error + } + const delayMs = Math.round(afterMinutes * 60_000) + const t = now() + let mode = delayMs > 0 ? 'scheduled' : 'immediate' + + await withState((current) => { + const job = getJob(current, jobId) + if (!job) { + const error = new Error('job not found') + error.code = 'NOT_FOUND' + throw error + } + if (identity) assertCanAccessJob(job, identity) + if (job.schedule?.kind !== 'at') { + const error = new Error('cron_retrigger is only for one-shot (at) jobs; use pause/resume for recurring jobs') + error.code = 'INVALID_RETRIGGER' + throw error + } + const runs = (current.runs || []).filter((run) => run?.jobId === job.id) + if (runs.some((run) => ACTIVE_RUN_STATUSES.has(run.status))) { + const error = new Error('job already has a pending or running instance') + error.code = 'ALREADY_RUNNING' + throw error + } + if (job.nextRunAt != null) { + const error = new Error('one-shot still has a pending next run; wait or edit the schedule instead of retrigger') + error.code = 'INVALID_RETRIGGER' + throw error + } + const terminal = job.lastStatus === 'succeeded' + || job.lastStatus === 'failed' + || job.lastStatus === 'skipped' + if (!terminal) { + const error = new Error('one-shot has not finished a run yet; nothing to retrigger') + error.code = 'INVALID_RETRIGGER' + throw error + } + + if (delayMs > 0) { + const atMs = t + delayMs + return upsertJob(current, { + ...job, + enabled: true, + schedule: { + kind: 'at', + at: new Date(atMs).toISOString(), + timezone: job.schedule?.timezone || current.settings?.timezone || 'Asia/Shanghai', + }, + nextRunAt: atMs, + updatedAt: t, + }) + } + if (job.enabled === false) { + return upsertJob(current, { ...job, enabled: true, updatedAt: t }) + } + return current + }) + + if (mode === 'scheduled') { + const state = await snapshot() + const job = getJob(state, jobId) + return { + ok: true, + mode, + job: jobView(job, state.runs), + run: null, + nextRunAt: job?.nextRunAt ?? null, + } + } + + const result = await dispatchRun(jobId, 'run-now') + return { + ok: true, + mode, + job: result.job, + run: result.run, + decision: result.decision, + nextRunAt: result.job?.nextRunAt ?? null, + } + } + async function deleteJob(jobId) { await withState((current) => { if (!getJob(current, jobId)) { @@ -1001,7 +1097,7 @@ export function createHostService(options = {}) { return } - const jobMatch = path.match(new RegExp(`^${API_PREFIX}/jobs/([^/]+)(/run|/pause|/resume)?$`)) + const jobMatch = path.match(new RegExp(`^${API_PREFIX}/jobs/([^/]+)(/run|/pause|/resume|/retrigger)?$`)) if (jobMatch) { const jobId = decodeURIComponent(jobMatch[1]) const rest = jobMatch[2] || '' @@ -1030,6 +1126,12 @@ export function createHostService(options = {}) { write(200, { ok: true, ...result }) return } + if (method === 'POST' && rest === '/retrigger') { + const body = await readJsonBody(req).catch(() => ({})) + const result = await retriggerJob(jobId, body || {}, identity) + write(200, result) + return + } if (method === 'POST' && rest === '/pause') { const job = await pauseJob(jobId, false, identity) write(200, { ok: true, job }) @@ -1182,7 +1284,10 @@ 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' || code === 'INVALID_WATCH' || code === 'INVALID_PROGRESS' || code === 'INVALID_REPORT' || code === 'INVALID_PERSIST_HISTORY' || code === 'INVALID_PROGRESS_PATH' || code === 'INVALID_ARCHIVE' || code === 'ARCHIVE_FAILED' || code === 'ARCHIVE_UNAVAILABLE') { + if (code === 'ALREADY_RUNNING') { + return write(409, { ok: false, error: error.message, code, message: error.message }) + } + 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' || code === 'INVALID_ARCHIVE' || code === 'ARCHIVE_FAILED' || code === 'ARCHIVE_UNAVAILABLE' || code === 'INVALID_RETRIGGER') { 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)) @@ -1354,6 +1459,7 @@ export function createHostService(options = {}) { createJob, updateJob, pauseJob, + retriggerJob, deleteJob, updateSettings, dispatchRun, diff --git a/lib/i18n.js b/lib/i18n.js index 25b105d..f696213 100644 --- a/lib/i18n.js +++ b/lib/i18n.js @@ -31,6 +31,7 @@ export const MESSAGES = { 'tool.list': '列出定时任务', 'tool.pause': '暂停定时任务', 'tool.resume': '恢复定时任务', + 'tool.retrigger': '重新触发一次性任务', 'tool.delete': '删除定时任务', }, en: { @@ -51,6 +52,7 @@ export const MESSAGES = { 'tool.list': 'List scheduled tasks', 'tool.pause': 'Pause scheduled task', 'tool.resume': 'Resume scheduled task', + 'tool.retrigger': 'Retrigger one-shot task', 'tool.delete': 'Delete scheduled task', }, } diff --git a/lib/index.js b/lib/index.js index 9bf7197..ad09f4a 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/query/runs/progress/pause/resume/delete registered') + tctx.logger?.info?.('[dsh-ops-cron] tools cron_create/list/query/runs/progress/pause/resume/retrigger/delete registered') }) ctx.inject(['skills'], (sctx) => { diff --git a/lib/monitor.js b/lib/monitor.js index a93633e..25329a5 100644 --- a/lib/monitor.js +++ b/lib/monitor.js @@ -226,6 +226,24 @@ export function persistHistoryLimit(policy, settingsLimit = 200) { return settingsLimit } +/** + * One-shot job whose schedule has been consumed (next cleared) and has a terminal lastStatus. + * Used by listeners to decide cron_retrigger vs wait. + * @param {object|null} job + * @param {object[]} [runs] + */ +export function isRetriggerable(job, runs = []) { + if (!job || job.schedule?.kind !== 'at') return false + if (job.nextRunAt != null) return false + const terminal = job.lastStatus === 'succeeded' + || job.lastStatus === 'failed' + || job.lastStatus === 'skipped' + if (!terminal) return false + const list = Array.isArray(runs) ? runs : [] + if (list.some((run) => run && ACTIVE_RUN_STATUSES.has(run.status))) return false + return true +} + /** * Project job + runs into a lifecycle state for listeners. * @param {object} job diff --git a/lib/tools.js b/lib/tools.js index df08ad6..5ed10dc 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -53,6 +53,7 @@ const JOB_SCHEMA = { lastRunAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, lastStatus: { oneOf: [{ type: 'string' }, { type: 'null' }] }, nextRunAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + retriggerable: { type: 'boolean' }, }, } @@ -685,7 +686,7 @@ export function cronToolDefinitions(service, deps = {}) { }, { name: 'cron_resume', - description: 'Resume a paused scheduled task by id.', + description: 'Resume a paused scheduled task by id. Only flips enabled=true; does NOT re-fire a consumed one-shot (next=n/a). Use cron_retrigger for that.', parameters: { type: 'object', additionalProperties: false, @@ -707,6 +708,56 @@ export function cronToolDefinitions(service, deps = {}) { return { job } }, }, + { + name: 'cron_retrigger', + description: 'Re-fire a consumed one-shot job (succeeded/failed/skipped with next=n/a) as a new run. Auto-enables if paused. Use after_minutes>0 to delay. For recurring jobs use pause/resume. Listeners: when retriggerable=true and progress done_flag is false, call this instead of creating a new job.', + parameters: { + type: 'object', + additionalProperties: false, + properties: { + task_id: { type: 'string', description: 'One-shot job id (alias: id).' }, + id: { type: 'string', description: 'Alias of task_id.' }, + after_minutes: { type: 'number', description: 'Delay before fire. Omit or 0 = immediate run-now.' }, + }, + }, + output: { + schema: { + type: 'object', + additionalProperties: true, + properties: { + ok: { type: 'boolean' }, + mode: { type: 'string' }, + job: JOB_SCHEMA, + run: { oneOf: [RUN_SCHEMA, { type: 'null' }] }, + nextRunAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + }, + }, + render: (_args, value) => { + const name = value.job?.name || value.job?.id || '' + if (value.mode === 'scheduled') { + const when = value.nextRunAt ? new Date(value.nextRunAt).toISOString() : 'pending' + return text(`Retrigger scheduled "${name}" → next ${when}.`) + } + return text(`Retriggered "${name}"${value.run?.id ? ` (run ${value.run.id})` : ''}.`) + }, + }, + presentCall: (args) => ({ + card: 'generic', + title: t('tool.retrigger'), + content: String(args?.task_id || args?.id || ''), + }), + async execute(args, exec) { + aborted(exec) + const peer = await resolveCallerPeer(exec, getDshIm()) + const identity = resolveToolIdentity(exec, service) + const id = String(args?.task_id || args?.id || '').trim() + if (!id) throw new Error('task_id is required') + await requireOwnedJob(service, id, peer, identity) + return service.retriggerJob(id, { + after_minutes: args?.after_minutes, + }, identity) + }, + }, { name: 'cron_delete', description: 'Permanently delete a scheduled task by id. History rows for that job remain until pruned.', @@ -761,8 +812,9 @@ export function cronGuidanceText(nowMs = Date.now(), timeZone = 'Asia/Shanghai') '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 / cron_query / cron_runs / cron_progress are in your tool list, call them.', + 'One-shot wake: cron_resume only unpauses. For a finished one-shot (retriggerable=true, next=n/a), call cron_retrigger to start a new run on the same job id — do not create a duplicate executor.', + 'When the user asks to look at, create, pause, resume, retrigger, or delete 定时任务 / scheduled tasks / cron jobs:', + '1. If cron_list / cron_create / cron_pause / cron_resume / cron_retrigger / 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" for usage guidance, then look again. skill_load does NOT inject tools — cron_* are registered by the dsh-ops-cron Host plugin at boot.', '3. If cron_* are still missing after a Host restart with the latest dsh-ops-cron, the plugin failed to register (check Host logs for unsupported JSON schema / tools.register). 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.', @@ -786,7 +838,7 @@ Tools: - 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 +- cron_pause / cron_resume / cron_retrigger / cron_delete — by id from cron_list. Use cron_retrigger (not resume) to re-fire a consumed one-shot worker. `, } } diff --git a/test/host.test.js b/test/host.test.js index 98e8642..9e2b873 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -285,6 +285,119 @@ test('overlap skip writes a skipped history row instead of a second session', as assert.equal(inflight, 1) }) +test('retriggerJob re-fires consumed oneshot; rejects cron / pending / in-flight', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) + t.after(() => rm(dir, { recursive: true, force: true })) + let clock = Date.parse('2026-09-10T08:00:00.000Z') + let fires = 0 + const service = createTestHost({ + filePath: join(dir, 'store.json'), + now: () => clock, + sessionPort: { + async createAndPrompt() { + fires += 1 + return { sessionId: `s-${fires}`, status: 'succeeded', summary: `ok-${fires}` } + }, + async archiveSession() {}, + }, + }) + + const cron = await service.createJob({ + name: 'cron-job', + prompt: 'loop', + schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'UTC' }, + }) + await assert.rejects(() => service.retriggerJob(cron.id), (err) => err.code === 'INVALID_RETRIGGER') + + const waiting = await service.createJob({ + name: 'waiting', + prompt: 'soon', + schedule: { kind: 'at', at: new Date(clock + 3600_000).toISOString(), timezone: 'UTC' }, + }) + await assert.rejects(() => service.retriggerJob(waiting.id), (err) => err.code === 'INVALID_RETRIGGER') + + const worker = await service.createJob({ + name: 'worker', + prompt: 'continue', + enabled: false, + cwd: '/tmp/user-workspaces/tester', + schedule: { kind: 'at', at: new Date(clock + 60_000).toISOString(), timezone: 'UTC' }, + }, { empNo: 'tester', displayName: 'Tester', permissions: { canViewAllSessions: false } }) + // Force consume: run-now while enabling via update, then clear next by settling through schedule fire path. + await service.pauseJob(worker.id, true) + const first = await service.dispatchRun(worker.id, 'run-now') + assert.equal(first.run.status, 'succeeded') + // Manually mark as consumed oneshot (run-now keeps nextRunAt). + await service.store.mutate((state) => { + const job = state.jobs.find((row) => row.id === worker.id) + return { + ...state, + jobs: state.jobs.map((row) => (row.id === worker.id + ? { ...job, nextRunAt: null, lastStatus: 'succeeded', enabled: false } + : row)), + } + }) + const listed = await service.listJobs() + const view = listed.find((row) => row.id === worker.id) + assert.equal(view.retriggerable, true) + assert.equal(view.enabled, false) + + const delayed = await service.retriggerJob(worker.id, { after_minutes: 5 }) + assert.equal(delayed.mode, 'scheduled') + assert.equal(delayed.job.enabled, true) + assert.ok(delayed.nextRunAt > clock) + assert.equal(delayed.job.retriggerable, false) + + // Consume again then immediate retrigger. + await service.store.mutate((state) => ({ + ...state, + jobs: state.jobs.map((row) => (row.id === worker.id + ? { ...row, nextRunAt: null, lastStatus: 'succeeded' } + : row)), + })) + const again = await service.retriggerJob(worker.id) + assert.equal(again.mode, 'immediate') + assert.ok(again.run) + assert.equal(again.run.status, 'succeeded') + assert.equal(fires, 2) + + // In-flight reject + await service.store.mutate((state) => ({ + ...state, + jobs: state.jobs.map((row) => (row.id === worker.id + ? { ...row, nextRunAt: null, lastStatus: 'succeeded' } + : row)), + runs: [ + { + id: 'inflight', + jobId: worker.id, + status: 'running', + scheduledAt: clock, + actualAt: clock, + stateEnteredAt: clock, + }, + ...(state.runs || []), + ], + })) + await assert.rejects(() => service.retriggerJob(worker.id), (err) => err.code === 'ALREADY_RUNNING') + + const http = await listen(service) + t.after(() => http.close()) + await service.store.mutate((state) => ({ + ...state, + runs: (state.runs || []).filter((run) => run.id !== 'inflight'), + jobs: state.jobs.map((row) => (row.id === worker.id + ? { ...row, nextRunAt: null, lastStatus: 'succeeded', enabled: true } + : row)), + })) + const httpHit = await jsonRequest(http.url, `/dsh-ops-cron/jobs/${worker.id}/retrigger`, { + method: 'POST', + body: JSON.stringify({}), + }) + assert.equal(httpHit.status, 200) + assert.equal(httpHit.body.mode, 'immediate') + assert.ok(httpHit.body.run?.id) +}) test('http api: anonymous list/run-now is rejected', async (t) => { const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) diff --git a/test/monitor.test.js b/test/monitor.test.js index 3192cf5..d5eccbf 100644 --- a/test/monitor.test.js +++ b/test/monitor.test.js @@ -4,6 +4,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import { test } from 'node:test' import { + isRetriggerable, labelsMatch, normalizeLabels, normalizePersistHistory, @@ -19,6 +20,27 @@ import { claimOccurrence, publicJob, settleRun } from '../lib/fire.js' import { createHostService } from '../lib/host.js' import { cronToolDefinitions } from '../lib/tools.js' +test('isRetriggerable / publicJob.retriggerable for consumed oneshots', () => { + const now = Date.parse('2026-09-10T08:00:00.000Z') + const base = { + id: 'j1', + name: 'worker', + enabled: true, + schedule: { kind: 'at', at: new Date(now - 60_000).toISOString(), timezone: 'UTC' }, + nextRunAt: null, + lastStatus: 'succeeded', + lastRunAt: now - 30_000, + timeoutMinutes: 10, + } + assert.equal(isRetriggerable(base, []), true) + assert.equal(publicJob(base, []).retriggerable, true) + + assert.equal(isRetriggerable({ ...base, schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'UTC' } }, []), false) + assert.equal(isRetriggerable({ ...base, nextRunAt: now + 60_000 }, []), false) + assert.equal(isRetriggerable({ ...base, lastStatus: null }, []), false) + assert.equal(isRetriggerable(base, [{ id: 'r1', jobId: 'j1', status: 'running' }]), false) +}) + test('settleRun keeps running enter time and stamps exitedAt separately', () => { const now = Date.parse('2026-09-10T07:00:00.000Z') const started = now - 60_000 diff --git a/test/tools.test.js b/test/tools.test.js index 8d92d43..85c5113 100644 --- a/test/tools.test.js +++ b/test/tools.test.js @@ -439,6 +439,7 @@ test('registerCronTools registers each definition and disposer unregisters', () 'cron_progress', 'cron_pause', 'cron_resume', + 'cron_retrigger', 'cron_delete', ]) off() @@ -454,6 +455,7 @@ test('cron tool output schemas never use type arrays (Host rejects them)', () => async getProgress() { return { snapshots: [], count: 0 } }, async pauseJob() { return {} }, async resumeJob() { return {} }, + async retriggerJob() { return {} }, async deleteJob() { return {} }, }) const bad = []