From 839660d12f7b4a25eae9d081b8b89fdf001b6875 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 15:46:53 +0800 Subject: [PATCH 01/10] Fix monitor test syntax after settleRun duration coverage. Co-authored-by: Cursor --- test/monitor.test.js | 1 + 1 file changed, 1 insertion(+) diff --git a/test/monitor.test.js b/test/monitor.test.js index 25a6bb6..3192cf5 100644 --- a/test/monitor.test.js +++ b/test/monitor.test.js @@ -49,6 +49,7 @@ test('settleRun keeps running enter time and stamps exitedAt separately', () => assert.ok(run.exitedAt > run.stateEnteredAt) }) +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) From 6724332114f6e13e510ecb997e536684c14e95a7 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 16:01:31 +0800 Subject: [PATCH 02/10] Show agent preset on the scheduled-task list and editor hint. List rows previously hid agentPreset; surface Host-default vs pinned id so selection is visible without opening each job. Co-authored-by: Cursor --- lib/client.js | 27 +++++++++++++++++++++++---- 1 file changed, 23 insertions(+), 4 deletions(-) diff --git a/lib/client.js b/lib/client.js index 2d0a9c1..fdd9a97 100644 --- a/lib/client.js +++ b/lib/client.js @@ -232,7 +232,7 @@ window.__ModuleLoader__.load({ } function jobStamp(rows) { - return (rows || []).map((row) => `${row.id}:${row.updatedAt}:${row.nextRunAt}:${row.enabled}:${row.lastStatus}:${row.state}:${row.stuck}:${row.model}:${row.provider}:${row.delivery?.kind}:${row.delivery?.targetId}:${JSON.stringify(row.labels || {})}`).join('|') + return (rows || []).map((row) => `${row.id}:${row.updatedAt}:${row.nextRunAt}:${row.enabled}:${row.lastStatus}:${row.state}:${row.stuck}:${row.model}:${row.provider}:${row.agentPreset || ''}:${row.delivery?.kind}:${row.delivery?.targetId}:${JSON.stringify(row.labels || {})}`).join('|') } function runStamp(rows) { @@ -391,6 +391,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba stateFailed: '失败', statePaused: '已暂停', stateStuck: '疑似卡死', + presetHostDefault: 'Host默认', labels: '标签', labelsHint: '每行 key=value,或 JSON 对象。用于批量监控与发现。', labelsPlaceholder: 'role=worker\ntask=theory', @@ -458,6 +459,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba stateFailed: 'Failed', statePaused: 'Paused', stateStuck: 'Stuck', + presetHostDefault: 'Host default', labels: 'Labels', labelsHint: 'One key=value per line, or a JSON object. Used for discovery and batch monitoring.', labelsPlaceholder: 'role=worker\ntask=theory', @@ -877,7 +879,15 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba return '' } - function JobList({ t, jobs, runs, selection, expanded, viewer, error, onSelectJob, onSelectRun, onNew, onRun, onToggle, onRemove, onToggleGroup }) { + function jobPresetLabel(job, t, presets) { + const id = typeof job?.agentPreset === 'string' ? job.agentPreset.trim() : '' + if (!id) return t('presetHostDefault') + const items = Array.isArray(presets?.items) ? presets.items : [] + const hit = items.find((row) => row.id === id) + return hit?.name || id + } + + function JobList({ t, jobs, runs, selection, expanded, viewer, presets, error, onSelectJob, onSelectRun, onNew, onRun, onToggle, onRemove, onToggleGroup }) { const skin = workspaceSkin() const [query, setQuery] = useState('') const [searchOn, setSearchOn] = useState(false) @@ -897,6 +907,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba job.ownerDisplayName, job.state, labelsInline(job), + job.agentPreset, ...(runs.filter((run) => run.jobId === job.id).map((run) => run.summary || run.status)), ].join(' ').toLowerCase() return hay.includes(needle) @@ -971,6 +982,11 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba 'data-state': jobStateKey(job), 'data-stuck': job.stuck ? 'true' : 'false', }, jobStateLabel(job, t)), + h('span', { className: 'dsh-ct-labels', title: t('agentPreset') }, + `preset=${jobPresetLabel(job, t, presets)}`), + jobModelLabel(job) + ? h('span', { className: 'dsh-ct-labels' }, jobModelLabel(job)) + : null, labelsInline(job) ? h('span', { className: 'dsh-ct-labels' }, labelsInline(job)) : null, @@ -1216,7 +1232,10 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba key: row.id, }, row.name || row.id)), ), - h('span', { className: 'dsh-ct-cwdHint' }, t('agentPresetHint')), + h('span', { className: 'dsh-ct-cwdHint' }, + selected + ? `${t('agentPreset')}: ${known ? (items.find((row) => row.id === selected)?.name || selected) : selected}` + : `${t('agentPreset')}: ${t('presetHostDefault')} — ${t('agentPresetHint')}`), ) } @@ -1761,7 +1780,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba return h(React.Fragment, null, h(JobList, { - t, jobs, runs, selection, viewer, expanded, error, + t, jobs, runs, selection, viewer, presets, expanded, error, onSelectJob: selectJob, onSelectRun: selectRun, onNew: selectNew, From c11221606a3a45814ac2b0edd1bf78926b129cdd Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 17:40:39 +0800 Subject: [PATCH 03/10] Add cron_retrigger for consumed one-shot jobs. Resume only unpauses; listeners need a native way to re-fire the same job id with a new run without creating duplicate executors. Co-authored-by: Cursor --- lib/fire.js | 2 + lib/host.js | 110 ++++++++++++++++++++++++++++++++++++++++- lib/i18n.js | 2 + lib/index.js | 2 +- lib/monitor.js | 18 +++++++ lib/tools.js | 60 +++++++++++++++++++++-- test/host.test.js | 113 +++++++++++++++++++++++++++++++++++++++++++ test/monitor.test.js | 22 +++++++++ test/tools.test.js | 2 + 9 files changed, 324 insertions(+), 7 deletions(-) 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 = [] From 3ce60b200ccab3a69cc1b73401671439994b5749 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 21:22:06 +0800 Subject: [PATCH 04/10] Allow cron_retrigger to immediately fire recurring jobs. Recurring tasks can start one extra run now without changing the cron schedule; delay remains oneshot-only. Co-authored-by: Cursor --- lib/host.js | 26 +++++++++++++++++++++----- lib/monitor.js | 16 +++++++++------- lib/tools.js | 10 +++++----- test/host.test.js | 16 +++++++++++++--- test/monitor.test.js | 11 ++++++++++- 5 files changed, 58 insertions(+), 21 deletions(-) diff --git a/lib/host.js b/lib/host.js index 3554d5e..45a5185 100644 --- a/lib/host.js +++ b/lib/host.js @@ -701,9 +701,10 @@ export function createHostService(options = {}) { } /** - * 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. + * Fire an extra run, or re-arm a consumed one-shot. + * - recurring (cron): after_minutes must be 0/omitted → enable + immediate run-now (keeps cron next) + * - one-shot (at): consumed only; 0 → run-now; >0 → reschedule next + * Auto-enables paused jobs. Rejects when a run is already pending/running. */ async function retriggerJob(jobId, opts = {}, identity = null) { const afterRaw = opts.afterMinutes ?? opts.after_minutes @@ -727,8 +728,9 @@ export function createHostService(options = {}) { 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') + const kind = job.schedule?.kind + if (kind !== 'at' && kind !== 'cron') { + const error = new Error('cron_retrigger requires a one-shot (at) or recurring (cron) job') error.code = 'INVALID_RETRIGGER' throw error } @@ -738,6 +740,20 @@ export function createHostService(options = {}) { error.code = 'ALREADY_RUNNING' throw error } + + if (kind === 'cron') { + if (delayMs > 0) { + const error = new Error('after_minutes delay is only for one-shot jobs; omit it to fire a recurring job immediately') + error.code = 'INVALID_RETRIGGER' + throw error + } + if (job.enabled === false) { + return upsertJob(current, { ...job, enabled: true, updatedAt: t }) + } + return current + } + + // one-shot 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' diff --git a/lib/monitor.js b/lib/monitor.js index 25329a5..430048c 100644 --- a/lib/monitor.js +++ b/lib/monitor.js @@ -227,21 +227,23 @@ export function persistHistoryLimit(policy, settingsLimit = 200) { } /** - * One-shot job whose schedule has been consumed (next cleared) and has a terminal lastStatus. - * Used by listeners to decide cron_retrigger vs wait. + * Whether cron_retrigger may start a new run now. + * - recurring (cron): idle (no pending/running run) + * - one-shot (at): consumed (next=n/a) + terminal lastStatus + idle * @param {object|null} job * @param {object[]} [runs] */ export function isRetriggerable(job, runs = []) { - if (!job || job.schedule?.kind !== 'at') return false + if (!job) return false + const list = Array.isArray(runs) ? runs : [] + if (list.some((run) => run && ACTIVE_RUN_STATUSES.has(run.status))) return false + if (job.schedule?.kind === 'cron') return true + if (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 + return terminal } /** diff --git a/lib/tools.js b/lib/tools.js index 5ed10dc..6c7b678 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -710,14 +710,14 @@ export function cronToolDefinitions(service, deps = {}) { }, { 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.', + description: 'Start an extra run now: for recurring (cron) jobs fires immediately without changing the schedule; for consumed one-shots (next=n/a) re-arms/fires a new run. Auto-enables if paused. after_minutes>0 only for one-shots. Rejects if already pending/running. Listeners: when retriggerable=true and progress done_flag is false, call this instead of creating a duplicate job.', parameters: { type: 'object', additionalProperties: false, properties: { - task_id: { type: 'string', description: 'One-shot job id (alias: id).' }, + task_id: { type: 'string', description: '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.' }, + after_minutes: { type: 'number', description: 'One-shot only: delay before fire. Omit/0 = immediate. Not allowed for recurring jobs.' }, }, }, output: { @@ -812,7 +812,7 @@ 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.', - '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.', + 'One-shot wake: cron_resume only unpauses. For a finished one-shot (retriggerable=true, next=n/a), call cron_retrigger. For a recurring job, cron_retrigger also starts one extra run immediately without changing the cron schedule.', '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.', @@ -838,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_retrigger / cron_delete — by id from cron_list. Use cron_retrigger (not resume) to re-fire a consumed one-shot worker. +- cron_pause / cron_resume / cron_retrigger / cron_delete — by id from cron_list. Use cron_retrigger to re-fire a consumed one-shot or to immediately start one extra run of a recurring job. `, } } diff --git a/test/host.test.js b/test/host.test.js index 9e2b873..71988b7 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -285,7 +285,7 @@ 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) => { +test('retriggerJob fires recurring immediately and re-arms consumed oneshot', 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') @@ -307,7 +307,17 @@ test('retriggerJob re-fires consumed oneshot; rejects cron / pending / in-flight prompt: 'loop', schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'UTC' }, }) - await assert.rejects(() => service.retriggerJob(cron.id), (err) => err.code === 'INVALID_RETRIGGER') + const cronNext = cron.nextRunAt + assert.equal(cron.retriggerable, true) + await assert.rejects( + () => service.retriggerJob(cron.id, { after_minutes: 5 }), + (err) => err.code === 'INVALID_RETRIGGER', + ) + const cronHit = await service.retriggerJob(cron.id) + assert.equal(cronHit.mode, 'immediate') + assert.equal(cronHit.run.status, 'succeeded') + assert.equal(cronHit.job.nextRunAt, cronNext) + assert.equal(fires, 1) const waiting = await service.createJob({ name: 'waiting', @@ -359,7 +369,7 @@ test('retriggerJob re-fires consumed oneshot; rejects cron / pending / in-flight assert.equal(again.mode, 'immediate') assert.ok(again.run) assert.equal(again.run.status, 'succeeded') - assert.equal(fires, 2) + assert.equal(fires, 3) // In-flight reject await service.store.mutate((state) => ({ diff --git a/test/monitor.test.js b/test/monitor.test.js index d5eccbf..7e3ec37 100644 --- a/test/monitor.test.js +++ b/test/monitor.test.js @@ -35,7 +35,16 @@ test('isRetriggerable / publicJob.retriggerable for consumed oneshots', () => { assert.equal(isRetriggerable(base, []), true) assert.equal(publicJob(base, []).retriggerable, true) - assert.equal(isRetriggerable({ ...base, schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'UTC' } }, []), false) + const cron = { + ...base, + schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'UTC' }, + nextRunAt: now + 60_000, + lastStatus: null, + } + assert.equal(isRetriggerable(cron, []), true) + assert.equal(publicJob(cron, []).retriggerable, true) + assert.equal(isRetriggerable(cron, [{ id: 'r1', jobId: 'j1', status: 'running' }]), 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) From 4472a9c74b869163ab349750c1672245686b77a2 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 21:46:16 +0800 Subject: [PATCH 05/10] Normalize doubled provider/model routes on cron create and fire. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Agents often pass zte/Qwen3-… as both provider and model; split combined routes so scheduled runs resolve a real model. Co-authored-by: Cursor --- lib/host.js | 6 ++++-- lib/index.js | 2 +- lib/store.js | 38 ++++++++++++++++++++++++++++++++++++-- lib/tools.js | 28 ++++++++++++++++++---------- test/tools.test.js | 36 ++++++++++++++++++++++++++++++++++++ 5 files changed, 95 insertions(+), 15 deletions(-) diff --git a/lib/host.js b/lib/host.js index 45a5185..c627106 100644 --- a/lib/host.js +++ b/lib/host.js @@ -36,6 +36,7 @@ import { listJobs, normalizeSettings, removeJob, + splitProviderModel, storePath, upsertJob, } from './store.js' @@ -1832,8 +1833,9 @@ export async function listPresetChoices(ctx) { } export async function resolveJobModel(ctx, job) { - const provider = typeof job?.provider === 'string' ? job.provider.trim() : '' - const model = typeof job?.model === 'string' ? job.model.trim() : '' + const split = splitProviderModel(job?.provider, job?.model) + const provider = split.provider + const model = split.model if (provider && model) { return { provider, diff --git a/lib/index.js b/lib/index.js index ad09f4a..96f32f8 100644 --- a/lib/index.js +++ b/lib/index.js @@ -48,7 +48,7 @@ export { workspaceVisibleIds, } from './isolation.js' export { wrapScheduledPrompt } from './prompt.js' -export { normalizeJobModel } from './store.js' +export { normalizeJobModel, splitProviderModel } from './store.js' export { assertCanAccessJob, canViewAllJobs, diff --git a/lib/store.js b/lib/store.js index d12b552..2f485ad 100644 --- a/lib/store.js +++ b/lib/store.js @@ -77,12 +77,46 @@ export function newId() { return randomUUID() } +/** + * Repair provider/model pairs when LLMs or hosts pass a combined "provider/model" + * route in one or both fields (e.g. provider=model="zte/Qwen3-…" → zte + Qwen3-…). + * Leaves legitimate model ids that contain "/" alone when provider is a distinct id. + * @param {unknown} provider + * @param {unknown} model + * @returns {{ provider: string, model: string }} + */ +export function splitProviderModel(provider, model) { + let p = String(provider || '').trim() + let m = String(model || '').trim() + if (!p && !m) return { provider: '', model: '' } + + if (p && m && p === m && p.includes('/')) { + const at = p.indexOf('/') + return { provider: p.slice(0, at).trim(), model: p.slice(at + 1).trim() } + } + if (!p && m.includes('/')) { + const at = m.indexOf('/') + return { provider: m.slice(0, at).trim(), model: m.slice(at + 1).trim() } + } + if (p.includes('/') && !m) { + const at = p.indexOf('/') + return { provider: p.slice(0, at).trim(), model: p.slice(at + 1).trim() } + } + if (p && m.startsWith(`${p}/`)) { + return { provider: p, model: m.slice(p.length + 1).trim() } + } + if (p.includes('/') && m && p.endsWith(`/${m}`)) { + const at = p.indexOf('/') + return { provider: p.slice(0, at).trim(), model: m } + } + return { provider: p, model: m } +} + export function normalizeJobModel(input = {}) { - const provider = typeof input.provider === 'string' ? input.provider.trim() : '' - const model = typeof input.model === 'string' ? input.model.trim() : '' const reasoningEffort = typeof input.reasoningEffort === 'string' ? input.reasoningEffort.trim() : (typeof input.reasoning_effort === 'string' ? input.reasoning_effort.trim() : '') + const { provider, model } = splitProviderModel(input.provider, input.model) if (!provider && !model) { return { provider: '', model: '', reasoningEffort: '' } } diff --git a/lib/tools.js b/lib/tools.js index 6c7b678..2f1fe32 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -20,6 +20,7 @@ import { jobVisibleToIdentity, UNASSIGNED_OWNER, } from './ownership.js' +import { splitProviderModel } from './store.js' const JOB_SCHEMA = { type: 'object', @@ -328,20 +329,27 @@ async function requireOwnedJob(service, id, peer, identity = null) { export function callerModelSelection(exec) { const opts = exec?.agent?.options || {} - const provider = String(opts.provider || '').trim() - const model = String(opts.model || '').trim() - if (!provider || !model) return { provider: '', model: '', reasoningEffort: '' } + const split = splitProviderModel(opts.provider, opts.model) + if (!split.provider || !split.model) return { provider: '', model: '', reasoningEffort: '' } const reasoningEffort = String(opts.reasoningEffort || '').trim() - return { provider, model, reasoningEffort } + return { provider: split.provider, model: split.model, reasoningEffort } } export function resolveCreateModel(args, exec) { - const provider = typeof args?.provider === 'string' ? args.provider.trim() : '' - const model = typeof args?.model === 'string' ? args.model.trim() : '' const reasoningEffort = typeof args?.reasoning_effort === 'string' ? args.reasoning_effort.trim() : (typeof args?.reasoningEffort === 'string' ? args.reasoningEffort.trim() : '') - if (provider && model) return { provider, model, reasoningEffort } + const explicit = splitProviderModel(args?.provider, args?.model) + if (explicit.provider && explicit.model) { + return { provider: explicit.provider, model: explicit.model, reasoningEffort } + } + // One of provider/model alone (after split) → ignore; inherit session instead of throwing. + if (explicit.provider || explicit.model) { + const fromSession = callerModelSelection(exec) + if (fromSession.provider && fromSession.model) { + return { ...fromSession, reasoningEffort: reasoningEffort || fromSession.reasoningEffort } + } + } const fromSession = callerModelSelection(exec) return { ...fromSession, reasoningEffort: reasoningEffort || fromSession.reasoningEffort } } @@ -379,8 +387,8 @@ export function cronToolDefinitions(service, deps = {}) { timezone: { type: 'string', description: 'IANA timezone for expr/at. Default Asia/Shanghai. Do not pass UTC unless the user asked for UTC.' }, time_zone: { type: 'string', description: 'Alias of timezone.' }, cwd: { type: 'string', description: 'Filesystem path of the workspace this job should run in. When the user is in a workspace conversation, pass THAT workspace path (current session cwd). If omitted, the current session cwd is used. WhatsApp/IM creates always use the current session cwd.' }, - provider: { type: 'string', description: 'Provider route for this job (e.g. minimax-cn, deepseek). Must be passed with model. If omitted, the current session model is stored so quota stays predictable.' }, - model: { type: 'string', description: 'Model id for this job. Must be passed with provider. Scheduled runs bill this model.' }, + provider: { type: 'string', description: 'Provider id only (e.g. zte, deepseek, minimax-cn). Do NOT pass "provider/model". Must pair with model, or omit both to inherit the current session.' }, + model: { type: 'string', description: 'Bare model id only (e.g. Qwen3-235B-A22B). Do NOT pass "provider/model" or repeat the provider. Must pair with provider, or omit both to inherit the current session.' }, reasoning_effort: { type: 'string', description: 'Optional reasoning effort for this job.' }, timeout_minutes: { type: 'integer', description: 'Per-run timeout in minutes, 1-240.' }, enabled: { type: 'boolean', description: 'If false, create paused. Default true.' }, @@ -807,7 +815,7 @@ export function cronGuidanceText(nowMs = Date.now(), timeZone = 'Asia/Shanghai') 'For "in N minutes / 一分钟后", call cron_create with after_minutes=N only (do not also pass at, hour, or expr).', 'For a clock time tonight, pass only hour and minute in 24h (晚上11点34 → hour=23, minute=34; 零点33 → hour=0, minute=33). Extra at/expr fields are ignored.', 'Working directory: if the user is chatting in a workspace, pass cwd as that workspace filesystem path (the current session working directory). If they name another workspace, use that path. If cwd is omitted, cron_create uses the current session cwd.', - 'Model: pass provider+model for the job. If omitted, cron_create stores the current session model. Scheduled runs consume that model\'s quota.', + 'Model: omit provider+model to inherit the current session. If you pass them, use bare ids only (provider=zte, model=Qwen3-235B-A22B) — never pass "zte/Qwen3-…" as either field.', '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.', diff --git a/test/tools.test.js b/test/tools.test.js index 85c5113..c7c24b4 100644 --- a/test/tools.test.js +++ b/test/tools.test.js @@ -116,6 +116,42 @@ test('resolveCreateCwd prefers explicit cwd then the calling session', () => { }) }) +test('resolveCreateModel / normalizeJobModel split doubled provider/model routes', async (t) => { + const { normalizeJobModel, splitProviderModel } = await import('../lib/store.js') + assert.deepEqual( + splitProviderModel('zte/Qwen3-235B-A22B', 'zte/Qwen3-235B-A22B'), + { provider: 'zte', model: 'Qwen3-235B-A22B' }, + ) + assert.deepEqual( + splitProviderModel('zte', 'zte/Qwen3-235B-A22B'), + { provider: 'zte', model: 'Qwen3-235B-A22B' }, + ) + assert.deepEqual( + splitProviderModel('', 'zte/Qwen3-235B-A22B'), + { provider: 'zte', model: 'Qwen3-235B-A22B' }, + ) + // Legitimate model id with slash under a distinct provider stays intact. + assert.deepEqual( + splitProviderModel('openai', 'org/custom-model'), + { provider: 'openai', model: 'org/custom-model' }, + ) + assert.deepEqual( + normalizeJobModel({ provider: 'zte/Qwen3-235B-A22B', model: 'zte/Qwen3-235B-A22B' }), + { provider: 'zte', model: 'Qwen3-235B-A22B', reasoningEffort: '' }, + ) + assert.deepEqual( + resolveCreateModel( + { provider: 'zte/Qwen3-235B-A22B', model: 'zte/Qwen3-235B-A22B' }, + { agent: { options: { provider: 'other', model: 'other-model' } } }, + ), + { provider: 'zte', model: 'Qwen3-235B-A22B', reasoningEffort: '' }, + ) + assert.deepEqual( + resolveCreateModel({}, { agent: { options: { provider: 'zte/Qwen3-235B-A22B', model: 'zte/Qwen3-235B-A22B' } } }), + { provider: 'zte', model: 'Qwen3-235B-A22B', reasoningEffort: '' }, + ) +}) + test('cron_create hour+minute uses today and rejects a guessed past calendar date', async (t) => { const dir = await mkdtemp(join(tmpdir(), 'dsh-cron-tools-')) t.after(() => rm(dir, { recursive: true, force: true })) From ad6e90f1978332183b758093e876e8f08678f777 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 21:54:31 +0800 Subject: [PATCH 06/10] Add sidebar-only agentAccess lock for scheduled jobs. Default remains allow; humans can deny agent pause/resume/retrigger/delete in the panel while tools never expose or set the field. Co-authored-by: Cursor --- lib/client.js | 21 +++++++++++++++++++ lib/fire.js | 1 + lib/host.js | 3 +++ lib/index.js | 2 +- lib/store.js | 22 ++++++++++++++++++++ lib/tools.js | 51 +++++++++++++++++++++++++++++++++++++--------- test/tools.test.js | 41 +++++++++++++++++++++++++++++++++++++ 7 files changed, 130 insertions(+), 11 deletions(-) diff --git a/lib/client.js b/lib/client.js index fdd9a97..1b42ca7 100644 --- a/lib/client.js +++ b/lib/client.js @@ -430,6 +430,10 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba deliveryHint: '选 WhatsApp/IM 后从下拉选择已保存的投递目标(仅超管)。普通用户请在 IM 聊天里用 cron_create 建渠道任务。目标在 IM 机器人 → 投递设置里创建。', mirrorToSession: '镜像回原会话', mirrorToSessionHint: '开启后把每次运行摘要写入创建时的 WhatsApp/Web 会话(不新开模型轮次)。默认关闭。', + agentAccess: 'Agent 操作', + agentAccessAllow: '允许 Agent 用工具管理(默认)', + agentAccessDeny: '仅人工(禁止 Agent 暂停/恢复/重触发/删除)', + agentAccessHint: '仅侧栏可改;Agent 看不到此开关。默认保持现状(允许)。关掉后 Agent 仍可查询进度,但无法改任务。', loginRequired: '登录后才能使用定时任务', runNoSession: '这次运行没有可打开的会话', runOpenFailed: '打不开这次对话,会话可能已被删除', @@ -498,6 +502,10 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba deliveryHint: 'Pick a saved IM delivery target (super_admin only). Regular users should create channel jobs from the IM chat via cron_create. Create targets under IM bot → Delivery settings.', mirrorToSession: 'Mirror into origin session', mirrorToSessionHint: 'When enabled, append each run summary into the creating WhatsApp/Web session (no new model turn). Off by default.', + agentAccess: 'Agent control', + agentAccessAllow: 'Allow agents to manage via tools (default)', + agentAccessDeny: 'Human only (block agent pause/resume/retrigger/delete)', + agentAccessHint: 'Sidebar only; agents never see this toggle. Default keeps current behavior (allow). When denied, agents can still query progress but cannot mutate the job.', loginRequired: 'Sign in to use scheduled tasks', runNoSession: 'This run has no session to open', runOpenFailed: 'Could not open this chat; the session may have been deleted', @@ -755,6 +763,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba imBotId: '', imTargetId: '', mirrorToSession: false, + agentAccess: 'allow', labelsText: '', persistKind: 'retain', archiveEndpoint: '', @@ -781,6 +790,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba imBotId: job.delivery?.botId || '', imTargetId: job.delivery?.targetId || '', mirrorToSession: job.mirrorToSession === true, + agentAccess: job.agentAccess === 'deny' ? 'deny' : 'allow', labelsText: formatLabelsText(job.labels), persistKind: persistHistoryKind(job), archiveEndpoint: job.persistHistory?.kind === 'archive' ? (job.persistHistory.endpoint || '') : '', @@ -1421,6 +1431,16 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba ? h(ImDeliveryFields, { t, form, setForm, imCatalog }) : null, h('span', { className: 'dsh-ct-cwdHint' }, t('deliveryHint')), + h('label', null, t('agentAccess'), + h('select', { + value: form.agentAccess === 'deny' ? 'deny' : 'allow', + onChange: (e) => setForm({ ...form, agentAccess: e.target.value }), + }, + h('option', { value: 'allow' }, t('agentAccessAllow')), + h('option', { value: 'deny' }, t('agentAccessDeny')), + ), + ), + h('span', { className: 'dsh-ct-cwdHint' }, t('agentAccessHint')), h('label', { className: 'dsh-ct-check' }, h('input', { type: 'checkbox', @@ -1689,6 +1709,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba ? { kind: 'im', botId: form.imBotId, targetId: form.imTargetId } : { kind: 'dsh' }, mirrorToSession: form.mirrorToSession === true, + agentAccess: form.agentAccess === 'deny' ? 'deny' : 'allow', labels: parseLabelsText(form.labelsText), persistHistory: form.persistKind === 'forever' ? { kind: 'forever' } diff --git a/lib/fire.js b/lib/fire.js index 7a1e751..846a8d2 100644 --- a/lib/fire.js +++ b/lib/fire.js @@ -108,6 +108,7 @@ export function publicJob(job, runs = null) { model: job.model || '', reasoningEffort: job.reasoningEffort || '', agentPreset: job.agentPreset || '', + agentAccess: job.agentAccess === 'deny' ? 'deny' : 'allow', ownerEmpNo: job.ownerEmpNo || '', ownerDisplayName: job.ownerDisplayName || '', delivery: job.delivery && job.delivery.kind === 'im' diff --git a/lib/host.js b/lib/host.js index c627106..999e6bb 100644 --- a/lib/host.js +++ b/lib/host.js @@ -630,6 +630,9 @@ export function createHostService(options = {}) { model: patch.model !== undefined ? patch.model : job.model, reasoningEffort: patch.reasoningEffort !== undefined ? patch.reasoningEffort : job.reasoningEffort, agentPreset: patch.agentPreset !== undefined ? patch.agentPreset : job.agentPreset, + agentAccess: patch.agentAccess !== undefined + ? patch.agentAccess + : (patch.agent_access !== undefined ? patch.agent_access : job.agentAccess), delivery: patch.delivery !== undefined ? mergeDeliveryMention(job.delivery, patch.delivery) : job.delivery, diff --git a/lib/index.js b/lib/index.js index 96f32f8..66e8314 100644 --- a/lib/index.js +++ b/lib/index.js @@ -48,7 +48,7 @@ export { workspaceVisibleIds, } from './isolation.js' export { wrapScheduledPrompt } from './prompt.js' -export { normalizeJobModel, splitProviderModel } from './store.js' +export { normalizeJobModel, normalizeAgentAccess, agentMayMutateJob, splitProviderModel } from './store.js' export { assertCanAccessJob, canViewAllJobs, diff --git a/lib/store.js b/lib/store.js index 2f485ad..a332c12 100644 --- a/lib/store.js +++ b/lib/store.js @@ -77,6 +77,26 @@ export function newId() { return randomUUID() } +/** + * Whether agents may mutate this job via cron_* tools. + * Default allow (= current behavior). Only the sidebar should set deny. + * @param {unknown} input job or raw value + * @returns {'allow'|'deny'} + */ +export function normalizeAgentAccess(input) { + const raw = input && typeof input === 'object' && !Array.isArray(input) + ? (input.agentAccess ?? input.agent_access) + : input + if (raw === false || raw === 0 || raw === 'deny' || raw === 'human' || raw === 'locked' || raw === 'off') { + return 'deny' + } + return 'allow' +} + +export function agentMayMutateJob(job) { + return normalizeAgentAccess(job) === 'allow' +} + /** * Repair provider/model pairs when LLMs or hosts pass a combined "provider/model" * route in one or both fields (e.g. provider=model="zte/Qwen3-…" → zte + Qwen3-…). @@ -169,6 +189,7 @@ export function createJobRecord(input, state, now) { input.persistHistory ?? input.persist_history, settings.historyLimit, ) + const agentAccess = normalizeAgentAccess(input) const job = { id: String(input.id || newId()), name, @@ -182,6 +203,7 @@ export function createJobRecord(input, state, now) { model: model.model, reasoningEffort: model.reasoningEffort, agentPreset, + agentAccess, delivery, mirrorToSession, ownerEmpNo, diff --git a/lib/tools.js b/lib/tools.js index 2f1fe32..22bf5a9 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -20,7 +20,7 @@ import { jobVisibleToIdentity, UNASSIGNED_OWNER, } from './ownership.js' -import { splitProviderModel } from './store.js' +import { agentMayMutateJob, splitProviderModel } from './store.js' const JOB_SCHEMA = { type: 'object', @@ -216,6 +216,20 @@ function jobLine(job) { return `${job.name} [${life}${stuck}] ${sched} tz=${tz} cwd=${job.cwd || '(recent workspace)'} model=${model} ${delivery}${labels} next=${when} id=${job.id}` } +/** Strip human-only fields so agents cannot see or game panel-only controls. */ +function forAgentView(job) { + if (!job || typeof job !== 'object') return job + const { agentAccess: _hidden, ...rest } = job + return rest +} + +function assertAgentMayMutate(job) { + if (agentMayMutateJob(job)) return + const error = new Error('this job is human-managed; change “Agent 操作” in the 定时任务 panel (agents cannot toggle it)') + error.code = 'AGENT_ACCESS_DENIED' + throw error +} + export function callerWorkingDirectory(exec) { const session = exec?.agent?.session return String(session?.header?.cwd || session?.cwd || exec?.agent?.cwd || '').trim() @@ -458,7 +472,7 @@ export function cronToolDefinitions(service, deps = {}) { agentPreset, ...monitorFieldsFromArgs(args), }, identity, { fromImPeer: !!peer?.botId }) - return { job } + return { job: forAgentView(job) } } catch (error) { const tz = args.timezone || args.time_zone || 'Asia/Shanghai' const nowText = formatInZone(Date.now(), tz) @@ -512,7 +526,8 @@ export function cronToolDefinitions(service, deps = {}) { 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 } + const visible = jobs.map(forAgentView) + return { jobs: visible, count: visible.length } }, }, { @@ -565,6 +580,12 @@ export function cronToolDefinitions(service, deps = {}) { result.items = (result.items || []).filter((item) => jobVisibleToPeer(item.job, peer)) result.count = result.jobs.length } + result.jobs = (result.jobs || []).map(forAgentView) + result.items = (result.items || []).map((item) => ( + item && typeof item === 'object' + ? { ...item, job: forAgentView(item.job) } + : item + )) return result }, }, @@ -665,6 +686,11 @@ export function cronToolDefinitions(service, deps = {}) { result.snapshots = (result.snapshots || []).filter((row) => jobVisibleToPeer(row.job, peer)) result.count = result.snapshots.length } + result.snapshots = (result.snapshots || []).map((row) => ( + row && typeof row === 'object' + ? { ...row, job: forAgentView(row.job) } + : row + )) return result }, }, @@ -687,9 +713,10 @@ export function cronToolDefinitions(service, deps = {}) { const peer = await resolveCallerPeer(exec, getDshIm()) const identity = resolveToolIdentity(exec, service) const id = requireId(args) - await requireOwnedJob(service, id, peer, identity) + const owned = await requireOwnedJob(service, id, peer, identity) + assertAgentMayMutate(owned) const job = await service.pauseJob(id, false, identity) - return { job } + return { job: forAgentView(job) } }, }, { @@ -711,9 +738,10 @@ export function cronToolDefinitions(service, deps = {}) { const peer = await resolveCallerPeer(exec, getDshIm()) const identity = resolveToolIdentity(exec, service) const id = requireId(args) - await requireOwnedJob(service, id, peer, identity) + const owned = await requireOwnedJob(service, id, peer, identity) + assertAgentMayMutate(owned) const job = await service.pauseJob(id, true, identity) - return { job } + return { job: forAgentView(job) } }, }, { @@ -760,10 +788,12 @@ export function cronToolDefinitions(service, deps = {}) { 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, { + const owned = await requireOwnedJob(service, id, peer, identity) + assertAgentMayMutate(owned) + const result = await service.retriggerJob(id, { after_minutes: args?.after_minutes, }, identity) + return { ...result, job: forAgentView(result.job) } }, }, { @@ -794,7 +824,8 @@ export function cronToolDefinitions(service, deps = {}) { const identity = resolveToolIdentity(exec, service) const id = requireId(args) try { - await requireOwnedJob(service, id, peer, identity) + const owned = await requireOwnedJob(service, id, peer, identity) + assertAgentMayMutate(owned) await service.deleteJob(id) return { id, deleted: true } } catch (error) { diff --git a/test/tools.test.js b/test/tools.test.js index c7c24b4..44e4daf 100644 --- a/test/tools.test.js +++ b/test/tools.test.js @@ -453,6 +453,47 @@ test('cron_create from IM peer stamps unassigned owner and forces session cwd', assert.equal(created.job.delivery?.kind, 'im') }) +test('agentAccess deny is hidden from tools and blocks pause/delete', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-cron-tools-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const service = createHostService({ + filePath: join(dir, 'store.json'), + now: () => Date.now(), + sessionPort: { + async createAndPrompt() { + return { sessionId: 'tool-sess', status: 'succeeded', summary: 'ok' } + }, + }, + }) + const tools = byName(cronToolDefinitions(service)) + const created = await tools.cron_create.execute({ + name: 'locked', + prompt: 'ping', + after_minutes: 30, + timezone: 'Asia/Shanghai', + }, {}) + assert.equal(created.job.agentAccess, undefined) + assert.equal(created.job.id != null, true) + + await service.updateJob(created.job.id, { agentAccess: 'deny' }) + const listed = await tools.cron_list.execute({}, {}) + const row = listed.jobs.find((job) => job.id === created.job.id) + assert.ok(row) + assert.equal(row.agentAccess, undefined) + + await assert.rejects( + () => tools.cron_pause.execute({ id: created.job.id }, {}), + (err) => err.code === 'AGENT_ACCESS_DENIED', + ) + await assert.rejects( + () => tools.cron_delete.execute({ id: created.job.id }, {}), + (err) => err.code === 'AGENT_ACCESS_DENIED', + ) + + const httpView = await service.getJob(created.job.id) + assert.equal(httpView.agentAccess, 'deny') +}) + test('registerCronTools registers each definition and disposer unregisters', () => { const registered = [] const ctx = { From 8cf624ce7f39f4141acad1ea3ebf787e927d11d1 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 10 Sep 2026 22:08:39 +0800 Subject: [PATCH 07/10] Add sidebar run permission preset for scheduled jobs. Default inherits Host new-session access mode; humans can raise to workspace-write or full access. Agents never see or set the field. Co-authored-by: Cursor --- lib/client.js | 27 +++++++++++++++++++++++++++ lib/fire.js | 1 + lib/host.js | 18 +++++++++++++++++- lib/index.js | 2 +- lib/store.js | 34 ++++++++++++++++++++++++++++++++++ lib/tools.js | 2 +- test/host.test.js | 45 +++++++++++++++++++++++++++++++++++++++++++++ test/tools.test.js | 15 ++++++++++++++- 8 files changed, 140 insertions(+), 4 deletions(-) diff --git a/lib/client.js b/lib/client.js index 1b42ca7..a62f296 100644 --- a/lib/client.js +++ b/lib/client.js @@ -434,6 +434,12 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba agentAccessAllow: '允许 Agent 用工具管理(默认)', agentAccessDeny: '仅人工(禁止 Agent 暂停/恢复/重触发/删除)', agentAccessHint: '仅侧栏可改;Agent 看不到此开关。默认保持现状(允许)。关掉后 Agent 仍可查询进度,但无法改任务。', + permissionPreset: '操作权限', + permissionPresetDefault: '跟随 Host 默认(当前行为)', + permissionPresetReadOnly: '仅可查看', + permissionPresetWorkspaceWrite: '可写入工作区', + permissionPresetFullAccess: '完全权限', + permissionPresetHint: '仅侧栏可改。默认不指定,触发时沿用 Host「新会话」默认权限;可手工提权到可写/完全权限。Agent 不可见也不可改。', loginRequired: '登录后才能使用定时任务', runNoSession: '这次运行没有可打开的会话', runOpenFailed: '打不开这次对话,会话可能已被删除', @@ -506,6 +512,12 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba agentAccessAllow: 'Allow agents to manage via tools (default)', agentAccessDeny: 'Human only (block agent pause/resume/retrigger/delete)', agentAccessHint: 'Sidebar only; agents never see this toggle. Default keeps current behavior (allow). When denied, agents can still query progress but cannot mutate the job.', + permissionPreset: 'Run permission', + permissionPresetDefault: 'Host default (current behavior)', + permissionPresetReadOnly: 'Read only', + permissionPresetWorkspaceWrite: 'Workspace write', + permissionPresetFullAccess: 'Full access', + permissionPresetHint: 'Sidebar only. Leave default to inherit the Host new-session permission preset; raise manually to workspace write or full access. Agents cannot see or change this.', loginRequired: 'Sign in to use scheduled tasks', runNoSession: 'This run has no session to open', runOpenFailed: 'Could not open this chat; the session may have been deleted', @@ -764,6 +776,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba imTargetId: '', mirrorToSession: false, agentAccess: 'allow', + permissionPreset: '', labelsText: '', persistKind: 'retain', archiveEndpoint: '', @@ -791,6 +804,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba imTargetId: job.delivery?.targetId || '', mirrorToSession: job.mirrorToSession === true, agentAccess: job.agentAccess === 'deny' ? 'deny' : 'allow', + permissionPreset: job.permissionPreset || '', labelsText: formatLabelsText(job.labels), persistKind: persistHistoryKind(job), archiveEndpoint: job.persistHistory?.kind === 'archive' ? (job.persistHistory.endpoint || '') : '', @@ -1441,6 +1455,18 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba ), ), h('span', { className: 'dsh-ct-cwdHint' }, t('agentAccessHint')), + h('label', null, t('permissionPreset'), + h('select', { + value: form.permissionPreset || '', + onChange: (e) => setForm({ ...form, permissionPreset: e.target.value }), + }, + h('option', { value: '' }, t('permissionPresetDefault')), + h('option', { value: 'read-only' }, t('permissionPresetReadOnly')), + h('option', { value: 'workspace-write' }, t('permissionPresetWorkspaceWrite')), + h('option', { value: 'danger-full-access' }, t('permissionPresetFullAccess')), + ), + ), + h('span', { className: 'dsh-ct-cwdHint' }, t('permissionPresetHint')), h('label', { className: 'dsh-ct-check' }, h('input', { type: 'checkbox', @@ -1710,6 +1736,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba : { kind: 'dsh' }, mirrorToSession: form.mirrorToSession === true, agentAccess: form.agentAccess === 'deny' ? 'deny' : 'allow', + permissionPreset: form.permissionPreset || '', labels: parseLabelsText(form.labelsText), persistHistory: form.persistKind === 'forever' ? { kind: 'forever' } diff --git a/lib/fire.js b/lib/fire.js index 846a8d2..ac875d8 100644 --- a/lib/fire.js +++ b/lib/fire.js @@ -109,6 +109,7 @@ export function publicJob(job, runs = null) { reasoningEffort: job.reasoningEffort || '', agentPreset: job.agentPreset || '', agentAccess: job.agentAccess === 'deny' ? 'deny' : 'allow', + permissionPreset: job.permissionPreset || '', ownerEmpNo: job.ownerEmpNo || '', ownerDisplayName: job.ownerDisplayName || '', delivery: job.delivery && job.delivery.kind === 'im' diff --git a/lib/host.js b/lib/host.js index 999e6bb..4443c98 100644 --- a/lib/host.js +++ b/lib/host.js @@ -633,6 +633,13 @@ export function createHostService(options = {}) { agentAccess: patch.agentAccess !== undefined ? patch.agentAccess : (patch.agent_access !== undefined ? patch.agent_access : job.agentAccess), + permissionPreset: patch.permissionPreset !== undefined + ? patch.permissionPreset + : (patch.permission_preset !== undefined + ? patch.permission_preset + : (patch.accessMode !== undefined + ? patch.accessMode + : (patch.access_mode !== undefined ? patch.access_mode : job.permissionPreset))), delivery: patch.delivery !== undefined ? mergeDeliveryMention(job.delivery, patch.delivery) : job.delivery, @@ -1307,7 +1314,7 @@ export function createHostService(options = {}) { 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') { + 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' || code === 'INVALID_PERMISSION_PRESET') { 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)) @@ -1968,6 +1975,15 @@ export function makeLiveSessionPort(ctx) { if (!agent || typeof agent.followup !== 'function') { throw new Error('created agent has no followup') } + try { + const wanted = typeof job?.permissionPreset === 'string' ? job.permissionPreset.trim() : '' + const permissionPresets = tryGet(ctx, 'permissionPresets') + if (wanted && permissionPresets && typeof permissionPresets.set === 'function' && agent.session) { + permissionPresets.set(agent.session, wanted) + } + } catch (error) { + console.warn?.(`[dsh-ops-cron] permission preset apply failed for job ${job?.id}: ${error instanceof Error ? error.message : error}`) + } const message = { id: randomUUID(), role: 'user', diff --git a/lib/index.js b/lib/index.js index 66e8314..aadfe65 100644 --- a/lib/index.js +++ b/lib/index.js @@ -48,7 +48,7 @@ export { workspaceVisibleIds, } from './isolation.js' export { wrapScheduledPrompt } from './prompt.js' -export { normalizeJobModel, normalizeAgentAccess, agentMayMutateJob, splitProviderModel } from './store.js' +export { normalizeJobModel, normalizeAgentAccess, agentMayMutateJob, normalizePermissionPreset, PERMISSION_PRESET_IDS, splitProviderModel } from './store.js' export { assertCanAccessJob, canViewAllJobs, diff --git a/lib/store.js b/lib/store.js index a332c12..46ca589 100644 --- a/lib/store.js +++ b/lib/store.js @@ -97,6 +97,38 @@ export function agentMayMutateJob(job) { return normalizeAgentAccess(job) === 'allow' } +/** Host permission-preset ids (DSH 访问模式). Empty = inherit Host default for new sessions. */ +export const PERMISSION_PRESET_IDS = Object.freeze([ + 'read-only', + 'workspace-write', + 'danger-full-access', +]) + +/** + * @param {unknown} input job or raw value + * @returns {''|'read-only'|'workspace-write'|'danger-full-access'} + */ +export function normalizePermissionPreset(input) { + const raw = input && typeof input === 'object' && !Array.isArray(input) + ? (input.permissionPreset ?? input.permission_preset ?? input.accessMode ?? input.access_mode) + : input + const value = String(raw || '').trim().toLowerCase() + if (!value || value === 'default' || value === 'inherit' || value === 'host') return '' + if (value === 'readonly' || value === 'read_only' || value === 'view' || value === 'read-only') { + return 'read-only' + } + if (value === 'workspace' || value === 'workspace_write' || value === 'write' || value === 'workspace-write') { + return 'workspace-write' + } + if (value === 'full' || value === 'fullaccess' || value === 'full_access' || value === 'danger-full-access') { + return 'danger-full-access' + } + if (PERMISSION_PRESET_IDS.includes(value)) return value + const error = new Error(`permissionPreset must be empty|read-only|workspace-write|danger-full-access (got ${JSON.stringify(raw)})`) + error.code = 'INVALID_PERMISSION_PRESET' + throw error +} + /** * Repair provider/model pairs when LLMs or hosts pass a combined "provider/model" * route in one or both fields (e.g. provider=model="zte/Qwen3-…" → zte + Qwen3-…). @@ -190,6 +222,7 @@ export function createJobRecord(input, state, now) { settings.historyLimit, ) const agentAccess = normalizeAgentAccess(input) + const permissionPreset = normalizePermissionPreset(input) const job = { id: String(input.id || newId()), name, @@ -204,6 +237,7 @@ export function createJobRecord(input, state, now) { reasoningEffort: model.reasoningEffort, agentPreset, agentAccess, + permissionPreset, delivery, mirrorToSession, ownerEmpNo, diff --git a/lib/tools.js b/lib/tools.js index 22bf5a9..2c65723 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -219,7 +219,7 @@ function jobLine(job) { /** Strip human-only fields so agents cannot see or game panel-only controls. */ function forAgentView(job) { if (!job || typeof job !== 'object') return job - const { agentAccess: _hidden, ...rest } = job + const { agentAccess: _a, permissionPreset: _p, ...rest } = job return rest } diff --git a/test/host.test.js b/test/host.test.js index 71988b7..0a0a0b6 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -596,6 +596,51 @@ test('makeLiveSessionPort creates with default model and setup', async () => { assert.equal(assembled.variables.model, 'deepseek-chat') }) +test('makeLiveSessionPort applies job permissionPreset when Host supports it', async () => { + const applied = [] + const fakeSession = { id: 'sess-perm' } + const ctx = { + get(name) { + if (name === 'agentDefaultModel') { + return { currentSelection: () => ({ provider: 'deepseek', model: 'deepseek-chat' }) } + } + if (name === 'permissionPresets') { + return { + set(session, preset) { + applied.push({ session, preset }) + }, + } + } + if (name === 'agents') { + return { + async create() { + const agent = createFakeAgent() + agent.session = fakeSession + return { agent, dispose: async () => {} } + }, + get() {}, + } + } + return undefined + }, + } + const port = makeLiveSessionPort(ctx) + await port.createAndPrompt({ + job: { name: 'elevated', timeoutMinutes: 1, permissionPreset: 'danger-full-access' }, + run: {}, + text: 'hi', + }) + assert.deepEqual(applied, [{ session: fakeSession, preset: 'danger-full-access' }]) + + applied.length = 0 + await port.createAndPrompt({ + job: { name: 'default', timeoutMinutes: 1, permissionPreset: '' }, + run: {}, + text: 'hi', + }) + assert.deepEqual(applied, []) +}) + test('resolveDefaultModel fails loud when Models has no selection', async () => { await assert.rejects( () => resolveDefaultModel({ get: () => undefined }, { waitMs: 0 }), diff --git a/test/tools.test.js b/test/tools.test.js index 44e4daf..49d297b 100644 --- a/test/tools.test.js +++ b/test/tools.test.js @@ -473,13 +473,15 @@ test('agentAccess deny is hidden from tools and blocks pause/delete', async (t) timezone: 'Asia/Shanghai', }, {}) assert.equal(created.job.agentAccess, undefined) + assert.equal(created.job.permissionPreset, undefined) assert.equal(created.job.id != null, true) - await service.updateJob(created.job.id, { agentAccess: 'deny' }) + await service.updateJob(created.job.id, { agentAccess: 'deny', permissionPreset: 'danger-full-access' }) const listed = await tools.cron_list.execute({}, {}) const row = listed.jobs.find((job) => job.id === created.job.id) assert.ok(row) assert.equal(row.agentAccess, undefined) + assert.equal(row.permissionPreset, undefined) await assert.rejects( () => tools.cron_pause.execute({ id: created.job.id }, {}), @@ -492,6 +494,17 @@ test('agentAccess deny is hidden from tools and blocks pause/delete', async (t) const httpView = await service.getJob(created.job.id) assert.equal(httpView.agentAccess, 'deny') + assert.equal(httpView.permissionPreset, 'danger-full-access') +}) + +test('normalizePermissionPreset maps UI aliases and rejects junk', async () => { + const { normalizePermissionPreset } = await import('../lib/store.js') + assert.equal(normalizePermissionPreset(''), '') + assert.equal(normalizePermissionPreset({ permissionPreset: 'inherit' }), '') + assert.equal(normalizePermissionPreset('read-only'), 'read-only') + assert.equal(normalizePermissionPreset('workspace-write'), 'workspace-write') + assert.equal(normalizePermissionPreset('full'), 'danger-full-access') + assert.throws(() => normalizePermissionPreset('nope'), (err) => err.code === 'INVALID_PERMISSION_PRESET') }) test('registerCronTools registers each definition and disposer unregisters', () => { From b22048d37727539bb953d8a60bd611fcf7d6301f Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 11 Sep 2026 07:00:29 +0800 Subject: [PATCH 08/10] Fix one-shot jobs that never auto-fire after restart or slow ticks. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit recover was clearing past-due at nextRunAt via nextFire→null; keep and re-arm them, and evaluate schedule misfire at tick start so a slow prior job cannot age peers out of grace. Co-authored-by: Cursor --- lib/host.js | 17 +++++++++--- test/host.test.js | 68 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 82 insertions(+), 3 deletions(-) diff --git a/lib/host.js b/lib/host.js index 4443c98..6544993 100644 --- a/lib/host.js +++ b/lib/host.js @@ -842,9 +842,9 @@ export function createHostService(options = {}) { return state.settings } - async function dispatchRun(jobId, trigger) { + async function dispatchRun(jobId, trigger, opts = {}) { let claimed - const t = now() + const t = Number.isFinite(opts.at) ? opts.at : now() await withState((current) => { claimed = claimOccurrence(current, jobId, t, trigger, current.settings) if (claimed.decision.action === 'wait') return current @@ -897,7 +897,9 @@ export function createHostService(options = {}) { misfirePolicy: state.settings.misfirePolicy, }) if (decision.action === 'wait') continue - const result = await dispatchRun(job.id, 'schedule') + // Use tick start time for claim/misfire so a slow earlier job cannot + // push later oneshots past the grace window. + const result = await dispatchRun(job.id, 'schedule', { at: t }) if (result.run) fired.push(result) } await tickWatcherReports() @@ -967,6 +969,15 @@ export function createHostService(options = {}) { try { if (job.nextRunAt === null) return job if (Number.isFinite(job.nextRunAt) && job.nextRunAt > t) return job + // Past-due one-shot: do NOT call nextFire (that returns null and + // permanently kills the job). Re-arm to now so the next tick fires + // a catch-up run instead of wiping or grace-misfiring after restart. + if (job.schedule?.kind === 'at') { + if (Number.isFinite(job.nextRunAt) && job.nextRunAt <= t) { + return { ...job, nextRunAt: t } + } + return job + } const nextRunAt = nextFire(job.schedule, t, job.schedule?.timezone || next.settings.timezone) return { ...job, nextRunAt } } catch { diff --git a/test/host.test.js b/test/host.test.js index 0a0a0b6..df7d367 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -718,6 +718,74 @@ test('one-shot at job fires once then later ticks wait instead of retriggering', assert.equal(misfires.length, 0) }) +test('recover re-arms past-due oneshot instead of clearing nextRunAt', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const at = Date.parse('2026-09-11T00:00:00.000Z') + let clock = at - 30_000 + let sessions = 0 + const service = createTestHost({ + filePath: join(dir, 'store.json'), + now: () => clock, + sessionPort: { + async createAndPrompt() { + sessions += 1 + return { sessionId: `catchup-${sessions}`, status: 'succeeded', summary: 'catch-up' } + }, + }, + }) + const job = await service.createJob({ + name: 'catchup-at', + prompt: 'ping', + schedule: { kind: 'at', at: new Date(at).toISOString(), timezone: 'UTC' }, + }) + assert.equal(job.nextRunAt, at) + + // Simulate Host restart after the due time (previously wiped nextRunAt via nextFire→null). + clock = at + 90_000 + await service.recover() + const after = await service.getJob(job.id) + assert.equal(after.nextRunAt, clock) + + const fired = await service.tick() + assert.equal(fired.length, 1) + assert.equal(fired[0].run.status, 'succeeded') + assert.equal(sessions, 1) +}) + +test('tick evaluates misfire against tick-start time so a slow prior job cannot age out oneshots', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const due = Date.parse('2026-09-11T01:00:00.000Z') + let clock = due + let sessions = 0 + const service = createTestHost({ + filePath: join(dir, 'store.json'), + now: () => clock, + sessionPort: { + async createAndPrompt({ job }) { + sessions += 1 + // First job stretches wall clock past oneshot grace (60s). + if (job.name === 'slow') clock = due + 90_000 + return { sessionId: `s-${sessions}`, status: 'succeeded', summary: 'ok' } + }, + }, + }) + await service.createJob({ + name: 'slow', + prompt: 'first', + schedule: { kind: 'at', at: new Date(due).toISOString(), timezone: 'UTC' }, + }) + await service.createJob({ + name: 'peer', + prompt: 'second', + schedule: { kind: 'at', at: new Date(due).toISOString(), timezone: 'UTC' }, + }) + const fired = await service.tick() + assert.equal(fired.filter((row) => row.run?.status === 'succeeded').length, 2) + assert.equal(sessions, 2) +}) + test('unarchiveSession drops the id from archivedSessionIds', async () => { const state = { workspaceIds: ['w'], archivedSessionIds: ['s1', 's2'], initialized: true } const ctx = { From 89ab5f59afb8894e8a41618fca923d8493d61127 Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 11 Sep 2026 19:57:22 +0800 Subject: [PATCH 09/10] Add cron_reschedule so pending one-shots can be moved without delete/recreate. Keeps retrigger for consumed jobs only; agents can advance or delay nextRunAt safely. Co-authored-by: Cursor --- docs/cron_reschedule原语详设.md | 105 +++++++++++++++++++++++++++++ lib/host.js | 113 +++++++++++++++++++++++++++++++- lib/i18n.js | 2 + lib/tools.js | 66 +++++++++++++++++-- test/host.test.js | 104 +++++++++++++++++++++++++++++ test/tools.test.js | 2 + 6 files changed, 382 insertions(+), 10 deletions(-) create mode 100644 docs/cron_reschedule原语详设.md diff --git a/docs/cron_reschedule原语详设.md b/docs/cron_reschedule原语详设.md new file mode 100644 index 0000000..ab7eaaf --- /dev/null +++ b/docs/cron_reschedule原语详设.md @@ -0,0 +1,105 @@ +# cron_reschedule 原语详设 + +**状态**:已在本包落地(`rescheduleJob` + `cron_reschedule` + `POST /jobs/:id/reschedule`) +**包**:`dsh-ops-cron`(本仓库可读源码,非混淆服务) +**背景**:一次性 `at` 任务在仍有 `nextRunAt`(pending)时,`cron_retrigger` 会拒绝;Agent 只能删重建,体验差。 + +--- + +## 1. 结论 + +**改 pending 可以做,且应做。** +语义是「把预定的那一次触发改期/提前」,不是「再开一次并发 run」。比放开 `cron_retrigger` 更精确、更安全。 + +| 方案 | 做法 | 评价 | +|------|------|------| +| A | 改 `cron_retrigger`:有 pending 时也接受,行为=抢占改期 | 可行,但混淆「已消费再触发」与「未到期改期」 | +| **B(推荐)** | **新增 `cron_reschedule`** | 不碰 retrigger 语义;最小侵入、向后兼容 | + +--- + +## 2. 现状(代码事实) + +- `isRetriggerable` / `retriggerJob`:one-shot 仅当 `nextRunAt == null` 且 `lastStatus` 终态才可 retrigger(`lib/monitor.js`、`lib/host.js`)。 +- pending 拒绝文案已写明:`one-shot still has a pending next run; wait or edit the schedule instead of retrigger`。 +- **侧栏 / HTTP 已能改期**:`PATCH /jobs/:id` → `updateJob({ schedule })`;若 patch 了 `schedule`,会经 `createJobRecord` **重算 `nextRunAt`**(`host.js` 约 698–700 行:仅当未改 schedule/enabled 才保留旧 next)。 +- **缺口**:Agent 工具面没有 `cron_update` / `cron_reschedule`,只有 create / list / pause / resume / retrigger / delete。 + +因此:不必「热改运行中混淆服务」;在本包加工具(可选再加 host 薄封装)即可落地。 + +--- + +## 3. 推荐 API:`cron_reschedule` + +### 3.1 签名 + +```text +cron_reschedule({ + id | task_id: string, // 必填 + after_minutes?: number, // ≥0;与 at 二选一(优先 after_minutes) + at?: string, // ISO 或与 cron_create 一致的本地时间语义 + timezone?: string // 可选;默认沿用 job.schedule.timezone +}) +``` + +### 3.2 行为 + +1. 仅 `schedule.kind === 'at'`。 +2. 仅当 **仍有 pending next**(`nextRunAt != null`)且 **无 active run**(queued/running)。 +3. 计算新触发时刻 `T`(`after_minutes=0` → 立即 due;`>0` → now+N;`at` → 解析后的绝对时间)。 +4. `T` 必须 ≥ now(允许 0 表示立刻进入 due,由现有 tick/`run-now` 路径消费;不要另开第二条 pending)。 +5. 写回: + - `schedule = { kind: 'at', at: ISO(T), timezone }` + - `nextRunAt = T` + - `enabled` 保持不变(paused 允许改 next,**不隐式 resume**) +6. **不**调用 `dispatchRun`(除非产品明确要 `fire_now=true`;默认不火,交给调度器)。 +7. 返回 `{ ok, job, nextRunAt, previousNextRunAt }`。 + +### 3.3 拒绝条件(稳定 error.code) + +| 条件 | code | +|------|------| +| 非 one-shot | `INVALID_RESCHEDULE` | +| `nextRunAt == null`(已消费) | `INVALID_RESCHEDULE` → 提示改用 `cron_retrigger` | +| 已有 queued/running | `ALREADY_RUNNING` | +| 时间非法 / 过去 | `INVALID_RESCHEDULE` | +| 周期 cron | `INVALID_RESCHEDULE`(周期改期另议,勿混进本原语) | + +### 3.4 与 retrigger 分工(写进 tool description + skill) + +| 状态 | 用哪个 | +|------|--------| +| one-shot 已跑完(`next=n/a`,`retriggerable=true`) | `cron_retrigger` | +| one-shot 仍在等(有 `nextRunAt`) | **`cron_reschedule`** | +| recurring 立刻多跑一次 | `cron_retrigger`(不改 cron 表达式) | + +--- + +## 4. 实现落点(本仓库) + +1. **`lib/host.js`**:新增 `rescheduleJob(jobId, opts, identity)`(或在 `updateJob` 外包一层校验);可复用 `scheduleFromArgs` / `nextFire`。 +2. **`lib/tools.js`**:注册 `cron_reschedule`;更新 `scheduled-tasks` skill 文案(禁止删重建来改期)。 +3. **HTTP(可选)**:`POST /jobs/:id/reschedule`,与 PATCH 并存;Agent 走工具即可。 +4. **测试**:`host.test.js` / `tools.test.js` + - pending → after_minutes=1 → `nextRunAt` 前移 + - pending + active run → 拒绝 + - consumed → 拒绝并指向 retrigger + - cron kind → 拒绝 + - paused pending → 只改 next,仍 `enabled=false` + +工作量小:调度内核已支持「改 schedule ⇒ 新 next」;主要是 **Agent 可发现原语 + 边界校验**。 + +--- + +## 5. 明确不做 + +- 不对 pending one-shot 放开无条件 `cron_retrigger`(避免与「已消费再触发」心智冲突)。 +- 不引入第二条并发 pending。 +- 不把 reschedule 做成隐式 start(paused 保持 paused)。 + +--- + +## 6. 临时绕过(原语未上线前) + +- **人**:侧栏编辑任务时间 → Save(已走 PATCH)。 +- **Agent**:无工具时只能 delete + create(应在 skill 里标明为 workaround,待 `cron_reschedule` 替换)。 diff --git a/lib/host.js b/lib/host.js index 6544993..1d0e631 100644 --- a/lib/host.js +++ b/lib/host.js @@ -766,7 +766,7 @@ export function createHostService(options = {}) { // one-shot if (job.nextRunAt != null) { - const error = new Error('one-shot still has a pending next run; wait or edit the schedule instead of retrigger') + const error = new Error('one-shot still has a pending next run; use cron_reschedule (or PATCH schedule) instead of retrigger') error.code = 'INVALID_RETRIGGER' throw error } @@ -822,6 +822,106 @@ export function createHostService(options = {}) { } } + /** + * Move a pending one-shot's next fire time (does not start a second concurrent run). + * Prefer this over delete+create when nextRunAt is still set. + * Does not change enabled (paused stays paused). Does not dispatch. + */ + async function rescheduleJob(jobId, opts = {}, identity = null) { + const afterRaw = opts.afterMinutes ?? opts.after_minutes + const hasAfter = afterRaw !== undefined && afterRaw !== null && afterRaw !== '' + const afterMinutes = hasAfter ? Number(afterRaw) : null + const atRaw = typeof opts.at === 'string' ? opts.at.trim() : '' + const tzOpt = typeof opts.timezone === 'string' && opts.timezone.trim() + ? opts.timezone.trim() + : (typeof opts.time_zone === 'string' && opts.time_zone.trim() ? opts.time_zone.trim() : '') + + if (hasAfter && (!Number.isFinite(afterMinutes) || afterMinutes < 0)) { + const error = new Error('after_minutes must be >= 0') + error.code = 'INVALID_RESCHEDULE' + throw error + } + if (!hasAfter && !atRaw) { + const error = new Error('provide after_minutes or at') + error.code = 'INVALID_RESCHEDULE' + throw error + } + + const t = now() + let previousNextRunAt = null + + 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_reschedule only applies to one-shot (at) jobs') + error.code = 'INVALID_RESCHEDULE' + throw error + } + if (job.nextRunAt == null) { + const error = new Error('one-shot has no pending next run; use cron_retrigger after it finishes') + error.code = 'INVALID_RESCHEDULE' + 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 + } + + const timezone = tzOpt + || job.schedule?.timezone + || current.settings?.timezone + || 'Asia/Shanghai' + let atMs + if (hasAfter) { + atMs = t + Math.round(afterMinutes * 60_000) + } else { + const validated = validateSchedule({ kind: 'at', at: atRaw, timezone }, timezone) + atMs = validated.at + } + if (!Number.isFinite(atMs)) { + const error = new Error('invalid reschedule time') + error.code = 'INVALID_RESCHEDULE' + throw error + } + // Align with createJobRecord: slight past (≤60s) → fire ASAP; older → reject. + if (atMs < t - 60_000) { + const error = new Error(`reschedule time is already in the past (now is ${new Date(t).toISOString()})`) + error.code = 'INVALID_RESCHEDULE' + throw error + } + if (atMs < t) atMs = t + + previousNextRunAt = job.nextRunAt + return upsertJob(current, { + ...job, + schedule: { + kind: 'at', + at: new Date(atMs).toISOString(), + timezone, + }, + nextRunAt: atMs, + updatedAt: t, + }) + }) + + const state = await snapshot() + const job = getJob(state, jobId) + return { + ok: true, + job: jobView(job, state.runs), + nextRunAt: job?.nextRunAt ?? null, + previousNextRunAt, + } + } + async function deleteJob(jobId) { await withState((current) => { if (!getJob(current, jobId)) { @@ -1135,7 +1235,7 @@ export function createHostService(options = {}) { return } - const jobMatch = path.match(new RegExp(`^${API_PREFIX}/jobs/([^/]+)(/run|/pause|/resume|/retrigger)?$`)) + const jobMatch = path.match(new RegExp(`^${API_PREFIX}/jobs/([^/]+)(/run|/pause|/resume|/retrigger|/reschedule)?$`)) if (jobMatch) { const jobId = decodeURIComponent(jobMatch[1]) const rest = jobMatch[2] || '' @@ -1170,6 +1270,12 @@ export function createHostService(options = {}) { write(200, result) return } + if (method === 'POST' && rest === '/reschedule') { + const body = await readJsonBody(req).catch(() => ({})) + const result = await rescheduleJob(jobId, body || {}, identity) + write(200, result) + return + } if (method === 'POST' && rest === '/pause') { const job = await pauseJob(jobId, false, identity) write(200, { ok: true, job }) @@ -1325,7 +1431,7 @@ export function createHostService(options = {}) { 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' || code === 'INVALID_PERMISSION_PRESET') { + 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' || code === 'INVALID_RESCHEDULE' || code === 'INVALID_PERMISSION_PRESET') { 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)) @@ -1498,6 +1604,7 @@ export function createHostService(options = {}) { updateJob, pauseJob, retriggerJob, + rescheduleJob, deleteJob, updateSettings, dispatchRun, diff --git a/lib/i18n.js b/lib/i18n.js index f696213..b580edc 100644 --- a/lib/i18n.js +++ b/lib/i18n.js @@ -32,6 +32,7 @@ export const MESSAGES = { 'tool.pause': '暂停定时任务', 'tool.resume': '恢复定时任务', 'tool.retrigger': '重新触发一次性任务', + 'tool.reschedule': '改期一次性任务', 'tool.delete': '删除定时任务', }, en: { @@ -53,6 +54,7 @@ export const MESSAGES = { 'tool.pause': 'Pause scheduled task', 'tool.resume': 'Resume scheduled task', 'tool.retrigger': 'Retrigger one-shot task', + 'tool.reschedule': 'Reschedule pending one-shot', 'tool.delete': 'Delete scheduled task', }, } diff --git a/lib/tools.js b/lib/tools.js index 2c65723..b0aa3fa 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -746,7 +746,7 @@ export function cronToolDefinitions(service, deps = {}) { }, { name: 'cron_retrigger', - description: 'Start an extra run now: for recurring (cron) jobs fires immediately without changing the schedule; for consumed one-shots (next=n/a) re-arms/fires a new run. Auto-enables if paused. after_minutes>0 only for one-shots. Rejects if already pending/running. Listeners: when retriggerable=true and progress done_flag is false, call this instead of creating a duplicate job.', + description: 'Start an extra run now: for recurring (cron) jobs fires immediately without changing the schedule; for consumed one-shots (next=n/a) re-arms/fires a new run. Auto-enables if paused. after_minutes>0 only for one-shots. Rejects if already pending/running OR if a one-shot still has a pending next (use cron_reschedule). Listeners: when retriggerable=true and progress done_flag is false, call this instead of creating a duplicate job.', parameters: { type: 'object', additionalProperties: false, @@ -796,6 +796,58 @@ export function cronToolDefinitions(service, deps = {}) { return { ...result, job: forAgentView(result.job) } }, }, + { + name: 'cron_reschedule', + description: 'Move a pending one-shot next fire time (still has next≠n/a). Prefer after_minutes (0=ASAP / next tick) or at. Does NOT start a second concurrent run; does NOT auto-resume paused jobs; does NOT replace cron_retrigger for consumed oneshots. Use this instead of delete+create when the user wants earlier/later.', + parameters: { + type: 'object', + additionalProperties: false, + properties: { + task_id: { type: 'string', description: 'Job id (alias: id).' }, + id: { type: 'string', description: 'Alias of task_id.' }, + after_minutes: { type: 'number', description: 'Delay from now before fire. 0 = due immediately. Preferred over at.' }, + at: { type: 'string', description: 'ISO timestamp or HH:mm in job timezone. Ignored if after_minutes is set.' }, + timezone: { type: 'string', description: 'Optional IANA tz; default keeps the job timezone.' }, + }, + }, + output: { + schema: { + type: 'object', + additionalProperties: true, + properties: { + ok: { type: 'boolean' }, + job: JOB_SCHEMA, + nextRunAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + previousNextRunAt: { oneOf: [{ type: 'number' }, { type: 'null' }] }, + }, + }, + render: (_args, value) => { + const name = value.job?.name || value.job?.id || '' + const when = value.nextRunAt ? new Date(value.nextRunAt).toISOString() : 'n/a' + return text(`Rescheduled "${name}" → next ${when}.`) + }, + }, + presentCall: (args) => ({ + card: 'generic', + title: t('tool.reschedule'), + 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') + const owned = await requireOwnedJob(service, id, peer, identity) + assertAgentMayMutate(owned) + const result = await service.rescheduleJob(id, { + after_minutes: args?.after_minutes, + at: args?.at, + timezone: args?.timezone, + }, identity) + return { ...result, job: forAgentView(result.job) } + }, + }, { name: 'cron_delete', description: 'Permanently delete a scheduled task by id. History rows for that job remain until pruned.', @@ -851,9 +903,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.', - 'One-shot wake: cron_resume only unpauses. For a finished one-shot (retriggerable=true, next=n/a), call cron_retrigger. For a recurring job, cron_retrigger also starts one extra run immediately without changing the cron schedule.', - '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.', + 'One-shot wake: cron_resume only unpauses. For a finished one-shot (retriggerable=true, next=n/a), call cron_retrigger. For a still-pending one-shot (has next), call cron_reschedule(after_minutes=N) instead of delete+create. For a recurring job, cron_retrigger starts one extra run immediately without changing the cron schedule.', + 'When the user asks to look at, create, pause, resume, retrigger, reschedule, or delete 定时任务 / scheduled tasks / cron jobs:', + '1. If cron_list / cron_create / cron_pause / cron_resume / cron_retrigger / cron_reschedule / 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.', @@ -865,8 +917,8 @@ export const CRON_GUIDANCE = cronGuidanceText() export function makeCronSkill() { return { name: 'scheduled-tasks', - description: '定时任务: 查看、创建、暂停、恢复、删除 DSH 侧栏定时任务(今天晚上几点、一次性执行、scheduled job、cron)。WhatsApp 里创建默认回投 IM;Web 里创建进侧栏会话。不要用系统 crontab,也不是会话内 reminder。', - whenToUse: 'User asks to list or create 定时任务 / scheduled tasks, schedule something for tonight/today, pause a job, or mentions cron in DeepSeek Harness.', + description: '定时任务: 查看、创建、暂停、恢复、改期、删除 DSH 侧栏定时任务(今天晚上几点、一次性执行、scheduled job、cron)。WhatsApp 里创建默认回投 IM;Web 里创建进侧栏会话。不要用系统 crontab,也不是会话内 reminder。', + whenToUse: 'User asks to list or create 定时任务 / scheduled tasks, schedule something for tonight/today, pause/reschedule a job, or mentions cron in DeepSeek Harness.', source: 'runtime', provider: 'runtime', content: `${cronGuidanceText()} @@ -877,7 +929,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_retrigger / cron_delete — by id from cron_list. Use cron_retrigger to re-fire a consumed one-shot or to immediately start one extra run of a recurring job. +- cron_pause / cron_resume / cron_retrigger / cron_reschedule / cron_delete — by id from cron_list. Use cron_retrigger to re-fire a consumed one-shot or to immediately start one extra run of a recurring job. Use cron_reschedule to move a still-pending one-shot (do not delete+create). `, } } diff --git a/test/host.test.js b/test/host.test.js index df7d367..df1879f 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -409,6 +409,110 @@ test('retriggerJob fires recurring immediately and re-arms consumed oneshot', as assert.ok(httpHit.body.run?.id) }) +test('rescheduleJob moves pending one-shot next without dispatch', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) + t.after(() => rm(dir, { recursive: true, force: true })) + let clock = Date.parse('2026-08-24T01:00:00.000Z') + let fires = 0 + const service = createTestHost({ + filePath: join(dir, 'store.json'), + now: () => clock, + sessionPort: { + async createAndPrompt() { + fires += 1 + return { sessionId: `sess-${fires}`, status: 'succeeded', summary: 'ok' } + }, + async archiveSession() {}, + }, + }) + + const pending = await service.createJob({ + name: 'later', + prompt: 'do it', + enabled: false, + schedule: { kind: 'at', at: new Date(clock + 3600_000).toISOString(), timezone: 'UTC' }, + }) + const previous = pending.nextRunAt + assert.ok(previous > clock) + + const moved = await service.rescheduleJob(pending.id, { after_minutes: 1 }) + assert.equal(moved.ok, true) + assert.equal(moved.previousNextRunAt, previous) + assert.equal(moved.nextRunAt, clock + 60_000) + assert.equal(moved.job.enabled, false) + assert.equal(moved.job.nextRunAt, clock + 60_000) + assert.equal(fires, 0) + + const asap = await service.rescheduleJob(pending.id, { after_minutes: 0 }) + assert.equal(asap.nextRunAt, clock) + assert.equal(asap.job.enabled, false) + + await assert.rejects( + () => service.rescheduleJob(pending.id, {}), + (err) => err.code === 'INVALID_RESCHEDULE', + ) + + const cron = await service.createJob({ + name: 'cron', + prompt: 'loop', + schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'UTC' }, + }) + await assert.rejects( + () => service.rescheduleJob(cron.id, { after_minutes: 1 }), + (err) => err.code === 'INVALID_RESCHEDULE', + ) + + await service.store.mutate((state) => ({ + ...state, + jobs: state.jobs.map((row) => (row.id === pending.id + ? { ...row, nextRunAt: null, lastStatus: 'succeeded' } + : row)), + })) + await assert.rejects( + () => service.rescheduleJob(pending.id, { after_minutes: 1 }), + (err) => err.code === 'INVALID_RESCHEDULE', + ) + + const waiting = await service.createJob({ + name: 'waiting2', + prompt: 'soon', + cwd: '/tmp/user-workspaces/tester', + schedule: { kind: 'at', at: new Date(clock + 7200_000).toISOString(), timezone: 'UTC' }, + }, { empNo: 'tester', displayName: 'Tester', permissions: { canViewAllSessions: false } }) + await service.store.mutate((state) => ({ + ...state, + runs: [ + { + id: 'inflight-rs', + jobId: waiting.id, + status: 'running', + scheduledAt: clock, + actualAt: clock, + stateEnteredAt: clock, + }, + ...(state.runs || []), + ], + })) + await assert.rejects( + () => service.rescheduleJob(waiting.id, { after_minutes: 1 }), + (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-rs'), + })) + const httpHit = await jsonRequest(http.url, `/dsh-ops-cron/jobs/${waiting.id}/reschedule`, { + method: 'POST', + body: JSON.stringify({ after_minutes: 2 }), + }) + assert.equal(httpHit.status, 200) + assert.equal(httpHit.body.nextRunAt, clock + 120_000) + assert.equal(httpHit.body.ok, true) +}) + test('http api: anonymous list/run-now is rejected', async (t) => { const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) t.after(() => rm(dir, { recursive: true, force: true })) diff --git a/test/tools.test.js b/test/tools.test.js index 49d297b..1b3e872 100644 --- a/test/tools.test.js +++ b/test/tools.test.js @@ -530,6 +530,7 @@ test('registerCronTools registers each definition and disposer unregisters', () 'cron_pause', 'cron_resume', 'cron_retrigger', + 'cron_reschedule', 'cron_delete', ]) off() @@ -546,6 +547,7 @@ test('cron tool output schemas never use type arrays (Host rejects them)', () => async pauseJob() { return {} }, async resumeJob() { return {} }, async retriggerJob() { return {} }, + async rescheduleJob() { return {} }, async deleteJob() { return {} }, }) const bad = [] From a41a2b53c96be75c07af135e5ec076c66039ca2e Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 11 Sep 2026 20:54:40 +0800 Subject: [PATCH 10/10] Make cron_retrigger fire-and-forget so callers do not wait for the run. Detached settle still finishes history and delivery in the background. Co-authored-by: Cursor --- lib/host.js | 62 +++++++++++++++++++++++++++++++++++++++++++---- lib/tools.js | 11 ++++++--- test/host.test.js | 38 +++++++++++++++++++++++++++++ 3 files changed, 102 insertions(+), 9 deletions(-) diff --git a/lib/host.js b/lib/host.js index 1d0e631..8cc7f7e 100644 --- a/lib/host.js +++ b/lib/host.js @@ -716,6 +716,7 @@ export function createHostService(options = {}) { * - recurring (cron): after_minutes must be 0/omitted → enable + immediate run-now (keeps cron next) * - one-shot (at): consumed only; 0 → run-now; >0 → reschedule next * Auto-enables paused jobs. Rejects when a run is already pending/running. + * Immediate run-now is fire-and-forget: returns once the session starts; settle/delivery continue in background. */ async function retriggerJob(jobId, opts = {}, identity = null) { const afterRaw = opts.afterMinutes ?? opts.after_minutes @@ -811,10 +812,12 @@ export function createHostService(options = {}) { } } - const result = await dispatchRun(jobId, 'run-now') + const result = await dispatchRun(jobId, 'run-now', { wait: false }) + const detached = result.detached === true + || (result.run && ACTIVE_RUN_STATUSES.has(result.run.status)) return { ok: true, - mode, + mode: detached ? 'started' : 'immediate', job: result.job, run: result.run, decision: result.decision, @@ -945,6 +948,7 @@ export function createHostService(options = {}) { async function dispatchRun(jobId, trigger, opts = {}) { let claimed const t = Number.isFinite(opts.at) ? opts.at : now() + const waitForCompletion = opts.wait !== false await withState((current) => { claimed = claimOccurrence(current, jobId, t, trigger, current.settings) if (claimed.decision.action === 'wait') return current @@ -964,6 +968,19 @@ export function createHostService(options = {}) { }) let run = (executed.runs || []).find((row) => row.id === claimed.run.id) if (run?.status === 'running' && typeof sessionPort?.waitForTurn === 'function') { + if (!waitForCompletion) { + const runId = run.id + const sessionId = run.sessionId + const decision = claimed.decision + void trackDetachedSettle(jobId, runId, sessionId, decision) + const state = await snapshot() + return { + job: jobView(getJob(state, jobId), state.runs), + run: runView(run), + decision, + detached: true, + } + } let terminal try { terminal = await sessionPort.waitForTurn(run.sessionId) @@ -978,6 +995,38 @@ export function createHostService(options = {}) { return settleAndNotify(jobId, claimed.run.id, claimed.decision) } + async function settleDetachedRun(jobId, runId, sessionId, claimedDecision) { + let terminal + try { + terminal = await sessionPort.waitForTurn(sessionId) + } catch (error) { + const message = error instanceof Error ? error.message : String(error) + terminal = { status: 'failed', error: message, summary: message } + } + try { + await store.mutate((current) => settleRun(current, runId, terminal, now())) + await settleAndNotify(jobId, runId, claimedDecision) + } catch (error) { + logger.warn?.(`[dsh-ops-cron] detached settle failed for ${runId}: ${error instanceof Error ? error.message : error}`) + } + } + + const detachedSettles = new Set() + + function trackDetachedSettle(jobId, runId, sessionId, claimedDecision) { + const task = settleDetachedRun(jobId, runId, sessionId, claimedDecision) + .catch(() => {}) + .finally(() => { detachedSettles.delete(task) }) + detachedSettles.add(task) + return task + } + + async function flushDetachedRuns() { + while (detachedSettles.size) { + await Promise.allSettled([...detachedSettles]) + } + } + async function tick() { if (ticking) return [] ticking = true @@ -1023,9 +1072,11 @@ export function createHostService(options = {}) { } function stopTimer() { - if (!timer) return - clearInterval(timer) - timer = null + if (timer) { + clearInterval(timer) + timer = null + } + return flushDetachedRuns() } async function concealKnownSessions() { @@ -1613,6 +1664,7 @@ export function createHostService(options = {}) { recover, startTimer, stopTimer, + flushDetachedRuns, handleRequest, snapshot, getUdsAuth, diff --git a/lib/tools.js b/lib/tools.js index b0aa3fa..7a70b32 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -746,14 +746,14 @@ export function cronToolDefinitions(service, deps = {}) { }, { name: 'cron_retrigger', - description: 'Start an extra run now: for recurring (cron) jobs fires immediately without changing the schedule; for consumed one-shots (next=n/a) re-arms/fires a new run. Auto-enables if paused. after_minutes>0 only for one-shots. Rejects if already pending/running OR if a one-shot still has a pending next (use cron_reschedule). Listeners: when retriggerable=true and progress done_flag is false, call this instead of creating a duplicate job.', + description: 'Start an extra run now (fire-and-forget): returns as soon as the run session starts — does NOT wait for the job to finish. Poll cron_runs / cron_progress for status. Recurring (cron): one extra run without changing the schedule. Consumed one-shots (next=n/a): re-arms/fires. Auto-enables if paused. after_minutes>0 only for one-shots (schedules next, no run yet). Rejects if already pending/running OR if a one-shot still has a pending next (use cron_reschedule). Listeners: when retriggerable=true and progress done_flag is false, call this instead of creating a duplicate job.', parameters: { type: 'object', additionalProperties: false, properties: { task_id: { type: 'string', description: 'Job id (alias: id).' }, id: { type: 'string', description: 'Alias of task_id.' }, - after_minutes: { type: 'number', description: 'One-shot only: delay before fire. Omit/0 = immediate. Not allowed for recurring jobs.' }, + after_minutes: { type: 'number', description: 'One-shot only: delay before fire. Omit/0 = start now (returns without waiting). Not allowed for recurring jobs.' }, }, }, output: { @@ -774,6 +774,9 @@ export function cronToolDefinitions(service, deps = {}) { const when = value.nextRunAt ? new Date(value.nextRunAt).toISOString() : 'pending' return text(`Retrigger scheduled "${name}" → next ${when}.`) } + if (value.mode === 'started') { + return text(`Started "${name}"${value.run?.id ? ` (run ${value.run.id})` : ''} — running in background.`) + } return text(`Retriggered "${name}"${value.run?.id ? ` (run ${value.run.id})` : ''}.`) }, }, @@ -903,7 +906,7 @@ 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.', - 'One-shot wake: cron_resume only unpauses. For a finished one-shot (retriggerable=true, next=n/a), call cron_retrigger. For a still-pending one-shot (has next), call cron_reschedule(after_minutes=N) instead of delete+create. For a recurring job, cron_retrigger starts one extra run immediately without changing the cron schedule.', + 'One-shot wake: cron_resume only unpauses. For a finished one-shot (retriggerable=true, next=n/a), call cron_retrigger (fire-and-forget — do not wait for it to finish; poll cron_runs). For a still-pending one-shot (has next), call cron_reschedule(after_minutes=N) instead of delete+create. For a recurring job, cron_retrigger starts one extra run immediately without changing the cron schedule.', 'When the user asks to look at, create, pause, resume, retrigger, reschedule, or delete 定时任务 / scheduled tasks / cron jobs:', '1. If cron_list / cron_create / cron_pause / cron_resume / cron_retrigger / cron_reschedule / 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.', @@ -929,7 +932,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_retrigger / cron_reschedule / cron_delete — by id from cron_list. Use cron_retrigger to re-fire a consumed one-shot or to immediately start one extra run of a recurring job. Use cron_reschedule to move a still-pending one-shot (do not delete+create). +- cron_pause / cron_resume / cron_retrigger / cron_reschedule / cron_delete — by id from cron_list. Use cron_retrigger to re-fire a consumed one-shot or to immediately start one extra run of a recurring job (returns when started; does not wait for completion). Use cron_reschedule to move a still-pending one-shot (do not delete+create). `, } } diff --git a/test/host.test.js b/test/host.test.js index df1879f..023befa 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -409,6 +409,44 @@ test('retriggerJob fires recurring immediately and re-arms consumed oneshot', as assert.ok(httpHit.body.run?.id) }) +test('retriggerJob fire-and-forget returns while live run is still running', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const archived = [] + const sessionPort = makeLiveSessionPort(fakeLiveCtx(archived)) + const service = createTestHost({ + filePath: join(dir, 'store.json'), + now: () => Date.parse('2026-08-24T01:00:00.000Z'), + sessionPort, + }) + + const job = await service.createJob({ + name: 'detach-me', + prompt: 'slow work', + enabled: false, + schedule: { kind: 'at', at: new Date(Date.parse('2026-08-24T01:00:00.000Z') + 60_000).toISOString(), timezone: 'UTC' }, + }) + await service.store.mutate((state) => ({ + ...state, + jobs: state.jobs.map((row) => (row.id === job.id + ? { ...row, nextRunAt: null, lastStatus: 'succeeded', enabled: false } + : row)), + })) + + const started = await service.retriggerJob(job.id) + assert.equal(started.ok, true) + assert.equal(started.mode, 'started') + assert.equal(started.run.status, 'running') + assert.ok(started.run.sessionId) + + await service.flushDetachedRuns() + const history = await service.listHistory(job.id) + const terminal = history.find((row) => row.id === started.run.id) + assert.ok(terminal) + assert.equal(terminal.status, 'succeeded') + assert.equal(terminal.summary, '测试成功') +}) + test('rescheduleJob moves pending one-shot next without dispatch', async (t) => { const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) t.after(() => rm(dir, { recursive: true, force: true }))