From 67238d746e572c6745bfb9e58bf93601d0cecafc Mon Sep 17 00:00:00 2001 From: oliver Date: Sun, 6 Sep 2026 14:45:37 +0800 Subject: [PATCH] Isolate IM cron tools and trust-gate read APIs. Scope list/pause/delete to the calling chat, bind IM delivery to that peer, ignore client job ids, require trusted Host for GETs, and wrap mirrored summaries as system reminders. Co-authored-by: Cursor --- lib/delivery.js | 70 +++++++++++++++++++++++++++++++++---------- lib/host.js | 18 ++++++++++- lib/index.d.ts | 2 ++ lib/index.js | 2 ++ lib/tools.js | 69 ++++++++++++++++++++++++++++++++++-------- test/delivery.test.js | 66 ++++++++++++++++++++++++++++++++++++++++ test/host.test.js | 24 +++++++++++++++ 7 files changed, 222 insertions(+), 29 deletions(-) diff --git a/lib/delivery.js b/lib/delivery.js index 5caff5c..cd85de2 100644 --- a/lib/delivery.js +++ b/lib/delivery.js @@ -149,9 +149,7 @@ export function matchTargetForPeer(targets, peer) { const jid = String(route.jid || route.chatId || route.openId || route.userId || '').trim().toLowerCase() if (!jid) continue if (wantGroup && !jid.endsWith('@g.us')) continue - if (candidates.some((peerJid) => ( - jid === peerJid || peerJid.endsWith(jid) || jid.endsWith(peerJid) - ))) { + if (candidates.some((peerJid) => jid === peerJid)) { return target } } @@ -365,25 +363,64 @@ export async function resolveCreateDelivery(args = {}, exec = {}, deps = {}) { const targetId = typeof args.im_target_id === 'string' ? args.im_target_id.trim() : (typeof args.imTargetId === 'string' ? args.imTargetId.trim() : '') + const dshIm = deps.dshIm + const peer = await resolveCallerPeer(exec, dshIm) + + // IM peers cannot retarget to an arbitrary catalog entry (confused deputy). + if (peer?.botId) { + if (explicit === 'dsh') return normalizeDelivery({ kind: 'dsh' }) + const ensured = await ensureImDeliveryForPeer(peer, dshIm) + if (ensured) return ensured + if (explicit === 'im' || (botId && targetId)) { + const error = new Error('IM sessions cannot pick a free im_bot_id/im_target_id; delivery is bound to this chat') + error.code = 'IM_TARGET_FORBIDDEN' + throw error + } + return normalizeDelivery({ kind: 'dsh' }) + } + if (explicit === 'dsh') return normalizeDelivery({ kind: 'dsh' }) if (explicit === 'im' || (botId && targetId)) { return normalizeDelivery({ kind: 'im', botId, targetId }) } - const dshIm = deps.dshIm - const sessionId = callerSessionId(exec) - if (sessionId && dshIm && typeof dshIm.resolveSessionPeer === 'function') { - try { - const peer = await dshIm.resolveSessionPeer(sessionId) - const ensured = await ensureImDeliveryForPeer(peer, dshIm) - if (ensured) return ensured - } catch { - // fall through to dsh - } - } return normalizeDelivery({ kind: 'dsh' }) } +/** + * Resolve the IM peer for the tool-calling session, if any. + * @param {object} exec + * @param {object|undefined} dshIm + */ +export async function resolveCallerPeer(exec, dshIm) { + const sessionId = callerSessionId(exec) + if (!sessionId || !dshIm || typeof dshIm.resolveSessionPeer !== 'function') return null + try { + const peer = await dshIm.resolveSessionPeer(sessionId) + return peer?.botId ? peer : null + } catch { + return null + } +} + +/** + * Whether a scheduled job was created from this IM conversation. + * Web/sidebar callers are unrestricted; IM peers only manage their own jobs. + * @param {object|null|undefined} job + * @param {object} peer + */ +export function jobVisibleToPeer(job, peer) { + if (!job || !peer?.botId) return false + const botId = String(peer.botId).trim() + const conversationKey = String(peer.conversationKey || '').trim() + const conversationId = String(peer.conversationId || '').trim() + const origin = normalizeOrigin(job.origin) + if (origin?.kind !== 'im' || origin.peer?.botId !== botId) return false + if (conversationKey && origin.peer.conversationKey === conversationKey) return true + if (conversationId && origin.peer.conversationId === conversationId) return true + return false +} + const IM_MAX_CHARS = 3500 /** @@ -522,14 +559,17 @@ export async function mirrorRunToSession(job, summary, deps = {}) { } const text = formatRunResultBody(job, summary) + const safe = text.replaceAll('', '<\\/system-reminder>') const plugin = typeof deps.pluginName === 'string' && deps.pluginName.trim() ? deps.pluginName.trim() : 'dsh-ops-cron' const id = typeof deps.newId === 'function' ? deps.newId() : `cron-mirror-${Date.now()}` + // Keep role=user + user/message (assistant/message requires model source + open turn). + // Wrap as system-reminder so the next model turn does not treat cron output as user intent. const message = { id, role: 'user', - content: [{ type: 'text', text }], + content: [{ type: 'text', text: `\n${safe}\n` }], source: { kind: 'plugin', plugin }, } diff --git a/lib/host.js b/lib/host.js index a1c1412..0eeb026 100644 --- a/lib/host.js +++ b/lib/host.js @@ -178,9 +178,11 @@ export function createHostService(options = {}) { async function createJob(input) { const t = now() + // Ignore client-supplied ids on create — otherwise POST/tools can overwrite. + const { id: _ignoredId, ...safeInput } = input && typeof input === 'object' ? input : {} let created await withState((current) => { - created = createJobRecord(input, current, t) + created = createJobRecord(safeInput, current, t) return upsertJob(current, created) }) return jobView(created) @@ -389,6 +391,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/settings` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const state = await snapshot() write(200, { ok: true, settings: state.settings }) return @@ -403,6 +406,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/models` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const catalog = typeof sessionPort?.listModels === 'function' ? await sessionPort.listModels() : { groups: [], current: null } @@ -411,6 +415,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/presets` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const catalog = typeof sessionPort?.listPresets === 'function' ? await sessionPort.listPresets() : { items: [], current: null } @@ -419,6 +424,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/workspaces` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const workspaces = typeof sessionPort?.listWorkspaces === 'function' ? await sessionPort.listWorkspaces() : [] @@ -427,6 +433,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/im-catalog` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const dshIm = getDshIm() if (!dshIm || typeof dshIm.listDeliveryCatalog !== 'function') { write(200, { @@ -456,6 +463,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/jobs` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const state = await snapshot() write(200, { ok: true, jobs: listJobs(state).map(jobView) }) return @@ -474,6 +482,7 @@ export function createHostService(options = {}) { const jobId = decodeURIComponent(jobMatch[1]) const rest = jobMatch[2] || '' if (method === 'GET' && !rest) { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const state = await snapshot() const job = getJob(state, jobId) if (!job) return write(404, { ok: false, error: 'job not found' }) @@ -545,6 +554,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/history` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const state = await snapshot() const jobId = url.searchParams.get('jobId') || undefined write(200, { ok: true, runs: listHistory(state, jobId).map(runView) }) @@ -552,6 +562,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/preview` && method === 'POST') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const body = await readJsonBody(req) const settings = (await snapshot()).settings const schedule = validateSchedule(body.schedule || body, body.timezone || settings.timezone) @@ -561,6 +572,7 @@ export function createHostService(options = {}) { } if (path === `${API_PREFIX}/workspace-visible` && method === 'GET') { + if (!isTrustedApiRequest(req)) return write(403, { ok: false, error: 'forbidden' }) const state = await snapshot() const listed = url.searchParams.getAll('id') write(200, { @@ -601,6 +613,10 @@ export function createHostService(options = {}) { async listJobs() { return listJobs(await snapshot()).map(jobView) }, + async getJob(jobId) { + const job = getJob(await snapshot(), jobId) + return job ? jobView(job) : null + }, async listHistory(jobId) { return listHistory(await snapshot(), jobId).map(runView) }, diff --git a/lib/index.d.ts b/lib/index.d.ts index 6794c52..3c3abfb 100644 --- a/lib/index.d.ts +++ b/lib/index.d.ts @@ -46,6 +46,8 @@ export function normalizeOrigin(input?: object): { kind: 'web' | 'im', sessionId export function resolveCreateDelivery(args?: object, exec?: object, deps?: object): Promise<{ kind: 'dsh' | 'im', botId?: string, targetId?: string }> export function resolveCreateOrigin(args?: object, exec?: object, deps?: object): Promise<{ kind: 'web' | 'im', sessionId?: string, peer?: object } | null> export function resolveMirrorSession(job: object, deps?: object): Promise<{ sessionId: string, via: string } | null> +export function resolveCallerPeer(exec?: object, dshIm?: object): Promise +export function jobVisibleToPeer(job: object, peer: object): boolean export function formatRunResultBody(job: object, summary?: string): string export function deliverRunToIm(job: object, summary?: string, deps?: object): Promise export function mirrorRunToSession(job: object, summary?: string, deps?: object): Promise diff --git a/lib/index.js b/lib/index.js index 3c2afb6..75312a3 100644 --- a/lib/index.js +++ b/lib/index.js @@ -52,9 +52,11 @@ export { claimOccurrence, executeClaimedRun, extractAssistantText, TITLE_PREFIX export { deliverRunToIm, formatRunResultBody, + jobVisibleToPeer, mirrorRunToSession, normalizeDelivery, normalizeOrigin, + resolveCallerPeer, resolveCreateDelivery, resolveCreateOrigin, resolveMirrorSession, diff --git a/lib/tools.js b/lib/tools.js index b88376b..79a8b29 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -4,7 +4,13 @@ */ import { formatInZone, resolveTodayAt } from './scheduler.js' -import { deliveryLine, resolveCreateDelivery, resolveCreateOrigin } from './delivery.js' +import { + deliveryLine, + jobVisibleToPeer, + resolveCallerPeer, + resolveCreateDelivery, + resolveCreateOrigin, +} from './delivery.js' import { resolveCreateAgentPreset } from './preset.js' const JOB_SCHEMA = { @@ -107,10 +113,35 @@ export function callerWorkingDirectory(exec) { return String(session?.header?.cwd || session?.cwd || exec?.agent?.cwd || '').trim() } -export function resolveCreateCwd(args, exec) { +export function resolveCreateCwd(args, exec, { peer = null } = {}) { + const sessionCwd = callerWorkingDirectory(exec) const passed = typeof args?.cwd === 'string' ? args.cwd.trim() : '' + // IM peers cannot point scheduled Agents at arbitrary host paths. + if (peer?.botId) return sessionCwd if (passed) return passed - return callerWorkingDirectory(exec) + return sessionCwd +} + +function forbidForeignJob(job, peer) { + if (!peer?.botId) return job + if (jobVisibleToPeer(job, peer)) return job + const error = new Error('job not found or not owned by this chat') + error.code = 'NOT_FOUND' + throw error +} + +async function requireOwnedJob(service, id, peer) { + const job = await service.getJob?.(id) + if (job) return forbidForeignJob(job, peer) + // Fallback when getJob is absent: list and find. + const jobs = await service.listJobs() + const hit = (jobs || []).find((row) => row.id === id) + if (!hit) { + const error = new Error('job not found') + error.code = 'NOT_FOUND' + throw error + } + return forbidForeignJob(hit, peer) } export function callerModelSelection(exec) { @@ -165,15 +196,15 @@ export function cronToolDefinitions(service, deps = {}) { minute: { type: 'integer', description: '0-59, used with hour.' }, 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. Only skip this to isolate from the project if the user asked for a custom folder.' }, + 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.' }, 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.' }, - delivery: { type: 'string', description: 'dsh (sidebar session) or im (proactive WhatsApp/IM via botId+targetId). Omit to auto-detect from the current session.' }, - im_bot_id: { type: 'string', description: 'Opaque botId from IM 投递设置 when delivery=im.' }, - im_target_id: { type: 'string', description: 'Opaque targetId from IM 投递设置 when delivery=im.' }, + delivery: { type: 'string', description: 'dsh (sidebar session) or im (proactive WhatsApp/IM via botId+targetId). Omit to auto-detect from the current session. On WhatsApp/IM, delivery is always bound to the current chat — free im_bot_id/im_target_id are ignored.' }, + im_bot_id: { type: 'string', description: 'Opaque botId from IM 投递设置 when delivery=im (Web/sidebar only; IM chats cannot retarget).' }, + im_target_id: { type: 'string', description: 'Opaque targetId from IM 投递设置 when delivery=im (Web/sidebar only; IM chats cannot retarget).' }, agent_preset: { type: 'string', description: 'Agent preset id for scheduled runs. Omit to inherit from the current WhatsApp/IM chat or Host default.' }, }, required: ['name', 'prompt'], @@ -203,18 +234,20 @@ export function cronToolDefinitions(service, deps = {}) { async execute(args, exec) { aborted(exec) const timeout = Number(args.timeout_minutes) + const dshIm = getDshIm() try { - const delivery = await resolveCreateDelivery(args, exec, { dshIm: getDshIm() }) - const origin = await resolveCreateOrigin(args, exec, { dshIm: getDshIm() }) + const peer = await resolveCallerPeer(exec, dshIm) + const delivery = await resolveCreateDelivery(args, exec, { dshIm }) + const origin = await resolveCreateOrigin(args, exec, { dshIm }) const agentPreset = await resolveCreateAgentPreset(args, exec, { - dshIm: getDshIm(), + dshIm, agentPresets: getAgentPresets(), }) const job = await service.createJob({ name: args.name, prompt: args.prompt, schedule: scheduleFromArgs(args, Date.now()), - cwd: resolveCreateCwd(args, exec), + cwd: resolveCreateCwd(args, exec, { peer }), ...resolveCreateModel(args, exec), timeoutMinutes: Number.isFinite(timeout) && timeout > 0 ? timeout : undefined, enabled: args.enabled !== false, @@ -263,7 +296,9 @@ export function cronToolDefinitions(service, deps = {}) { presentCall: () => ({ card: 'generic', title: '列出定时任务' }), async execute(args, exec) { aborted(exec) + const peer = await resolveCallerPeer(exec, getDshIm()) let jobs = await service.listJobs() + 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 } }, @@ -284,7 +319,10 @@ export function cronToolDefinitions(service, deps = {}) { presentCall: (args) => ({ card: 'generic', title: '暂停定时任务', content: String(args?.id || '') }), async execute(args, exec) { aborted(exec) - const job = await service.pauseJob(requireId(args), false) + const peer = await resolveCallerPeer(exec, getDshIm()) + const id = requireId(args) + await requireOwnedJob(service, id, peer) + const job = await service.pauseJob(id, false) return { job } }, }, @@ -304,7 +342,10 @@ export function cronToolDefinitions(service, deps = {}) { presentCall: (args) => ({ card: 'generic', title: '恢复定时任务', content: String(args?.id || '') }), async execute(args, exec) { aborted(exec) - const job = await service.pauseJob(requireId(args), true) + const peer = await resolveCallerPeer(exec, getDshIm()) + const id = requireId(args) + await requireOwnedJob(service, id, peer) + const job = await service.pauseJob(id, true) return { job } }, }, @@ -332,8 +373,10 @@ export function cronToolDefinitions(service, deps = {}) { presentCall: (args) => ({ card: 'generic', title: '删除定时任务', content: String(args?.id || '') }), async execute(args, exec) { aborted(exec) + const peer = await resolveCallerPeer(exec, getDshIm()) const id = requireId(args) try { + await requireOwnedJob(service, id, peer) await service.deleteJob(id) return { id, deleted: true } } catch (error) { diff --git a/test/delivery.test.js b/test/delivery.test.js index 631bce9..f973c10 100644 --- a/test/delivery.test.js +++ b/test/delivery.test.js @@ -10,6 +10,7 @@ import { formatRunResultBody, deliverRunToIm, mirrorRunToSession, + jobVisibleToPeer, } from '../lib/delivery.js' test('normalizeDelivery defaults to dsh', () => { @@ -307,10 +308,75 @@ test('mirrorRunToSession appends plugin notice without waking a turn', async () assert.equal(appended[0][0], 'user/message') assert.match(appended[0][1].content[0].text, /日报/) assert.match(appended[0][1].content[0].text, /line one/) + assert.match(appended[0][1].content[0].text, //) assert.equal(appended[0][1].source.kind, 'plugin') assert.deepEqual(appended[0][2], { surfaceOp: 'append' }) }) +test('matchTargetForPeer requires exact JID equality', () => { + const targets = [{ + targetId: 'suffix-trap', + kind: 'user', + route: { jid: '8613800138000@s.whatsapp.net' }, + }] + const peer = { + kind: 'direct', + conversationId: '13800138000@s.whatsapp.net', + phone: '13800138000', + } + assert.equal(matchTargetForPeer(targets, peer), null) +}) + +test('jobVisibleToPeer scopes IM ownership to origin conversation', () => { + const peer = { + botId: 'bot-a', + conversationKey: 'group:120363@g.us:user:1@lid', + conversationId: '120363@g.us', + } + assert.equal(jobVisibleToPeer({ + origin: { + kind: 'im', + peer: { botId: 'bot-a', conversationKey: 'group:120363@g.us:user:1@lid', conversationId: '120363@g.us' }, + }, + }, peer), true) + assert.equal(jobVisibleToPeer({ + origin: { + kind: 'im', + peer: { botId: 'bot-a', conversationKey: 'group:other@g.us', conversationId: 'other@g.us' }, + }, + }, peer), false) + assert.equal(jobVisibleToPeer({ + origin: { kind: 'web', sessionId: 'web-1' }, + }, peer), false) +}) + +test('resolveCreateDelivery binds IM peers to the current chat', async () => { + const peer = { + botId: 'bot-a', + conversationKey: 'direct:86138@s.whatsapp.net', + conversationId: '86138@s.whatsapp.net', + kind: 'direct', + phone: '86138', + } + const delivery = await resolveCreateDelivery( + { delivery: 'im', im_bot_id: 'other-bot', im_target_id: 'other-target' }, + { agent: { session: { id: 'sess-im' } } }, + { + dshIm: { + resolveSessionPeer: async () => peer, + listTargets: async () => [{ + targetId: 'auto-dm', + kind: 'user', + route: { jid: '86138@s.whatsapp.net' }, + }], + }, + }, + ) + assert.equal(delivery.kind, 'im') + assert.equal(delivery.botId, 'bot-a') + assert.equal(delivery.targetId, 'auto-dm') +}) + test('mirrorRunToSession resumes a cold origin session instead of create', async () => { const appended = [] const resumed = [] diff --git a/test/host.test.js b/test/host.test.js index 89c263a..a77bd24 100644 --- a/test/host.test.js +++ b/test/host.test.js @@ -237,6 +237,30 @@ test('shipped HTTP handler: create, list, run-now, history', async (t) => { assert.deepEqual(visible.body.visible, ['other']) }) +test('createJob ignores client-supplied id to prevent overwrite', async (t) => { + const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) + t.after(() => rm(dir, { recursive: true, force: true })) + const service = createHostService({ + filePath: join(dir, 'store.json'), + now: () => Date.parse('2026-08-24T01:00:00.000Z'), + }) + const first = await service.createJob({ + name: 'keep-me', + prompt: 'first', + schedule: { kind: 'cron', expr: '0 9 * * *', timezone: 'Asia/Shanghai' }, + }) + const second = await service.createJob({ + id: first.id, + name: 'attacker', + prompt: 'overwrite?', + schedule: { kind: 'cron', expr: '0 10 * * *', timezone: 'Asia/Shanghai' }, + }) + assert.notEqual(second.id, first.id) + const listed = await service.listJobs() + assert.equal(listed.find((job) => job.id === first.id)?.name, 'keep-me') + assert.equal(listed.find((job) => job.id === second.id)?.name, 'attacker') +}) + test('POST /runs/:id/open reveals the run session id', async (t) => { const dir = await mkdtemp(join(tmpdir(), 'dsh-ops-cron-')) t.after(() => rm(dir, { recursive: true, force: true }))