diff --git a/FORK.md b/FORK.md index d23094f..96ee8e8 100644 --- a/FORK.md +++ b/FORK.md @@ -5,6 +5,7 @@ - Agent 自建定时任务(`cron_create` / `cron_list` / …) - 按创建来源默认投递:**WhatsApp 会话 → IM**;**Web/DSH 会话 → 侧栏新会话** - 与 [dsh-im-ops](https://github.com/hansjone/dsh-im-ops) 的 `ctx.dshIm.send(botId, targetId, text)` 打通 +- 跑完后把摘要**镜像**进创建时的有效会话(Web 来源会话 / WhatsApp 对端绑定会话),便于继续追问 不要与上游 `@dsh-external/dsh-cron-tasks` 装进同一 profile(工具名冲突)。 @@ -27,7 +28,18 @@ dsh plugin --profile web add -w "D:\project\chatgpt\dsh-ops-cron" 1. 包名 / cordis id / HTTP 前缀仍为 `dsh-ops-cron`;侧栏产品文案为「定时任务」 2. Job 增加 `delivery: { kind: 'dsh'|'im', botId?, targetId? }` -3. `cron_create` 支持 `delivery` / `im_bot_id` / `im_target_id`;未显式指定时,若当前会话能 `resolveSessionPeer` 且已有匹配投递目标 → 默认 IM -4. 开火后若 `delivery.kind === 'im'`,把 assistant 摘要经 `ctx.dshIm.send` 投回 +3. Job 增加 `origin: { kind: 'web'|'im', sessionId?, peer? }`(创建时钉死;IM 开火时优先按 `conversationKey` 解析当前绑定会话) +4. `cron_create` 支持 `delivery` / `im_bot_id` / `im_target_id`;未显式指定时,若当前会话能 `resolveSessionPeer` 且已有匹配投递目标 → 默认 IM +5. 开火后若 `delivery.kind === 'im'`,把 assistant 摘要经 `ctx.dshIm.send` 投回 +6. 开火终态后把同构摘要镜像进有效会话(`agents.get` 或 `agents.resume` 后 `session.append`,不另开一轮模型;不用 `agents.create`,避免覆盖已有会话) + +## 执行会话 vs 有效会话 + +| 角色 | 是什么 | 用途 | +|------|--------|------| +| 执行会话 | 每次开火新建的 root Session(侧栏历史可打开) | 跑 Agent、保留工具轨迹 | +| 有效会话 | 创建时的 Web 会话,或 WhatsApp/IM `conversationKey` 当前绑定的 Session | 用户日常追问所在处;摘要镜像写入这里 | + +`dshIm.send` 只保证平台侧可见,不会自动写 Harness 会话日志;镜像是 cron 层额外一步。 关键告警推送仍在 **netxops**,不经过本插件。 diff --git a/README.md b/README.md index cad842e..3c2f14f 100644 --- a/README.md +++ b/README.md @@ -5,6 +5,7 @@ Scheduled-task plugin for DeepSeek Harness (fork of [dsh-cron-tasks](https://git - Sidebar **定时任务** under New Session - Agent tools: `cron_create` / `cron_list` / `cron_pause` / `cron_resume` / `cron_delete` - Per-job delivery: **DSH** (new root session) or **IM** (`ctx.dshIm.send`) +- After each fire, mirrors the summary into the **origin effective session** (creator Web chat or WhatsApp peer binding) so follow-ups share context See [FORK.md](FORK.md) for install notes. Do **not** install upstream `@dsh-external/dsh-cron-tasks` in the same profile. @@ -18,14 +19,16 @@ dsh plugin --profile web add -w "github:hansjone/dsh-ops-cron" ## Delivery defaults -| Created from | Default delivery | -|--------------|------------------| -| WhatsApp / IM session | `im` → reuse or **auto-create** a 投递目标 for that chat (group→group, DM→DM) | -| Web / plain DSH session | `dsh` → sidebar history session | -| Explicit `delivery` / `im_bot_id`+`im_target_id` | as specified | +| Created from | Default delivery | Effective-session mirror | +|--------------|------------------|--------------------------| +| WhatsApp / IM session | `im` → reuse or **auto-create** a 投递目标 for that chat (group→group, DM→DM) | Live session for that `conversationKey` (latest binding) | +| Web / plain DSH session | `dsh` → sidebar history session | Creator session pinned as `origin.sessionId` | +| Explicit `delivery` / `im_bot_id`+`im_target_id` | as specified | Same mirror rules when `origin` / peer can be resolved | You usually do **not** need to paste botId/targetId when creating from WhatsApp — omit `delivery` and the job binds to the current chat. +**Execution session ≠ effective session.** Each fire still opens a fresh root Session for the Agent turn (visible under 定时任务 history). The mirrored notice is what makes the WhatsApp/Web chat you actually use able to continue from the result. + ## Agent preset | Source | Behavior | diff --git a/lib/client.js b/lib/client.js index eeebe55..c782fa5 100644 --- a/lib/client.js +++ b/lib/client.js @@ -302,7 +302,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba saveFailed: '本部署没有接受这些值,已保留供你修改。', enabled: '启用调度器', timezone: '默认时区', historyLimit: '历史保留条数', overlapPolicy: '重叠策略', misfirePolicy: '漏跑策略', policySkip: '跳过', - hint: '任务在左侧「新会话」下方进入。运行结果会出现在历史行里,点击记录会打开和普通会话一样的对话页,可以继续聊。', + hint: '每次执行会新开会话跑任务;摘要会镜像回创建时的有效会话(WhatsApp/Web),便于继续追问。完整工具轨迹仍可从历史打开。', entry: '定时任务', entryLabel: '打开定时任务', backToWorkspace: '返回工作区', backLabel: '返回工作区', newJob: '新建任务', name: '名称', prompt: '提示词', kind: '日程', @@ -322,7 +322,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba agentPreset: 'Agent Preset', agentPresetDefault: '每次运行用当时的 Host 默认', agentPresetHint: '指定后每次触发都挂载该 Preset;留空则跟随 Host 默认(创建时若从 WhatsApp 继承会自动写入)。', - editorLead: '保存后按日程在新会话里执行;若投递为 IM,结束后还会发摘要到 WhatsApp/IM。点击下方某次运行,右侧会打开那次对话,可以继续。', + editorLead: '到点仍会新开会话执行。结束后摘要会镜像回创建时的有效会话;IM 投递还会经 WhatsApp/IM 发出。日常追问请用原会话;完整工具轨迹可从下方运行记录打开。', lastOutput: '上次输出', noOutput: '还没有输出。先立即运行一次。', delivery: '投递', deliveryDsh: 'DSH 侧栏会话', @@ -343,7 +343,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba saveFailed: 'The deployment did not accept these values; they were left for you to correct.', enabled: 'Enable scheduler', timezone: 'Default time zone', historyLimit: 'History retention', overlapPolicy: 'Overlap policy', misfirePolicy: 'Misfire policy', policySkip: 'Skip', - hint: 'Open from under New Session. Run output appears in history; click a run to continue in the normal chat view.', + hint: 'Each fire runs in a fresh session; the summary is mirrored into the origin WhatsApp/Web session so you can follow up there. Full tool traces stay in run history.', entry: 'Scheduled tasks', entryLabel: 'Open scheduled tasks', backToWorkspace: 'Back to workspace', backLabel: 'Back to workspace', newJob: 'New job', name: 'Name', prompt: 'Prompt', kind: 'Schedule', @@ -363,7 +363,7 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba agentPreset: 'Agent preset', agentPresetDefault: 'Use the Host default at fire time', agentPresetHint: 'Pin a preset for every fire. Leave empty to follow the Host default (WhatsApp-created jobs inherit the chat preset automatically).', - editorLead: 'Runs start a new session. IM delivery also sends a summary via dshIm after the turn. Click a run below to open that conversation.', + editorLead: 'Fires still run in a new session. Afterward the summary is mirrored into the origin session; IM delivery also sends via WhatsApp/IM. Follow up in the origin chat; open a run below for the full tool trace.', lastOutput: 'Last output', noOutput: 'No output yet. Run it once.', delivery: 'Delivery', deliveryDsh: 'DSH sidebar session', @@ -1230,6 +1230,16 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba if (opened === false) setError('打不开这次对话,会话可能已被删除') } + function currentSessionId(faces) { + try { + const snap = faces?.sessions?.list?.getSnapshot?.() + const current = snap?.current + return typeof current === 'string' && current.trim() ? current.trim() : '' + } catch { + return '' + } + } + async function save() { try { const atValue = form.kind === 'at' ? localInputToAt(form.at, form.timezone) : form.at @@ -1248,6 +1258,8 @@ body>.dsh-ct-main{position:fixed;top:0;right:0;bottom:0;left:var(--dsh-ct-sideba if (selection.type === 'job' && selection.jobId) { await api(`/jobs/${selection.jobId}`, { method: 'PATCH', body: JSON.stringify(payload) }) } else { + const originSessionId = currentSessionId(faces) + if (originSessionId) payload.origin = { kind: 'web', sessionId: originSessionId } const created = await api('/jobs', { method: 'POST', body: JSON.stringify(payload) }) setSelection({ type: 'job', jobId: created.job.id }) } diff --git a/lib/delivery.js b/lib/delivery.js index ad077f5..5caff5c 100644 --- a/lib/delivery.js +++ b/lib/delivery.js @@ -266,6 +266,91 @@ export function callerSessionId(exec) { ).trim() } +/** + * Normalize a durable origin binding for result mirroring. + * @param {object|null|undefined} input + * @returns {{ kind: 'web' | 'im', sessionId?: string, peer?: object } | null} + */ +export function normalizeOrigin(input) { + if (!input || typeof input !== 'object') return null + const sessionId = typeof input.sessionId === 'string' ? input.sessionId.trim() : '' + const peerSrc = input.peer && typeof input.peer === 'object' ? input.peer : null + const peerBotId = typeof peerSrc?.botId === 'string' ? peerSrc.botId.trim() : '' + const conversationKey = typeof peerSrc?.conversationKey === 'string' + ? peerSrc.conversationKey.trim() + : '' + const kindRaw = String(input.kind || '').trim().toLowerCase() + const kind = kindRaw === 'im' || peerBotId + ? 'im' + : (kindRaw === 'web' || sessionId ? 'web' : '') + if (!kind) return null + if (!sessionId && !(kind === 'im' && peerBotId && conversationKey)) return null + const origin = { kind } + if (sessionId) origin.sessionId = sessionId + if (kind === 'im' && peerBotId && conversationKey) { + origin.peer = { + botId: peerBotId, + conversationKey, + ...typeof peerSrc.conversationId === 'string' && peerSrc.conversationId.trim() + ? { conversationId: peerSrc.conversationId.trim() } + : {}, + ...peerSrc.kind === 'group' || peerSrc.kind === 'direct' + ? { kind: peerSrc.kind } + : {}, + } + } + return origin +} + +function slimPeer(peer) { + if (!peer || typeof peer !== 'object') return null + const botId = typeof peer.botId === 'string' ? peer.botId.trim() : '' + const conversationKey = typeof peer.conversationKey === 'string' + ? peer.conversationKey.trim() + : '' + if (!botId || !conversationKey) return null + return { + botId, + conversationKey, + ...typeof peer.conversationId === 'string' && peer.conversationId.trim() + ? { conversationId: peer.conversationId.trim() } + : {}, + ...peer.kind === 'group' || peer.kind === 'direct' ? { kind: peer.kind } : {}, + } +} + +/** + * Resolve origin for cron_create / sidebar create. + * Prefer an explicit origin; else capture the caller session (+ IM peer when present). + * @param {object} args + * @param {object} exec + * @param {{ dshIm?: object }} deps + */ +export async function resolveCreateOrigin(args = {}, exec = {}, deps = {}) { + const explicit = normalizeOrigin(args.origin) + if (explicit) return explicit + + const sessionId = typeof args.origin_session_id === 'string' + ? args.origin_session_id.trim() + : (typeof args.originSessionId === 'string' ? args.originSessionId.trim() : '') + const callerId = sessionId || callerSessionId(exec) + if (!callerId) return null + + const dshIm = deps.dshIm + if (dshIm && typeof dshIm.resolveSessionPeer === 'function') { + try { + const peer = await dshIm.resolveSessionPeer(callerId) + const slim = slimPeer(peer) + if (slim) { + return normalizeOrigin({ kind: 'im', sessionId: callerId, peer: slim }) + } + } catch { + // fall through to web origin + } + } + return normalizeOrigin({ kind: 'web', sessionId: callerId }) +} + /** * Resolve delivery for cron_create: explicit args win; else IM peer → matching * or auto-created 投递目标; else dsh. @@ -301,6 +386,17 @@ export async function resolveCreateDelivery(args = {}, exec = {}, deps = {}) { const IM_MAX_CHARS = 3500 +/** + * Shared body for IM delivery and session mirror (without @mention prefix). + * @param {object} job + * @param {string} [summary] + */ +export function formatRunResultBody(job, summary) { + const text = String(summary || '').trim() || `(定时任务「${job?.name || ''}」已完成,无文本摘要)` + const clipped = text.length > IM_MAX_CHARS ? `${text.slice(0, IM_MAX_CHARS)}…` : text + return `【定时任务 · ${job?.name || ''}】\n${clipped}` +} + /** * Send run summary to IM when job.delivery.kind === 'im'. * Group deliveries @mention the creator when mentionJid was captured at create time. @@ -314,17 +410,165 @@ export async function deliverRunToIm(job, summary, deps = {}) { error.code = 'IM_UNAVAILABLE' throw error } - const text = String(summary || '').trim() || `(定时任务「${job.name}」已完成,无文本摘要)` - const clipped = text.length > IM_MAX_CHARS ? `${text.slice(0, IM_MAX_CHARS)}…` : text - const header = `【定时任务 · ${job.name}】\n` + const headerAndBody = formatRunResultBody(job, summary) const mentionJid = typeof delivery.mentionJid === 'string' ? delivery.mentionJid.trim() : '' - let body = `${header}${clipped}` + let body = headerAndBody const options = {} if (mentionJid && /@(s\.whatsapp\.net|lid)$/i.test(mentionJid)) { const token = mentionJid.slice(0, mentionJid.indexOf('@')) - body = `@${token}\n${header}${clipped}` + body = `@${token}\n${headerAndBody}` options.mentions = [mentionJid] } await dshIm.send(delivery.botId, delivery.targetId, body, options) return { sent: true, mentioned: Boolean(options.mentions) } } + +/** + * Resolve which live/effective session should receive a mirrored run summary. + * Prefer the current IM conversation binding over a stale origin.sessionId. + * @param {object} job + * @param {{ dshIm?: object }} deps + * @returns {Promise<{ sessionId: string, via: string } | null>} + */ +export async function resolveMirrorSession(job, deps = {}) { + const dshIm = deps.dshIm + const origin = normalizeOrigin(job?.origin) + const delivery = job?.delivery + + if (origin?.peer?.botId && origin.peer.conversationKey + && dshIm && typeof dshIm.resolveConversationSession === 'function') { + try { + const hit = await dshIm.resolveConversationSession( + origin.peer.botId, + origin.peer.conversationKey, + ) + const sessionId = typeof hit?.sessionId === 'string' ? hit.sessionId.trim() : '' + if (sessionId) return { sessionId, via: 'conversation' } + } catch { + // try target / origin fallbacks + } + } + + if (delivery?.kind === 'im' && delivery.botId && delivery.targetId + && dshIm && typeof dshIm.resolveTargetSession === 'function') { + try { + const hit = await dshIm.resolveTargetSession(delivery.botId, delivery.targetId) + const sessionId = typeof hit?.sessionId === 'string' ? hit.sessionId.trim() : '' + if (sessionId) return { sessionId, via: 'target' } + } catch { + // fall through + } + } + + if (origin?.sessionId) return { sessionId: origin.sessionId, via: 'origin' } + return null +} + +/** + * Append the run summary into the effective origin/IM session without waking a turn. + * Prefer a live agent; otherwise resume the persisted session (never create — that + * id already belongs to the origin chat). + * @param {object} job + * @param {string} [summary] + * @param {{ + * dshIm?: object, + * getAgents?: () => object | undefined, + * pluginName?: string, + * newId?: () => string, + * }} deps + */ +export async function mirrorRunToSession(job, summary, deps = {}) { + const target = await resolveMirrorSession(job, deps) + if (!target?.sessionId) { + return { mirrored: false, skipped: true, reason: 'no_target' } + } + const agents = typeof deps.getAgents === 'function' ? deps.getAgents() : undefined + if (!agents) { + return { + mirrored: false, + skipped: true, + reason: 'agents_unavailable', + sessionId: target.sessionId, + via: target.via, + } + } + + let agent = typeof agents.get === 'function' ? agents.get(target.sessionId) : null + let resumed = false + if (!agent && typeof agents.resume === 'function') { + try { + const handle = await agents.resume({ resumeSessionId: target.sessionId }) + agent = handle?.agent || null + resumed = Boolean(agent) + } catch (error) { + return { + mirrored: false, + skipped: true, + reason: 'session_unavailable', + sessionId: target.sessionId, + via: target.via, + error: error instanceof Error ? error.message : String(error), + } + } + } + if (!agent) { + return { + mirrored: false, + skipped: true, + reason: 'session_unavailable', + sessionId: target.sessionId, + via: target.via, + } + } + + const text = formatRunResultBody(job, summary) + 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()}` + const message = { + id, + role: 'user', + content: [{ type: 'text', text }], + source: { kind: 'plugin', plugin }, + } + + try { + if (typeof agent.session?.append === 'function') { + agent.session.append('user/message', message, { surfaceOp: 'append' }) + return { + mirrored: true, + sessionId: target.sessionId, + via: target.via, + method: 'append', + resumed, + } + } + if (typeof agent.inject === 'function') { + agent.inject(message) + return { + mirrored: true, + sessionId: target.sessionId, + via: target.via, + method: 'inject', + resumed, + } + } + } catch (error) { + return { + mirrored: false, + skipped: true, + reason: 'mirror_failed', + sessionId: target.sessionId, + via: target.via, + error: error instanceof Error ? error.message : String(error), + } + } + return { + mirrored: false, + skipped: true, + reason: 'no_append', + sessionId: target.sessionId, + via: target.via, + } +} diff --git a/lib/fire.js b/lib/fire.js index 6bb7431..624b654 100644 --- a/lib/fire.js +++ b/lib/fire.js @@ -65,6 +65,7 @@ export function publicJob(job) { ...job.delivery.mentionName ? { mentionName: job.delivery.mentionName } : {}, } : { kind: 'dsh' }, + ...job.origin ? { origin: job.origin } : {}, schedule: job.schedule, createdAt: job.createdAt, updatedAt: job.updatedAt, diff --git a/lib/host.js b/lib/host.js index ed4d235..a1c1412 100644 --- a/lib/host.js +++ b/lib/host.js @@ -8,7 +8,7 @@ import { mkdir } from 'node:fs/promises' import { homedir } from 'node:os' import { basename, join } from 'node:path' import { claimOccurrence, executeClaimedRun, extractAssistantText, interruptActiveRuns, publicJob, settleRun, TITLE_PREFIX } from './fire.js' -import { deliverRunToIm, mergeDeliveryMention } from './delivery.js' +import { deliverRunToIm, mergeDeliveryMention, mirrorRunToSession, normalizeOrigin } from './delivery.js' import { workspaceVisibleIds } from './isolation.js' import { decideDispatch, nextFire, validateSchedule } from './scheduler.js' import { @@ -112,6 +112,7 @@ function runView(run) { * @param {string} [options.filePath] * @param {object} [options.sessionPort] { createAndPrompt, archiveSession, waitForTurn } * @param {() => object|undefined} [options.getDshIm] soft-injected proactive IM API + * @param {() => object|undefined} [options.getAgents] soft-injected agents service for session mirror * @param {{ warn?: Function, info?: Function }} [options.logger] * @param {number} [options.tickIntervalMs] */ @@ -121,6 +122,7 @@ export function createHostService(options = {}) { const store = options.store || createStore({ filePath }) const sessionPort = options.sessionPort || null const getDshIm = typeof options.getDshIm === 'function' ? options.getDshIm : () => undefined + const getAgents = typeof options.getAgents === 'function' ? options.getAgents : () => undefined const logger = options.logger || console const tickIntervalMs = Number(options.tickIntervalMs) > 0 ? Number(options.tickIntervalMs) : 15_000 let timer = null @@ -144,6 +146,36 @@ export function createHostService(options = {}) { } } + async function maybeMirrorResult(job, run) { + if (!job || !run || run.status === 'skipped') return + try { + const result = await mirrorRunToSession(job, run.summary || run.error || '', { + dshIm: getDshIm(), + getAgents, + pluginName: PLUGIN_NAME, + newId: () => randomUUID(), + }) + if (result?.mirrored) { + logger.info?.(`[dsh-ops-cron] mirrored run ${run.id} → session ${result.sessionId} via ${result.via} (${result.method}${result.resumed ? ', resumed' : ''})`) + } else if (result?.reason && result.reason !== 'no_target') { + logger.warn?.(`[dsh-ops-cron] mirror skipped for run ${run.id}: ${result.reason}${result.sessionId ? ` session=${result.sessionId}` : ''}${result.error ? ` (${result.error})` : ''}`) + } else if (!result?.skipped) { + logger.warn?.(`[dsh-ops-cron] mirror incomplete for run ${run.id}: ${result?.reason || 'unknown'}`) + } + } catch (error) { + logger.warn?.(`[dsh-ops-cron] mirror failed for job ${job.id}: ${error instanceof Error ? error.message : error}`) + } + } + + async function settleAndNotify(jobId, runId, claimedDecision) { + const state = await snapshot() + const run = (state.runs || []).find((row) => row.id === runId) + const job = getJob(state, jobId) + await maybeDeliverIm(job, run) + await maybeMirrorResult(job, run) + return { job: jobView(job), run: runView(run), decision: claimedDecision } + } + async function createJob(input) { const t = now() let created @@ -178,11 +210,15 @@ export function createHostService(options = {}) { delivery: patch.delivery !== undefined ? mergeDeliveryMention(job.delivery, patch.delivery) : job.delivery, + origin: patch.origin !== undefined + ? (normalizeOrigin(patch.origin) || undefined) + : job.origin, } const record = createJobRecord(nextInput, current, t) record.createdAt = job.createdAt record.lastRunAt = job.lastRunAt record.lastStatus = job.lastStatus + if (patch.origin === undefined && job.origin) record.origin = job.origin if (patch.schedule === undefined && patch.enabled === undefined) { record.nextRunAt = job.nextRunAt } @@ -246,13 +282,9 @@ export function createHostService(options = {}) { } const settled = await store.mutate((current) => settleRun(current, run.id, terminal, now())) run = (settled.runs || []).find((row) => row.id === claimed.run.id) - const job = getJob(settled, jobId) - await maybeDeliverIm(job, run) - return { job: jobView(job), run: runView(run), decision: claimed.decision } + return settleAndNotify(jobId, run.id, claimed.decision) } - const job = getJob(executed, jobId) - await maybeDeliverIm(job, run) - return { job: jobView(job), run: runView(run), decision: claimed.decision } + return settleAndNotify(jobId, claimed.run.id, claimed.decision) } async function tick() { diff --git a/lib/index.d.ts b/lib/index.d.ts index f117b0a..6794c52 100644 --- a/lib/index.d.ts +++ b/lib/index.d.ts @@ -42,5 +42,10 @@ export const TITLE_PREFIX: string export function claimOccurrence(state: object, jobId: string, now: number, trigger: string, policies?: object): object export function executeClaimedRun(state: object, runId: string, deps: object): Promise export function normalizeDelivery(input?: object): { kind: 'dsh' | 'im', botId?: string, targetId?: string } +export function normalizeOrigin(input?: object): { kind: 'web' | 'im', sessionId?: string, peer?: object } | null 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 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 74ee89e..3c2afb6 100644 --- a/lib/index.js +++ b/lib/index.js @@ -51,8 +51,13 @@ export { normalizeJobModel } from './store.js' export { claimOccurrence, executeClaimedRun, extractAssistantText, TITLE_PREFIX } from './fire.js' export { deliverRunToIm, + formatRunResultBody, + mirrorRunToSession, normalizeDelivery, + normalizeOrigin, resolveCreateDelivery, + resolveCreateOrigin, + resolveMirrorSession, } from './delivery.js' function loadPkg(id) { @@ -131,10 +136,12 @@ function registerWebRoute(ctx, handler) { export function apply(ctx, config = {}) { const entry = resolveConfig(config) const getDshIm = () => tryGet(ctx, 'dshIm') + const getAgents = () => tryGet(ctx, 'agents') const getAgentPresets = () => tryGet(ctx, 'agentPresets') const service = createHostService({ sessionPort: makeLiveSessionPort(ctx), getDshIm, + getAgents, logger: ctx.logger, }) diff --git a/lib/store.js b/lib/store.js index 5cdbd71..b82d7c2 100644 --- a/lib/store.js +++ b/lib/store.js @@ -8,7 +8,7 @@ import { dirname, join } from 'node:path' import { randomUUID } from 'node:crypto' import { applyRunIsolation, recordHiddenSession } from './isolation.js' import { nextFire, validateSchedule } from './scheduler.js' -import { normalizeDelivery } from './delivery.js' +import { normalizeDelivery, normalizeOrigin } from './delivery.js' import { normalizeAgentPresetId } from './preset.js' export const STORE_VERSION = 1 @@ -95,6 +95,7 @@ export function createJobRecord(input, state, now) { const timeoutMinutes = Number(input.timeoutMinutes) const model = normalizeJobModel(input) const delivery = normalizeDelivery(input.delivery) + const origin = normalizeOrigin(input.origin) const agentPreset = normalizeAgentPresetId(input.agentPreset ?? input.agent_preset) const job = { id: String(input.id || newId()), @@ -110,6 +111,7 @@ export function createJobRecord(input, state, now) { reasoningEffort: model.reasoningEffort, agentPreset, delivery, + ...origin ? { origin } : {}, schedule: schedule.kind === 'cron' ? { kind: 'cron', expr: schedule.expr, timezone: schedule.timezone } : { kind: 'at', at: new Date(schedule.at).toISOString(), timezone: schedule.timezone }, diff --git a/lib/tools.js b/lib/tools.js index e9d30b9..b88376b 100644 --- a/lib/tools.js +++ b/lib/tools.js @@ -4,7 +4,7 @@ */ import { formatInZone, resolveTodayAt } from './scheduler.js' -import { deliveryLine, resolveCreateDelivery } from './delivery.js' +import { deliveryLine, resolveCreateDelivery, resolveCreateOrigin } from './delivery.js' import { resolveCreateAgentPreset } from './preset.js' const JOB_SCHEMA = { @@ -22,6 +22,7 @@ const JOB_SCHEMA = { reasoningEffort: { type: 'string' }, agentPreset: { type: 'string' }, delivery: { type: 'object', additionalProperties: true }, + origin: { type: 'object', additionalProperties: true }, schedule: { type: 'object', additionalProperties: true }, createdAt: { type: 'number' }, updatedAt: { type: 'number' }, @@ -204,6 +205,7 @@ export function cronToolDefinitions(service, deps = {}) { const timeout = Number(args.timeout_minutes) try { const delivery = await resolveCreateDelivery(args, exec, { dshIm: getDshIm() }) + const origin = await resolveCreateOrigin(args, exec, { dshIm: getDshIm() }) const agentPreset = await resolveCreateAgentPreset(args, exec, { dshIm: getDshIm(), agentPresets: getAgentPresets(), @@ -217,6 +219,7 @@ export function cronToolDefinitions(service, deps = {}) { timeoutMinutes: Number.isFinite(timeout) && timeout > 0 ? timeout : undefined, enabled: args.enabled !== false, delivery, + ...origin ? { origin } : {}, agentPreset, }) return { job } diff --git a/test/delivery.test.js b/test/delivery.test.js index 1075d77..631bce9 100644 --- a/test/delivery.test.js +++ b/test/delivery.test.js @@ -3,8 +3,13 @@ import test from 'node:test' import { matchTargetForPeer, normalizeDelivery, + normalizeOrigin, resolveCreateDelivery, + resolveCreateOrigin, + resolveMirrorSession, + formatRunResultBody, deliverRunToIm, + mirrorRunToSession, } from '../lib/delivery.js' test('normalizeDelivery defaults to dsh', () => { @@ -19,6 +24,26 @@ test('normalizeDelivery requires botId+targetId for im', () => { ) }) +test('normalizeOrigin pins web and im shapes', () => { + assert.equal(normalizeOrigin(null), null) + assert.deepEqual( + normalizeOrigin({ kind: 'web', sessionId: 'sess-web' }), + { kind: 'web', sessionId: 'sess-web' }, + ) + assert.deepEqual( + normalizeOrigin({ + kind: 'im', + sessionId: 'sess-im', + peer: { botId: 'bot_1', conversationKey: 'direct:86138@s.whatsapp.net', kind: 'direct' }, + }), + { + kind: 'im', + sessionId: 'sess-im', + peer: { botId: 'bot_1', conversationKey: 'direct:86138@s.whatsapp.net', kind: 'direct' }, + }, + ) +}) + test('matchTargetForPeer matches whatsapp jid from conversationKey', () => { const peer = { botId: 'bot_1', conversationKey: 'direct:8613800000000@s.whatsapp.net' } const targets = [ @@ -55,6 +80,37 @@ test('resolveCreateDelivery auto-matches peer target', async () => { assert.deepEqual(delivery, { kind: 'im', botId: 'bot_1', targetId: 'wa-dm' }) }) +test('resolveCreateOrigin captures im peer from caller session', async () => { + const dshIm = { + resolveSessionPeer: async () => ({ + botId: 'bot_1', + conversationKey: 'group:120363@g.us:user:86138@s.whatsapp.net', + conversationId: '120363@g.us', + kind: 'group', + }), + } + const origin = await resolveCreateOrigin({}, { + agent: { session: { id: 'sess-group' } }, + }, { dshIm }) + assert.deepEqual(origin, { + kind: 'im', + sessionId: 'sess-group', + peer: { + botId: 'bot_1', + conversationKey: 'group:120363@g.us:user:86138@s.whatsapp.net', + conversationId: '120363@g.us', + kind: 'group', + }, + }) +}) + +test('resolveCreateOrigin falls back to web when no peer', async () => { + const origin = await resolveCreateOrigin({}, { + agent: { session: { id: 'sess-web' } }, + }, { dshIm: { resolveSessionPeer: async () => null } }) + assert.deepEqual(origin, { kind: 'web', sessionId: 'sess-web' }) +}) + test('resolveCreateDelivery auto-creates group target when none exists', async () => { const created = [] const dshIm = { @@ -151,7 +207,6 @@ test('resolveCreateDelivery pins group creator mention from peer', async () => { kind: 'im', botId: 'bot_wa', targetId: 'ops-group', - // Prefer phone JID over LID so WhatsApp can show nickname + notify. mentionJid: '8613800000000@s.whatsapp.net', mentionName: 'Alice', }) @@ -196,3 +251,123 @@ test('deliverRunToIm sends via dshIm', async () => { assert.match(sent[0][2], /定时任务/) assert.match(sent[0][2], /hello world/) }) + +test('resolveMirrorSession prefers live conversation binding over stale origin.sessionId', async () => { + const hit = await resolveMirrorSession({ + origin: { + kind: 'im', + sessionId: 'stale-sess', + peer: { botId: 'bot_1', conversationKey: 'direct:86138@s.whatsapp.net' }, + }, + delivery: { kind: 'im', botId: 'bot_1', targetId: 'wa-dm' }, + }, { + dshIm: { + resolveConversationSession: async () => ({ + sessionId: 'live-sess', + botId: 'bot_1', + conversationKey: 'direct:86138@s.whatsapp.net', + }), + }, + }) + assert.deepEqual(hit, { sessionId: 'live-sess', via: 'conversation' }) +}) + +test('resolveMirrorSession falls back to origin.sessionId', async () => { + const hit = await resolveMirrorSession({ + origin: { kind: 'web', sessionId: 'origin-web' }, + delivery: { kind: 'dsh' }, + }, {}) + assert.deepEqual(hit, { sessionId: 'origin-web', via: 'origin' }) +}) + +test('mirrorRunToSession appends plugin notice without waking a turn', async () => { + const appended = [] + const result = await mirrorRunToSession( + { + name: '日报', + origin: { kind: 'web', sessionId: 'origin-1' }, + delivery: { kind: 'dsh' }, + }, + 'line one', + { + getAgents: () => ({ + get: () => ({ + session: { + append: (...args) => { appended.push(args) }, + }, + }), + }), + newId: () => 'msg-1', + pluginName: 'dsh-ops-cron', + }, + ) + assert.equal(result.mirrored, true) + assert.equal(result.method, 'append') + assert.equal(appended.length, 1) + 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.equal(appended[0][1].source.kind, 'plugin') + assert.deepEqual(appended[0][2], { surfaceOp: 'append' }) +}) + +test('mirrorRunToSession resumes a cold origin session instead of create', async () => { + const appended = [] + const resumed = [] + const result = await mirrorRunToSession( + { + name: '冷会话', + origin: { kind: 'web', sessionId: 'cold-1' }, + }, + 'hello', + { + getAgents: () => ({ + get: () => null, + create: async () => { + throw new Error('create must not be used for origin mirror') + }, + resume: async (opts) => { + resumed.push(opts) + return { + agent: { + session: { + append: (...args) => { appended.push(args) }, + }, + }, + } + }, + }), + newId: () => 'msg-2', + }, + ) + assert.equal(result.mirrored, true) + assert.equal(result.resumed, true) + assert.deepEqual(resumed, [{ resumeSessionId: 'cold-1' }]) + assert.equal(appended[0][0], 'user/message') +}) + +test('mirrorRunToSession skips when origin session is unavailable', async () => { + const result = await mirrorRunToSession( + { + name: 'gone', + origin: { kind: 'web', sessionId: 'missing' }, + }, + 'x', + { + getAgents: () => ({ + get: () => null, + resume: async () => { + throw new Error('not found') + }, + }), + }, + ) + assert.equal(result.mirrored, false) + assert.equal(result.reason, 'session_unavailable') +}) + +test('formatRunResultBody matches IM delivery header', () => { + const body = formatRunResultBody({ name: '晨报' }, 'ok') + assert.match(body, /^【定时任务 · 晨报】\n/) + assert.match(body, /ok$/) +}) diff --git a/test/package.test.js b/test/package.test.js index 7e700e9..c76b7b6 100644 --- a/test/package.test.js +++ b/test/package.test.js @@ -44,6 +44,8 @@ test('installable bundle declares host apply, client half, unique id, and no @de const host = await readFile(join(root, 'lib/host.js'), 'utf8') assert.match(host, /agentDefaultModel/) assert.match(host, /deliverRunToIm/) + assert.match(host, /mirrorRunToSession/) + assert.match(host, /getAgents/) assert.match(host, /PLUGIN_NAME = 'dsh-ops-cron'/) assert.match(host, /API_PREFIX = '\/dsh-ops-cron'/) assert.match(host, /resolveSessionPlacement/) @@ -68,7 +70,10 @@ test('installable bundle declares host apply, client half, unique id, and no @de assert.match(preset, /resolveCreateAgentPreset/) const delivery = await readFile(join(root, 'lib/delivery.js'), 'utf8') assert.match(delivery, /resolveCreateDelivery/) + assert.match(delivery, /resolveCreateOrigin/) + assert.match(delivery, /mirrorRunToSession/) assert.match(delivery, /deliverRunToIm/) + assert.match(client, /origin/) const prd = await readFile(join(root, 'docs/PRD.md'), 'utf8') assert.match(prd, /定时任务/) assert.match(prd, /dsh-schedule/)