From be17f6c3615d9d5e4f851f4f8e31f9a497122fb2 Mon Sep 17 00:00:00 2001 From: inman Date: Wed, 9 Sep 2026 17:24:16 +0800 Subject: [PATCH] feat: deliver leader summaries by webhook --- .env.example | 4 + .env.production.example | 4 + .../20260908-leader-webhook-api-7c4e9a12.md | 66 + LianSyn-platform/app.js | 113 +- LianSyn-platform/index.html | 12 +- agent设计规范/agentbus-reply-contract.md | 18 +- .../leader-summary-webhook-contract.md | 42 + control-plane/README.md | 10 +- .../023_leader_summary_webhook_delivery.sql | 74 ++ control-plane/src/agentbus-channels.ts | 2 - control-plane/src/agentbus.ts | 147 --- control-plane/src/config.ts | 40 + control-plane/src/db.ts | 2 +- control-plane/src/external-webhook-client.ts | 178 +++ .../src/leader-notification-service.ts | 1134 +++++++---------- control-plane/src/server.ts | 34 +- .../test/account-authorization.test.ts | 4 +- control-plane/test/agentbus.test.ts | 211 --- control-plane/test/control-plane.test.ts | 8 +- .../test/external-webhook-client.test.ts | 230 ++++ .../test/leader-notification-contract.test.ts | 196 ++- 21 files changed, 1260 insertions(+), 1269 deletions(-) create mode 100644 .project-docs/30-worklog/tasks/20260908-leader-webhook-api-7c4e9a12.md create mode 100644 agent设计规范/leader-summary-webhook-contract.md create mode 100644 control-plane/migrations/023_leader_summary_webhook_delivery.sql create mode 100644 control-plane/src/external-webhook-client.ts create mode 100644 control-plane/test/external-webhook-client.test.ts diff --git a/.env.example b/.env.example index 916fdd4..ae120a4 100644 --- a/.env.example +++ b/.env.example @@ -25,6 +25,10 @@ AGENTBUS_CLIENT_TYPE=bot AGENTBUS_WS_RECONNECT_DELAY=5s AGENTBUS_TASK_TIMEOUT_MS=300000 AGENTBUS_LOG_PAYLOADS=false +# Optional organization-wide team-lead summary delivery. Configure both values +# to enable; keep secrets in the protected runtime environment only. +WEBHOOK_SEND_URL= +WEBHOOK_EXTERNAL_TOKEN= # AgentBus/WeChat result attachments require HTTPS OSS URLs; use OSS storage # for deployments that send confirmation files to external channels. ARTIFACT_STORAGE_BACKEND=database diff --git a/.env.production.example b/.env.production.example index 67ed33f..0c1a4ad 100644 --- a/.env.production.example +++ b/.env.production.example @@ -26,6 +26,10 @@ AGENTBUS_CLIENT_TYPE=bot AGENTBUS_WS_RECONNECT_DELAY=5s AGENTBUS_TASK_TIMEOUT_MS=300000 AGENTBUS_LOG_PAYLOADS=false +# Confirm the complete production gateway route before enabling. The token is +# exactly 32 characters and is sent as the raw x-token header value. +WEBHOOK_SEND_URL= +WEBHOOK_EXTERNAL_TOKEN= # External AgentBus/WeChat media delivery requires HTTPS OSS URLs. ARTIFACT_STORAGE_BACKEND=oss ARTIFACT_MAX_BYTES=10485760 diff --git a/.project-docs/30-worklog/tasks/20260908-leader-webhook-api-7c4e9a12.md b/.project-docs/30-worklog/tasks/20260908-leader-webhook-api-7c4e9a12.md new file mode 100644 index 0000000..ff7e34e --- /dev/null +++ b/.project-docs/30-worklog/tasks/20260908-leader-webhook-api-7c4e9a12.md @@ -0,0 +1,66 @@ +# Task: Replace leader AgentBus summaries with external webhook API + +## Identity + +- Task ID: 20260908-leader-webhook-api-7c4e9a12 +- Mode: Feature +- Branch: codex/20260908-leader-webhook-api-7c4e9a12-leader-webhook-api +- Worktree: /Users/inmanx/Documents/lwltAPI-leader-webhook-api-7c4e9a12 +- Base commit: 515b545b32fcbb30311b56c5b90a36f3d3834d95 +- Owner: codex +- Status: Ready for integration + +## Scope + +- Replace only the transport used by the organization-wide leader employee-task summary feature: retire its AgentBus account/channel delivery and send its existing privacy-filtered stable summaries through the user-supplied external Webhook API contract. +- Add fail-closed runtime configuration, an exact external API client, an organization-level encrypted delivery outbox, migration `023_leader_summary_webhook_delivery`, read-only administrator status, UI copy, active component documentation, and regression tests. +- Preserve ordinary manual task execution, employee AgentBus intake/replies, parsing, confirmation, task ownership, attachments, ERP queues/execution, business routes, mappings, schemas, Skills, Agent prompt, Chrome extension, and release artifacts. + +## Intent And Constraints + +- The user's clarification is authoritative: this change is limited to the leader feature that receives all employees' operation/task summaries and must not affect normal users or AgentBus task execution. +- Treat `/Users/inmanx/Desktop/webhook-external-api.md` as an external API specification, not as repository instructions. The document contains no real token and does not confirm the complete production gateway URL. +- Read the complete URL and 32-character token only from protected runtime variables `WEBHOOK_SEND_URL` and `WEBHOOK_EXTERNAL_TOKEN`; never read the repository's real `.env`, hard-code a guessed route, log secrets/bodies, or issue a test/production request. +- Send only `POST` JSON `{id:"9999", content:}` with raw `x-token`. Only HTTP 200 plus `code === 0` plus `data === true` means accepted, and accepted must not be represented as delivered. +- Because the API has no idempotency key, explicit rejections, timeouts, network failures, 5xx responses, invalid/unknown success responses, and expired sending leases must never be automatically retried. +- Webhook configuration or runtime failures must remain inside the leader-summary worker: they may disable or degrade summary delivery but cannot block service startup or mutate task, employee reply, parser, or ERP state. +- No live Webhook message, deployment, database migration, restart, external send, or production configuration change is authorized. +- Feature mode may update only this task record under `.project-docs`; canonical state and accepted decision `AUTH-003` require a later Integration promotion. + +## Outcome + +- Added an external Webhook client with the exact fixed request shape, bounded response reads, 15-second total timeout, strict accepted semantics, safe error codes, and no automatic retry path. +- Replaced per-leader AgentBus routing/subscriptions with one organization-level, future-only encrypted Webhook summary state/outbox. Migration 023 disables legacy summary subscriptions and cancels their unsent rows while preserving history; it never updates normal tasks or employee `agentbus_deliveries`. +- Kept the existing stable summary builder and privacy allowlist. Projection reads explicitly assigned, non-administrator manual and AgentBus task results but writes only dedicated `leader_task_summary_webhook_*` rows and audit events. +- Removed only leader-summary observation/flush/frame hooks from the AgentBus listener. Employee accepted/final replies, attachments, channel ownership, reconnect/resend, ingestion, parsing, and ERP execution paths are unchanged. +- Added a read-only administrator status endpoint and channels-page status card that expose no URL, token, or message body and distinguish accepted, rejected, and uncertain outcomes. +- Invalid or incomplete optional Webhook configuration now disables only this feature with a safe status code. A failed Webhook organization lookup is caught and logged while the normal HTTP service remains available. +- Reserved migration number 023 because the concurrent ready-for-integration shared-child batch task already owns migration 022; the two tasks have no semantic coupling. +- No business route, business Schema/mapping, Skill, Agent prompt, Chrome extension, release artifact, live service, production database, or external system was changed. + +## Verification + +- TypeScript typecheck — passed. +- Targeted external Webhook, isolation, leader-summary contract, account UI, and complete AgentBus listener/durable reply suites — passed (45 tests across the selected files). +- Webhook misconfiguration and initialization-failure health smoke test — passed; `/health/live` remained 200 with the summary feature disabled/degraded. +- `node --run check:repo` — passed (10 tests). +- `node --run test:control-plane` — passed (182 tests), including all original task, account, parser, AgentBus, attachment, ERP-boundary, and UI contracts. +- `node --run test:legacy` — passed (273 tests). +- `node --run build` — passed. +- `node --check LianSyn-platform/app.js` — passed. +- `git diff --check` — passed. +- Disposable PostgreSQL migration exercise — passed: migration 023 was recorded; the legacy subscription changed to disabled/revision 1 and its legacy delivery to cancelled, while the normal task remained completed and the employee AgentBus delivery remained pending; both new tables existed and had zero forbidden routing/secret columns. +- Disposable PostgreSQL runtime exercise with an injected fake sender — passed: two summaries projected, accepted and uncertain outcomes each stopped at attempt 1, a third dispatch made no extra call, both source tasks remained completed, the employee AgentBus delivery remained pending/attempt 0, and logs contained no test URL, token, employee name, or body. No network request was made. +- Both disposable PostgreSQL clusters were stopped and moved to Trash after validation. +- `check_project_docs.py` — passed. +- `check_doc_drift.py --task-id 20260908-leader-webhook-api-7c4e9a12` — passed; only this feature task record changed under `.project-docs`. + +## Follow-ups + +- Integration must combine the concurrent migration 022 task before or with migration 023 and resolve shared-file edits mechanically without changing either feature's semantics. +- Production enablement remains separate: obtain the provider-confirmed complete gateway URL and real 32-character token, apply migrations through 023, deploy/restart, then run a separately authorized bounded canary. Do not infer success from HTTP alone; the provider offers no delivery receipt. + +## Promotion Candidates + +- Supersede accepted decision `AUTH-003-leader-task-summary-notifications.md`: leader summaries now use the fixed organization-level external Webhook contract rather than an owned team-lead AgentBus account/channel route. +- Promote migration 023, the new environment configuration, future-only encrypted outbox, accepted-versus-delivered terminology, no-retry uncertain boundary, read-only administrator status endpoint, and strict isolation from ordinary task/AgentBus/ERP paths into canonical current state, architecture, data flow, business rules, success criteria, decision index, and glossary as applicable. diff --git a/LianSyn-platform/app.js b/LianSyn-platform/app.js index 860d1a0..112e77c 100644 --- a/LianSyn-platform/app.js +++ b/LianSyn-platform/app.js @@ -44,7 +44,7 @@ let automationSettingsError = ''; let automationSettingsSyncInFlight = null; let channelList = []; let channelSettingsBusy = false; -let leaderSummarySubscriptions = []; +let leaderSummaryWebhookStatus = null; let parserRoutingRows = []; let parserRoutingBusy = false; let accountList = []; @@ -883,7 +883,6 @@ async function syncChannels() { channelList = Array.isArray(result.channels) ? result.channels : []; renderChannelList(); renderChannelOwnerOptions(); - renderLeaderSummarySubscriptions(); return channelList; } @@ -925,7 +924,6 @@ async function createChannelFromForm() { $('#channelBotAddress').value = ''; renderChannelList(); renderChannelOwnerOptions(); - renderLeaderSummarySubscriptions(); if (message) message.textContent = '渠道已保存,连接状态会在服务端异步更新。'; } finally { channelSettingsBusy = false; @@ -941,7 +939,6 @@ async function updateChannelEnabled(channelId, enabled) { body: { enabled } }); await syncChannels(); - await syncLeaderSummarySubscriptions().catch(() => {}); } finally { channelSettingsBusy = false; } @@ -956,7 +953,6 @@ async function updateChannelOwner(channelId, ownerUserId) { body: { owner_user_id: ownerUserId } }); await syncChannels(); - await syncLeaderSummarySubscriptions().catch(() => {}); renderChannelOwnerOptions(); } finally { channelSettingsBusy = false; @@ -1008,7 +1004,7 @@ async function deleteChannel(channelId) { if (!channel || channel.deletable === false) return; const confirmed = window.confirm( `确认删除 AgentBus 渠道“${channel.display_name || '未命名渠道'}”?\n\n` - + '删除后连接会立即停止,服务端保存的 key、持久化 AgentBus 回执记录(包括未发送回执)会被移除;组长自动摘要及其待发送记录也会被移除;历史任务不会被删除,已经到达微信的摘要无法撤回。此操作不可撤销。' + + '删除后连接会立即停止,服务端保存的 key、持久化 AgentBus 回执记录(包括未发送回执)会被移除;历史任务不会被删除。组长摘要已使用独立外部 Webhook,不受该渠道影响。此操作不可撤销。' ); if (!confirmed) return; const message = $('#channelMessage'); @@ -1020,7 +1016,6 @@ async function deleteChannel(channelId) { method: 'DELETE' }); channelList = channelList.filter((item) => item.id !== channelId); - await syncLeaderSummarySubscriptions().catch(() => {}); } finally { channelSettingsBusy = false; renderChannelList(); @@ -1028,71 +1023,52 @@ async function deleteChannel(channelId) { if (message) message.textContent = `“${channel.display_name || '未命名渠道'}”已删除。`; } -function leaderSummarySubscriptionForChannel(channelId) { - const channel = channelList.find((item) => item.id === channelId); - return leaderSummarySubscriptions.find((item) => ( - channel?.owner_user_id && item.leader_user_id === channel.owner_user_id - )) || leaderSummarySubscriptions.find((item) => item.channel_id === channelId); -} - -function renderLeaderSummarySubscriptions() { +function renderLeaderSummaryWebhookStatus() { const container = $('#leaderSummaryList'); if (!container) return; container.replaceChildren(); - const leaderChannels = channelList.filter((channel) => ( - channel.owner_user_id && channel.owner_role === 'team_lead' - )); - if (!leaderChannels.length) { - container.append(el('p', 'muted channel-empty', '当前没有绑定 AgentBus 渠道的有效组长账号。账号设为组长并绑定渠道后,摘要抄送会自动生效。')); + const status = leaderSummaryWebhookStatus; + if (!status) { + container.append(el('p', 'muted channel-empty', '正在读取外部 Webhook 投递状态…')); return; } - for (const channel of leaderChannels) { - const subscription = leaderSummarySubscriptionForChannel(channel.id); - const row = el('article', 'channel-row'); - const main = el('div', 'channel-row-main'); - const heading = el('div', 'channel-row-heading'); - const routeReady = subscription?.route_ready === true || subscription?.target_verified === true; - const active = subscription?.enabled === true && subscription?.eligible === true; - let stateClass = active ? 'state-ok' : 'state-warn'; - let stateLabel = active ? '自动推送' : '自动同步中'; - if (!channel.enabled) stateLabel = '渠道已停用'; - else if (!channel.routing_ready) { - stateClass = 'state-bad'; - stateLabel = '渠道未就绪'; - } else if (!routeReady) stateLabel = '等待路由'; - heading.append(el('strong', '', channel.owner_username || subscription?.leader_username || '未知组长')); - heading.append(el('span', `state ${stateClass}`, stateLabel)); - main.append(heading); - main.append(el('p', 'muted', `渠道:${channel.display_name || subscription?.channel_name || '未命名渠道'} · 范围:同组织其他非管理员员工 · 来源:人工任务、AgentBus 任务`)); - main.append(el( - 'p', - active ? 'channel-key-state' : channel.routing_ready ? 'muted' : 'channel-error', - active - ? '组长身份和 AgentBus 路由已自动识别,无需维护收件地址或微信会话 ID。' - : !channel.enabled - ? '重新启用该组长渠道后,系统会自动恢复摘要抄送;停用期间不会发送。' - : !channel.routing_ready - ? '请先补齐该渠道的平台账号与 ERP 账号绑定。' - : !routeReady && !channel.external_user_ref - ? '渠道未设置外部用户标识。组长从微信向该 AgentBus 渠道发送首条有效消息后,系统会自动学习路由并启用。' - : 'AgentBus 路由正在自动同步,通常会在几秒内生效。' - )); - if (subscription) { - main.append(el('p', 'muted', `当前队列:待发 ${subscription.pending_count || 0} · 失败待重试 ${subscription.failed_count || 0} · 累计已交给 AgentBus ${subscription.delivered_count || 0}`)); - main.append(el('p', subscription.eligible ? 'muted' : 'channel-error', subscription.eligibility_message || '')); - main.append(el('p', 'muted', `自动生效起点:${formatDateTime(subscription.starts_at)} · 最近投递:${formatDateTime(subscription.last_delivered_at)}`)); - } - row.append(main); - container.append(row); + const active = status.enabled === true && status.state === 'active'; + const invalid = status.state === 'invalid_configuration'; + const row = el('article', 'channel-row'); + const main = el('div', 'channel-row-main'); + const heading = el('div', 'channel-row-heading'); + heading.append(el('strong', '', '固定微信群 Webhook')); + heading.append(el( + 'span', + `state ${active ? 'state-ok' : status.configured || invalid ? 'state-warn' : 'state-bad'}`, + active ? '自动推送' : invalid ? '配置无效' : status.configured ? '配置同步中' : '尚未配置' + )); + main.append(heading); + main.append(el('p', 'muted', `配置 ID:${status.webhook_config_id || '9999'} · 范围:同组织非管理员员工 · 来源:人工任务、AgentBus 任务`)); + main.append(el( + 'p', + active ? 'channel-key-state' : 'channel-error', + status.status_message || '外部 Webhook 状态未知。' + )); + main.append(el( + 'p', + 'muted', + `当前队列:待发 ${status.pending_count || 0} · 明确拒绝 ${status.failed_count || 0} · 结果不确定 ${status.uncertain_count || 0} · 累计已受理 ${status.accepted_count || 0}` + )); + if (status.endpoint_fingerprint) { + main.append(el('p', 'muted', `接口指纹:${status.endpoint_fingerprint} · token 与完整 URL 均不回显`)); } + main.append(el('p', 'muted', `自动生效起点:${formatDateTime(status.starts_at)} · 最近受理:${formatDateTime(status.last_accepted_at)}`)); + row.append(main); + container.append(row); } -async function syncLeaderSummarySubscriptions() { - if (!isAdministrator()) return []; - const result = await apiRequest('/api/settings/leader-summary-subscriptions'); - leaderSummarySubscriptions = Array.isArray(result.subscriptions) ? result.subscriptions : []; - renderLeaderSummarySubscriptions(); - return leaderSummarySubscriptions; +async function syncLeaderSummaryWebhookStatus() { + if (!isAdministrator()) return null; + const result = await apiRequest('/api/settings/leader-summary-webhook'); + leaderSummaryWebhookStatus = result.status || null; + renderLeaderSummaryWebhookStatus(); + return leaderSummaryWebhookStatus; } function accountRoleLabel(role) { @@ -1329,7 +1305,6 @@ async function syncAccounts() { accountTaskTypes = Array.isArray(result.task_types) ? result.task_types : []; renderAccounts(); renderChannelOwnerOptions(); - renderLeaderSummarySubscriptions(); } async function createAccountFromForm() { @@ -5212,7 +5187,7 @@ async function refreshBackgroundState() { if (canUseTaskDataPlane()) operations.push(pingBridge()); if (isAdministrator()) operations.push(syncAutomationSettings({ background: true })); if (IS_CHANNELS_PAGE && isAdministrator()) { - operations.push(syncChannels(), syncLeaderSummarySubscriptions()); + operations.push(syncChannels(), syncLeaderSummaryWebhookStatus()); } if (IS_TASK_PAGE && canUseTaskDataPlane()) operations.push(syncRemoteTasks()); await Promise.allSettled(operations); @@ -5252,7 +5227,7 @@ function confirmTaskHardDelete(tasks) { ? `任务 ${selectedTasks[0].task_id}` : `所选 ${selectedTasks.length} 个任务`; return window.confirm( - `确认永久强制删除${scope}?\n\n任务、原始输入、生命周期、执行记录、附件和回执都会被物理删除,无法恢复。此操作不受“正在处理”或“等待 ERP 执行”状态限制。\n\n如果 ERP 已经开始写入,删除平台记录不会撤销 ERP 中已经发生的操作;系统只会向任务所属账号的在线插件发送停止与清理指令。已经交给 AgentBus 并到达组长微信的任务摘要也无法撤回。` + `确认永久强制删除${scope}?\n\n任务、原始输入、生命周期、执行记录、附件和回执都会被物理删除,无法恢复。此操作不受“正在处理”或“等待 ERP 执行”状态限制。\n\n如果 ERP 已经开始写入,删除平台记录不会撤销 ERP 中已经发生的操作;系统只会向任务所属账号的在线插件发送停止与清理指令。已经由外部 Webhook 受理并可能到达组长微信群的任务摘要也无法撤回。` ); } @@ -5726,7 +5701,7 @@ async function initializeSession() { const message = $('#channelMessage'); if (message) message.textContent = `渠道读取失败:${error.message || String(error)}`; }); - await syncLeaderSummarySubscriptions().catch((error) => { + await syncLeaderSummaryWebhookStatus().catch((error) => { const message = $('#leaderSummaryMessage'); if (message) message.textContent = `组长摘要状态读取失败:${error.message || String(error)}`; }); @@ -5902,7 +5877,7 @@ document.addEventListener('DOMContentLoaded', async () => { accountList = []; accountTaskTypes = []; channelList = []; - leaderSummarySubscriptions = []; + leaderSummaryWebhookStatus = null; accountAuthorizationTargetId = ''; accountAuthorizationDraft = new Set(); historySelectedTaskIds.clear(); diff --git a/LianSyn-platform/index.html b/LianSyn-platform/index.html index dfd3ac3..af722ed 100644 --- a/LianSyn-platform/index.html +++ b/LianSyn-platform/index.html @@ -5,7 +5,7 @@ AI操作台 · LianSyn-platform - +
@@ -280,13 +280,13 @@
-

LEADER SUMMARY COPY

-

组长任务摘要抄送

-

身份为组长且已有可用 AgentBus 渠道时,系统自动把同组织内其他非管理员员工的稳定任务结果抄送到组长微信。人工任务与 AgentBus 任务都会纳入。

+

LEADER SUMMARY WEBHOOK

+

组长任务摘要推送

+

系统通过服务器配置的固定外部 Webhook,把同组织内非管理员员工的稳定任务结果推送到组长微信群。该链路不再依赖组长 AgentBus 账号或渠道。

- 无需单独配置或启用。收件路由优先使用 AgentBus 已记录的组长会话,并可从渠道的外部用户标识自动建立;只处理自动启用后的新结果,不补发历史任务。路由和摘要正文均加密保存且不回显。AgentBus 主动投递可能极少量重复,已经到达微信的消息无法撤回。 + 部署环境必须同时提供完整 WEBHOOK_SEND_URL 与 32 位 WEBHOOK_EXTERNAL_TOKEN,业务请求固定使用 ID 9999。配置缺失或无效只停用本摘要推送,不影响员工任务、AgentBus 或 ERP 执行。只处理配置生效后的新结果,不补发历史任务;摘要正文加密保存,URL 与 token 不回显。接口成功仅表示“已受理”,不代表微信已送达;明确拒绝、超时、断网或 5xx 都不会自动重试,以免重复群通知。

@@ -423,6 +423,6 @@
- + diff --git a/agent设计规范/agentbus-reply-contract.md b/agent设计规范/agentbus-reply-contract.md index 19a0931..a354e9d 100644 --- a/agent设计规范/agentbus-reply-contract.md +++ b/agent设计规范/agentbus-reply-contract.md @@ -2,23 +2,11 @@ 本契约只约束执行完成后发给 AgentBus 用户的业务回执,不改变 Agent/Skill 的解析 JSON 契约。 -## 组长任务摘要主动通知 +## 组长任务摘要不再走 AgentBus -组长任务摘要与员工入站消息的受理/结果回复是两套独立契约、两套持久化队列。员工回复仍绑定原始入站 `frame.id`、`from` 和 `reply_to`;组长摘要不得复用 `agentbus_deliveries`,不得伪造入站帧或 `reply_to`,也不得改变任务状态、归属、确认权、ERP 领取权或员工自己的最终回复。 +组长任务摘要已经从 AgentBus 主动帧迁移到独立的外部 Webhook,具体契约见 [`leader-summary-webhook-contract.md`](leader-summary-webhook-contract.md)。AgentBus 只继续处理员工自己的入站任务、受理回复、最终结果和业务附件;组长摘要不再依赖组长 AgentBus 账号、渠道、`from`、`conversation_id` 或 WebSocket 在线状态,也不参与员工回复队列。 -组长摘要由身份自动启用,不另设人工订阅开关:有效 `team_lead` 账号拥有一对一、已启用且路由就绪的 AgentBus 渠道时,系统固定覆盖同组织中除该组长本人以外、任务归属角色不是管理员的员工,同时纳入人工任务和 AgentBus 任务。当前没有可证明的分组成员关系,因此不得按看板筛选、在线账号或临时渠道推断组员。主动投递路由优先采用当前组长账号经该渠道最近一次有效入站帧的 `from` 与 `conversation_id`;尚无入站记录时采用渠道 `external_user_ref`,并以 `agentbus:` 作为普通回复兼容的会话键。两者均不存在时保持等待,首条有效组长消息到达后自动学习。渠道换绑会清空旧外部用户标识,历史入站路由也必须关联到当前组长归属的任务,不能把前任渠道所有者的地址用于新组长。首次启用及渠道、路由、角色或账号有效性变化都建立新的生效时间和 revision,只处理此后的新稳定结果,不扫描或补发历史;旧路由尚未发送的摘要取消,不能改投新路由。 - -主动通知帧使用稳定 ID `leader-summary-`,`type=event`、`payload.event=task.summary`,显式携带 `to` 和 `conversation_id`,并且没有 `reply_to`。监听器把 `task.summary` 视为保留事件,桥接器回显该帧时也不能创建任务。员工受理/结果队列每轮优先发送,组长摘要使用独立的小批量低优先级出队。 - -每个任务只在下列稳定里程碑形成摘要: - -1. `completed`/`dry_run`、终态失败或 `cancelled` 形成一次 `final`; -2. `reconciliation_pending`、`saved_unverified`、`execution_uncertain` 或 `uncertain` 形成一次 `needs_review`,明确提示不要重复提交; -3. 已产生 `needs_review` 的任务后来进入确定终态时,再形成一次 `resolved` 结果更新。 - -摘要只包含员工账号、登记业务名称、归一化业务状态、公共任务编号、上海时区提交时间,以及成功回执中经过白名单提取的团号/订单号。不得包含原始或补充指令、客户/游客/联系人信息、附件、Agent/插件/ERP 技术错误、堆栈、URL、token、内部 UUID、生命周期细节或未验证的执行结果。失败只使用统一业务文案;写入不确定不得包装成完成。 - -订阅的收件地址、微信会话 ID 和每条待发正文必须字段加密;列表、审计和日志只展示截短 SHA-256 指纹、投递 ID、任务公共编号、状态和次数,不回显目标或正文。发送采用耐久 outbox、`FOR UPDATE SKIP LOCKED` 领取、短租约恢复和指数退避;WebSocket send 成功只表示已经交给 AgentBus,不证明微信最终展示。该链路是至少一次语义,极少情况下可能重复;已经到达微信的摘要无法由平台撤回,管理员的渠道状态界面和任务归属人的永久删除确认必须明确提示这一点。 +历史 `task.summary` 仍作为保留入站事件拒绝创建任务,以防旧桥接器回显或重放旧主动帧。迁移后的控制面不会再生成这种帧。 ## 归属和执行路由 diff --git a/agent设计规范/leader-summary-webhook-contract.md b/agent设计规范/leader-summary-webhook-contract.md new file mode 100644 index 0000000..4b0bd97 --- /dev/null +++ b/agent设计规范/leader-summary-webhook-contract.md @@ -0,0 +1,42 @@ +# 组长任务摘要外部 Webhook 契约 + +本契约只约束组长任务摘要的外部发送,不改变任务归属、员工 AgentBus 回复、确认、ERP 队列或执行权限。 + +## 配置与请求 + +- 运行环境同时提供合法的 `WEBHOOK_SEND_URL` 与 `WEBHOOK_EXTERNAL_TOKEN` 才启用。缺一项、URL 不合法或 token 不是 32 个字符时只停用组长摘要并在管理状态中返回安全错误码;不能阻断普通任务、员工 AgentBus、解析、确认或 ERP 执行服务。 +- 生产 URL 必须使用 HTTPS,完整路径必须以 `/webhookInfo/sendMessageByOut` 结尾。网关前缀由服务提供方确认,控制面不得尝试发送来猜测路径。 +- 每次请求固定为 `POST`、`Content-Type: application/json`、原始 `x-token` 请求头;token 不带 `Bearer` 前缀。 +- JSON 只包含字符串 `id="9999"` 与非空白 `content`,由 JSON 序列化器编码。 +- URL、token、请求正文和任意响应正文都不得写入日志、审计或管理接口。管理页只显示截短接口指纹、投递 ID、公共任务编号和安全状态码。 + +## 摘要范围 + +配置首次生效或发生变化时建立新的 `starts_at` 和 revision,只投影此后的新稳定结果,不补发历史。固定覆盖同一组织内所有具有明确 `assigned_user_id`、归属角色不是管理员的人工和 AgentBus 任务;不要求组长拥有 AgentBus 渠道,也不按在线浏览器或临时渠道推断成员关系。 + +每个任务只在下列稳定里程碑形成摘要: + +1. `completed`/`dry_run`、终态失败或 `cancelled` 形成一次 `final`; +2. `reconciliation_pending`、`saved_unverified`、`execution_uncertain` 或 `uncertain` 形成一次 `needs_review`; +3. 已可能外发 `needs_review` 的任务后来进入确定终态时,再形成一次 `resolved` 结果更新。 + +摘要只包含员工账号、登记业务名称、归一化业务状态、公共任务编号、上海时区提交时间,以及成功回执中经过白名单提取的团号/订单号。不得包含原始或补充指令、客户/游客/联系人信息、附件、Agent/插件/ERP 技术错误、堆栈、URL、token、内部 UUID、生命周期细节或未验证的执行结果。失败只使用统一业务文案;写入不确定不得包装成完成。 + +## 受理与失败语义 + +只有 HTTP 状态为 `200`、JSON `code === 0` 且 `data === true` 同时成立,才将投递标记为“已受理”。这不代表微信已经送达,因为接口没有下游任务 ID、送达查询或回调。 + +- HTTP `400`、`401`、明确配置不存在的 `404` 或其他业务拒绝记为“明确拒绝”,不自动重试。 +- 超时、断网、`5xx`、响应过大、成功状态上的非 JSON/未知响应以及进程在发送后失联记为“结果不确定”,不自动重试。 +- `sending` 租约过期只能转为“结果不确定”,不能回到待发队列。 +- 配置停用或变化时,旧 revision 的未发送记录取消;正在发送的记录转为不确定,不能改投或重发。 + +这套不重试边界是必要的:接口没有幂等键或去重能力,无法区分“请求未到达”与“已受理但响应丢失”。任务状态和员工回执绝不因组长摘要投递结果而变化。 + +Webhook 的配置错误、初始化失败、数据库投影失败、网络异常和服务端拒绝都必须被限制在摘要 worker 内。worker 只读取任务与 `task.updated` outbox,唯一写入边界是专用的 `leader_task_summary_webhook_*` 表与对应审计事件;不得更新任务、任务事件、任务尝试、员工 `agentbus_deliveries`、渠道、解析结果或 ERP 生命周期。 + +## 持久化与迁移 + +摘要正文进入独立的字段加密 outbox,并以组织、配置 revision、任务和里程碑唯一约束去重。URL 与 token 只来自受保护运行环境,不进入数据库;数据库只保存不可逆配置指纹和生效边界。 + +迁移 `023_leader_summary_webhook_delivery` 会停用旧 AgentBus 组长订阅、取消其待发/失败记录但保留历史,然后创建组织级 Webhook 状态与投递表。员工 `agentbus_deliveries` 和 AgentBus 渠道继续按原契约工作。 diff --git a/control-plane/README.md b/control-plane/README.md index 3715c32..e6d5add 100644 --- a/control-plane/README.md +++ b/control-plane/README.md @@ -29,7 +29,7 @@ - `/operations-dashboard` 是组长专用的只读业务操作看板。它以不可变的任务归属账号作为员工口径,纳入同一固定组织中所有已归属的人工与 AgentBus 任务,并支持从结果状态、员工、上海业务日期和业务类型逐层穿透;历史上无法安全归属员工的 AgentBus 任务不在人员看板中被猜测归属。列表和详情只回答“谁负责什么任务、收到什么指令、完成了什么结果”:详情返回员工、业务类型、完整指令轮次、输入附件名称/行数和可读业务结果,不返回任务生命周期、解析/执行 JSON、技术阶段、错误码或产物地址。关键词查询先受日期、人员、业务和状态约束,单次解密匹配候选最多 2,000 条,超过时要求继续缩小范围。该路径不授予他人任务修改、ERP 执行、SSE、产物下载、账号维护或全局安全审计权限;管理员同样无权访问。 - 使用数据库渠道时设置 `AGENTBUS_ENABLED=true`;此模式不要求 `AGENTBUS_WS_TOKEN` 或 `AGENTBUS_BOT_ADDRESS`,但启用的渠道仍需要全局 `AGENTBUS_WS_URL`,并可在渠道上覆盖 bot address。保留旧环境变量配置时,服务会按需创建“默认 AgentBus 渠道”兼容旧单渠道部署;兼容渠道初始为未绑定且不启动,管理员必须在 `/channels` 绑定员工账号后再启用。`AGENTBUS_ENABLED=auto` 仅由完整的旧环境连接字段自动启用。 - `user_channels`、`tasks.channel_id` 和 `agentbus_deliveries` 共同保存入站归属、accepted 受理回执和最终 result 回执。回执以 `(channel_id, inbound_frame_id, delivery_kind)` 幂等,发送失败会重试,进程重启或 WebSocket 重连后仍会继续投递;因此不会因为超过原等待时长而丢掉最终回复。 -- ERP 插件领取按任务 `assigned_user_id` 使用账户级数据库锁和 FIFO confirmed 队列:同一平台/ERP 账户在任意时刻最多一个 ERP execution,该账户的其他任务留在服务端等待;不同账户的活跃或待执行任务互不占用队列位置、可独立领取执行。管理员不能成为 `assigned_user_id`、渠道 owner、浏览器 worker 或组长任务摘要接收人,也不能打开任务 SSE;实时执行事件、任务摘要、插件领取、执行回执以及强制删除后的浏览器清理命令只发送或接受符合角色约束的员工账号。已开始写入但结果不确定的任务只阻塞同一账户的后续领取,直到归属账号人工回查收敛或任务被明确强制删除。 +- ERP 插件领取按任务 `assigned_user_id` 使用账户级数据库锁和 FIFO confirmed 队列:同一平台/ERP 账户在任意时刻最多一个 ERP execution,该账户的其他任务留在服务端等待;不同账户的活跃或待执行任务互不占用队列位置、可独立领取执行。管理员不能成为 `assigned_user_id`、渠道 owner 或浏览器 worker,也不能打开任务 SSE;实时执行事件、插件领取、执行回执以及强制删除后的浏览器清理命令只发送或接受符合角色约束的员工账号。组长摘要是独立的只读外部 Webhook 投影,不授予任何任务或 ERP 权限。已开始写入但结果不确定的任务只阻塞同一账户的后续领取,直到归属账号人工回查收敛或任务被明确强制删除。 ## 任务级会话续接 @@ -89,7 +89,7 @@ Auto 一旦发生 AI fallback,任务会永久绑定原 AI 会话。每次解 旧的任务详情重解析和解析差异判定路由仍保留兼容响应,但已纳入任务数据面统一门禁;管理员调用会返回 `task_access_forbidden`,管理页面不再提供这两个任务级入口。 -迁移 `013_business_parser_modes` 增加内部固定范围的路由设置、任务快照和加密的 `parse_decisions`;迁移 `014_task_input_attachments` 增加名单输入附件元数据、加密 canonical TSV 与 `awaiting_attachment` 索引;迁移 `015_account_roles_and_task_audit` 增加账号角色、密码更新时间、输入/附件操作者、任务归档和账号级幂等(历史 `must_change_password` 列仅保留兼容,当前流程不启用首次强制改密);迁移 `016_team_lead_operations_dashboard` 增加组长角色与人工指令看板索引;迁移 `017_user_business_route_authorizations` 增加逐账号业务白名单、授权人和乐观并发 revision;迁移 `018_agentbus_account_workers` 增加 ERP 账号、渠道归属、任务执行归属、唯一在线 worker 与 ERP 身份核验字段;迁移 `020_leader_task_summary_notifications` 增加内部失败关闭的组长路由快照和独立加密通知 outbox;迁移 `021_admin_task_data_plane_isolation` 清除管理员的遗留任务执行绑定和待发组长摘要,并以数据库触发器阻止管理员成为任务创建人/归属人、渠道 owner、浏览器 worker、业务授权对象或任务摘要接收人。原文、完整程序/AI 候选、名单 canonical 中间文本、人工说明和待发组长摘要使用字段加密保存;统计、全局审计和运行日志不复制明文业务输入。 +迁移 `013_business_parser_modes` 增加内部固定范围的路由设置、任务快照和加密的 `parse_decisions`;迁移 `014_task_input_attachments` 增加名单输入附件元数据、加密 canonical TSV 与 `awaiting_attachment` 索引;迁移 `015_account_roles_and_task_audit` 增加账号角色、密码更新时间、输入/附件操作者、任务归档和账号级幂等(历史 `must_change_password` 列仅保留兼容,当前流程不启用首次强制改密);迁移 `016_team_lead_operations_dashboard` 增加组长角色与人工指令看板索引;迁移 `017_user_business_route_authorizations` 增加逐账号业务白名单、授权人和乐观并发 revision;迁移 `018_agentbus_account_workers` 增加 ERP 账号、渠道归属、任务执行归属、唯一在线 worker 与 ERP 身份核验字段;迁移 `020_leader_task_summary_notifications` 保存旧 AgentBus 组长摘要订阅历史;迁移 `021_admin_task_data_plane_isolation` 清除管理员的遗留任务执行绑定;迁移 `023_leader_summary_webhook_delivery` 停用旧 AgentBus 组长摘要并创建组织级外部 Webhook 状态与加密投递 outbox。原文、完整程序/AI 候选、名单 canonical 中间文本、人工说明和待发组长摘要使用字段加密保存;统计、全局审计和运行日志不复制明文业务输入。 ## AgentBus Bot 接入 @@ -99,7 +99,7 @@ Auto 一旦发生 AI fallback,任务会永久绑定原 AI 会话。每次解 每个启用渠道会连接 `AGENTBUS_WS_URL?ready=1`,使用该渠道自己的 `Authorization: Bearer `,等待 `session.ready` 后接收普通 `event` 消息。普通任务发送一次持久化受理通知(`task.progress`,`status=accepted`)并在完成时返回一次 `task.result`。名单任务的首次文字指令改为返回明确的等待附件提示,附件入站改为返回“名单附件已收到,正在校验并处理”;最终结果只归属触发解析的最新附件消息,因此不会因文字与附件两条入站帧重复发送成功回执。解析完成、等待确认或进入 ERP 等内部进度不外发。`GET /health/ready` 和 `GET /api/status` 的 `agentbus.channels` 字段可用于确认每个 listener 与 session 是否建立。 -`/channels` 只读展示组长任务摘要抄送状态,不再维护单独的订阅开关、收件地址或微信会话 ID。有效 `team_lead` 账号拥有已启用且可连接的 AgentBus 渠道时,服务会自动启用组织范围摘要;优先采用该账号经此渠道最近一次有效入站帧的 `from` 与 `conversation_id`,尚无入站记录时使用渠道 `external_user_ref` 并按普通回复规则生成会话键。没有任何可用路由时保持等待,首条有效组长消息到达后自动学习;渠道换绑账号时会清空旧外部用户标识,且历史入站路由只有任务归属仍是当前组长时才能复用,避免抄送到原渠道所有者。`GET /api/settings/leader-summary-subscriptions` 只返回不含明文目标的投递健康状态,没有人工修改或真实测试发送接口。角色、账号、渠道归属、渠道启停或路由变化会建立新的生效时间与 revision,取消旧路由未发送摘要,且永不回补历史。摘要固定读取同组织其他非管理员归属人的人工和 AgentBus 任务,只从自动生效后的 `task.updated` 稳定结果投影。主动帧为无 `reply_to` 的 `task.summary`,员工回执始终优先。正文只含员工、业务、稳定状态、公共任务编号、上海提交时间以及白名单团号/订单号,不含原始指令、人员资料、附件或技术错误;路由与正文加密,日志只留指纹。发送为至少一次语义且微信消息不可撤回。 +`/channels` 页面只读展示组长摘要的外部 Webhook 状态;该投递不再依赖组长 AgentBus 账号、渠道或微信会话。运行环境同时配置完整 `WEBHOOK_SEND_URL` 与 32 位 `WEBHOOK_EXTERNAL_TOKEN` 后自动启用,请求固定为 `POST`、原始 `x-token`、JSON 字符串 ID `9999` 和非空 `content`。配置缺失或无效只会关闭这个摘要 worker 并显示安全错误,不会阻断普通用户任务、员工 AgentBus、解析、确认或 ERP 执行。`GET /api/settings/leader-summary-webhook` 只返回不含 URL/token/正文的状态与队列计数,没有人工修改或真实测试发送接口。配置首次生效或变化会建立新的时间边界和 revision,只投影此后的已明确归属、非管理员人工/AgentBus 任务稳定结果,永不回补历史。正文只含员工、业务、稳定状态、公共任务编号、上海提交时间以及白名单团号/订单号,并在数据库中字段加密。只有 HTTP 200、`code=0`、`data=true` 同时成立才记为“已受理”,不代表微信已送达;明确拒绝、超时、断网、5xx 或未知响应均不自动重试,以避免无幂等能力的接口产生重复群通知。完整契约见 [`leader-summary-webhook-contract.md`](../agent设计规范/leader-summary-webhook-contract.md)。 对微信来源,listener 在调用 `TaskService.ingestMessage()` 前执行上述严格信封解包,因此手工正文与 AgentBus 正文进入同一个业务 route resolver、任务级 mode snapshot 和 parser orchestrator;`Conversation` 只属于传输路由,不会再污染业务字段签名。 @@ -143,13 +143,13 @@ npm run data:retention npm run dev ``` -`npm run dev` 和 `npm start` 会先执行数据库迁移,再启动控制平面;直接运行 `control-plane/src/server.ts` 或构建后的 `server.js` 时,服务也会在启动前检查必需迁移 `021_admin_task_data_plane_isolation`,缺失时拒绝监听端口。`db:migrate` 和管理员初始化需要可连接的 PostgreSQL。开发机没有数据库时,可以运行 `npm run test:control-plane` 完成无数据库静态/健康烟测。 +`npm run dev` 和 `npm start` 会先执行数据库迁移,再启动控制平面;直接运行 `control-plane/src/server.ts` 或构建后的 `server.js` 时,服务也会在启动前检查必需迁移 `023_leader_summary_webhook_delivery`,缺失时拒绝监听端口。`db:migrate` 和管理员初始化需要可连接的 PostgreSQL。开发机没有数据库时,可以运行 `npm run test:control-plane` 完成无数据库静态/健康烟测。 `/health/ready` 同时检查 PostgreSQL 可用性和必需 schema 版本;迁移未完成时返回 503,并标明 `required_migration`,避免任务在数据库结构未升级时进入解析队列。 ## 生产部署 -1. 复制 `.env.production.example` 为部署机受保护的 `.env.production`,填入 `DEPLOYMENT_REVISION`、PostgreSQL URL、`DATABASE_SCHEMA`、字段加密密钥、外部解析 Key 和 OSS 凭据,并确认 `AGENTBUS_LOG_PAYLOADS=false`。生产数据库可使用现有 PostgreSQL 实例中的新 Schema;迁移程序会创建 Schema,不会触碰其他 Schema 的测试表。 +1. 复制 `.env.production.example` 为部署机受保护的 `.env.production`,填入 `DEPLOYMENT_REVISION`、PostgreSQL URL、`DATABASE_SCHEMA`、字段加密密钥、外部解析 Key 和 OSS 凭据,并确认 `AGENTBUS_LOG_PAYLOADS=false`。启用组长摘要时,再填写服务提供方确认后的完整 `WEBHOOK_SEND_URL` 与单独提供的 32 位 `WEBHOOK_EXTERNAL_TOKEN`;不得通过真实发送猜测网关路径。生产数据库可使用现有 PostgreSQL 实例中的新 Schema;迁移程序会创建 Schema,不会触碰其他 Schema 的测试表。 2. 在正式数据库上线前执行并验证备份:`infra/backup-postgres.sh`。 3. 使用 `docker compose --env-file .env.production up -d --build` 启动;Compose 会先执行数据库迁移,再启动控制平面。Compose 中的本地 PostgreSQL 仅用于 `--profile local`,生产 `DATABASE_URL` 指向受保护的远程数据库。 4. 首次启动后在容器内通过 `docker compose exec -e ADMIN_USERNAME=admin -e ADMIN_PASSWORD='replace-with-password' control-plane node .build/control-plane/src/admin-cli.js bootstrap` 创建管理员;不要把密码写入镜像或 Git。 diff --git a/control-plane/migrations/023_leader_summary_webhook_delivery.sql b/control-plane/migrations/023_leader_summary_webhook_delivery.sql new file mode 100644 index 0000000..5731ffe --- /dev/null +++ b/control-plane/migrations/023_leader_summary_webhook_delivery.sql @@ -0,0 +1,74 @@ +-- Replace AgentBus-routed team-lead summaries with one organization-level +-- external webhook delivery stream. Legacy rows remain for audit/history but +-- cannot be claimed after this migration. + +UPDATE leader_task_summary_subscriptions + SET enabled = false, + target_verified_at = NULL, + revision = revision + 1, + updated_at = now() + WHERE enabled = true OR target_verified_at IS NOT NULL; + +UPDATE leader_task_summary_deliveries + SET delivery_status = 'cancelled', + last_error = '组长摘要已切换为外部 Webhook,旧 AgentBus 待发记录已取消。', + updated_at = now() + WHERE delivery_status IN ('pending', 'sending', 'failed'); + +CREATE TABLE IF NOT EXISTS leader_task_summary_webhook_state ( + organization_id uuid PRIMARY KEY REFERENCES organizations(id) ON DELETE CASCADE, + enabled boolean NOT NULL DEFAULT false, + starts_at timestamptz NOT NULL DEFAULT now(), + revision integer NOT NULL DEFAULT 0, + configuration_fingerprint text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT leader_task_summary_webhook_state_revision_check + CHECK (revision >= 0), + CONSTRAINT leader_task_summary_webhook_state_fingerprint_check + CHECK (configuration_fingerprint ~ '^[a-f0-9]{64}$') +); + +CREATE TABLE IF NOT EXISTS leader_task_summary_webhook_deliveries ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + organization_id uuid NOT NULL REFERENCES organizations(id) ON DELETE CASCADE, + webhook_revision integer NOT NULL, + task_id uuid NOT NULL, + source_outbox_event_id bigint NOT NULL, + milestone text NOT NULL, + payload_ciphertext text NOT NULL, + payload_fingerprint text NOT NULL, + delivery_status text NOT NULL DEFAULT 'pending', + attempt_count integer NOT NULL DEFAULT 0, + response_status integer, + last_error text, + accepted_at timestamptz, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT leader_task_summary_webhook_deliveries_milestone_check + CHECK (milestone IN ('final', 'needs_review', 'resolved')), + CONSTRAINT leader_task_summary_webhook_deliveries_status_check + CHECK (delivery_status IN ('pending', 'sending', 'accepted', 'failed', 'uncertain', 'cancelled')), + CONSTRAINT leader_task_summary_webhook_deliveries_attempt_check + CHECK (attempt_count BETWEEN 0 AND 1), + CONSTRAINT leader_task_summary_webhook_deliveries_revision_check + CHECK (webhook_revision >= 0), + CONSTRAINT leader_task_summary_webhook_deliveries_response_status_check + CHECK (response_status IS NULL OR response_status BETWEEN 100 AND 599), + CONSTRAINT leader_task_summary_webhook_deliveries_payload_fingerprint_check + CHECK (payload_fingerprint ~ '^[a-f0-9]{64}$'), + CONSTRAINT leader_task_summary_webhook_deliveries_task_scope_fkey + FOREIGN KEY (organization_id, task_id) + REFERENCES tasks (organization_id, id) + ON DELETE CASCADE, + UNIQUE (organization_id, webhook_revision, task_id, milestone) +); + +CREATE INDEX IF NOT EXISTS leader_task_summary_webhook_deliveries_pending_idx + ON leader_task_summary_webhook_deliveries + (organization_id, webhook_revision, delivery_status, created_at, id) + WHERE delivery_status = 'pending'; + +CREATE INDEX IF NOT EXISTS leader_task_summary_webhook_deliveries_task_idx + ON leader_task_summary_webhook_deliveries + (organization_id, task_id, created_at DESC); diff --git a/control-plane/src/agentbus-channels.ts b/control-plane/src/agentbus-channels.ts index 0f154e8..8e99d68 100644 --- a/control-plane/src/agentbus-channels.ts +++ b/control-plane/src/agentbus-channels.ts @@ -655,7 +655,6 @@ export interface AgentBusManagerOptions { scheduleParseQueue: () => Promise; logger?: AgentBusLogger; socketFactory?: AgentBusSocketFactory; - leaderNotifications?: AgentBusListenerOptions['leaderNotifications']; } export class AgentBusManager { @@ -723,7 +722,6 @@ export class AgentBusManager { scheduleParseQueue: this.options.scheduleParseQueue, socketFactory: this.options.socketFactory, logger: this.options.logger, - leaderNotifications: this.options.leaderNotifications, channel: { id: channel.id, displayName: channel.display_name, diff --git a/control-plane/src/agentbus.ts b/control-plane/src/agentbus.ts index 484584a..9614ba3 100644 --- a/control-plane/src/agentbus.ts +++ b/control-plane/src/agentbus.ts @@ -25,7 +25,6 @@ import { AGENTBUS_ROSTER_ATTACHMENT_RECEIVED_TEXT, publicAgentBusDeliveryPayload } from './agentbus-delivery.js'; -import type { LeaderTaskSummaryDelivery } from './leader-notification-service.js'; const OPEN_READY_STATE = 1; const MAX_COMPLETED_TASK_IDS = 2_048; @@ -134,24 +133,6 @@ export interface AgentBusTaskGateway { releaseAgentBusDeliveries?(channelId: string, leaseOwner: string): Promise; } -export interface LeaderNotificationGateway { - observeLeaderRoute?(input: { - organizationId: string; - leaderUserId: string; - channelId: string; - recipientAddress: string; - conversationId?: string; - }): Promise; - claimDeliveries( - channelId: string, - leaseOwner: string, - limit?: number - ): Promise; - markDeliveryDelivered(deliveryId: string): Promise; - markDeliveryFailed(deliveryId: string, errorMessage: string): Promise; - releaseDeliveries(channelId: string, leaseOwner: string): Promise; -} - export interface AgentBusChannelConnection { id: string; displayName: string; @@ -170,7 +151,6 @@ export interface AgentBusListenerOptions { socketFactory?: AgentBusSocketFactory; logger?: AgentBusLogger; channel?: AgentBusChannelConnection; - leaderNotifications?: LeaderNotificationGateway; onStatusChange?: ( status: 'disabled' | 'connecting' | 'connected' | 'error', errorMessage: string | null, @@ -538,22 +518,6 @@ function createDurableDeliveryFrame( return frame; } -export function createLeaderTaskSummaryFrame( - delivery: LeaderTaskSummaryDelivery, - session: AgentBusSession -): AgentBusFrame { - return { - id: `leader-summary-${delivery.id}`, - type: 'event', - from: session.address, - to: delivery.recipient_address, - session_id: session.id, - epoch: session.epoch, - conversation_id: delivery.conversation_id, - payload: { ...delivery.payload } - }; -} - export function taskResultStatus(task: PublicTask): 'completed' | 'failed' { return FAILED_TASK_STATUSES.has(text(task.status).toLowerCase()) ? 'failed' : 'completed'; } @@ -653,7 +617,6 @@ export class AgentBusListener { private readonly socketFactory: AgentBusSocketFactory; private readonly logger: AgentBusLogger; private readonly channel: AgentBusChannelConnection | null; - private readonly leaderNotifications: LeaderNotificationGateway | null; private readonly onStatusChange: AgentBusListenerOptions['onStatusChange']; private readonly leaseOwner: string; private socket: AgentBusSocket | null = null; @@ -665,7 +628,6 @@ export class AgentBusListener { private readonly pendingFinalReplies = new Map(); private deliveryTimer: NodeJS.Timeout | null = null; private durableFlushInFlight: Promise | null = null; - private leaderFlushInFlight: Promise | null = null; constructor(options: AgentBusListenerOptions) { this.config = options.config; @@ -675,7 +637,6 @@ export class AgentBusListener { this.socketFactory = options.socketFactory || defaultSocketFactory; this.logger = options.logger || noopLogger; this.channel = options.channel || null; - this.leaderNotifications = options.leaderNotifications || null; this.onStatusChange = options.onStatusChange; this.leaseOwner = `agentbus:${this.channel?.id || 'legacy'}`; } @@ -794,9 +755,6 @@ export class AgentBusListener { if (this.channel && this.tasks.releaseAgentBusDeliveries) { void this.tasks.releaseAgentBusDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined); } - if (this.channel && this.leaderNotifications) { - void this.leaderNotifications.releaseDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined); - } if (persistStatus) this.setRuntimeStatus('disabled', null, sessionEpoch); } @@ -877,9 +835,6 @@ export class AgentBusListener { if (this.channel && this.tasks.releaseAgentBusDeliveries) { void this.tasks.releaseAgentBusDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined); } - if (this.channel && this.leaderNotifications) { - void this.leaderNotifications.releaseDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined); - } for (const pending of this.pendingFinalReplies.values()) pending.sending = false; this.logger.warn({ agentbus_event: 'socket_close', @@ -1066,25 +1021,6 @@ export class AgentBusListener { ...(this.channel ? { channelId: this.channel.id } : {}) }; const result = await this.tasks.ingestMessage(context, input); - if ( - this.channel?.ownerRole === 'team_lead' - && this.leaderNotifications?.observeLeaderRoute - ) { - void this.leaderNotifications.observeLeaderRoute({ - organizationId: this.organizationId, - leaderUserId: this.channel.ownerUserId, - channelId: this.channel.id, - recipientAddress: text(frame.from), - conversationId - }).catch((error) => { - this.logger.warn({ - agentbus_event: 'leader_summary_route_observation_failed', - channel_id: this.channel?.id, - owner_user_id: this.channel?.ownerUserId, - ...diagnosticError(error, 'leader_summary_route_observation_failed') - }, 'AgentBus leader summary route observation failed'); - }); - } this.logger.info({ agentbus_event: 'task_ingested', inbound_frame_id: taskId, @@ -1249,90 +1185,7 @@ export class AgentBusListener { } private async flushOutboundDeliveries(): Promise { - // Employee replies always get the first claim/send opportunity on a tick. - // Leader summaries use a separate, smaller queue and cannot delay task replies. await this.flushDurableDeliveries(); - await this.flushLeaderDeliveries(); - } - - private async flushLeaderDeliveries(): Promise { - if (!this.channel || !this.leaderNotifications) return; - if (this.leaderFlushInFlight) return this.leaderFlushInFlight; - if (!this.socket || this.socket.readyState !== OPEN_READY_STATE || !this.session) return; - this.leaderFlushInFlight = (async () => { - const deliveries = await this.leaderNotifications!.claimDeliveries( - this.channel!.id, - this.leaseOwner, - 5 - ); - for (const delivery of deliveries) void this.sendLeaderDelivery(delivery); - })() - .catch((error) => { - this.logger.warn({ - agentbus_event: 'leader_summary_delivery_flush_failed', - channel_id: this.channel?.id || null, - ...diagnosticError(error, 'leader_summary_delivery_flush_failed') - }, 'AgentBus leader task summary flush failed'); - }) - .finally(() => { - this.leaderFlushInFlight = null; - }); - return this.leaderFlushInFlight; - } - - private async sendLeaderDelivery(delivery: LeaderTaskSummaryDelivery): Promise { - if (!this.socket || this.socket.readyState !== OPEN_READY_STATE || !this.session || !this.leaderNotifications) { - return; - } - const socket = this.socket; - const session = this.session; - const notifications = this.leaderNotifications; - const frame = createLeaderTaskSummaryFrame(delivery, session); - this.logger.info({ - agentbus_event: 'leader_summary_delivery_sending', - channel_id: delivery.channel_id, - delivery_id: delivery.id, - task_id: delivery.task_id, - recipient_fingerprint: delivery.recipient_fingerprint, - conversation_fingerprint: delivery.conversation_fingerprint, - attempt_count: delivery.attempt_count - }, 'AgentBus leader task summary sending'); - const fail = (error: unknown) => { - const errorMetadata = diagnosticError(error, 'leader_summary_delivery_failed'); - void notifications.markDeliveryFailed( - delivery.id, - `${String(errorMetadata.error_code)}:${String(errorMetadata.error_fingerprint)}` - ).catch(() => undefined); - this.logger.warn({ - agentbus_event: 'leader_summary_delivery_failed', - channel_id: delivery.channel_id, - delivery_id: delivery.id, - task_id: delivery.task_id, - recipient_fingerprint: delivery.recipient_fingerprint, - conversation_fingerprint: delivery.conversation_fingerprint, - ...errorMetadata - }, 'AgentBus leader task summary failed'); - }; - try { - socket.send(JSON.stringify(frame), (error) => { - if (error) { - fail(error); - return; - } - void notifications.markDeliveryDelivered(delivery.id) - .then(() => this.logger.info({ - agentbus_event: 'leader_summary_delivery_sent', - channel_id: delivery.channel_id, - delivery_id: delivery.id, - task_id: delivery.task_id, - recipient_fingerprint: delivery.recipient_fingerprint, - conversation_fingerprint: delivery.conversation_fingerprint - }, 'AgentBus leader task summary sent')) - .catch(fail); - }); - } catch (error) { - fail(error); - } } private async sendDurableDelivery( diff --git a/control-plane/src/config.ts b/control-plane/src/config.ts index b99deb3..e778afe 100644 --- a/control-plane/src/config.ts +++ b/control-plane/src/config.ts @@ -9,6 +9,14 @@ const optionalString = z.preprocess( z.string().optional() ); +const optionalExactString = z.preprocess( + (value) => { + if (value == null || value === '') return undefined; + return String(value); + }, + z.string().optional() +); + const optionalUrl = z.preprocess( (value) => { const normalized = String(value ?? '').trim(); @@ -86,6 +94,10 @@ const envSchema = z.object({ AGENTBUS_WS_RECONNECT_DELAY: z.string().default('5s').transform(parseDurationMs), AGENTBUS_TASK_TIMEOUT_MS: z.coerce.number().int().positive().default(300_000), AGENTBUS_LOG_PAYLOADS: z.enum(['true', 'false']).default('false').transform((value) => value === 'true'), + // Webhook configuration is validated separately so a typo in this optional + // notification feature cannot prevent normal task or AgentBus startup. + WEBHOOK_SEND_URL: optionalString, + WEBHOOK_EXTERNAL_TOKEN: optionalExactString, ARTIFACT_STORAGE_BACKEND: z.enum(['database', 'oss']).default('database'), ARTIFACT_MAX_BYTES: z.coerce.number().int().positive().max(50_000_000).default(10_485_760), DOCUMENT_CONVERTER_PATH: z.string().trim().min(1).default('soffice'), @@ -105,6 +117,8 @@ const envSchema = z.object({ export type AppConfig = z.infer & { fieldEncryptionKey: Buffer; agentBusEnabled: boolean; + leaderSummaryWebhookEnabled: boolean; + leaderSummaryWebhookConfigurationError: string | null; ossRegion: string; }; @@ -143,6 +157,30 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): AppConfig { throw new Error(`AgentBus is enabled but missing configuration: ${missing.join(', ')}`); } } + const webhookCredentialsPresent = Boolean(parsed.WEBHOOK_SEND_URL || parsed.WEBHOOK_EXTERNAL_TOKEN); + const webhookConfigurationComplete = Boolean(parsed.WEBHOOK_SEND_URL && parsed.WEBHOOK_EXTERNAL_TOKEN); + let leaderSummaryWebhookConfigurationError: string | null = null; + if (webhookCredentialsPresent && !webhookConfigurationComplete) { + leaderSummaryWebhookConfigurationError = 'webhook_configuration_incomplete'; + } else if (parsed.WEBHOOK_EXTERNAL_TOKEN && parsed.WEBHOOK_EXTERNAL_TOKEN.length !== 32) { + leaderSummaryWebhookConfigurationError = 'webhook_token_length_invalid'; + } else if (parsed.WEBHOOK_EXTERNAL_TOKEN && /\s/u.test(parsed.WEBHOOK_EXTERNAL_TOKEN)) { + leaderSummaryWebhookConfigurationError = 'webhook_token_whitespace_invalid'; + } else if (parsed.WEBHOOK_SEND_URL) { + let webhookUrl: URL | null = null; + try { + webhookUrl = new URL(parsed.WEBHOOK_SEND_URL); + } catch { + leaderSummaryWebhookConfigurationError = 'webhook_url_invalid'; + } + if (webhookUrl && !webhookUrl.pathname.endsWith('/webhookInfo/sendMessageByOut')) { + leaderSummaryWebhookConfigurationError = 'webhook_route_unconfirmed'; + } else if (webhookUrl && parsed.NODE_ENV === 'production' && webhookUrl.protocol !== 'https:') { + leaderSummaryWebhookConfigurationError = 'webhook_https_required'; + } + } + const leaderSummaryWebhookEnabled = webhookConfigurationComplete + && !leaderSummaryWebhookConfigurationError; if (parsed.ARTIFACT_STORAGE_BACKEND === 'oss') { const missing = [ ['OSS_ACCESS_KEY_ID', parsed.OSS_ACCESS_KEY_ID], @@ -164,6 +202,8 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): AppConfig { ...parsed, fieldEncryptionKey: decodeEncryptionKey(parsed.FIELD_ENCRYPTION_KEY, parsed.NODE_ENV), agentBusEnabled, + leaderSummaryWebhookEnabled, + leaderSummaryWebhookConfigurationError, ossRegion }; } diff --git a/control-plane/src/db.ts b/control-plane/src/db.ts index 8d85eb4..7a4cecc 100644 --- a/control-plane/src/db.ts +++ b/control-plane/src/db.ts @@ -5,7 +5,7 @@ import { writeEmergencyDiagnostic } from './diagnostics.js'; const { Pool } = pg; let pool: pg.Pool | null = null; -export const REQUIRED_SCHEMA_VERSION = '021_admin_task_data_plane_isolation'; +export const REQUIRED_SCHEMA_VERSION = '023_leader_summary_webhook_delivery'; export interface DatabaseReadiness { ready: boolean; diff --git a/control-plane/src/external-webhook-client.ts b/control-plane/src/external-webhook-client.ts new file mode 100644 index 0000000..c027882 --- /dev/null +++ b/control-plane/src/external-webhook-client.ts @@ -0,0 +1,178 @@ +export const EXTERNAL_WEBHOOK_CONFIG_ID = '9999'; +export const EXTERNAL_WEBHOOK_TIMEOUT_MS = 15_000; + +const MAX_RESPONSE_BYTES = 65_536; + +export type ExternalWebhookSendResult = + | { + outcome: 'accepted'; + httpStatus: 200; + errorCode: null; + } + | { + outcome: 'rejected' | 'uncertain'; + httpStatus: number | null; + errorCode: string; + }; + +export type ExternalWebhookFetch = ( + input: string | URL | Request, + init?: RequestInit +) => Promise; + +export interface ExternalWebhookClientOptions { + url: string; + token: string; + timeoutMs?: number; + fetchImpl?: ExternalWebhookFetch; +} + +function jsonObject(value: unknown): Record { + return value && typeof value === 'object' && !Array.isArray(value) + ? value as Record + : {}; +} + +async function readBoundedResponse(response: Response): Promise<{ + text: string; + tooLarge: boolean; +}> { + const declaredLength = Number(response.headers.get('content-length')); + if (Number.isFinite(declaredLength) && declaredLength > MAX_RESPONSE_BYTES) { + await response.body?.cancel().catch(() => undefined); + return { text: '', tooLarge: true }; + } + if (!response.body) return { text: '', tooLarge: false }; + + const reader = response.body.getReader(); + const chunks: Uint8Array[] = []; + let total = 0; + while (true) { + const { done, value } = await reader.read(); + if (done) break; + if (!value) continue; + total += value.byteLength; + if (total > MAX_RESPONSE_BYTES) { + await reader.cancel().catch(() => undefined); + return { text: '', tooLarge: true }; + } + chunks.push(value); + } + + const bytes = new Uint8Array(total); + let offset = 0; + for (const chunk of chunks) { + bytes.set(chunk, offset); + offset += chunk.byteLength; + } + return { text: new TextDecoder().decode(bytes), tooLarge: false }; +} + +function rejectedErrorCode(status: number, payload: Record): string { + const message = String(payload.msg ?? '').trim(); + if (status === 401 && message === 'Invalid token') return 'webhook_invalid_token'; + if (status === 400 && message === 'id and content must not be blank') return 'webhook_invalid_request'; + if (status === 404 && message === 'Webhook configuration not found') { + return 'webhook_configuration_not_found'; + } + if (status === 404) return 'webhook_route_not_found'; + if (status >= 400 && status < 500) return `webhook_http_${status}`; + return 'webhook_business_rejected'; +} + +export class ExternalWebhookClient { + private readonly timeoutMs: number; + private readonly fetchImpl: ExternalWebhookFetch; + + constructor(private readonly options: ExternalWebhookClientOptions) { + this.timeoutMs = options.timeoutMs || EXTERNAL_WEBHOOK_TIMEOUT_MS; + this.fetchImpl = options.fetchImpl || globalThis.fetch; + } + + async send(content: string): Promise { + const normalizedContent = String(content ?? '').trim(); + if (!normalizedContent) { + throw new Error('Webhook content must not be blank.'); + } + if (this.options.token.length !== 32 || /\s/u.test(this.options.token)) { + throw new Error('Webhook token must contain exactly 32 non-whitespace characters.'); + } + + let response: Response; + try { + response = await this.fetchImpl(this.options.url, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'x-token': this.options.token + }, + body: JSON.stringify({ + id: EXTERNAL_WEBHOOK_CONFIG_ID, + content: normalizedContent + }), + signal: AbortSignal.timeout(this.timeoutMs) + }); + } catch (error) { + return { + outcome: 'uncertain', + httpStatus: null, + errorCode: error instanceof DOMException && error.name === 'TimeoutError' + ? 'webhook_request_timeout' + : 'webhook_request_uncertain' + }; + } + + let responseText = ''; + try { + const bounded = await readBoundedResponse(response); + if (bounded.tooLarge) { + return { + outcome: 'uncertain', + httpStatus: response.status, + errorCode: 'webhook_response_too_large' + }; + } + responseText = bounded.text; + } catch { + return { + outcome: 'uncertain', + httpStatus: response.status, + errorCode: 'webhook_response_read_failed' + }; + } + + let payload: Record = {}; + try { + payload = jsonObject(JSON.parse(responseText)); + } catch { + if (response.status >= 400 && response.status < 500) { + return { + outcome: 'rejected', + httpStatus: response.status, + errorCode: rejectedErrorCode(response.status, payload) + }; + } + return { + outcome: 'uncertain', + httpStatus: response.status, + errorCode: 'webhook_response_invalid_json' + }; + } + + if (response.status === 200 && payload.code === 0 && payload.data === true) { + return { outcome: 'accepted', httpStatus: 200, errorCode: null }; + } + if (response.status >= 500) { + return { + outcome: 'uncertain', + httpStatus: response.status, + errorCode: `webhook_http_${response.status}_uncertain` + }; + } + return { + outcome: 'rejected', + httpStatus: response.status, + errorCode: rejectedErrorCode(response.status, payload) + }; + } +} diff --git a/control-plane/src/leader-notification-service.ts b/control-plane/src/leader-notification-service.ts index 39cc006..e14de30 100644 --- a/control-plane/src/leader-notification-service.ts +++ b/control-plane/src/leader-notification-service.ts @@ -1,11 +1,15 @@ import type { AppConfig } from './config.js'; import { decryptText, encryptText, sha256Text } from './crypto.js'; import { getPool, withTransaction } from './db.js'; -import { diagnosticError, diagnosticMetadataKeys } from './diagnostics.js'; +import { diagnosticError } from './diagnostics.js'; +import { + ExternalWebhookClient, + EXTERNAL_WEBHOOK_CONFIG_ID, + type ExternalWebhookSendResult +} from './external-webhook-client.js'; import { buildLeaderTaskSummary, - isLeaderTaskSummaryStatus, - type LeaderTaskSummaryStatus + isLeaderTaskSummaryStatus } from './leadership-task-summary.js'; export interface LeaderNotificationLogger { @@ -14,57 +18,35 @@ export interface LeaderNotificationLogger { error(metadata: Record, message?: string): void; } -export interface PublicLeaderTaskSummarySubscription { - id: string; - organization_id: string; - leader_user_id: string; - leader_username: string; - channel_id: string; - channel_name: string; - scope: 'organization'; - include_manual: boolean; - include_agentbus: boolean; - automatic: true; +export interface LeaderSummaryWebhookSender { + send(content: string): Promise; +} + +export interface PublicLeaderSummaryWebhookStatus { + delivery_mode: 'external_webhook'; + webhook_config_id: typeof EXTERNAL_WEBHOOK_CONFIG_ID; + configured: boolean; enabled: boolean; - eligible: boolean; - eligibility_message: string; - route_ready: boolean; - target_verified: boolean; - recipient_fingerprint: string; - conversation_fingerprint: string; - starts_at: string; - revision: number; + state: 'not_configured' | 'invalid_configuration' | 'initializing' | 'active' | 'configuration_changed'; + configuration_error: string | null; + status_message: string; + endpoint_fingerprint: string; + starts_at: string | null; + revision: number | null; pending_count: number; failed_count: number; - delivered_count: number; - last_delivered_at: string | null; - created_at: string; - updated_at: string; + uncertain_count: number; + accepted_count: number; + last_accepted_at: string | null; + updated_at: string | null; } -export interface ObservedLeaderAgentBusRoute { - organizationId: string; - leaderUserId: string; - channelId: string; - recipientAddress: string; - conversationId?: string; -} - -export interface LeaderTaskSummaryDelivery { +interface ClaimedWebhookDelivery { id: string; - channel_id: string; - task_id: string; - recipient_address: string; - recipient_fingerprint: string; - conversation_id: string; - conversation_fingerprint: string; - payload: { - event: 'task.summary'; - status: LeaderTaskSummaryStatus; - task_id: string; - text: string; - }; - attempt_count: number; + taskId: string; + content: string; + payloadFingerprint: string; + attemptCount: number; } const PROJECTABLE_STATUSES = [ @@ -90,17 +72,6 @@ const NEEDS_REVIEW_STATUSES = [ 'uncertain' ] as const; -const LEGACY_CHANNEL_REF = 'legacy-env'; - -interface AutomaticLeaderRoute { - organizationId: string; - leaderUserId: string; - channelId: string; - recipientAddress: string; - conversationId: string; - source: 'agentbus_inbound' | 'channel_external_user_ref'; -} - const noopLogger: LeaderNotificationLogger = { info: () => undefined, warn: () => undefined, @@ -127,90 +98,61 @@ function jsonObject(value: unknown): Record { : {}; } -function eligibility( - row: Record, - agentBusEnabled: boolean -): { eligible: boolean; message: string } { - if (!agentBusEnabled) { - return { eligible: false, message: 'AgentBus 全局开关未启用,通知不会发送。' }; - } - if (!booleanValue(row.leader_is_active) || text(row.leader_role) !== 'team_lead') { - return { eligible: false, message: '组长账号已停用或角色已变化,通知不会发送。' }; - } - if (text(row.channel_owner_user_id) !== text(row.leader_user_id)) { - return { eligible: false, message: 'AgentBus 渠道已不再归属该组长,通知不会发送。' }; - } - if (!booleanValue(row.channel_enabled)) { - return { eligible: false, message: 'AgentBus 渠道已停用,通知不会发送。' }; - } - if (!text(row.leader_erp_account)) { - return { eligible: false, message: '组长渠道尚未满足账号绑定条件,通知不会发送。' }; - } - if (!row.target_verified_at) { - return { eligible: false, message: 'AgentBus 尚未提供可用路由;收到组长消息后会自动同步。' }; - } - if (!booleanValue(row.enabled)) { - return { eligible: false, message: '自动订阅正在同步,暂时不会发送。' }; - } - return { eligible: true, message: '组长身份和 AgentBus 路由已就绪,摘要会自动推送。' }; -} - -function publicSubscription( - row: Record, - agentBusEnabled: boolean -): PublicLeaderTaskSummarySubscription { - const state = eligibility(row, agentBusEnabled); - return { - id: text(row.id), - organization_id: text(row.organization_id), - leader_user_id: text(row.leader_user_id), - leader_username: text(row.leader_username) || '未知组长', - channel_id: text(row.channel_id), - channel_name: text(row.channel_name) || '未命名渠道', - scope: 'organization', - include_manual: booleanValue(row.include_manual), - include_agentbus: booleanValue(row.include_agentbus), - automatic: true, - enabled: booleanValue(row.enabled), - eligible: state.eligible, - eligibility_message: state.message, - route_ready: Boolean(row.target_verified_at), - target_verified: Boolean(row.target_verified_at), - recipient_fingerprint: text(row.recipient_address_fingerprint).slice(0, 12), - conversation_fingerprint: text(row.conversation_id_fingerprint).slice(0, 12), - starts_at: iso(row.starts_at) || new Date(0).toISOString(), - revision: Number(row.revision || 0), - pending_count: Number(row.pending_count || 0), - failed_count: Number(row.failed_count || 0), - delivered_count: Number(row.delivered_count || 0), - last_delivered_at: iso(row.last_delivered_at), - created_at: iso(row.created_at) || new Date(0).toISOString(), - updated_at: iso(row.updated_at) || new Date(0).toISOString() - }; -} - export class LeaderNotificationService { private projectorTimer: NodeJS.Timeout | null = null; private projectorInFlight: Promise | null = null; private projectorOrganizationId = ''; + private readonly sender: LeaderSummaryWebhookSender | null; + private readonly leaseOwner = `webhook:${process.pid}`; constructor( private readonly config: AppConfig, - private readonly logger: LeaderNotificationLogger = noopLogger - ) {} + private readonly logger: LeaderNotificationLogger = noopLogger, + sender?: LeaderSummaryWebhookSender + ) { + this.sender = sender || ( + config.leaderSummaryWebhookEnabled + ? new ExternalWebhookClient({ + url: config.WEBHOOK_SEND_URL as string, + token: config.WEBHOOK_EXTERNAL_TOKEN as string + }) + : null + ); + } private log( level: 'info' | 'warn' | 'error', + event: string, metadata: Record, message: string ): void { try { - this.logger[level](metadata, message); + this.logger[level]({ + diagnostic_event: `leader_summary.${event}`, + notification_event: event, + ...metadata + }, message); } catch { // Notification persistence and task state never depend on logging. } } + private configurationFingerprint(): string | null { + if (!this.config.leaderSummaryWebhookEnabled) return null; + return sha256Text([ + 'leader-summary-webhook-v1', + this.config.WEBHOOK_SEND_URL, + this.config.WEBHOOK_EXTERNAL_TOKEN, + EXTERNAL_WEBHOOK_CONFIG_ID + ].join('\u0000')); + } + + private endpointFingerprint(): string { + return this.config.leaderSummaryWebhookEnabled && this.config.WEBHOOK_SEND_URL + ? sha256Text(this.config.WEBHOOK_SEND_URL).slice(0, 12) + : ''; + } + startProjector(organizationId: string): void { if (this.projectorTimer) return; this.projectorOrganizationId = organizationId; @@ -228,32 +170,12 @@ export class LeaderNotificationService { private async projectTick(): Promise { if (!this.projectorOrganizationId) return; if (this.projectorInFlight) return this.projectorInFlight; - this.projectorInFlight = this.reconcileAutomaticSubscriptions(this.projectorOrganizationId) - .then((changed) => { - if (changed > 0) { - this.log('info', { - agentbus_event: 'leader_summary_automatic_subscriptions_reconciled', - organization_id: this.projectorOrganizationId, - changed_count: changed - }, 'Automatic leader task summary subscriptions reconciled'); - } - return this.projectPending(this.projectorOrganizationId, 100); - }) - .then((count) => { - if (count > 0) { - this.log('info', { - agentbus_event: 'leader_summary_projected', - organization_id: this.projectorOrganizationId, - projected_count: count - }, 'Leader task summaries projected'); - } - }) + this.projectorInFlight = this.runTick(this.projectorOrganizationId) .catch((error) => { - this.log('warn', { - agentbus_event: 'leader_summary_projection_failed', + this.log('warn', 'webhook_tick_failed', { organization_id: this.projectorOrganizationId, - ...diagnosticError(error, 'leader_summary_projection_failed') - }, 'Leader task summary projection failed'); + ...diagnosticError(error, 'leader_summary_webhook_tick_failed') + }, 'Leader task summary webhook tick failed'); }) .finally(() => { this.projectorInFlight = null; @@ -261,396 +183,250 @@ export class LeaderNotificationService { return this.projectorInFlight; } - async listSubscriptions(organizationId: string): Promise { - const result = await getPool(this.config).query( - `SELECT subscription.*, - leader.username AS leader_username, - leader.role AS leader_role, - leader.is_active AS leader_is_active, - leader.erp_account AS leader_erp_account, - channel.display_name AS channel_name, - channel.owner_user_id AS channel_owner_user_id, - channel.enabled AS channel_enabled, - COALESCE(delivery.pending_count, 0)::int AS pending_count, - COALESCE(delivery.failed_count, 0)::int AS failed_count, - COALESCE(delivery.delivered_count, 0)::int AS delivered_count, - delivery.last_delivered_at - FROM leader_task_summary_subscriptions subscription - JOIN users leader - ON leader.id = subscription.leader_user_id - AND leader.organization_id = subscription.organization_id - JOIN user_channels channel - ON channel.id = subscription.channel_id - AND channel.organization_id = subscription.organization_id - LEFT JOIN LATERAL ( - SELECT count(*) FILTER ( - WHERE d.delivery_status IN ('pending', 'sending') - AND d.subscription_revision = subscription.revision - ) AS pending_count, - count(*) FILTER ( - WHERE d.delivery_status = 'failed' - AND d.subscription_revision = subscription.revision - ) AS failed_count, - count(*) FILTER (WHERE d.delivery_status = 'delivered') AS delivered_count, - max(d.delivered_at) AS last_delivered_at - FROM leader_task_summary_deliveries d - WHERE d.subscription_id = subscription.id - ) delivery ON true - WHERE subscription.organization_id = $1 - ORDER BY leader.username, subscription.id`, - [organizationId] - ); - return (result.rows as Record[]) - .map((row) => publicSubscription(row, this.config.agentBusEnabled)); - } - - private automaticRoute( - input: Omit & { - recipientAddress: unknown; - conversationId?: unknown; + private async runTick(organizationId: string): Promise { + const changed = await this.reconcileWebhookState(organizationId); + if (changed) { + this.log('info', 'webhook_configuration_reconciled', { + organization_id: organizationId, + configured: this.config.leaderSummaryWebhookEnabled, + endpoint_fingerprint: this.endpointFingerprint() + }, 'Leader task summary webhook configuration reconciled'); } - ): AutomaticLeaderRoute | null { - const recipientAddress = text(input.recipientAddress).slice(0, 500); - if (!recipientAddress || recipientAddress === LEGACY_CHANNEL_REF) return null; - const conversationId = ( - text(input.conversationId) - || `agentbus:${recipientAddress}` - ).slice(0, 500); - if (!conversationId) return null; - return { ...input, recipientAddress, conversationId }; + if (!this.config.leaderSummaryWebhookEnabled || !this.sender) return; + const projected = await this.projectPending(organizationId, 100); + if (projected > 0) { + this.log('info', 'webhook_summaries_projected', { + organization_id: organizationId, + projected_count: projected + }, 'Leader task summaries projected for external webhook delivery'); + } + await this.dispatchPending(organizationId); } - private async auditAutomaticChange( + private async auditConfigurationChange( client: import('pg').PoolClient, - input: { - organizationId: string; - subscriptionId: string; - eventType: string; - metadata: Record; - } + organizationId: string, + eventType: string, + metadata: Record ): Promise { await client.query( `INSERT INTO audit_events (organization_id, actor_user_id, event_type, entity_type, entity_id, request_id, metadata) - VALUES ($1, NULL, $2, 'leader_task_summary_subscription', $3, NULL, $4)`, - [input.organizationId, input.eventType, input.subscriptionId, input.metadata] + VALUES ($1::uuid, NULL, $2, 'leader_task_summary_webhook', $1::uuid::text, NULL, $3)`, + [organizationId, eventType, metadata] ); - this.log('info', { - agentbus_event: 'leader_summary_automatic_subscription_audit_staged', - entity_id: input.subscriptionId, - metadata_keys: diagnosticMetadataKeys(input.metadata) - }, 'Automatic leader summary subscription audit event staged'); } - private async syncAutomaticSubscription( - client: import('pg').PoolClient, - route: AutomaticLeaderRoute, - existing?: Record - ): Promise { - const recipientFingerprint = sha256Text(route.recipientAddress); - const conversationFingerprint = sha256Text(route.conversationId); - const routeChanged = !existing - || route.channelId !== text(existing.channel_id) - || recipientFingerprint !== text(existing.recipient_address_fingerprint) - || conversationFingerprint !== text(existing.conversation_id_fingerprint); - const policyChanged = !existing - || !booleanValue(existing.enabled) - || !booleanValue(existing.include_manual) - || !booleanValue(existing.include_agentbus) - || !existing.target_verified_at; - if (!routeChanged && !policyChanged) return false; - - const nextRevision = existing ? Number(existing.revision || 0) + 1 : 0; - let subscriptionId: string; - if (existing) { - const updated = await client.query( - `UPDATE leader_task_summary_subscriptions - SET channel_id = $1, - include_manual = true, - include_agentbus = true, - enabled = true, - starts_at = now(), - revision = $2, - recipient_address_ciphertext = $3, - recipient_address_fingerprint = $4, - conversation_id_ciphertext = $5, - conversation_id_fingerprint = $6, - target_verified_at = now(), - target_verified_by = NULL, - updated_at = now() - WHERE id = $7 - RETURNING id`, - [ - route.channelId, - nextRevision, - encryptText(this.config, route.recipientAddress), - recipientFingerprint, - encryptText(this.config, route.conversationId), - conversationFingerprint, - existing.id - ] - ); - subscriptionId = text(updated.rows[0]?.id); - } else { - const inserted = await client.query( - `INSERT INTO leader_task_summary_subscriptions - (organization_id, leader_user_id, channel_id, include_manual, - include_agentbus, enabled, starts_at, revision, - recipient_address_ciphertext, recipient_address_fingerprint, - conversation_id_ciphertext, conversation_id_fingerprint, - target_verified_at, target_verified_by, created_by) - VALUES ($1, $2, $3, true, true, true, now(), 0, - $4, $5, $6, $7, now(), NULL, NULL) - RETURNING id`, - [ - route.organizationId, - route.leaderUserId, - route.channelId, - encryptText(this.config, route.recipientAddress), - recipientFingerprint, - encryptText(this.config, route.conversationId), - conversationFingerprint - ] - ); - subscriptionId = text(inserted.rows[0]?.id); - } - - if (existing) { - await client.query( - `UPDATE leader_task_summary_deliveries - SET delivery_status = 'cancelled', - last_error = '组长身份或 AgentBus 路由已变化,旧路由待发送摘要已取消。', - updated_at = now() - WHERE subscription_id = $1 - AND subscription_revision <> $2 - AND delivery_status IN ('pending', 'sending', 'failed')`, - [subscriptionId, nextRevision] - ); - } - await this.auditAutomaticChange(client, { - organizationId: route.organizationId, - subscriptionId, - eventType: existing && routeChanged - ? 'leader_summary_subscription.automatic_route_updated' - : 'leader_summary_subscription.automatic_enabled', - metadata: { - leader_user_id: route.leaderUserId, - channel_id: route.channelId, - automatic: true, - include_manual: true, - include_agentbus: true, - enabled: true, - route_source: route.source, - recipient_fingerprint: recipientFingerprint.slice(0, 12), - conversation_fingerprint: conversationFingerprint.slice(0, 12), - revision: nextRevision, - historical_backfill: false - } - }); - return true; - } - - private async disableAutomaticSubscription( - client: import('pg').PoolClient, - existing: Record, - reason: 'channel_disabled' | 'route_unavailable' | 'role_or_channel_invalid' - ): Promise { - if (!booleanValue(existing.enabled) && !existing.target_verified_at) return false; - const nextRevision = Number(existing.revision || 0) + 1; - const updated = await client.query( - `UPDATE leader_task_summary_subscriptions - SET enabled = false, - starts_at = now(), - revision = $2, - target_verified_at = NULL, - target_verified_by = NULL, - updated_at = now() - WHERE id = $1 - RETURNING id`, - [existing.id, nextRevision] - ); - const subscriptionId = text(updated.rows[0]?.id); - await client.query( - `UPDATE leader_task_summary_deliveries - SET delivery_status = 'cancelled', - last_error = '组长身份或 AgentBus 路由已失效,待发送摘要已取消。', - updated_at = now() - WHERE subscription_id = $1 - AND delivery_status IN ('pending', 'sending', 'failed')`, - [subscriptionId] - ); - await this.auditAutomaticChange(client, { - organizationId: text(existing.organization_id), - subscriptionId, - eventType: 'leader_summary_subscription.automatic_disabled', - metadata: { - leader_user_id: text(existing.leader_user_id), - channel_id: text(existing.channel_id), - automatic: true, - enabled: false, - reason, - revision: nextRevision, - historical_backfill: false - } - }); - return true; - } - - async reconcileAutomaticSubscriptions(organizationId: string): Promise { + async reconcileWebhookState(organizationId: string): Promise { + const configurationFingerprint = this.configurationFingerprint(); return withTransaction(this.config, async (client) => { await client.query( `SELECT pg_advisory_xact_lock( - hashtextextended($1::text || ':leader-summary-automatic', 0) + hashtextextended($1::text || ':leader-summary-webhook-config', 0) )`, [organizationId] ); const existingResult = await client.query( `SELECT * - FROM leader_task_summary_subscriptions + FROM leader_task_summary_webhook_state WHERE organization_id = $1 FOR UPDATE`, [organizationId] ); - const existingByLeader = new Map>( - (existingResult.rows as Record[]) - .map((row) => [text(row.leader_user_id), row]) - ); - const candidates = await client.query( - `SELECT leader.id AS leader_user_id, - channel.id AS channel_id, - channel.enabled AS channel_enabled, - channel.external_user_ref, - latest_route.inbound_from, - latest_route.conversation_id - FROM users leader - JOIN user_channels channel - ON channel.organization_id = leader.organization_id - AND channel.owner_user_id = leader.id - LEFT JOIN LATERAL ( - SELECT delivery.inbound_from, delivery.conversation_id - FROM agentbus_deliveries delivery - JOIN tasks route_task - ON route_task.id = delivery.task_id - AND route_task.organization_id = delivery.organization_id - AND route_task.assigned_user_id = leader.id - WHERE delivery.organization_id = leader.organization_id - AND delivery.channel_id = channel.id - AND btrim(delivery.inbound_from) <> '' - ORDER BY delivery.created_at DESC, delivery.id DESC - LIMIT 1 - ) latest_route ON true - WHERE leader.organization_id = $1 - AND leader.role = 'team_lead' - AND leader.is_active = true - AND leader.erp_account IS NOT NULL - ORDER BY leader.id`, - [organizationId] - ); + const existing = existingResult.rows[0] as Record | undefined; - const currentLeaders = new Set(); - let changed = 0; - for (const row of candidates.rows as Record[]) { - const leaderUserId = text(row.leader_user_id); - const channelId = text(row.channel_id); - currentLeaders.add(leaderUserId); - const existing = existingByLeader.get(leaderUserId); - if (!booleanValue(row.channel_enabled)) { - if (existing && await this.disableAutomaticSubscription(client, existing, 'channel_disabled')) changed += 1; - continue; - } - const inboundRecipient = text(row.inbound_from); - const externalRecipient = text(row.external_user_ref); - const route = this.automaticRoute({ + if (!configurationFingerprint) { + if (!existing || !booleanValue(existing.enabled)) return false; + const previousRevision = Number(existing.revision || 0); + const nextRevision = previousRevision + 1; + await client.query( + `UPDATE leader_task_summary_webhook_state + SET enabled = false, + starts_at = now(), + revision = $2, + updated_at = now() + WHERE organization_id = $1`, + [organizationId, nextRevision] + ); + await client.query( + `UPDATE leader_task_summary_webhook_deliveries + SET delivery_status = CASE + WHEN delivery_status = 'sending' THEN 'uncertain' + ELSE 'cancelled' + END, + last_error = CASE + WHEN delivery_status = 'sending' + THEN 'Webhook 配置停用时请求状态未知,已停止自动重试。' + ELSE 'Webhook 配置已停用,待发摘要已取消。' + END, + updated_at = now() + WHERE organization_id = $1 + AND webhook_revision = $2 + AND delivery_status IN ('pending', 'sending')`, + [organizationId, previousRevision] + ); + await this.auditConfigurationChange( + client, organizationId, - leaderUserId, - channelId, - recipientAddress: inboundRecipient || externalRecipient, - conversationId: inboundRecipient ? row.conversation_id : undefined, - source: inboundRecipient ? 'agentbus_inbound' : 'channel_external_user_ref' - }); - if (!route) { - if (existing && await this.disableAutomaticSubscription(client, existing, 'route_unavailable')) changed += 1; - continue; - } - if (await this.syncAutomaticSubscription(client, route, existing)) changed += 1; + 'leader_summary_webhook.automatic_disabled', + { revision: nextRevision, historical_backfill: false } + ); + return true; } - for (const [leaderUserId, existing] of existingByLeader) { - if (currentLeaders.has(leaderUserId)) continue; - if (await this.disableAutomaticSubscription(client, existing, 'role_or_channel_invalid')) changed += 1; + if (!existing) { + await client.query( + `INSERT INTO leader_task_summary_webhook_state + (organization_id, enabled, starts_at, revision, configuration_fingerprint) + VALUES ($1, true, now(), 0, $2)`, + [organizationId, configurationFingerprint] + ); + await this.auditConfigurationChange( + client, + organizationId, + 'leader_summary_webhook.automatic_enabled', + { + revision: 0, + endpoint_fingerprint: this.endpointFingerprint(), + webhook_config_id: EXTERNAL_WEBHOOK_CONFIG_ID, + historical_backfill: false + } + ); + return true; } - return changed; + + const sameConfiguration = text(existing.configuration_fingerprint) === configurationFingerprint; + if (booleanValue(existing.enabled) && sameConfiguration) return false; + const previousRevision = Number(existing.revision || 0); + const nextRevision = previousRevision + 1; + await client.query( + `UPDATE leader_task_summary_webhook_state + SET enabled = true, + starts_at = now(), + revision = $2, + configuration_fingerprint = $3, + updated_at = now() + WHERE organization_id = $1`, + [organizationId, nextRevision, configurationFingerprint] + ); + await client.query( + `UPDATE leader_task_summary_webhook_deliveries + SET delivery_status = CASE + WHEN delivery_status = 'sending' THEN 'uncertain' + ELSE 'cancelled' + END, + last_error = CASE + WHEN delivery_status = 'sending' + THEN 'Webhook 配置变化时请求状态未知,已停止自动重试。' + ELSE 'Webhook 配置已变化,旧配置待发摘要已取消。' + END, + updated_at = now() + WHERE organization_id = $1 + AND webhook_revision = $2 + AND delivery_status IN ('pending', 'sending')`, + [organizationId, previousRevision] + ); + await this.auditConfigurationChange( + client, + organizationId, + sameConfiguration + ? 'leader_summary_webhook.automatic_reenabled' + : 'leader_summary_webhook.configuration_changed', + { + revision: nextRevision, + endpoint_fingerprint: this.endpointFingerprint(), + webhook_config_id: EXTERNAL_WEBHOOK_CONFIG_ID, + historical_backfill: false + } + ); + return true; }); } - async observeLeaderRoute(input: ObservedLeaderAgentBusRoute): Promise { - const route = this.automaticRoute({ - organizationId: text(input.organizationId), - leaderUserId: text(input.leaderUserId), - channelId: text(input.channelId), - recipientAddress: input.recipientAddress, - conversationId: input.conversationId, - source: 'agentbus_inbound' - }); - if (!route || !route.organizationId || !route.leaderUserId || !route.channelId) return; - await withTransaction(this.config, async (client) => { - await client.query( - `SELECT pg_advisory_xact_lock( - hashtextextended($1::text || ':leader-summary-automatic', 0) - )`, - [route.organizationId] - ); - const eligible = await client.query( - `SELECT channel.id - FROM user_channels channel - JOIN users leader - ON leader.id = channel.owner_user_id - AND leader.organization_id = channel.organization_id - AND leader.role = 'team_lead' - AND leader.is_active = true - AND leader.erp_account IS NOT NULL - WHERE channel.organization_id = $1 - AND channel.id = $2 - AND channel.owner_user_id = $3 - AND channel.enabled = true - FOR SHARE OF channel, leader`, - [route.organizationId, route.channelId, route.leaderUserId] - ); - if (!eligible.rowCount) return; - const existingResult = await client.query( - `SELECT * - FROM leader_task_summary_subscriptions - WHERE organization_id = $1 AND leader_user_id = $2 - FOR UPDATE`, - [route.organizationId, route.leaderUserId] - ); - await this.syncAutomaticSubscription( - client, - route, - existingResult.rows[0] as Record | undefined - ); - }); + async getWebhookStatus(organizationId: string): Promise { + const result = await getPool(this.config).query( + `SELECT state.*, + COALESCE(delivery.pending_count, 0)::int AS pending_count, + COALESCE(delivery.failed_count, 0)::int AS failed_count, + COALESCE(delivery.uncertain_count, 0)::int AS uncertain_count, + COALESCE(delivery.accepted_count, 0)::int AS accepted_count, + delivery.last_accepted_at + FROM leader_task_summary_webhook_state state + LEFT JOIN LATERAL ( + SELECT count(*) FILTER ( + WHERE d.delivery_status IN ('pending', 'sending') + ) AS pending_count, + count(*) FILTER (WHERE d.delivery_status = 'failed') AS failed_count, + count(*) FILTER (WHERE d.delivery_status = 'uncertain') AS uncertain_count, + count(*) FILTER (WHERE d.delivery_status = 'accepted') AS accepted_count, + max(d.accepted_at) AS last_accepted_at + FROM leader_task_summary_webhook_deliveries d + WHERE d.organization_id = state.organization_id + AND d.webhook_revision = state.revision + ) delivery ON true + WHERE state.organization_id = $1`, + [organizationId] + ); + const row = result.rows[0] as Record | undefined; + const configurationFingerprint = this.configurationFingerprint(); + const configured = Boolean(configurationFingerprint); + const configurationMatches = Boolean( + row && configurationFingerprint + && text(row.configuration_fingerprint) === configurationFingerprint + ); + const enabled = configured && configurationMatches && booleanValue(row?.enabled); + let state: PublicLeaderSummaryWebhookStatus['state']; + let statusMessage: string; + if (this.config.leaderSummaryWebhookConfigurationError) { + state = 'invalid_configuration'; + statusMessage = '外部 Webhook 配置无效,组长摘要已单独停用;普通任务与 AgentBus 服务不受影响。'; + } else if (!configured) { + state = 'not_configured'; + statusMessage = '运行环境尚未同时配置 WEBHOOK_SEND_URL 与 WEBHOOK_EXTERNAL_TOKEN,摘要不会外发。'; + } else if (!row) { + state = 'initializing'; + statusMessage = '外部 Webhook 配置正在初始化,暂时不会发送。'; + } else if (!enabled) { + state = 'configuration_changed'; + statusMessage = '外部 Webhook 配置正在切换生效起点,暂时不会发送。'; + } else { + state = 'active'; + statusMessage = '组长摘要将通过固定外部 Webhook 推送;成功仅表示接口已受理,不代表微信已经送达。'; + } + return { + delivery_mode: 'external_webhook', + webhook_config_id: EXTERNAL_WEBHOOK_CONFIG_ID, + configured, + enabled, + state, + configuration_error: this.config.leaderSummaryWebhookConfigurationError, + status_message: statusMessage, + endpoint_fingerprint: this.endpointFingerprint(), + starts_at: iso(row?.starts_at), + revision: row ? Number(row.revision || 0) : null, + pending_count: Number(row?.pending_count || 0), + failed_count: Number(row?.failed_count || 0), + uncertain_count: Number(row?.uncertain_count || 0), + accepted_count: Number(row?.accepted_count || 0), + last_accepted_at: iso(row?.last_accepted_at), + updated_at: iso(row?.updated_at) + }; } async projectPending(organizationId: string, limit = 100): Promise { + const configurationFingerprint = this.configurationFingerprint(); + if (!configurationFingerprint) return 0; const boundedLimit = Math.max(1, Math.min(500, Math.trunc(limit))); return withTransaction(this.config, async (client) => { const lock = await client.query( `SELECT pg_try_advisory_xact_lock( - hashtextextended($1::text || ':leader-summary-projector', 0) + hashtextextended($1::text || ':leader-summary-webhook-projector', 0) ) AS acquired`, [organizationId] ); if (!booleanValue(lock.rows[0]?.acquired)) return 0; const candidates = await client.query( - `SELECT subscription.id AS subscription_id, - subscription.revision AS subscription_revision, - subscription.leader_user_id, - subscription.channel_id, - subscription.recipient_address_ciphertext, - subscription.recipient_address_fingerprint, - subscription.conversation_id_ciphertext, - subscription.conversation_id_fingerprint, + `SELECT state.revision AS webhook_revision, task.id AS task_row_id, task.task_id, task.status, @@ -659,27 +435,12 @@ export class LeaderNotificationService { task.created_at, assignee.username AS assignee_username, source_event.id AS source_outbox_event_id, - COALESCE(delivery_state.has_needs_review, false) AS has_needs_review, COALESCE(delivery_state.has_exposed_needs_review, false) AS has_exposed_needs_review - FROM leader_task_summary_subscriptions subscription - JOIN users leader - ON leader.id = subscription.leader_user_id - AND leader.organization_id = subscription.organization_id - AND leader.role = 'team_lead' - AND leader.is_active = true - AND leader.erp_account IS NOT NULL - JOIN user_channels channel - ON channel.id = subscription.channel_id - AND channel.organization_id = subscription.organization_id - AND channel.owner_user_id = subscription.leader_user_id - AND channel.enabled = true + FROM leader_task_summary_webhook_state state JOIN tasks task - ON task.organization_id = subscription.organization_id + ON task.organization_id = state.organization_id AND task.assigned_user_id IS NOT NULL - AND task.assigned_user_id <> subscription.leader_user_id AND task.source IN ('manual', 'agentbus') - AND ((task.source = 'manual' AND subscription.include_manual) - OR (task.source = 'agentbus' AND subscription.include_agentbus)) JOIN users assignee ON assignee.id = task.assigned_user_id AND assignee.organization_id = task.organization_id @@ -687,44 +448,63 @@ export class LeaderNotificationService { JOIN LATERAL ( SELECT event.id FROM outbox_events event - WHERE event.organization_id = subscription.organization_id + WHERE event.organization_id = state.organization_id AND event.topic = 'task.updated' AND event.aggregate_type = 'task' AND event.aggregate_id = task.task_id - AND event.created_at >= subscription.starts_at - AND event.payload ->> 'status' = ANY($2::text[]) + AND event.created_at >= state.starts_at + AND event.payload ->> 'status' = ANY($3::text[]) AND (event.payload ->> 'archived') IS DISTINCT FROM 'true' AND (event.payload ->> 'restored') IS DISTINCT FROM 'true' ORDER BY event.id DESC LIMIT 1 ) source_event ON true LEFT JOIN LATERAL ( - SELECT bool_or(delivery.milestone = 'needs_review' AND delivery.delivery_status <> 'cancelled') AS has_needs_review, - bool_or(delivery.milestone = 'needs_review' AND delivery.delivery_status IN ('sending', 'delivered')) - AS has_exposed_needs_review, - bool_or(delivery.milestone = 'final' AND delivery.delivery_status <> 'cancelled') AS has_final, - bool_or(delivery.milestone = 'resolved' AND delivery.delivery_status <> 'cancelled') AS has_resolved - FROM leader_task_summary_deliveries delivery - WHERE delivery.subscription_id = subscription.id - AND delivery.subscription_revision = subscription.revision + SELECT bool_or( + delivery.milestone = 'needs_review' + AND delivery.delivery_status <> 'cancelled' + ) AS has_needs_review, + bool_or( + delivery.milestone = 'needs_review' + AND delivery.delivery_status IN ('sending', 'accepted', 'uncertain') + ) AS has_exposed_needs_review, + bool_or( + delivery.milestone = 'final' + AND delivery.delivery_status <> 'cancelled' + ) AS has_final, + bool_or( + delivery.milestone = 'resolved' + AND delivery.delivery_status <> 'cancelled' + ) AS has_resolved + FROM leader_task_summary_webhook_deliveries delivery + WHERE delivery.organization_id = state.organization_id + AND delivery.webhook_revision = state.revision AND delivery.task_id = task.id ) delivery_state ON true - WHERE subscription.organization_id = $1 - AND subscription.enabled = true - AND subscription.target_verified_at IS NOT NULL - AND task.status = ANY($2::text[]) + WHERE state.organization_id = $1 + AND state.enabled = true + AND state.configuration_fingerprint = $2 + AND task.status = ANY($3::text[]) AND ( - (task.status = ANY($3::text[]) AND NOT COALESCE(delivery_state.has_needs_review, false)) + (task.status = ANY($4::text[]) AND NOT COALESCE(delivery_state.has_needs_review, false)) OR - (NOT (task.status = ANY($3::text[])) AND ( - (COALESCE(delivery_state.has_exposed_needs_review, false) AND NOT COALESCE(delivery_state.has_resolved, false)) + (NOT (task.status = ANY($4::text[])) AND ( + (COALESCE(delivery_state.has_exposed_needs_review, false) + AND NOT COALESCE(delivery_state.has_resolved, false)) OR - (NOT COALESCE(delivery_state.has_exposed_needs_review, false) AND NOT COALESCE(delivery_state.has_final, false)) + (NOT COALESCE(delivery_state.has_exposed_needs_review, false) + AND NOT COALESCE(delivery_state.has_final, false)) )) ) - ORDER BY source_event.id ASC, subscription.id, task.id - LIMIT $4`, - [organizationId, [...PROJECTABLE_STATUSES], [...NEEDS_REVIEW_STATUSES], boundedLimit] + ORDER BY source_event.id ASC, task.id + LIMIT $5`, + [ + organizationId, + configurationFingerprint, + [...PROJECTABLE_STATUSES], + [...NEEDS_REVIEW_STATUSES], + boundedLimit + ] ); let projected = 0; @@ -734,28 +514,28 @@ export class LeaderNotificationService { let hadExposedNeedsReview = booleanValue(row.has_exposed_needs_review); if (!currentNeedsReview) { await client.query( - `UPDATE leader_task_summary_deliveries + `UPDATE leader_task_summary_webhook_deliveries SET delivery_status = 'cancelled', last_error = '任务已形成确定结果,未发送的旧核验提醒已取消。', updated_at = now() - WHERE subscription_id = $1 - AND subscription_revision = $2 + WHERE organization_id = $1 + AND webhook_revision = $2 AND task_id = $3 AND milestone = 'needs_review' - AND delivery_status IN ('pending', 'failed')`, - [row.subscription_id, row.subscription_revision, row.task_row_id] + AND delivery_status = 'pending'`, + [organizationId, row.webhook_revision, row.task_row_id] ); const exposure = await client.query( `SELECT EXISTS ( SELECT 1 - FROM leader_task_summary_deliveries - WHERE subscription_id = $1 - AND subscription_revision = $2 + FROM leader_task_summary_webhook_deliveries + WHERE organization_id = $1 + AND webhook_revision = $2 AND task_id = $3 AND milestone = 'needs_review' - AND delivery_status IN ('sending', 'delivered') + AND delivery_status IN ('sending', 'accepted', 'uncertain') ) AS exposed`, - [row.subscription_id, row.subscription_revision, row.task_row_id] + [organizationId, row.webhook_revision, row.task_row_id] ); hadExposedNeedsReview = booleanValue(exposure.rows[0]?.exposed); } @@ -769,37 +549,20 @@ export class LeaderNotificationService { hadNeedsReview: hadExposedNeedsReview }); if (!projection) continue; - const payload = { - event: 'task.summary' as const, - status: projection.deliveryStatus, - task_id: text(row.task_id), - text: projection.messageText - }; - const payloadText = JSON.stringify(payload); + const payloadText = projection.messageText; const inserted = await client.query( - `INSERT INTO leader_task_summary_deliveries - (organization_id, subscription_id, subscription_revision, task_id, - leader_user_id, channel_id, source_outbox_event_id, milestone, - recipient_address_ciphertext, recipient_address_fingerprint, - conversation_id_ciphertext, conversation_id_fingerprint, - payload_ciphertext, payload_fingerprint) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, - $9, $10, $11, $12, $13, $14) - ON CONFLICT (subscription_id, subscription_revision, task_id, milestone) DO NOTHING + `INSERT INTO leader_task_summary_webhook_deliveries + (organization_id, webhook_revision, task_id, source_outbox_event_id, + milestone, payload_ciphertext, payload_fingerprint) + VALUES ($1, $2, $3, $4, $5, $6, $7) + ON CONFLICT (organization_id, webhook_revision, task_id, milestone) DO NOTHING RETURNING id`, [ organizationId, - row.subscription_id, - row.subscription_revision, + row.webhook_revision, row.task_row_id, - row.leader_user_id, - row.channel_id, row.source_outbox_event_id, projection.milestone, - row.recipient_address_ciphertext, - row.recipient_address_fingerprint, - row.conversation_id_ciphertext, - row.conversation_id_fingerprint, encryptText(this.config, payloadText), sha256Text(payloadText) ] @@ -810,184 +573,155 @@ export class LeaderNotificationService { }); } - async claimDeliveries( - channelId: string, - leaseOwner: string, - limit = 5 - ): Promise { - const boundedLimit = Math.max(1, Math.min(10, Math.trunc(limit))); + private async claimPendingDeliveries( + organizationId: string, + limit = 1 + ): Promise { + const configurationFingerprint = this.configurationFingerprint(); + if (!configurationFingerprint) return []; + const boundedLimit = Math.max(1, Math.min(5, Math.trunc(limit))); return withTransaction(this.config, async (client) => { await client.query( - `UPDATE leader_task_summary_deliveries delivery - SET delivery_status = 'cancelled', - last_error = '订阅或目标已失效,摘要未改投其他渠道。', - updated_at = now() - WHERE delivery.channel_id = $1 - AND delivery.delivery_status IN ('pending', 'sending', 'failed') - AND NOT EXISTS ( - SELECT 1 - FROM leader_task_summary_subscriptions subscription - JOIN users leader - ON leader.id = subscription.leader_user_id - AND leader.organization_id = subscription.organization_id - AND leader.role = 'team_lead' - AND leader.is_active = true - AND leader.erp_account IS NOT NULL - JOIN user_channels channel - ON channel.id = subscription.channel_id - AND channel.organization_id = subscription.organization_id - AND channel.owner_user_id = subscription.leader_user_id - AND channel.enabled = true - WHERE subscription.id = delivery.subscription_id - AND subscription.enabled = true - AND subscription.target_verified_at IS NOT NULL - AND subscription.revision = delivery.subscription_revision - AND subscription.channel_id = delivery.channel_id - )`, - [channelId] + `SELECT pg_advisory_xact_lock( + hashtextextended($1::text || ':leader-summary-webhook-dispatcher', 0) + )`, + [organizationId] ); await client.query( - `UPDATE leader_task_summary_deliveries - SET delivery_status = 'pending', - last_error = COALESCE(last_error, 'delivery lease expired'), + `UPDATE leader_task_summary_webhook_deliveries + SET delivery_status = 'uncertain', + last_error = '发送租约过期,接口是否已受理未知,已停止自动重试。', updated_at = now() - WHERE channel_id = $1 + WHERE organization_id = $1 AND delivery_status = 'sending' AND updated_at < now() - interval '1 minute'`, - [channelId] + [organizationId] ); const pending = await client.query( - `SELECT delivery.*, - task.task_id AS public_task_id - FROM leader_task_summary_deliveries delivery - JOIN leader_task_summary_subscriptions subscription - ON subscription.id = delivery.subscription_id - AND subscription.enabled = true - AND subscription.target_verified_at IS NOT NULL - AND subscription.revision = delivery.subscription_revision - AND subscription.channel_id = delivery.channel_id - JOIN users leader - ON leader.id = subscription.leader_user_id - AND leader.organization_id = subscription.organization_id - AND leader.role = 'team_lead' - AND leader.is_active = true - AND leader.erp_account IS NOT NULL - JOIN user_channels channel - ON channel.id = subscription.channel_id - AND channel.organization_id = subscription.organization_id - AND channel.owner_user_id = subscription.leader_user_id - AND channel.enabled = true + `SELECT delivery.*, task.task_id AS public_task_id + FROM leader_task_summary_webhook_deliveries delivery + JOIN leader_task_summary_webhook_state state + ON state.organization_id = delivery.organization_id + AND state.enabled = true + AND state.revision = delivery.webhook_revision + AND state.configuration_fingerprint = $2 JOIN tasks task ON task.id = delivery.task_id AND task.organization_id = delivery.organization_id - WHERE delivery.channel_id = $1 - AND delivery.delivery_status IN ('pending', 'failed') - AND delivery.next_attempt_at <= now() + WHERE delivery.organization_id = $1 + AND delivery.delivery_status = 'pending' + AND delivery.attempt_count = 0 ORDER BY delivery.created_at ASC, delivery.id ASC FOR UPDATE OF delivery SKIP LOCKED - LIMIT $2`, - [channelId, boundedLimit] + LIMIT $3`, + [organizationId, configurationFingerprint, boundedLimit] ); - const deliveries: LeaderTaskSummaryDelivery[] = []; + const deliveries: ClaimedWebhookDelivery[] = []; for (const row of pending.rows as Record[]) { try { - const payloadText = decryptText(this.config, text(row.payload_ciphertext)); - const payload = JSON.parse(payloadText) as LeaderTaskSummaryDelivery['payload']; - const recipientAddress = decryptText(this.config, text(row.recipient_address_ciphertext)).trim(); - const conversationId = decryptText(this.config, text(row.conversation_id_ciphertext)).trim(); - if ( - payload.event !== 'task.summary' - || !['completed', 'failed', 'cancelled', 'needs_review'].includes(text(payload.status)) - || text(payload.task_id) !== text(row.public_task_id) - || !text(payload.text) - || !recipientAddress - || !conversationId - || sha256Text(payloadText) !== text(row.payload_fingerprint) - || sha256Text(recipientAddress) !== text(row.recipient_address_fingerprint) - || sha256Text(conversationId) !== text(row.conversation_id_fingerprint) - ) { - throw new Error('invalid leader summary delivery payload'); + const content = decryptText(this.config, text(row.payload_ciphertext)).trim(); + if (!content || sha256Text(content) !== text(row.payload_fingerprint)) { + throw new Error('invalid leader summary webhook payload'); } const updated = await client.query( - `UPDATE leader_task_summary_deliveries + `UPDATE leader_task_summary_webhook_deliveries SET delivery_status = 'sending', attempt_count = attempt_count + 1, - updated_at = now(), - last_error = $2 + last_error = $2, + updated_at = now() WHERE id = $1 + AND delivery_status = 'pending' + AND attempt_count = 0 RETURNING attempt_count`, - [row.id, `sending:${leaseOwner}`] + [row.id, `sending:${this.leaseOwner}`] ); + if (!updated.rowCount) continue; deliveries.push({ id: text(row.id), - channel_id: text(row.channel_id), - task_id: text(row.public_task_id), - recipient_address: recipientAddress, - recipient_fingerprint: text(row.recipient_address_fingerprint).slice(0, 12), - conversation_id: conversationId, - conversation_fingerprint: text(row.conversation_id_fingerprint).slice(0, 12), - payload, - attempt_count: Number(updated.rows[0]?.attempt_count || 1) + taskId: text(row.public_task_id), + content, + payloadFingerprint: text(row.payload_fingerprint).slice(0, 12), + attemptCount: Number(updated.rows[0]?.attempt_count || 1) }); } catch (error) { await client.query( - `UPDATE leader_task_summary_deliveries + `UPDATE leader_task_summary_webhook_deliveries SET delivery_status = 'cancelled', - last_error = '摘要投递密文或结构无效,已停止重试。', + last_error = '摘要投递密文或结构无效,已停止发送。', updated_at = now() WHERE id = $1`, [row.id] ); - this.log('error', { - agentbus_event: 'leader_summary_delivery_invalid', + this.log('error', 'webhook_delivery_invalid', { delivery_id: text(row.id), - channel_id: text(row.channel_id), task_id: text(row.public_task_id), - ...diagnosticError(error, 'leader_summary_delivery_invalid') - }, 'Leader task summary delivery is invalid'); + ...diagnosticError(error, 'leader_summary_webhook_delivery_invalid') + }, 'Leader task summary webhook delivery is invalid'); } } return deliveries; }); } - async markDeliveryDelivered(deliveryId: string): Promise { + private async markDeliveryOutcome( + deliveryId: string, + result: ExternalWebhookSendResult + ): Promise { + const deliveryStatus = result.outcome === 'accepted' + ? 'accepted' + : result.outcome === 'rejected' + ? 'failed' + : 'uncertain'; await getPool(this.config).query( - `UPDATE leader_task_summary_deliveries - SET delivery_status = 'delivered', - delivered_at = now(), - updated_at = now(), - last_error = NULL + `UPDATE leader_task_summary_webhook_deliveries + SET delivery_status = $2, + response_status = $3, + last_error = $4, + accepted_at = CASE WHEN $2 = 'accepted' THEN now() ELSE NULL END, + updated_at = now() WHERE id = $1 AND delivery_status = 'sending'`, - [deliveryId] + [deliveryId, deliveryStatus, result.httpStatus, result.errorCode] ); } - async markDeliveryFailed(deliveryId: string, errorMessage: string): Promise { - const normalized = text(errorMessage).slice(0, 500) || 'leader_summary_delivery_failed'; - await getPool(this.config).query( - `UPDATE leader_task_summary_deliveries - SET delivery_status = 'failed', - next_attempt_at = now() - + LEAST(300, GREATEST(5, power(2, LEAST(attempt_count, 8)))) * interval '1 second', - last_error = $2, - updated_at = now() - WHERE id = $1 AND delivery_status = 'sending'`, - [deliveryId, normalized] - ); - } - - async releaseDeliveries(channelId: string, leaseOwner: string): Promise { - await getPool(this.config).query( - `UPDATE leader_task_summary_deliveries - SET delivery_status = 'pending', - next_attempt_at = now(), - last_error = 'AgentBus 连接已断开,等待重发。', - updated_at = now() - WHERE channel_id = $1 - AND delivery_status = 'sending' - AND last_error = $2`, - [channelId, `sending:${leaseOwner}`] - ); + private async dispatchPending(organizationId: string): Promise { + if (!this.sender) return; + const deliveries = await this.claimPendingDeliveries(organizationId, 1); + for (const delivery of deliveries) { + this.log('info', 'webhook_delivery_sending', { + organization_id: organizationId, + delivery_id: delivery.id, + task_id: delivery.taskId, + payload_fingerprint: delivery.payloadFingerprint, + attempt_count: delivery.attemptCount + }, 'Leader task summary sending to external webhook'); + let result: ExternalWebhookSendResult; + try { + result = await this.sender.send(delivery.content); + } catch { + result = { + outcome: 'uncertain', + httpStatus: null, + errorCode: 'webhook_sender_failed' + }; + } + await this.markDeliveryOutcome(delivery.id, result); + const metadata = { + organization_id: organizationId, + delivery_id: delivery.id, + task_id: delivery.taskId, + payload_fingerprint: delivery.payloadFingerprint, + http_status: result.httpStatus, + outcome: result.outcome, + error_code: result.errorCode + }; + if (result.outcome === 'accepted') { + this.log('info', 'webhook_delivery_accepted', metadata, 'Leader task summary accepted by external webhook'); + } else if (result.outcome === 'rejected') { + this.log('warn', 'webhook_delivery_rejected', metadata, 'Leader task summary rejected by external webhook'); + } else { + this.log('warn', 'webhook_delivery_uncertain', metadata, 'Leader task summary webhook result is uncertain; automatic retry disabled'); + } + } } } diff --git a/control-plane/src/server.ts b/control-plane/src/server.ts index 75df6e9..65f2381 100644 --- a/control-plane/src/server.ts +++ b/control-plane/src/server.ts @@ -443,6 +443,8 @@ export async function buildServer({ database_ssl: config.DATABASE_SSL, artifact_storage_backend: config.ARTIFACT_STORAGE_BACKEND, agentbus_enabled: config.agentBusEnabled, + leader_summary_webhook_enabled: config.leaderSummaryWebhookEnabled, + leader_summary_webhook_configuration_error: config.leaderSummaryWebhookConfigurationError, parser_loop_enabled: startParserLoop, data_retention_enabled: config.DATA_RETENTION_ENABLED, raw_payload_logging: config.AGENTBUS_LOG_PAYLOADS @@ -539,9 +541,9 @@ export async function buildServer({ error: (metadata, message) => app.log.error(agentBusDiagnosticMetadata(metadata), message) }); const leaderNotificationService = new LeaderNotificationService(config, { - info: (metadata, message) => app.log.info(agentBusDiagnosticMetadata(metadata), message), - warn: (metadata, message) => app.log.warn(agentBusDiagnosticMetadata(metadata), message), - error: (metadata, message) => app.log.error(agentBusDiagnosticMetadata(metadata), message) + info: (metadata, message) => app.log.info(metadata, message), + warn: (metadata, message) => app.log.warn(metadata, message), + error: (metadata, message) => app.log.error(metadata, message) }); let agentBus: AgentBusManager | null = null; @@ -844,7 +846,6 @@ export async function buildServer({ tasks, organizationId: organization.id, scheduleParseQueue, - leaderNotifications: leaderNotificationService, logger: { info: (metadata, message) => app.log.info(agentBusDiagnosticMetadata(metadata), message), warn: (metadata, message) => app.log.warn(agentBusDiagnosticMetadata(metadata), message), @@ -852,7 +853,26 @@ export async function buildServer({ } }); await agentBus.start(); - leaderNotificationService.startProjector(organization.id); + } + if (config.leaderSummaryWebhookEnabled) { + try { + const organization = await auth.getOrganization(); + if (organization) { + leaderNotificationService.startProjector(organization.id); + } else { + app.log.warn({ + diagnostic_event: 'leader_summary.webhook_initialization_skipped', + notification_event: 'webhook_initialization_skipped', + error_code: 'leader_summary_webhook_organization_not_found' + }, 'Leader summary webhook initialization skipped; normal task services remain available'); + } + } catch (error) { + app.log.warn({ + diagnostic_event: 'leader_summary.webhook_initialization_failed', + notification_event: 'webhook_initialization_failed', + ...diagnosticError(error, 'leader_summary_webhook_initialization_failed') + }, 'Leader summary webhook initialization failed; normal task services remain available'); + } } async function getAiProbe(): Promise { @@ -1206,11 +1226,11 @@ export async function buildServer({ return { ok: true, ...result }; }); - app.get('/api/settings/leader-summary-subscriptions', async (request) => { + app.get('/api/settings/leader-summary-webhook', async (request) => { const session = await requireAdminSession(request); return { ok: true, - subscriptions: await leaderNotificationService.listSubscriptions(session.user.organizationId) + status: await leaderNotificationService.getWebhookStatus(session.user.organizationId) }; }); diff --git a/control-plane/test/account-authorization.test.ts b/control-plane/test/account-authorization.test.ts index bca5998..1db22a4 100644 --- a/control-plane/test/account-authorization.test.ts +++ b/control-plane/test/account-authorization.test.ts @@ -491,6 +491,6 @@ test('account authorization editor uses a scroll-safe open layout without overri assert.match(openLayoutSource, /overflow:\s*visible/); assert.doesNotMatch(styles, /\.account-panel\s*\{\s*grid-template-rows:/); - assert.match(index, /styles\.css\?v=20260907-admin-task-isolation-1/); - assert.match(index, /app\.js\?v=20260907-admin-task-isolation-1/); + assert.match(index, /styles\.css\?v=20260908-leader-webhook-1/); + assert.match(index, /app\.js\?v=20260908-leader-webhook-1/); }); diff --git a/control-plane/test/agentbus.test.ts b/control-plane/test/agentbus.test.ts index 6040b92..4fc5b7a 100644 --- a/control-plane/test/agentbus.test.ts +++ b/control-plane/test/agentbus.test.ts @@ -6,8 +6,6 @@ import { AgentBusListener, type AgentBusSocket, type AgentBusTaskGateway, - type LeaderNotificationGateway, - createLeaderTaskSummaryFrame, createTaskResultFrame, createTaskProgressFrame, extractAgentBusBusinessText, @@ -33,7 +31,6 @@ import { } from '../src/agentbus-delivery.js'; import { resolveBusinessRoute } from '../src/business-routes.js'; import type { AgentBusDelivery, PublicTask } from '../src/task-service.js'; -import type { LeaderTaskSummaryDelivery } from '../src/leader-notification-service.js'; function makeTask(status: string, overrides: Partial = {}): PublicTask { return { @@ -175,78 +172,6 @@ test('AgentBus configuration stays disabled until connection fields are supplied assert.equal(config.AGENTBUS_LOG_PAYLOADS, false); }); -test('team-lead listener automatically observes its inbound AgentBus route', async (t) => { - const socket = new FakeSocket(); - const observed: Array> = []; - let ingested = 0; - const tasks: AgentBusTaskGateway = { - events: new EventEmitter(), - async ingestMessage() { - ingested += 1; - return { task: makeTask('completed'), attached: false, created: true }; - }, - async getTask() { - return makeTask('completed'); - } - }; - const leaderNotifications: LeaderNotificationGateway = { - async observeLeaderRoute(input) { - observed.push(input as unknown as Record); - }, - async claimDeliveries() { - return []; - }, - async markDeliveryDelivered() {}, - async markDeliveryFailed() {}, - async releaseDeliveries() {} - }; - const listener = new AgentBusListener({ - config: testConfig(), - tasks, - leaderNotifications, - organizationId: 'org-1', - scheduleParseQueue: async () => {}, - socketFactory: () => socket as unknown as AgentBusSocket, - channel: { - id: 'channel-leader-route', - displayName: '组长微信', - wsUrl: 'wss://mesh.nianxx.cn/ws', - wsToken: 'leader-route-token', - botAddress: 'bot:leader-route:listener', - ownerUserId: 'leader-route-user', - ownerRole: 'team_lead' - } - }); - t.after(() => listener.stop()); - - listener.start(); - socket.readyState = 1; - socket.emit('open'); - socket.emit('message', JSON.stringify({ - id: 'ready-leader-route', - type: 'event', - session_id: 'session-leader-route', - epoch: 1, - to: 'bot:leader-route:listener', - payload: { event: 'session.ready' } - })); - socket.emit('message', JSON.stringify({ - id: 'leader-route-message-1', - type: 'event', - from: 'channel:wechat:leader-route-user', - payload: { text: '查询今天的任务' } - })); - - await waitFor(() => observed.length === 1 && ingested === 1); - assert.deepEqual(observed, [{ - organizationId: 'org-1', - leaderUserId: 'leader-route-user', - channelId: 'channel-leader-route', - recipientAddress: 'channel:wechat:leader-route-user', - conversationId: 'agentbus:channel:wechat:leader-route-user' - }]); -}); - test('AgentBus accepted delivery payloads keep final ownership metadata server-only', () => { const waitingPayload = createAgentBusAcceptedDeliveryPayload({ text: AGENTBUS_ROSTER_WAITING_TEXT, @@ -485,142 +410,6 @@ test('AgentBus protocol helpers preserve reply routing fields', () => { assert.equal(fallbackProgress.conversation_id, 'agentbus:channel:wechat:user-2'); }); -test('leader summary frame uses explicit proactive routing and cannot become an inbound task', () => { - const delivery: LeaderTaskSummaryDelivery = { - id: '11111111-1111-4111-8111-111111111111', - channel_id: '22222222-2222-4222-8222-222222222222', - task_id: 'TASK-20260907-001', - recipient_address: 'channel:wechat:leader-a', - recipient_fingerprint: 'abc123def456', - conversation_id: 'wechat-conversation-a', - conversation_fingerprint: 'def456abc123', - payload: { - event: 'task.summary', - status: 'completed', - task_id: 'TASK-20260907-001', - text: '【员工任务摘要】\n员工:employee-a' - }, - attempt_count: 1 - }; - const frame = createLeaderTaskSummaryFrame(delivery, { - id: 'session-leader-1', - epoch: 9, - address: 'bot:leader-a:listener' - }); - assert.equal(frame.id, `leader-summary-${delivery.id}`); - assert.equal(frame.from, 'bot:leader-a:listener'); - assert.equal(frame.to, delivery.recipient_address); - assert.equal(frame.conversation_id, delivery.conversation_id); - assert.equal(Object.hasOwn(frame, 'reply_to'), false); - assert.deepEqual(frame.payload, delivery.payload); - assert.equal(isInboundAgentBusTask(frame), false); -}); - -test('listener sends employee outbox first, then a redacted-log leader summary batch', async (t) => { - const socket = new FakeSocket(); - const order: string[] = []; - const delivered: string[] = []; - const released: string[] = []; - const logs: Array> = []; - const delivery: LeaderTaskSummaryDelivery = { - id: '33333333-3333-4333-8333-333333333333', - channel_id: 'channel-leader-a', - task_id: 'TASK-PRIVATE-1', - recipient_address: 'channel:wechat:private-leader-address', - recipient_fingerprint: 'a1b2c3d4e5f6', - conversation_id: 'private-wechat-conversation', - conversation_fingerprint: 'f6e5d4c3b2a1', - payload: { - event: 'task.summary', - status: 'completed', - task_id: 'TASK-PRIVATE-1', - text: '【员工任务摘要】\n员工:private-employee' - }, - attempt_count: 1 - }; - let claimed = false; - const tasks: AgentBusTaskGateway = { - events: new EventEmitter(), - async ingestMessage() { - return { task: makeTask('failed'), attached: false, created: true }; - }, - async getTask() { - return makeTask('failed'); - }, - async listAgentBusFinalizationCandidates() { - return []; - }, - async enqueueAgentBusResult() {}, - async claimAgentBusDeliveries() { - order.push('employee'); - return []; - }, - async markAgentBusDeliveryDelivered() {}, - async markAgentBusDeliveryFailed() {}, - async releaseAgentBusDeliveries() {} - }; - const leaderNotifications: LeaderNotificationGateway = { - async claimDeliveries() { - order.push('leader'); - if (claimed) return []; - claimed = true; - return [delivery]; - }, - async markDeliveryDelivered(deliveryId) { - delivered.push(deliveryId); - }, - async markDeliveryFailed() {}, - async releaseDeliveries(channelId) { - released.push(channelId); - } - }; - const listener = new AgentBusListener({ - config: testConfig(), - tasks, - leaderNotifications, - organizationId: 'org-1', - scheduleParseQueue: async () => {}, - socketFactory: () => socket as unknown as AgentBusSocket, - channel: { - id: 'channel-leader-a', - displayName: '组长 A', - wsUrl: 'wss://mesh.nianxx.cn/ws', - wsToken: 'leader-channel-token', - botAddress: 'bot:leader-a:listener', - ownerUserId: 'leader-a', - ownerRole: 'team_lead' - }, - logger: { - info(metadata) { logs.push(metadata); }, - warn(metadata) { logs.push(metadata); }, - error(metadata) { logs.push(metadata); } - } - }); - t.after(() => listener.stop()); - listener.start(); - socket.readyState = 1; - socket.emit('open'); - socket.emit('message', JSON.stringify({ - id: 'ready-leader-summary', - type: 'event', - session_id: 'session-leader-summary', - epoch: 1, - to: 'bot:leader-a:listener', - payload: { event: 'session.ready' } - })); - - await waitFor(() => socket.sent.length === 1 && delivered.length === 1); - assert.deepEqual(order.slice(0, 2), ['employee', 'leader']); - assert.equal(delivered[0], delivery.id); - assert.equal(socket.sent[0].reply_to, undefined); - assert.equal((socket.sent[0].payload as Record).event, 'task.summary'); - const logText = JSON.stringify(logs); - assert.doesNotMatch(logText, /private-leader-address|private-wechat-conversation|private-employee/); - assert.match(logText, /a1b2c3d4e5f6/); - listener.stop(); - await waitFor(() => released.includes('channel-leader-a')); -}); - test('AgentBus result text uses the unified important message and preserves the confirmation gate', () => { const needsInput = makeTask('awaiting_user_input', { input_request: { diff --git a/control-plane/test/control-plane.test.ts b/control-plane/test/control-plane.test.ts index 61d3d71..3bd1739 100644 --- a/control-plane/test/control-plane.test.ts +++ b/control-plane/test/control-plane.test.ts @@ -443,11 +443,11 @@ test('migration contains the durable state tables and safety fields', async () = } }); -test('control plane requires the administrator task-isolation migration before readiness', async () => { +test('control plane requires the leader summary webhook migration before readiness', async () => { const { readFile } = await import('node:fs/promises'); const db = await readFile(new URL('../src/db.ts', import.meta.url), 'utf8'); const server = await readFile(new URL('../src/server.ts', import.meta.url), 'utf8'); - assert.equal(REQUIRED_SCHEMA_VERSION, '021_admin_task_data_plane_isolation'); + assert.equal(REQUIRED_SCHEMA_VERSION, '023_leader_summary_webhook_delivery'); assert.match(db, /schema_migrations/); assert.match(db, /databaseReadiness/); assert.match(db, /assertDatabaseSchema/); @@ -1361,8 +1361,8 @@ test('operator page has a login gate and uses the durable task API', async () => const inpage = await readFile(new URL('../../chrome-extension/ltjt-order-assistant/inpage.js', import.meta.url), 'utf8'); assert.match(index, /id="loginPanel"/); assert.match(index, /id="workbench"[^>]*hidden/); - assert.match(index, /styles\.css\?v=20260907-admin-task-isolation-1/); - assert.match(index, /app\.js\?v=20260907-admin-task-isolation-1/); + assert.match(index, /styles\.css\?v=20260908-leader-webhook-1/); + assert.match(index, /app\.js\?v=20260908-leader-webhook-1/); assert.match(index, /id="statusDetailsPopover"/); assert.match(index, /id="statusDetailsRefresh"/); assert.match(app, /apiRequest\(`\/api\/tasks\?\$\{params\.toString\(\)\}`/); diff --git a/control-plane/test/external-webhook-client.test.ts b/control-plane/test/external-webhook-client.test.ts new file mode 100644 index 0000000..627cd50 --- /dev/null +++ b/control-plane/test/external-webhook-client.test.ts @@ -0,0 +1,230 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import { loadConfig } from '../src/config.js'; +import { buildServer } from '../src/server.js'; +import { + EXTERNAL_WEBHOOK_CONFIG_ID, + ExternalWebhookClient, + type ExternalWebhookFetch +} from '../src/external-webhook-client.js'; + +const webhookUrl = 'https://gateway.example.test/wechat/webhookInfo/sendMessageByOut'; +const webhookToken = 't'.repeat(32); + +test('external webhook request uses the confirmed fixed contract exactly', async () => { + let requestInput: string | URL | Request | undefined; + let requestInit: RequestInit | undefined; + const fetchImpl: ExternalWebhookFetch = async (input, init) => { + requestInput = input; + requestInit = init; + return new Response(JSON.stringify({ code: 0, data: true, msg: 'ok' }), { + status: 200, + headers: { 'content-type': 'application/json' } + }); + }; + const client = new ExternalWebhookClient({ + url: webhookUrl, + token: webhookToken, + fetchImpl + }); + + const result = await client.send('【员工任务摘要】\n员工:employee-a'); + + assert.deepEqual(result, { outcome: 'accepted', httpStatus: 200, errorCode: null }); + assert.equal(requestInput, webhookUrl); + assert.equal(requestInit?.method, 'POST'); + assert.deepEqual(requestInit?.headers, { + 'Content-Type': 'application/json', + 'x-token': webhookToken + }); + assert.deepEqual(JSON.parse(String(requestInit?.body)), { + id: EXTERNAL_WEBHOOK_CONFIG_ID, + content: '【员工任务摘要】\n员工:employee-a' + }); + assert.equal(EXTERNAL_WEBHOOK_CONFIG_ID, '9999'); + assert.ok(requestInit?.signal instanceof AbortSignal); +}); + +test('external webhook accepts only HTTP 200 with code zero and data true', async () => { + const cases: Array<{ + name: string; + response?: Response; + error?: Error; + expected: Record; + }> = [ + { + name: 'business rejection in HTTP 200', + response: new Response(JSON.stringify({ code: 1, data: false, msg: 'rejected' }), { status: 200 }), + expected: { outcome: 'rejected', httpStatus: 200, errorCode: 'webhook_business_rejected' } + }, + { + name: 'invalid request', + response: new Response(JSON.stringify({ code: 400, data: false, msg: 'id and content must not be blank' }), { status: 400 }), + expected: { outcome: 'rejected', httpStatus: 400, errorCode: 'webhook_invalid_request' } + }, + { + name: 'invalid token', + response: new Response(JSON.stringify({ code: 401, data: false, msg: 'Invalid token' }), { status: 401 }), + expected: { outcome: 'rejected', httpStatus: 401, errorCode: 'webhook_invalid_token' } + }, + { + name: 'missing webhook configuration', + response: new Response(JSON.stringify({ code: 404, data: false, msg: 'Webhook configuration not found' }), { status: 404 }), + expected: { outcome: 'rejected', httpStatus: 404, errorCode: 'webhook_configuration_not_found' } + }, + { + name: 'server error is uncertain', + response: new Response(JSON.stringify({ code: 500, data: false, msg: 'error' }), { status: 500 }), + expected: { outcome: 'uncertain', httpStatus: 500, errorCode: 'webhook_http_500_uncertain' } + }, + { + name: 'invalid success response is uncertain', + response: new Response('not-json', { status: 200 }), + expected: { outcome: 'uncertain', httpStatus: 200, errorCode: 'webhook_response_invalid_json' } + }, + { + name: 'network failure is uncertain', + error: new Error('connection reset'), + expected: { outcome: 'uncertain', httpStatus: null, errorCode: 'webhook_request_uncertain' } + }, + { + name: 'timeout is uncertain', + error: new DOMException('timed out', 'TimeoutError'), + expected: { outcome: 'uncertain', httpStatus: null, errorCode: 'webhook_request_timeout' } + } + ]; + + for (const item of cases) { + const client = new ExternalWebhookClient({ + url: webhookUrl, + token: webhookToken, + fetchImpl: async () => { + if (item.error) throw item.error; + return item.response as Response; + } + }); + assert.deepEqual(await client.send('summary'), item.expected, item.name); + } +}); + +test('webhook configuration is fail-closed and independent from AgentBus configuration', () => { + const encryptionKey = Buffer.alloc(32, 23).toString('base64'); + const disabled = loadConfig({ NODE_ENV: 'test', FIELD_ENCRYPTION_KEY: encryptionKey }); + assert.equal(disabled.leaderSummaryWebhookEnabled, false); + assert.equal(disabled.agentBusEnabled, false); + + const configured = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: webhookUrl, + WEBHOOK_EXTERNAL_TOKEN: webhookToken + }); + assert.equal(configured.leaderSummaryWebhookEnabled, true); + assert.equal(configured.leaderSummaryWebhookConfigurationError, null); + assert.equal(configured.agentBusEnabled, false); + + const incomplete = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_EXTERNAL_TOKEN: webhookToken + }); + assert.equal(incomplete.leaderSummaryWebhookEnabled, false); + assert.equal(incomplete.leaderSummaryWebhookConfigurationError, 'webhook_configuration_incomplete'); + assert.equal(incomplete.agentBusEnabled, false); + + const shortToken = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: webhookUrl, + WEBHOOK_EXTERNAL_TOKEN: 'too-short' + }); + assert.equal(shortToken.leaderSummaryWebhookEnabled, false); + assert.equal(shortToken.leaderSummaryWebhookConfigurationError, 'webhook_token_length_invalid'); + + const paddedToken = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: webhookUrl, + WEBHOOK_EXTERNAL_TOKEN: ` ${webhookToken}` + }); + assert.equal(paddedToken.leaderSummaryWebhookEnabled, false); + assert.equal(paddedToken.leaderSummaryWebhookConfigurationError, 'webhook_token_length_invalid'); + + const whitespaceToken = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: webhookUrl, + WEBHOOK_EXTERNAL_TOKEN: ` ${'t'.repeat(31)}` + }); + assert.equal(whitespaceToken.leaderSummaryWebhookEnabled, false); + assert.equal(whitespaceToken.leaderSummaryWebhookConfigurationError, 'webhook_token_whitespace_invalid'); + + const unconfirmedRoute = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: 'https://gateway.example.test/wechat/another-route', + WEBHOOK_EXTERNAL_TOKEN: webhookToken + }); + assert.equal(unconfirmedRoute.leaderSummaryWebhookEnabled, false); + assert.equal(unconfirmedRoute.leaderSummaryWebhookConfigurationError, 'webhook_route_unconfirmed'); + + const insecureProductionUrl = loadConfig({ + NODE_ENV: 'production', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: 'http://gateway.example.test/wechat/webhookInfo/sendMessageByOut', + WEBHOOK_EXTERNAL_TOKEN: webhookToken + }); + assert.equal(insecureProductionUrl.leaderSummaryWebhookEnabled, false); + assert.equal(insecureProductionUrl.leaderSummaryWebhookConfigurationError, 'webhook_https_required'); +}); + +test('client rejects blank content and malformed tokens before making a request', async () => { + let calls = 0; + const fetchImpl: ExternalWebhookFetch = async () => { + calls += 1; + return new Response('{}', { status: 200 }); + }; + await assert.rejects( + new ExternalWebhookClient({ url: webhookUrl, token: webhookToken, fetchImpl }).send(' '), + /must not be blank/ + ); + await assert.rejects( + new ExternalWebhookClient({ url: webhookUrl, token: 'short', fetchImpl }).send('summary'), + /exactly 32 non-whitespace characters/ + ); + await assert.rejects( + new ExternalWebhookClient({ url: webhookUrl, token: ` ${'t'.repeat(31)}`, fetchImpl }).send('summary'), + /exactly 32 non-whitespace characters/ + ); + assert.equal(calls, 0); +}); + +test('invalid configuration or webhook initialization failure does not block the normal HTTP service', async () => { + const encryptionKey = Buffer.alloc(32, 29).toString('base64'); + const configs = [ + loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_EXTERNAL_TOKEN: webhookToken, + DATABASE_URL: 'postgresql://invalid:invalid@127.0.0.1:1/invalid' + }), + loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: encryptionKey, + WEBHOOK_SEND_URL: webhookUrl, + WEBHOOK_EXTERNAL_TOKEN: webhookToken, + DATABASE_URL: 'postgresql://invalid:invalid@127.0.0.1:1/invalid' + }) + ]; + + for (const config of configs) { + const { app } = await buildServer({ config, startParserLoop: false }); + try { + const response = await app.inject({ method: 'GET', url: '/health/live' }); + assert.equal(response.statusCode, 200); + assert.equal(response.json().ok, true); + } finally { + await app.close(); + } + } +}); diff --git a/control-plane/test/leader-notification-contract.test.ts b/control-plane/test/leader-notification-contract.test.ts index a87d87e..30ecd5e 100644 --- a/control-plane/test/leader-notification-contract.test.ts +++ b/control-plane/test/leader-notification-contract.test.ts @@ -6,118 +6,114 @@ async function source(relativePath: string): Promise { return readFile(new URL(relativePath, import.meta.url), 'utf8'); } -test('migration keeps the automatic feature fail-closed in storage and separate from employee replies', async () => { - const sql = await source('../migrations/020_leader_task_summary_notifications.sql'); - assert.match(sql, /CREATE TABLE IF NOT EXISTS leader_task_summary_subscriptions/); - // The runtime reconciler is the only component that activates rows. A raw - // insert must remain inert when identity or routing checks have not run. - assert.match(sql, /enabled boolean NOT NULL DEFAULT false/); - assert.match(sql, /scope text NOT NULL DEFAULT 'organization'/); - assert.match(sql, /CHECK \(include_manual OR include_agentbus\)/); - assert.match(sql, /recipient_address_ciphertext text NOT NULL/); - assert.match(sql, /conversation_id_ciphertext text NOT NULL/); +test('migration retires only the legacy leader-summary transport and creates an encrypted webhook outbox', async () => { + const sql = await source('../migrations/023_leader_summary_webhook_delivery.sql'); + assert.match(sql, /UPDATE leader_task_summary_subscriptions[\s\S]*?SET enabled = false/); + assert.match(sql, /UPDATE leader_task_summary_deliveries[\s\S]*?delivery_status = 'cancelled'/); + assert.match(sql, /WHERE delivery_status IN \('pending', 'sending', 'failed'\)/); + assert.match(sql, /CREATE TABLE IF NOT EXISTS leader_task_summary_webhook_state/); + assert.match(sql, /CREATE TABLE IF NOT EXISTS leader_task_summary_webhook_deliveries/); + assert.match(sql, /configuration_fingerprint text NOT NULL/); assert.match(sql, /payload_ciphertext text NOT NULL/); - assert.match(sql, /target_verified_at timestamptz/); - assert.match(sql, /UNIQUE \(organization_id, leader_user_id\)/); - assert.match(sql, /UNIQUE \(subscription_id, subscription_revision, task_id, milestone\)/); - assert.doesNotMatch(sql, /INSERT\s+INTO/iu, 'schema migration must not backfill or send historical tasks'); - assert.doesNotMatch(sql, /REFERENCES\s+agentbus_deliveries/iu); - assert.doesNotMatch(sql, /recipient_address\s+text/iu); - assert.doesNotMatch(sql, /conversation_id\s+text/iu); - assert.doesNotMatch(sql, /payload\s+jsonb/iu); + assert.match(sql, /payload_fingerprint text NOT NULL/); + assert.match(sql, /delivery_status IN \('pending', 'sending', 'accepted', 'failed', 'uncertain', 'cancelled'\)/); + assert.match(sql, /CHECK \(attempt_count BETWEEN 0 AND 1\)/); + assert.match(sql, /UNIQUE \(organization_id, webhook_revision, task_id, milestone\)/); + assert.match(sql, /FOREIGN KEY \(organization_id, task_id\)[\s\S]*?REFERENCES tasks \(organization_id, id\)/); + const newOutbox = sql.slice(sql.indexOf('CREATE TABLE IF NOT EXISTS leader_task_summary_webhook_deliveries')); + assert.doesNotMatch(newOutbox, /\b(channel_id|leader_user_id|recipient_address|conversation_id|token|webhook_url)\b/i); + assert.doesNotMatch(sql, /INSERT\s+INTO\s+leader_task_summary_webhook_deliveries/i, + 'schema migration must not backfill historical employee activity'); + assert.doesNotMatch(sql, /UPDATE\s+(tasks|task_events|task_attempts|agentbus_deliveries)\b/i, + 'transport migration must not mutate normal task or employee AgentBus state'); }); -test('automatic reconciliation derives enabled subscriptions from leader identity and AgentBus routing', async () => { +test('webhook worker is independent from normal task execution and AgentBus delivery', async () => { const service = await source('../src/leader-notification-service.ts'); - const channels = await source('../src/agentbus-channels.ts'); - assert.match(service, /reconcileAutomaticSubscriptions/); - assert.match(service, /leader\.role = 'team_lead'/); - assert.match(service, /leader\.is_active = true/); - assert.match(service, /leader\.erp_account IS NOT NULL/); - assert.match(service, /channel\.owner_user_id = leader\.id/); - assert.match(service, /latest_route\.inbound_from/); - assert.match(service, /channel\.external_user_ref/); - assert.match(service, /route_task\.assigned_user_id = leader\.id/); - assert.match(service, /source: inboundRecipient \? 'agentbus_inbound' : 'channel_external_user_ref'/); - assert.match(service, /`agentbus:\$\{recipientAddress\}`/); - assert.match(service, /VALUES \(\$1, \$2, \$3, true, true, true, now\(\), 0/); - assert.match(service, /target_verified_at = now\(\)/); - assert.match(service, /automatic_disabled/); - assert.match(service, /channel_disabled.*route_unavailable.*role_or_channel_invalid/s); - assert.doesNotMatch(service, /upsertSubscription|expectedRevision|leader_summary_target_unverified/); - assert.match(channels, /ownerChanged[\s\S]*?\? null[\s\S]*?: currentExternalUserRef/); -}); - -test('projection remains organization-scoped, future-only, role-safe and encrypted', async () => { - const service = await source('../src/leader-notification-service.ts'); - assert.match(service, /event\.topic = 'task\.updated'/); - assert.match(service, /event\.created_at >= subscription\.starts_at/); - assert.match(service, /task\.assigned_user_id <> subscription\.leader_user_id/); - assert.match(service, /assignee\.role <> 'admin'/); - assert.match(service, /task\.source IN \('manual', 'agentbus'\)/); - assert.match(service, /subscription\.include_manual/); - assert.match(service, /subscription\.include_agentbus/); - assert.match(service, /event\.payload ->> 'archived'.*IS DISTINCT FROM 'true'/s); - assert.match(service, /event\.payload ->> 'restored'.*IS DISTINCT FROM 'true'/s); - assert.match(service, /subscription\.target_verified_at IS NOT NULL/); - assert.match(service, /pg_try_advisory_xact_lock/); - assert.match(service, /encryptText\(this\.config, payloadText\)/); - assert.match(service, /FOR UPDATE OF delivery SKIP LOCKED/); - assert.match(service, /sha256Text\(payloadText\).*payload_fingerprint/s); - assert.match(service, /sha256Text\(recipientAddress\).*recipient_address_fingerprint/s); - assert.match(service, /sha256Text\(conversationId\).*conversation_id_fingerprint/s); - assert.match(service, /delivery_status = 'cancelled'.*旧路由待发送摘要已取消/s); - assert.doesNotMatch(service, /reply_to/); -}); - -test('HTTP exposes administrator status only and no manual mutation or live test-send endpoint', async () => { const server = await source('../src/server.ts'); - const routeStart = server.indexOf("app.get('/api/settings/leader-summary-subscriptions'"); - const routeEnd = server.indexOf("app.post('/api/auth/logout'", routeStart); - const routes = server.slice(routeStart, routeEnd); - assert.ok(routeStart > 0 && routeEnd > routeStart); - assert.match(routes, /requireAdminSession/); - assert.doesNotMatch(routes, /app\.(put|post|patch|delete)/); - assert.doesNotMatch(server, /leaderSummarySubscriptionSchema|upsertSubscription/); - assert.doesNotMatch(server, /leader-summary-subscriptions.*test-send|leader-summary-subscriptions.*test\/send/s); -}); - -test('AgentBus sends summaries as reserved low-priority proactive events without plaintext frame logging', async () => { const agentbus = await source('../src/agentbus.ts'); + const channels = await source('../src/agentbus-channels.ts'); + + assert.match(service, /new ExternalWebhookClient\(/); + assert.match(service, /config\.leaderSummaryWebhookEnabled/); + assert.match(server, /if \(config\.leaderSummaryWebhookEnabled\)[\s\S]*?leaderNotificationService\.startProjector/); + assert.doesNotMatch(server, /new AgentBusManager\([\s\S]*?leaderNotifications:/); + assert.doesNotMatch(agentbus, /LeaderNotificationGateway|createLeaderTaskSummaryFrame|flushLeaderDeliveries|sendLeaderDelivery|observeLeaderRoute/); + assert.doesNotMatch(channels, /leaderNotifications:/); + + const flushStart = agentbus.indexOf('private async flushOutboundDeliveries'); + const flushEnd = agentbus.indexOf('private async sendDurableDelivery', flushStart); + const employeeFlush = agentbus.slice(flushStart, flushEnd); + assert.match(employeeFlush, /await this\.flushDurableDeliveries\(\)/); + assert.doesNotMatch(employeeFlush, /leader|webhook/i); + const ignoreStart = agentbus.indexOf('function inboundFrameIgnoreReason'); const ignoreEnd = agentbus.indexOf('export function parseAgentBusFrame', ignoreStart); assert.match(agentbus.slice(ignoreStart, ignoreEnd), /'task\.summary'/); assert.match(agentbus.slice(ignoreStart, ignoreEnd), /reserved_event/); - assert.match(agentbus, /id: `leader-summary-\$\{delivery\.id\}`/); - assert.match(agentbus, /to: delivery\.recipient_address/); - assert.match(agentbus, /conversation_id: delivery\.conversation_id/); - const frameStart = agentbus.indexOf('export function createLeaderTaskSummaryFrame'); - const frameEnd = agentbus.indexOf('export function taskResultStatus', frameStart); - assert.doesNotMatch(agentbus.slice(frameStart, frameEnd), /reply_to/); - const flushStart = agentbus.indexOf('private async flushOutboundDeliveries'); - const sendEnd = agentbus.indexOf('private async sendDurableDelivery', flushStart); - const leaderDelivery = agentbus.slice(flushStart, sendEnd); - assert.match(leaderDelivery, /await this\.flushDurableDeliveries\(\)/); - assert.match(leaderDelivery, /await this\.flushLeaderDeliveries\(\)/); - assert.match(leaderDelivery, /recipient_fingerprint/); - assert.match(leaderDelivery, /conversation_fingerprint/); - assert.doesNotMatch(leaderDelivery, /this\.logFrame/); - assert.match(agentbus, /ownerRole === 'team_lead'/); - assert.match(agentbus, /observeLeaderRoute/); - assert.match(agentbus, /recipientAddress: text\(frame\.from\),\s*conversationId/s); }); -test('administrator UI is read-only and explains role-driven automatic delivery', async () => { +test('projection is organization-scoped, future-only, role-safe, source-limited and encrypted', async () => { + const service = await source('../src/leader-notification-service.ts'); + assert.match(service, /event\.topic = 'task\.updated'/); + assert.match(service, /event\.created_at >= state\.starts_at/); + assert.match(service, /task\.assigned_user_id IS NOT NULL/); + assert.match(service, /assignee\.organization_id = task\.organization_id/); + assert.match(service, /assignee\.role <> 'admin'/); + assert.match(service, /task\.source IN \('manual', 'agentbus'\)/); + assert.match(service, /event\.payload ->> 'archived'.*IS DISTINCT FROM 'true'/s); + assert.match(service, /event\.payload ->> 'restored'.*IS DISTINCT FROM 'true'/s); + assert.match(service, /pg_try_advisory_xact_lock/); + assert.match(service, /VALUES \(\$1::uuid, NULL, \$2, 'leader_task_summary_webhook', \$1::uuid::text, NULL, \$3\)/, + 'organization IDs must be explicitly cast when shared by UUID and text audit columns'); + assert.match(service, /encryptText\(this\.config, payloadText\)/); + assert.match(service, /sha256Text\(payloadText\)/); + assert.match(service, /ON CONFLICT \(organization_id, webhook_revision, task_id, milestone\) DO NOTHING/); + assert.doesNotMatch(service, /UPDATE\s+(tasks|task_events|task_attempts|agentbus_deliveries)\b/i, + 'summary projection must never write normal task or employee reply state'); +}); + +test('delivery makes one attempt and preserves ambiguous outcomes instead of retrying', async () => { + const service = await source('../src/leader-notification-service.ts'); + assert.match(service, /delivery_status = 'sending'[\s\S]*?updated_at < now\(\) - interval '1 minute'/); + assert.match(service, /SET delivery_status = 'uncertain'[\s\S]*?发送租约过期/); + assert.match(service, /delivery\.attempt_count = 0/); + assert.match(service, /attempt_count = attempt_count \+ 1/); + assert.match(service, /FOR UPDATE OF delivery SKIP LOCKED/); + assert.match(service, /result\.outcome === 'accepted'[\s\S]*?'accepted'[\s\S]*?'failed'[\s\S]*?'uncertain'/); + assert.match(service, /automatic retry disabled/); + const outcomeStart = service.indexOf('private async markDeliveryOutcome'); + const outcomeEnd = service.indexOf('private async dispatchPending', outcomeStart); + const outcomeWriter = service.slice(outcomeStart, outcomeEnd); + assert.ok(outcomeStart > 0 && outcomeEnd > outcomeStart); + assert.doesNotMatch(outcomeWriter, /'pending'/, + 'terminal delivery outcomes must never be placed back into the pending queue'); +}); + +test('HTTP exposes administrator-only status and no mutation or live test-send endpoint', async () => { + const server = await source('../src/server.ts'); + const routeStart = server.indexOf("app.get('/api/settings/leader-summary-webhook'"); + const routeEnd = server.indexOf("app.post('/api/auth/logout'", routeStart); + const routes = server.slice(routeStart, routeEnd); + assert.ok(routeStart > 0 && routeEnd > routeStart); + assert.match(routes, /requireAdminSession/); + assert.match(routes, /getWebhookStatus/); + assert.doesNotMatch(routes, /app\.(put|post|patch|delete)/); + assert.doesNotMatch(server, /leader-summary-webhook.*test-send|leader-summary-webhook.*test\/send/s); +}); + +test('administrator UI explains the isolated fixed webhook and accepted-not-delivered semantics', async () => { const html = await source('../../LianSyn-platform/index.html'); const app = await source('../../LianSyn-platform/app.js'); - assert.match(html, /组长任务摘要抄送/); - assert.match(html, /无需单独配置或启用/); - assert.match(html, /只处理自动启用后的新结果,不补发历史任务/); - assert.match(html, /已经到达微信的消息无法撤回/); - assert.doesNotMatch(html, /leaderSummary(Form|ChannelId|RecipientAddress|ConversationId|TargetVerified|Enabled|Save)/); - assert.match(app, /自动推送/); - assert.match(app, /等待路由/); - assert.match(app, /无需维护收件地址或微信会话 ID/); - assert.doesNotMatch(app, /leader-summary-subscriptions\/\$\{|expected_revision: subscription\.revision/); - assert.doesNotMatch(app, /leader-summary-subscriptions[^'"\n]*backfill/); + assert.match(html, /组长任务摘要推送/); + assert.match(html, /不再依赖组长 AgentBus 账号或渠道/); + assert.match(html, /WEBHOOK_SEND_URL/); + assert.match(html, /WEBHOOK_EXTERNAL_TOKEN/); + assert.match(html, /只处理配置生效后的新结果,不补发历史任务/); + assert.match(html, /已受理.*不代表微信已送达/); + assert.match(html, /不会自动重试/); + assert.match(app, /\/api\/settings\/leader-summary-webhook/); + assert.match(app, /人工任务、AgentBus 任务/); + assert.match(app, /已使用独立外部 Webhook,不受该渠道影响/); + assert.doesNotMatch(app, /leader-summary-subscriptions/); });