mirror of
https://github.com/hansjone/dsh-ops-cron.git
synced 2026-10-09 00:43:22 +08:00
Compare commits
10 commits
8cc3bed5f1
...
a41a2b53c9
| Author | SHA1 | Date | |
|---|---|---|---|
| a41a2b53c9 | |||
| 89ab5f59af | |||
| b22048d377 | |||
| 8cf624ce7f | |||
| ad6e90f197 | |||
| 4472a9c74b | |||
| 3ce60b200c | |||
| c11221606a | |||
| 6724332114 | |||
| 839660d12f |
12 changed files with 1294 additions and 41 deletions
105
docs/cron_reschedule原语详设.md
Normal file
105
docs/cron_reschedule原语详设.md
Normal file
|
|
@ -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` 替换)。
|
||||
|
|
@ -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',
|
||||
|
|
@ -429,6 +430,16 @@ 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 仍可查询进度,但无法改任务。',
|
||||
permissionPreset: '操作权限',
|
||||
permissionPresetDefault: '跟随 Host 默认(当前行为)',
|
||||
permissionPresetReadOnly: '仅可查看',
|
||||
permissionPresetWorkspaceWrite: '可写入工作区',
|
||||
permissionPresetFullAccess: '完全权限',
|
||||
permissionPresetHint: '仅侧栏可改。默认不指定,触发时沿用 Host「新会话」默认权限;可手工提权到可写/完全权限。Agent 不可见也不可改。',
|
||||
loginRequired: '登录后才能使用定时任务',
|
||||
runNoSession: '这次运行没有可打开的会话',
|
||||
runOpenFailed: '打不开这次对话,会话可能已被删除',
|
||||
|
|
@ -458,6 +469,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',
|
||||
|
|
@ -496,6 +508,16 @@ 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.',
|
||||
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',
|
||||
|
|
@ -753,6 +775,8 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba
|
|||
imBotId: '',
|
||||
imTargetId: '',
|
||||
mirrorToSession: false,
|
||||
agentAccess: 'allow',
|
||||
permissionPreset: '',
|
||||
labelsText: '',
|
||||
persistKind: 'retain',
|
||||
archiveEndpoint: '',
|
||||
|
|
@ -779,6 +803,8 @@ 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',
|
||||
permissionPreset: job.permissionPreset || '',
|
||||
labelsText: formatLabelsText(job.labels),
|
||||
persistKind: persistHistoryKind(job),
|
||||
archiveEndpoint: job.persistHistory?.kind === 'archive' ? (job.persistHistory.endpoint || '') : '',
|
||||
|
|
@ -877,7 +903,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 +931,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 +1006,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 +1256,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')}`),
|
||||
)
|
||||
}
|
||||
|
||||
|
|
@ -1402,6 +1445,28 @@ 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', 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',
|
||||
|
|
@ -1670,6 +1735,8 @@ 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',
|
||||
permissionPreset: form.permissionPreset || '',
|
||||
labels: parseLabelsText(form.labelsText),
|
||||
persistHistory: form.persistKind === 'forever'
|
||||
? { kind: 'forever' }
|
||||
|
|
@ -1761,7 +1828,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,
|
||||
|
|
|
|||
|
|
@ -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'
|
||||
|
|
@ -107,6 +108,8 @@ export function publicJob(job, runs = null) {
|
|||
model: job.model || '',
|
||||
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'
|
||||
|
|
@ -127,6 +130,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,
|
||||
|
|
|
|||
333
lib/host.js
333
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,
|
||||
|
|
@ -35,6 +36,7 @@ import {
|
|||
listJobs,
|
||||
normalizeSettings,
|
||||
removeJob,
|
||||
splitProviderModel,
|
||||
storePath,
|
||||
upsertJob,
|
||||
} from './store.js'
|
||||
|
|
@ -628,6 +630,16 @@ 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),
|
||||
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,
|
||||
|
|
@ -699,6 +711,220 @@ export function createHostService(options = {}) {
|
|||
return updateJob(jobId, { enabled }, identity)
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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
|
||||
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)
|
||||
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
|
||||
}
|
||||
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 (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; use cron_reschedule (or PATCH 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', { wait: false })
|
||||
const detached = result.detached === true
|
||||
|| (result.run && ACTIVE_RUN_STATUSES.has(result.run.status))
|
||||
return {
|
||||
ok: true,
|
||||
mode: detached ? 'started' : 'immediate',
|
||||
job: result.job,
|
||||
run: result.run,
|
||||
decision: result.decision,
|
||||
nextRunAt: result.job?.nextRunAt ?? null,
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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)) {
|
||||
|
|
@ -719,9 +945,10 @@ 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()
|
||||
const waitForCompletion = opts.wait !== false
|
||||
await withState((current) => {
|
||||
claimed = claimOccurrence(current, jobId, t, trigger, current.settings)
|
||||
if (claimed.decision.action === 'wait') return current
|
||||
|
|
@ -741,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)
|
||||
|
|
@ -755,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
|
||||
|
|
@ -774,7 +1046,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()
|
||||
|
|
@ -798,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() {
|
||||
|
|
@ -844,6 +1120,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 {
|
||||
|
|
@ -1001,7 +1286,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|/reschedule)?$`))
|
||||
if (jobMatch) {
|
||||
const jobId = decodeURIComponent(jobMatch[1])
|
||||
const rest = jobMatch[2] || ''
|
||||
|
|
@ -1030,6 +1315,18 @@ 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 === '/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 })
|
||||
|
|
@ -1182,7 +1479,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' || 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))
|
||||
|
|
@ -1354,6 +1654,8 @@ export function createHostService(options = {}) {
|
|||
createJob,
|
||||
updateJob,
|
||||
pauseJob,
|
||||
retriggerJob,
|
||||
rescheduleJob,
|
||||
deleteJob,
|
||||
updateSettings,
|
||||
dispatchRun,
|
||||
|
|
@ -1362,6 +1664,7 @@ export function createHostService(options = {}) {
|
|||
recover,
|
||||
startTimer,
|
||||
stopTimer,
|
||||
flushDetachedRuns,
|
||||
handleRequest,
|
||||
snapshot,
|
||||
getUdsAuth,
|
||||
|
|
@ -1710,8 +2013,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,
|
||||
|
|
@ -1841,6 +2145,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',
|
||||
|
|
|
|||
|
|
@ -31,6 +31,8 @@ export const MESSAGES = {
|
|||
'tool.list': '列出定时任务',
|
||||
'tool.pause': '暂停定时任务',
|
||||
'tool.resume': '恢复定时任务',
|
||||
'tool.retrigger': '重新触发一次性任务',
|
||||
'tool.reschedule': '改期一次性任务',
|
||||
'tool.delete': '删除定时任务',
|
||||
},
|
||||
en: {
|
||||
|
|
@ -51,6 +53,8 @@ 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.reschedule': 'Reschedule pending one-shot',
|
||||
'tool.delete': 'Delete scheduled task',
|
||||
},
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,7 +48,7 @@ export {
|
|||
workspaceVisibleIds,
|
||||
} from './isolation.js'
|
||||
export { wrapScheduledPrompt } from './prompt.js'
|
||||
export { normalizeJobModel } from './store.js'
|
||||
export { normalizeJobModel, normalizeAgentAccess, agentMayMutateJob, normalizePermissionPreset, PERMISSION_PRESET_IDS, splitProviderModel } from './store.js'
|
||||
export {
|
||||
assertCanAccessJob,
|
||||
canViewAllJobs,
|
||||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -226,6 +226,26 @@ export function persistHistoryLimit(policy, settingsLimit = 200) {
|
|||
return settingsLimit
|
||||
}
|
||||
|
||||
/**
|
||||
* 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) 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'
|
||||
return terminal
|
||||
}
|
||||
|
||||
/**
|
||||
* Project job + runs into a lifecycle state for listeners.
|
||||
* @param {object} job
|
||||
|
|
|
|||
94
lib/store.js
94
lib/store.js
|
|
@ -77,12 +77,98 @@ 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'
|
||||
}
|
||||
|
||||
/** 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-…).
|
||||
* 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: '' }
|
||||
}
|
||||
|
|
@ -135,6 +221,8 @@ export function createJobRecord(input, state, now) {
|
|||
input.persistHistory ?? input.persist_history,
|
||||
settings.historyLimit,
|
||||
)
|
||||
const agentAccess = normalizeAgentAccess(input)
|
||||
const permissionPreset = normalizePermissionPreset(input)
|
||||
const job = {
|
||||
id: String(input.id || newId()),
|
||||
name,
|
||||
|
|
@ -148,6 +236,8 @@ export function createJobRecord(input, state, now) {
|
|||
model: model.model,
|
||||
reasoningEffort: model.reasoningEffort,
|
||||
agentPreset,
|
||||
agentAccess,
|
||||
permissionPreset,
|
||||
delivery,
|
||||
mirrorToSession,
|
||||
ownerEmpNo,
|
||||
|
|
|
|||
192
lib/tools.js
192
lib/tools.js
|
|
@ -20,6 +20,7 @@ import {
|
|||
jobVisibleToIdentity,
|
||||
UNASSIGNED_OWNER,
|
||||
} from './ownership.js'
|
||||
import { agentMayMutateJob, splitProviderModel } from './store.js'
|
||||
|
||||
const JOB_SCHEMA = {
|
||||
type: 'object',
|
||||
|
|
@ -53,6 +54,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' },
|
||||
},
|
||||
}
|
||||
|
||||
|
|
@ -214,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: _a, permissionPreset: _p, ...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()
|
||||
|
|
@ -327,20 +343,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 }
|
||||
}
|
||||
|
|
@ -378,8 +401,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.' },
|
||||
|
|
@ -449,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)
|
||||
|
|
@ -503,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 }
|
||||
},
|
||||
},
|
||||
{
|
||||
|
|
@ -556,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
|
||||
},
|
||||
},
|
||||
|
|
@ -656,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
|
||||
},
|
||||
},
|
||||
|
|
@ -678,14 +713,15 @@ 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) }
|
||||
},
|
||||
},
|
||||
{
|
||||
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,
|
||||
|
|
@ -702,9 +738,117 @@ 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) }
|
||||
},
|
||||
},
|
||||
{
|
||||
name: 'cron_retrigger',
|
||||
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 = start now (returns without waiting). Not allowed for recurring jobs.' },
|
||||
},
|
||||
},
|
||||
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}.`)
|
||||
}
|
||||
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})` : ''}.`)
|
||||
},
|
||||
},
|
||||
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')
|
||||
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) }
|
||||
},
|
||||
},
|
||||
{
|
||||
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) }
|
||||
},
|
||||
},
|
||||
{
|
||||
|
|
@ -735,7 +879,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) {
|
||||
|
|
@ -756,13 +901,14 @@ 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.',
|
||||
'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 (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.',
|
||||
'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.',
|
||||
|
|
@ -774,8 +920,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()}
|
||||
|
|
@ -786,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_delete — by id from cron_list
|
||||
- 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).
|
||||
`,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -285,6 +285,271 @@ test('overlap skip writes a skipped history row instead of a second session', as
|
|||
assert.equal(inflight, 1)
|
||||
})
|
||||
|
||||
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')
|
||||
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' },
|
||||
})
|
||||
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',
|
||||
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, 3)
|
||||
|
||||
// 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('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 }))
|
||||
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-'))
|
||||
|
|
@ -473,6 +738,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 }),
|
||||
|
|
@ -550,6 +860,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 = {
|
||||
|
|
|
|||
|
|
@ -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,36 @@ 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)
|
||||
|
||||
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)
|
||||
})
|
||||
|
||||
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
|
||||
|
|
@ -49,6 +80,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)
|
||||
|
|
|
|||
|
|
@ -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 }))
|
||||
|
|
@ -417,6 +453,60 @@ 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.permissionPreset, undefined)
|
||||
assert.equal(created.job.id != null, true)
|
||||
|
||||
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 }, {}),
|
||||
(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')
|
||||
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', () => {
|
||||
const registered = []
|
||||
const ctx = {
|
||||
|
|
@ -439,6 +529,8 @@ test('registerCronTools registers each definition and disposer unregisters', ()
|
|||
'cron_progress',
|
||||
'cron_pause',
|
||||
'cron_resume',
|
||||
'cron_retrigger',
|
||||
'cron_reschedule',
|
||||
'cron_delete',
|
||||
])
|
||||
off()
|
||||
|
|
@ -454,6 +546,8 @@ 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 rescheduleJob() { return {} },
|
||||
async deleteJob() { return {} },
|
||||
})
|
||||
const bad = []
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue