mirror of
https://github.com/hansjone/dsh-ops-cron.git
synced 2026-10-08 22:00:46 +08:00
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 <cursoragent@cursor.com>
This commit is contained in:
parent
89ab5f59af
commit
a41a2b53c9
3 changed files with 102 additions and 9 deletions
62
lib/host.js
62
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)
|
* - 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
|
* - one-shot (at): consumed only; 0 → run-now; >0 → reschedule next
|
||||||
* Auto-enables paused jobs. Rejects when a run is already pending/running.
|
* 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) {
|
async function retriggerJob(jobId, opts = {}, identity = null) {
|
||||||
const afterRaw = opts.afterMinutes ?? opts.after_minutes
|
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 {
|
return {
|
||||||
ok: true,
|
ok: true,
|
||||||
mode,
|
mode: detached ? 'started' : 'immediate',
|
||||||
job: result.job,
|
job: result.job,
|
||||||
run: result.run,
|
run: result.run,
|
||||||
decision: result.decision,
|
decision: result.decision,
|
||||||
|
|
@ -945,6 +948,7 @@ export function createHostService(options = {}) {
|
||||||
async function dispatchRun(jobId, trigger, opts = {}) {
|
async function dispatchRun(jobId, trigger, opts = {}) {
|
||||||
let claimed
|
let claimed
|
||||||
const t = Number.isFinite(opts.at) ? opts.at : now()
|
const t = Number.isFinite(opts.at) ? opts.at : now()
|
||||||
|
const waitForCompletion = opts.wait !== false
|
||||||
await withState((current) => {
|
await withState((current) => {
|
||||||
claimed = claimOccurrence(current, jobId, t, trigger, current.settings)
|
claimed = claimOccurrence(current, jobId, t, trigger, current.settings)
|
||||||
if (claimed.decision.action === 'wait') return current
|
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)
|
let run = (executed.runs || []).find((row) => row.id === claimed.run.id)
|
||||||
if (run?.status === 'running' && typeof sessionPort?.waitForTurn === 'function') {
|
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
|
let terminal
|
||||||
try {
|
try {
|
||||||
terminal = await sessionPort.waitForTurn(run.sessionId)
|
terminal = await sessionPort.waitForTurn(run.sessionId)
|
||||||
|
|
@ -978,6 +995,38 @@ export function createHostService(options = {}) {
|
||||||
return settleAndNotify(jobId, claimed.run.id, claimed.decision)
|
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() {
|
async function tick() {
|
||||||
if (ticking) return []
|
if (ticking) return []
|
||||||
ticking = true
|
ticking = true
|
||||||
|
|
@ -1023,9 +1072,11 @@ export function createHostService(options = {}) {
|
||||||
}
|
}
|
||||||
|
|
||||||
function stopTimer() {
|
function stopTimer() {
|
||||||
if (!timer) return
|
if (timer) {
|
||||||
clearInterval(timer)
|
clearInterval(timer)
|
||||||
timer = null
|
timer = null
|
||||||
|
}
|
||||||
|
return flushDetachedRuns()
|
||||||
}
|
}
|
||||||
|
|
||||||
async function concealKnownSessions() {
|
async function concealKnownSessions() {
|
||||||
|
|
@ -1613,6 +1664,7 @@ export function createHostService(options = {}) {
|
||||||
recover,
|
recover,
|
||||||
startTimer,
|
startTimer,
|
||||||
stopTimer,
|
stopTimer,
|
||||||
|
flushDetachedRuns,
|
||||||
handleRequest,
|
handleRequest,
|
||||||
snapshot,
|
snapshot,
|
||||||
getUdsAuth,
|
getUdsAuth,
|
||||||
|
|
|
||||||
11
lib/tools.js
11
lib/tools.js
|
|
@ -746,14 +746,14 @@ export function cronToolDefinitions(service, deps = {}) {
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: 'cron_retrigger',
|
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: {
|
parameters: {
|
||||||
type: 'object',
|
type: 'object',
|
||||||
additionalProperties: false,
|
additionalProperties: false,
|
||||||
properties: {
|
properties: {
|
||||||
task_id: { type: 'string', description: 'Job id (alias: id).' },
|
task_id: { type: 'string', description: 'Job id (alias: id).' },
|
||||||
id: { type: 'string', description: 'Alias of task_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: {
|
output: {
|
||||||
|
|
@ -774,6 +774,9 @@ export function cronToolDefinitions(service, deps = {}) {
|
||||||
const when = value.nextRunAt ? new Date(value.nextRunAt).toISOString() : 'pending'
|
const when = value.nextRunAt ? new Date(value.nextRunAt).toISOString() : 'pending'
|
||||||
return text(`Retrigger scheduled "${name}" → next ${when}.`)
|
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})` : ''}.`)
|
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.',
|
'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.',
|
'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.',
|
'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:',
|
'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.',
|
'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.',
|
'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_runs — run history for a task_id
|
||||||
- cron_progress — read progress.channel.file snapshots
|
- 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_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).
|
||||||
`,
|
`,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -409,6 +409,44 @@ test('retriggerJob fires recurring immediately and re-arms consumed oneshot', as
|
||||||
assert.ok(httpHit.body.run?.id)
|
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) => {
|
test('rescheduleJob moves pending one-shot next without dispatch', async (t) => {
|
||||||
const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-'))
|
const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-'))
|
||||||
t.after(() => rm(dir, { recursive: true, force: true }))
|
t.after(() => rm(dir, { recursive: true, force: true }))
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue