From fd3934760378ae41c83cb92a6b8aa4451c46138b Mon Sep 17 00:00:00 2001 From: inman Date: Mon, 7 Sep 2026 20:19:25 +0800 Subject: [PATCH] feat: isolate administrators from task data plane --- LianSyn-platform/README.md | 2 +- LianSyn-platform/app.js | 165 +++++++++++--- LianSyn-platform/index.html | 12 +- control-plane/README.md | 22 +- .../021_admin_task_data_plane_isolation.sql | 213 ++++++++++++++++++ control-plane/src/agentbus.ts | 2 +- control-plane/src/auth.ts | 4 +- control-plane/src/db.ts | 2 +- control-plane/src/server.ts | 78 +++++-- control-plane/src/task-service.ts | 60 +++-- .../test/account-authorization.test.ts | 44 +++- control-plane/test/agentbus.test.ts | 2 +- control-plane/test/control-plane.test.ts | 150 +++++++++++- 13 files changed, 644 insertions(+), 112 deletions(-) create mode 100644 control-plane/migrations/021_admin_task_data_plane_isolation.sql diff --git a/LianSyn-platform/README.md b/LianSyn-platform/README.md index 2b78ecf..b826a6f 100644 --- a/LianSyn-platform/README.md +++ b/LianSyn-platform/README.md @@ -4,7 +4,7 @@ 操作台首页只展示按创建时间倒序排列的最近 10 条任务;任务数量超过 10 条时,点击“查看更多”进入 `/history` 历史任务目录。历史目录仍使用同一登录会话和任务详情面板,支持按任务编号/摘要搜索、按状态筛选、分页,以及补充信息、确认提交、插件回查和彻底删除等现有操作。历史目录不新增归档状态,服务端持久化任务仍是唯一事实源。 -操作台不会在页面刷新、插件重连或状态轮询时自动重交 ERP 任务。只有管理员的首次明确确认会申请服务端唯一执行权;领取成功后才向插件下发。领取后的投递超时或 ERP 回执不确定都会进入“待回查”,禁止自动重试。 +操作台不会在页面刷新、插件重连或状态轮询时自动重交 ERP 任务。只有任务所属员工账号的首次明确确认会申请服务端唯一执行权;领取成功后才向该账号绑定的插件下发。领取后的投递超时或 ERP 回执不确定都会进入“待回查”,禁止自动重试。管理员只使用账号、渠道、解析策略和审计等管理功能,不进入任务数据面。 生产控制平面说明见 [../control-plane/README.md](../control-plane/README.md)。 diff --git a/LianSyn-platform/app.js b/LianSyn-platform/app.js index 6b2e4c9..ff2734b 100644 --- a/LianSyn-platform/app.js +++ b/LianSyn-platform/app.js @@ -163,12 +163,16 @@ function isTeamLeader() { return authUser?.role === 'team_lead'; } +function canUseTaskDataPlane(user = authUser) { + return user?.role === 'team_lead' || user?.role === 'user'; +} + function canViewOperationsDashboard() { - return isAdministrator() || isTeamLeader(); + return isTeamLeader(); } function connectionIdForUser(user) { - if (!user?.id) return ''; + if (!user?.id || !canUseTaskDataPlane(user)) return ''; const storageKey = `liansyn_platform_browser_connection_id:${user.id}`; let connectionId = localStorage.getItem(storageKey) || ''; if (!/^platform-browser:[a-f0-9-]{16,80}$/i.test(connectionId)) { @@ -237,6 +241,11 @@ async function encodeRosterAttachment(file) { } async function apiRequest(path, options = {}) { + if (authUser && !canUseTaskDataPlane() && isTaskDataPlaneApiPath(path)) { + const accessError = new Error('管理员账号仅用于平台管理,不能访问业务任务。'); + accessError.code = 'task_access_forbidden'; + throw accessError; + } const method = String(options.method || 'GET').toUpperCase(); const headers = new Headers(options.headers || {}); const timeoutMs = Number(options.timeoutMs) || requestTimeoutForPath(path); @@ -273,6 +282,60 @@ async function apiRequest(path, options = {}) { return payload || {}; } +function isTaskDataPlaneApiPath(path) { + const pathname = String(path || '').split(/[?#]/, 1)[0]; + return pathname === '/api/messages' + || pathname === '/api/connections/heartbeat' + || pathname === '/api/events' + || pathname === '/api/tasks' + || pathname.startsWith('/api/tasks/') + || pathname.startsWith('/api/parser-decisions/') + || pathname === '/api/operations-dashboard' + || pathname.startsWith('/api/operations-dashboard/'); +} + +function hideAuthenticatedNavigation() { + for (const selector of [ + '#workbenchNav', + '#historyNav', + '#operationsDashboardNav', + '#channelsNav', + '#parserRoutingNav', + '#accountsNav', + '#auditNav', + '#bridgeState' + ]) { + const element = $(selector); + if (element) element.hidden = true; + } +} + +function clearTaskDataPlaneClientState() { + currentTaskId = ''; + taskStore = []; + runtimeTaskStore.clear(); + remoteTaskStore.clear(); + taskDetailStore.clear(); + taskDetailRequests.clear(); + taskInputHistoryStore.clear(); + taskInputHistoryRequests.clear(); + localTaskOverlayStore.clear(); + extensionResultPersistQueues.clear(); + pendingExtensionResults.clear(); + persistedExtensionResultVersions.clear(); + taskArchiveStates.clear(); + taskDeleteStates.clear(); + taskReplyStates.clear(); + taskReplyDrafts.clear(); + taskAttachmentDrafts.clear(); + historySelectedTaskIds.clear(); + locallyDeletedTaskIds.clear(); + sessionStorage.removeItem('liansyn_platform_current_task_id'); + stopPolling(); + if (eventStream) eventStream.close(); + eventStream = null; +} + function showLoginPanel(message = '') { authUser = null; browserConnectionId = ''; @@ -314,6 +377,7 @@ function showLoginPanel(message = '') { if (automationButton) automationButton.hidden = true; if (logoutButton) logoutButton.hidden = true; if (changePasswordButton) changePasswordButton.hidden = true; + hideAuthenticatedNavigation(); const submitButton = $('#loginForm button[type="submit"]'); if (submitButton) submitButton.disabled = false; const error = $('#loginError'); @@ -354,6 +418,7 @@ function showAuthChecking() { if (automationButton) automationButton.hidden = true; if (logoutButton) logoutButton.hidden = true; if (changePasswordButton) changePasswordButton.hidden = true; + hideAuthenticatedNavigation(); const submitButton = $('#loginForm button[type="submit"]'); if (submitButton) submitButton.disabled = true; const error = $('#loginError'); @@ -362,14 +427,22 @@ function showAuthChecking() { function showAuthenticatedApp(user) { authUser = user; - browserConnectionId = connectionIdForUser(user); + browserConnectionId = canUseTaskDataPlane() ? connectionIdForUser(user) : ''; + if (isAdministrator()) { + clearTaskDataPlaneClientState(); + if (user?.id) localStorage.removeItem(`liansyn_platform_browser_connection_id:${user.id}`); + } + if (isAdministrator() && (IS_TASK_PAGE || IS_OPERATIONS_DASHBOARD_PAGE)) { + window.location.replace('/accounts'); + return false; + } if (!isAdministrator() && IS_ADMIN_PAGE) { window.location.replace('/'); - return; + return false; } if (IS_OPERATIONS_DASHBOARD_PAGE && !canViewOperationsDashboard()) { window.location.replace('/'); - return; + return false; } const panel = $('#loginPanel'); const workbench = $('#workbench'); @@ -381,11 +454,12 @@ function showAuthenticatedApp(user) { const passwordPanel = $('#passwordChangePanel'); const authState = $('#authState'); const automationButton = $('#automationToggleButton'); + const bridgeState = $('#bridgeState'); const logoutButton = $('#logoutButton'); const changePasswordButton = $('#changePasswordButton'); if (panel) panel.hidden = true; if (passwordPanel) passwordPanel.hidden = true; - if (workbench) workbench.hidden = IS_MANAGEMENT_PAGE; + if (workbench) workbench.hidden = IS_MANAGEMENT_PAGE || !canUseTaskDataPlane(); if (channelsPage) channelsPage.hidden = !IS_CHANNELS_PAGE || !isAdministrator(); if (parserRoutingPage) parserRoutingPage.hidden = !IS_PARSER_ROUTING_PAGE || !isAdministrator(); if (accountsPage) accountsPage.hidden = !IS_ACCOUNTS_PAGE || !isAdministrator(); @@ -405,12 +479,18 @@ function showAuthenticatedApp(user) { } const operationsDashboardNav = $('#operationsDashboardNav'); if (operationsDashboardNav) operationsDashboardNav.hidden = !canViewOperationsDashboard(); + for (const selector of ['#workbenchNav', '#historyNav']) { + const nav = $(selector); + if (nav) nav.hidden = !canUseTaskDataPlane(); + } + if (bridgeState) bridgeState.hidden = !canUseTaskDataPlane(); if (automationButton) { automationButton.hidden = !isAdministrator() || IS_MANAGEMENT_PAGE; renderAutomationToggle(); } const submitButton = $('#loginForm button[type="submit"]'); if (submitButton) submitButton.disabled = false; + return true; } function showPasswordChangePanel() { @@ -1140,7 +1220,7 @@ function renderAccounts() { revoke.dataset.accountAction = 'revoke-sessions'; revoke.dataset.accountId = account.id; revoke.disabled = accountSettingsBusy; - const authorizations = el('button', 'secondary-button', account.role === 'admin' ? '全部任务' : '任务权限'); + const authorizations = el('button', 'secondary-button', account.role === 'admin' ? '不参与任务' : '任务权限'); authorizations.type = 'button'; authorizations.dataset.accountAction = 'business-authorizations'; authorizations.dataset.accountId = account.id; @@ -1280,7 +1360,7 @@ async function createAccountFromForm() { $('#accountErpAccount').value = ''; if (message) { message.textContent = result.account?.role === 'admin' - ? '管理员账号已创建;该角色固定拥有全部任务权限,初始密码不会再次显示。' + ? '管理员账号已创建;该角色仅能使用平台管理功能,不具备任何任务权限。' : '账号已创建,当前默认不能执行任何业务;请点击“任务权限”完成授权。'; } } finally { @@ -1859,7 +1939,7 @@ async function refreshCsrfToken() { } async function syncRemoteTasks() { - if (!authUser) return; + if (!canUseTaskDataPlane()) return; if (remoteSyncInProgress) { syncRequested = true; return; @@ -1907,12 +1987,12 @@ async function syncRemoteTasks() { if (currentTaskId) sessionStorage.setItem('liansyn_platform_current_task_id', currentTaskId); else sessionStorage.removeItem('liansyn_platform_current_task_id'); renderTaskCards(); - } while (syncRequested && authUser); + } while (syncRequested && canUseTaskDataPlane()); await syncRuntimeTasks(); await autoDispatchReadyTasks(); } finally { remoteSyncInProgress = false; - if (syncRequested && authUser) { + if (syncRequested && canUseTaskDataPlane()) { syncRequested = false; queueMicrotask(() => syncRemoteTasks().catch((error) => { setOutput({ status: 'sync_error', message: error.message }); @@ -1923,6 +2003,8 @@ async function syncRemoteTasks() { function startRemoteEventStream() { if (eventStream) eventStream.close(); + eventStream = null; + if (!canUseTaskDataPlane()) return; const lastEventId = runtimeTasks().filter(taskAssignedToCurrentAccount).reduce((highest, task) => ( Math.max(highest, Number(task?.last_event_id || 0)) ), 0); @@ -1980,7 +2062,7 @@ function hasPollableRuntimeTasks() { } async function syncRuntimeTasks() { - if (!authUser) return; + if (!canUseTaskDataPlane()) return; const result = await apiRequest('/api/tasks?status=active&limit=200&include_total=false&executable_by=me'); const activeTasks = (Array.isArray(result.tasks) ? result.tasks : []) .map(mergeRemoteTask) @@ -2269,12 +2351,12 @@ function taskStatusText(task) { if (status === 'parse_blocked') return '拆解失败'; if (status === 'dry_run') return '已规划,未写入 ERP'; if (status === 'operation_blocked') return '业务未接入'; - if (status === 'execution_uncertain') return '待管理员核验,请勿重复提交'; - if (status === 'reconciliation_pending') return '待管理员核验,请勿重复提交'; + if (status === 'execution_uncertain') return '待当前账号只读核验,请勿重复提交'; + if (status === 'reconciliation_pending') return '待当前账号只读核验,请勿重复提交'; if (status === 'cancelled') return '已取消'; if (status === 'post_save_recovery_required') return '需要恢复导出'; if (status === 'blocked') return '已阻断'; - if (status === 'saved_unverified') return '待管理员核验,请勿重复提交'; + if (status === 'saved_unverified') return '待当前账号只读核验,请勿重复提交'; return stage || status; } @@ -2306,9 +2388,9 @@ function taskStatusLabel(task) { waiting_extension: '等待插件', dry_run: '已规划', operation_blocked: '已阻断', - execution_uncertain: '待管理员核验', - reconciliation_pending: '待管理员核验', - saved_unverified: '待管理员核验', + execution_uncertain: '待人工核验', + reconciliation_pending: '待人工核验', + saved_unverified: '待人工核验', cancelled: '已取消', post_save_recovery_required: '待恢复' }; @@ -2332,7 +2414,7 @@ function requiresManualConfirmation(task) { } function taskAssignedToCurrentAccount(task) { - return Boolean(authUser?.id && task?.assignee?.id === authUser.id); + return Boolean(canUseTaskDataPlane() && authUser?.id && task?.assignee?.id === authUser.id); } function extensionTaskBelongsToCurrentAccount(taskId) { @@ -2676,7 +2758,7 @@ function isAutomaticTask(task) { } async function autoDispatchReadyTasks({ force = false } = {}) { - if (!authUser || autoHandoffInProgress) return; + if (!canUseTaskDataPlane() || autoHandoffInProgress) return; autoHandoffInProgress = true; try { const result = await apiRequest('/api/tasks?status=confirmed&limit=200&include_total=false&executable_by=me'); @@ -3389,7 +3471,7 @@ function taskStageSnapshot(task) { business = { mode: 'needs', current: 'processing', - currentLabel: '待管理员核验' + currentLabel: '待人工核验' }; } else if (['blocked', 'failed', 'cancelled', 'extension_error', 'batch_fallback_incomplete', 'live_submit_blocked', 'preflight_blocked'].includes(status)) { business = { @@ -3517,7 +3599,7 @@ function renderTaskStages(task) { function parserEngineLabel(parser = {}) { if (parser.fallback_reason && parser.fallback_reason !== 'manual_ai_reparse') return 'AI 兜底'; - if (parser.fallback_reason === 'manual_ai_reparse') return '管理员 AI 重解析'; + if (parser.fallback_reason === 'manual_ai_reparse') return '人工 AI 重解析'; return parser.authoritative_engine === 'program' ? '程序解析' : parser.authoritative_engine === 'ai' ? 'AI 解析' : '尚未确定'; } @@ -3689,7 +3771,7 @@ function renderTaskInputAudit(task) { const item = el('article', 'detail-block'); const sourceLabel = message.source === 'manual' ? '人工输入' : message.source === 'agentbus' ? 'AgentBus' - : message.source === 'reparse' ? '管理员重解析' : '历史/系统'; + : message.source === 'reparse' ? '人工重解析' : '历史/系统'; item.append(el( 'p', 'muted', @@ -4622,6 +4704,9 @@ function makeRequestId() { } function sendToExtension(type, payload = {}, timeoutMs = 2500) { + if (!canUseTaskDataPlane()) { + return Promise.reject(new Error('管理员账号不能连接或调用任务执行插件。')); + } const requestId = makeRequestId(); const promise = new Promise((resolve, reject) => { const timer = setTimeout(() => { @@ -4643,6 +4728,7 @@ window.addEventListener('message', (event) => { if (event.source !== window) return; const message = event.data || {}; if (message.source !== EXTENSION_SOURCE) return; + if (!canUseTaskDataPlane()) return; if (message.type === 'BRIDGE_READY') { applyBridgePayload(message.payload || { ok: true }); return; @@ -4999,6 +5085,7 @@ async function confirmAndSubmitToErpPlugin(task) { } async function pingBridge() { + if (!canUseTaskDataPlane()) return false; try { const result = await sendToExtension('PING', { expected_erp_account: authUser?.erp_account || '' @@ -5067,6 +5154,10 @@ async function createTask() { showLoginPanel('请先登录。'); return; } + if (!canUseTaskDataPlane()) { + window.location.replace('/accounts'); + return; + } if (taskCreateInProgress) return; const rawText = $('#rawInstruction').value; const fileInput = $('#rosterAttachment'); @@ -5117,12 +5208,13 @@ async function refreshBackgroundState() { if (!authUser || backgroundRefreshInProgress) return; backgroundRefreshInProgress = true; try { - const operations = [pingAi(), pingBridge()]; + const operations = [pingAi()]; + if (canUseTaskDataPlane()) operations.push(pingBridge()); if (isAdministrator()) operations.push(syncAutomationSettings({ background: true })); if (IS_CHANNELS_PAGE && isAdministrator()) { operations.push(syncChannels(), syncLeaderSummarySubscriptions()); } - if (IS_TASK_PAGE) operations.push(syncRemoteTasks()); + if (IS_TASK_PAGE && canUseTaskDataPlane()) operations.push(syncRemoteTasks()); await Promise.allSettled(operations); } finally { backgroundRefreshInProgress = false; @@ -5160,7 +5252,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 中已经发生的操作;系统只会向任务所属账号的在线插件发送停止与清理指令。已经交给 AgentBus 并到达组长微信的任务摘要也无法撤回。` ); } @@ -5540,7 +5632,7 @@ async function resumePrewriteTaskExecution(task) { } async function pollAllTaskResults() { - if (!bridgeConnected || pollInProgress) return; + if (!canUseTaskDataPlane() || !bridgeConnected || pollInProgress) return; pollInProgress = true; try { const activeTasks = runtimeTasks().filter((task) => taskAssignedToCurrentAccount(task) && isTaskPollable(task)); @@ -5567,6 +5659,7 @@ function stopPolling() { function startPolling() { stopPolling(); + if (!canUseTaskDataPlane()) return; pollAllTaskResults().catch((error) => { setTaskState('轮询失败'); setOutput({ status: 'poll_error', message: error.message }); @@ -5616,7 +5709,7 @@ async function initializeSession() { return false; } - showAuthenticatedApp(me.user); + if (!showAuthenticatedApp(me.user)) return false; if (isAdministrator()) { try { await syncAutomationSettings(); @@ -5656,7 +5749,7 @@ async function initializeSession() { if (message) message.textContent = operationsDashboardSafeError(error); }); } - if (IS_TASK_PAGE) { + if (IS_TASK_PAGE && canUseTaskDataPlane()) { try { await syncRemoteTasks(); } catch (error) { @@ -5720,7 +5813,7 @@ document.addEventListener('DOMContentLoaded', async () => { const result = await response.json(); if (!response.ok) throw new Error(result.message || '登录失败。'); csrfToken = result.csrf_token || ''; - showAuthenticatedApp(result.user); + if (!showAuthenticatedApp(result.user)) return; $('#loginPassword').value = ''; if (isAdministrator()) { await syncAutomationSettings().catch(() => setTaskState('全自动化设置读取失败')); @@ -5735,12 +5828,12 @@ document.addEventListener('DOMContentLoaded', async () => { } if (IS_ACCOUNTS_PAGE && isAdministrator()) await syncAccounts(); if (IS_AUDIT_PAGE && isAdministrator()) await syncAuditEvents(); - if (IS_TASK_PAGE) { + if (IS_TASK_PAGE && canUseTaskDataPlane()) { await syncRemoteTasks(); startRemoteEventStream(); } await pingAi(); - await pingBridge(); + if (canUseTaskDataPlane()) await pingBridge(); } catch (error) { if (errorNode) errorNode.textContent = error.message || String(error); } finally { @@ -6227,9 +6320,9 @@ document.addEventListener('DOMContentLoaded', async () => { void copyTaskLifecycle(button); }); if (await initializeSession()) { - if (IS_TASK_PAGE) renderTaskCards(); + if (IS_TASK_PAGE && canUseTaskDataPlane()) renderTaskCards(); pingAi().catch(() => {}); - pingBridge().catch(() => {}); + if (canUseTaskDataPlane()) pingBridge().catch(() => {}); // The control plane caches this health probe; keep the browser refresh // interval conservative so status checks cannot pressure the Agent API. setInterval(() => { @@ -6246,7 +6339,7 @@ document.addEventListener('DOMContentLoaded', async () => { window.addEventListener('focus', () => { refreshBackgroundState().catch(() => {}); }); - if (IS_TASK_PAGE && currentTaskId) renderTaskDetail(); - if (IS_TASK_PAGE && hasPollableRuntimeTasks()) startPolling(); + if (IS_TASK_PAGE && canUseTaskDataPlane() && currentTaskId) renderTaskDetail(); + if (IS_TASK_PAGE && canUseTaskDataPlane() && hasPollableRuntimeTasks()) startPolling(); } }); diff --git a/LianSyn-platform/index.html b/LianSyn-platform/index.html index 530d569..dfd3ac3 100644 --- a/LianSyn-platform/index.html +++ b/LianSyn-platform/index.html @@ -5,7 +5,7 @@ AI操作台 · LianSyn-platform - +
@@ -14,8 +14,8 @@
- + diff --git a/control-plane/README.md b/control-plane/README.md index 1ad874b..7464791 100644 --- a/control-plane/README.md +++ b/control-plane/README.md @@ -7,9 +7,9 @@ - PostgreSQL 是任务、事件、幂等、审计和回查状态的唯一事实源。 - 平台账号使用 `admin`(管理员)、`team_lead`(组长)和 `user`(普通用户)三种固定角色登录;服务端会话使用 HttpOnly/Secure/SameSite Cookie。管理员维护账号、角色、状态、密码重置、会话撤销和可执行任务类型;平台不提供组织或租户选择。 - 登录会话默认跨浏览器重启持续有效,不按空闲时间或绝对时长自动失效;用户退出、修改密码、账号停用或管理员强制撤销时由服务端立即撤销。 -- 业务页面通过 REST 创建任务;Agent 返回结构化结果后,任务归属人在自己的任务上一次点击“确认并提交到 ERP 插件”,再通过 SSE 或轮询读取服务端状态。普通用户与组长的任务 API、SSE、附件和插件执行权限覆盖分配给本人的人工任务与 AgentBus 任务;管理员可查看固定部署范围内的全部任务,但确认、插件领取和执行回执仍必须来自任务归属账号,查看权限不等于执行权限。 +- 业务页面通过 REST 创建任务;Agent 返回结构化结果后,任务归属人在自己的任务上一次点击“确认并提交到 ERP 插件”,再通过 SSE 或轮询读取服务端状态。普通用户与组长的任务 API、SSE、附件和插件执行权限只覆盖分配给本人的人工任务与 AgentBus 任务;管理员是纯管理面身份,不能创建、查看、修改、领取、接收事件或提交任何业务任务。 - 手工与 AgentBus 的每个新任务都会按同一组织、同一 18 项业务路由固化解析策略及配置 revision;来源不能覆盖模式。其中两项名单 route 固定为 `program_only`,其余 16 项可配置 `ai / shadow / auto / program`。全消息中的唯一已登记指令可以确定 route;没有指令时,只有全部标签都属于唯一 route、至少两个不同标签且业务定位必填项完整的字段签名才可确定 route。未知、冲突或多 route 输入不猜测。后续补充轮次沿用原任务快照,设置变化只影响新任务和新会话。 -- 管理员天然拥有全部 18 类人工业务。组长和普通用户使用逐账号白名单,新账号默认没有任何可执行任务类型;管理员在 `/accounts` 逐项授权后才能提交对应的新任务、补充指令或名单附件。已登记但未授权的业务返回 `business_not_authorized`;无法唯一确认 route 的非管理员输入返回 `business_type_unresolved`。两种拒绝都发生在解析器和 ERP 插件之前,并记录不含明文的授权拒绝审计。权限在补充输入、人工确认、全自动确认和插件领取前再次校验;取消授权后的任务不会进入 ERP 队列。AgentBus 入站使用渠道绑定员工的同一白名单,管理员账号不得成为员工渠道归属人。 +- 组长和普通用户使用逐账号白名单,新账号默认没有任何可执行任务类型;管理员自身的任务白名单固定为空,只能在 `/accounts` 为员工逐项授权。已登记但未授权的业务返回 `business_not_authorized`;无法唯一确认 route 的员工输入返回 `business_type_unresolved`。两种拒绝都发生在解析器和 ERP 插件之前,并记录不含明文的授权拒绝审计。权限在补充输入、人工确认、全自动确认和插件领取前再次校验;取消授权后的任务不会进入 ERP 队列。AgentBus 入站使用渠道绑定员工的同一白名单,管理员账号不得成为员工渠道归属人。 - 名单 route 创建后先进入 `awaiting_attachment`,只接收一份 `.xls/.xlsx`;控制面在前 100 行中自动定位唯一的 ERP 名单字段表头,按精确字段语义从任意表头行和列顺序中只选择 12 个必需源字段,并把表头下方连续名单数据在内存中规范化为 13 列 canonical TSV。身份证、年龄、源证件类型和其他普通额外列不参与名单行识别、字段校验或输出;必需字段缺失/重复、多个候选表头、必需字段中的非法公式及宏、外链等主动内容仍失败关闭。原始工作簿不入库,文件名和 canonical TSV 使用字段加密,解析通过或终止后清除 canonical 中间文本。附件到齐前不会领取解析任务,也不会进入 ERP。独立团初始 16 行、散拼子单初始 31 行仅为 ERP 动态扩行基线,5000 为技术上限。 - 平台的任务 ID、Agent 会话 ID、确认状态、重要摘要、事件和用户通讯内容只存在于平台任务/会话/结果信封中,不回写到 Agent `operation`,也不成为 ERP 业务字段。 - Agent `operation` 只保存解析态业务事实。插件领取已确认任务后先做无需 ERP 查询的前门禁,再在 ERP 内只读唯一解析对象、资源和当前状态,最后对内部 execution operation 执行严格写前门禁。 @@ -25,11 +25,11 @@ - 微信侧的 `[WeChat attachment: 文件名]` 只是一段传输占位文字,不代表控制面已经收到文件。若同一帧没有符合契约的 `payload.attachments[]`,listener 会在进入任务服务前失败关闭、保留原名单任务的等待状态,并返回“附件内容未传到平台”;不会把占位文字创建成新业务任务。附件元数据、HTTPS URL、DNS、大小或摘要校验失败时返回对应的安全摘要,仍不回显 URL、文件字节或名单内容。当前生产部署位于受信内网,入站附件 URL 可以使用内网域名、私网 IPv4/IPv6 或 localhost;因此 AgentBus 渠道和上游桥接器必须被视为受信输入边界。 - AgentBus 入站消息会复用 `TaskService` 的任务/会话/解析队列,并以渠道绑定员工写入 `created_by` 与不可变的 `assigned_user_id`,解析完成后通过同一 WebSocket 返回一次 `task.result`。组织级“全自动化”关闭时,手工与 AgentBus 新任务都需要人工确认;开启后,两种来源的合法解析结果都自动进入 ERP 队列,不再按来源或创建、名单、安排、修改、取消/恢复、导出等业务类型保留人工例外。历史未归属 AgentBus 任务不会自动执行。操作台在 EventSource 建连/重连、30 秒后台刷新以及页面重新可见或聚焦时重新读取数据库权威开关。缺资料、解析失败、歧义、插件校验失败或 ERP 回查不确定时仍会停止,不会绕过校验或重试不确定写入。 - 单一部署范围可以维护多个“用户渠道”。每个渠道代表一个外部 AgentBus 用户身份,并且必须一对一绑定一个有效的非管理员平台账号;一个平台账号也只能绑定一个渠道。管理员在 `/channels` 创建、绑定、停用、启用、轮换或删除渠道。只有绑定账号有效且已配置 ERP 账号的启用渠道才启动 listener;未绑定渠道失败关闭。删除会停止对应 listener、移除服务端保存的 key 和该渠道尚存的持久化回执;历史任务本体保留,其 `channel_id` 置空而 `assigned_user_id` 不变。每个渠道独立保存加密后的 AgentBus key,同一 key 不能被多个渠道复用;列表和日志都不会回显 key。`AGENTBUS_WS_URL`、重连策略和客户端类型仍是全局连接配置,`AGENTBUS_BOT_ADDRESS` 可作为渠道 bot address 的默认值。 -- `/history` 同时提供可恢复的归档/恢复与显式的永久强制删除。`POST /api/tasks/:taskId/archive`、`POST /api/tasks/:taskId/restore` 和 `POST /api/tasks/bulk-archive` 保留归档语义及运行状态门禁;`DELETE /api/tasks/:taskId` 与 `POST /api/tasks/bulk-delete` 会绕过任务状态门禁并物理删除任务及其输入、事件、尝试、会话、投递和附件记录,同时清理任务 outbox,并在提交后尽力清理 OSS 对象。普通用户和组长只能操作本人任务,管理员可以处理全部授权任务;不可逆删除仍保留最小化的删除审计事件。 -- `/operations-dashboard` 是组长和管理员专用的只读业务操作看板。它以不可变的任务归属账号作为员工口径,纳入同一固定组织中所有已归属的人工与 AgentBus 任务,并支持从结果状态、员工、上海业务日期和业务类型逐层穿透;历史上无法安全归属员工的 AgentBus 任务不在人员看板中被猜测归属。列表和详情只回答“谁负责什么任务、收到什么指令、完成了什么结果”:详情返回员工、业务类型、完整指令轮次、输入附件名称/行数和可读业务结果,不返回任务生命周期、解析/执行 JSON、技术阶段、错误码或产物地址。关键词查询先受日期、人员、业务和状态约束,单次解密匹配候选最多 2,000 条,超过时要求继续缩小范围。该路径不授予他人任务修改、ERP 执行、SSE、产物下载、账号维护或全局安全审计权限。 +- `/history` 同时提供可恢复的归档/恢复与显式的永久强制删除。`POST /api/tasks/:taskId/archive`、`POST /api/tasks/:taskId/restore` 和 `POST /api/tasks/bulk-archive` 保留归档语义及运行状态门禁;`DELETE /api/tasks/:taskId` 与 `POST /api/tasks/bulk-delete` 会绕过任务状态门禁并物理删除任务及其输入、事件、尝试、会话、投递和附件记录,同时清理任务 outbox,并在提交后尽力清理 OSS 对象。普通用户和组长只能操作本人任务;管理员不能进入该页面或调用这些接口。不可逆删除仍保留最小化的删除审计事件。 +- `/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` 对应的登录账号,管理员不会因为能查看员工任务而收到或处理该员工的插件任务。已开始写入但结果不确定的任务只阻塞同一账户的后续领取,直到人工回查收敛或任务被明确强制删除。 +- ERP 插件领取按任务 `assigned_user_id` 使用账户级数据库锁和 FIFO confirmed 队列:同一平台/ERP 账户在任意时刻最多一个 ERP execution,该账户的其他任务留在服务端等待;不同账户的活跃或待执行任务互不占用队列位置、可独立领取执行。管理员不能成为 `assigned_user_id`、渠道 owner、浏览器 worker 或组长任务摘要接收人,也不能打开任务 SSE;实时执行事件、任务摘要、插件领取、执行回执以及强制删除后的浏览器清理命令只发送或接受符合角色约束的员工账号。已开始写入但结果不确定的任务只阻塞同一账户的后续领取,直到归属账号人工回查收敛或任务被明确强制删除。 ## 任务级会话续接 @@ -77,19 +77,19 @@ Content-Type: application/json - `auto`:程序优先;仅 `program_unsupported_syntax`、`program_unknown_field`、`program_internal_error`、`program_timeout`、`program_contract_invalid` 可整体切到 AI。缺字段、非法值、冲突、歧义、多动作和业务规则阻断不会调用 AI。 - `program`:只运行程序;失败时追问或阻断,不调用 AI。 -Auto 一旦发生 AI fallback,任务会永久绑定原 AI 会话。每次解析只有一个完整权威结果,不拼接候选;任务进入确认、插件领取或 ERP 链路后不再自动重解析。尚未进入插件且明确 `no_plugin_dispatch=true / no_erp_write=true` 的解析失败任务,可由管理员授权一次 AI 重解析。 +Auto 一旦发生 AI fallback,任务会永久绑定原 AI 会话。每次解析只有一个完整权威结果,不拼接候选;任务进入确认、插件领取或 ERP 链路后不再自动重解析。管理员任务隔离后不再从任务详情执行一次性 AI 重解析或解析差异判定。 操作台根路径 `/parser-routing` 提供逐业务观察统计、即时切换检查、revision 乐观并发控制和事务级“全部切回 AI”。连续天数、任务数量和 fallback 比率不作为等待门槛:18 项已实现业务均可独立进入 Shadow,多个 Shadow 可并行测试且不受固定前序约束;已出现的差异全部判定且本业务程序关键、契约、运行错误为零即可进入 Auto;程序下游失败和上述错误为零即可进入 Program。同一时间最多一个业务处于 Auto,晋级仍只能 `AI → Shadow → Auto → Program` 逐级进行,回退可直接执行。7/30 日统计仅供观察。组织已开启全自动化时,切入 Auto/Program 不增加人工确认期。 -主要接口: +主要管理接口: - `GET /api/settings/parser-routing` - `PUT /api/settings/parser-routing/:routeId` - `POST /api/settings/parser-routing/emergency-ai` -- `POST /api/tasks/:taskId/reparse` -- `PUT /api/parser-decisions/:decisionId/review` -迁移 `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,是否启用由运行时角色与渠道核对自动维护。原文、完整程序/AI 候选、名单 canonical 中间文本、人工说明和待发组长摘要使用字段加密保存;统计、全局审计和运行日志不复制明文业务输入。 +旧的任务详情重解析和解析差异判定路由仍保留兼容响应,但已纳入任务数据面统一门禁;管理员调用会返回 `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 中间文本、人工说明和待发组长摘要使用字段加密保存;统计、全局审计和运行日志不复制明文业务输入。 ## AgentBus Bot 接入 @@ -143,7 +143,7 @@ npm run data:retention npm run dev ``` -`npm run dev` 和 `npm start` 会先执行数据库迁移,再启动控制平面;直接运行 `control-plane/src/server.ts` 或构建后的 `server.js` 时,服务也会在启动前检查必需迁移 `020_leader_task_summary_notifications`,缺失时拒绝监听端口。`db:migrate` 和管理员初始化需要可连接的 PostgreSQL。开发机没有数据库时,可以运行 `npm run test:control-plane` 完成无数据库静态/健康烟测。 +`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` 完成无数据库静态/健康烟测。 `/health/ready` 同时检查 PostgreSQL 可用性和必需 schema 版本;迁移未完成时返回 503,并标明 `required_migration`,避免任务在数据库结构未升级时进入解析队列。 diff --git a/control-plane/migrations/021_admin_task_data_plane_isolation.sql b/control-plane/migrations/021_admin_task_data_plane_isolation.sql new file mode 100644 index 0000000..1a23a3a --- /dev/null +++ b/control-plane/migrations/021_admin_task_data_plane_isolation.sql @@ -0,0 +1,213 @@ +-- Administrators are management-plane identities only. Remove any legacy +-- execution bindings and reject every future database write that would make +-- an administrator a task creator, assignee, channel owner, browser worker, +-- or business-route grantee. + +UPDATE tasks task + SET assigned_user_id = NULL, + lease_owner = NULL, + lease_expires_at = NULL, + updated_at = now() + FROM users account + WHERE account.organization_id = task.organization_id + AND account.id = task.assigned_user_id + AND account.role = 'admin'; + +UPDATE user_channels channel + SET owner_user_id = NULL, + enabled = false, + status = 'disabled', + last_error = '管理员账号已与任务执行面隔离,渠道绑定已解除。', + updated_at = now() + FROM users account + WHERE account.organization_id = channel.organization_id + AND account.id = channel.owner_user_id + AND account.role = 'admin'; + +UPDATE browser_connections connection + SET status = 'superseded', + erp_account_verified = false + FROM users account + WHERE account.organization_id = connection.organization_id + AND account.id = connection.user_id + AND account.role = 'admin' + AND connection.status = 'connected'; + +DELETE FROM user_business_route_authorizations route_grant + USING users account + WHERE account.organization_id = route_grant.organization_id + AND account.id = route_grant.user_id + AND account.role = 'admin'; + +UPDATE leader_task_summary_subscriptions subscription + SET enabled = false, + target_verified_at = NULL, + revision = revision + 1, + updated_at = now() + FROM users account + WHERE account.organization_id = subscription.organization_id + AND account.id = subscription.leader_user_id + AND account.role = 'admin'; + +UPDATE leader_task_summary_deliveries delivery + SET delivery_status = 'cancelled', + last_error = '管理员账号已与任务数据面隔离,组长任务摘要投递已取消。', + updated_at = now() + FROM users account + WHERE account.organization_id = delivery.organization_id + AND account.id = delivery.leader_user_id + AND account.role = 'admin' + AND delivery.delivery_status IN ('pending', 'sending', 'failed'); + +CREATE OR REPLACE FUNCTION reject_admin_task_principal() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +DECLARE + principal_id uuid; + principal_role text; +BEGIN + principal_id := NULLIF(to_jsonb(NEW) ->> TG_ARGV[0], '')::uuid; + IF principal_id IS NULL THEN + RETURN NEW; + END IF; + + SELECT role + INTO principal_role + FROM users + WHERE organization_id = NEW.organization_id + AND id = principal_id + FOR SHARE; + + IF principal_role = 'admin' THEN + RAISE EXCEPTION USING + ERRCODE = '23514', + CONSTRAINT = TG_ARGV[1], + MESSAGE = 'administrator identities cannot enter the task data plane'; + END IF; + RETURN NEW; +END; +$$; + +DROP TRIGGER IF EXISTS tasks_admin_assignee_forbidden_trigger ON tasks; +CREATE TRIGGER tasks_admin_assignee_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, assigned_user_id ON tasks +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'assigned_user_id', + 'tasks_admin_assignee_forbidden' +); + +DROP TRIGGER IF EXISTS tasks_admin_creator_forbidden_trigger ON tasks; +CREATE TRIGGER tasks_admin_creator_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, created_by ON tasks +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'created_by', + 'tasks_admin_creator_forbidden' +); + +DROP TRIGGER IF EXISTS user_channels_admin_owner_forbidden_trigger ON user_channels; +CREATE TRIGGER user_channels_admin_owner_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, owner_user_id ON user_channels +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'owner_user_id', + 'user_channels_admin_owner_forbidden' +); + +DROP TRIGGER IF EXISTS browser_connections_admin_worker_forbidden_trigger ON browser_connections; +CREATE TRIGGER browser_connections_admin_worker_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, user_id ON browser_connections +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'user_id', + 'browser_connections_admin_worker_forbidden' +); + +DROP TRIGGER IF EXISTS user_business_routes_admin_grantee_forbidden_trigger + ON user_business_route_authorizations; +CREATE TRIGGER user_business_routes_admin_grantee_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, user_id ON user_business_route_authorizations +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'user_id', + 'user_business_routes_admin_grantee_forbidden' +); + +DROP TRIGGER IF EXISTS leader_task_summary_admin_subscriber_forbidden_trigger + ON leader_task_summary_subscriptions; +CREATE TRIGGER leader_task_summary_admin_subscriber_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, leader_user_id ON leader_task_summary_subscriptions +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'leader_user_id', + 'leader_task_summary_admin_subscriber_forbidden' +); + +DROP TRIGGER IF EXISTS leader_task_summary_admin_recipient_forbidden_trigger + ON leader_task_summary_deliveries; +CREATE TRIGGER leader_task_summary_admin_recipient_forbidden_trigger +BEFORE INSERT OR UPDATE OF organization_id, leader_user_id ON leader_task_summary_deliveries +FOR EACH ROW EXECUTE FUNCTION reject_admin_task_principal( + 'leader_user_id', + 'leader_task_summary_admin_recipient_forbidden' +); + +CREATE OR REPLACE FUNCTION isolate_administrator_from_task_runtime() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF NEW.role <> 'admin' THEN + RETURN NEW; + END IF; + + UPDATE tasks + SET assigned_user_id = NULL, + lease_owner = NULL, + lease_expires_at = NULL, + updated_at = now() + WHERE organization_id = NEW.organization_id + AND assigned_user_id = NEW.id; + + UPDATE user_channels + SET owner_user_id = NULL, + enabled = false, + status = 'disabled', + last_error = '绑定账号已变更为管理员,渠道已解除绑定。', + updated_at = now() + WHERE organization_id = NEW.organization_id + AND owner_user_id = NEW.id; + + UPDATE browser_connections + SET status = 'superseded', + erp_account_verified = false + WHERE organization_id = NEW.organization_id + AND user_id = NEW.id + AND status = 'connected'; + + DELETE FROM user_business_route_authorizations + WHERE organization_id = NEW.organization_id + AND user_id = NEW.id; + + UPDATE leader_task_summary_subscriptions + SET enabled = false, + target_verified_at = NULL, + revision = revision + 1, + updated_at = now() + WHERE organization_id = NEW.organization_id + AND leader_user_id = NEW.id; + + UPDATE leader_task_summary_deliveries + SET delivery_status = 'cancelled', + last_error = '账号已变更为管理员,组长任务摘要投递已取消。', + updated_at = now() + WHERE organization_id = NEW.organization_id + AND leader_user_id = NEW.id + AND delivery_status IN ('pending', 'sending', 'failed'); + + RETURN NEW; +END; +$$; + +DROP TRIGGER IF EXISTS users_admin_task_runtime_isolation_trigger ON users; +CREATE TRIGGER users_admin_task_runtime_isolation_trigger +AFTER UPDATE OF role ON users +FOR EACH ROW +WHEN (NEW.role = 'admin' AND OLD.role IS DISTINCT FROM NEW.role) +EXECUTE FUNCTION isolate_administrator_from_task_runtime(); diff --git a/control-plane/src/agentbus.ts b/control-plane/src/agentbus.ts index ec82cd2..484584a 100644 --- a/control-plane/src/agentbus.ts +++ b/control-plane/src/agentbus.ts @@ -575,7 +575,7 @@ function conciseErrorText(value: unknown): string { if (!normalized) return '任务处理失败,请检查任务信息后重试。'; const firstMessage = normalized.split(/[;;]/u, 1)[0]?.trim() || normalized; if (/ERP 写入结果不确定|reconciliation_pending|saved_unverified|execution_uncertain/iu.test(normalized)) { - return 'ERP 写入已发起,但系统尚未确认最终结果。请勿重复提交同一任务,等待管理员只读核验。'; + return 'ERP 写入已发起,但系统尚未确认最终结果。请勿重复提交同一任务,等待任务所属账号只读核验。'; } if (/required fields still blank|必填字段/iu.test(firstMessage)) { return '任务缺少必要信息,未写入 ERP,请补充后重试。'; diff --git a/control-plane/src/auth.ts b/control-plane/src/auth.ts index 7123182..b30ab76 100644 --- a/control-plane/src/auth.ts +++ b/control-plane/src/auth.ts @@ -148,7 +148,7 @@ function mapAccount(row: Record): PublicAccount { role, erp_account: normalizeErpAccount(row.erp_account), is_active: row.is_active === true || String(row.is_active) === 'true', - authorized_business_route_ids: role === 'admin' ? [...ALL_BUSINESS_ROUTE_IDS] : storedRouteIds, + authorized_business_route_ids: role === 'admin' ? [] : storedRouteIds, business_authorization_revision: Math.max(0, Number(row.business_authorization_revision || 0)), last_login_at: isoOrNull(row.last_login_at), created_at: isoOrNull(row.created_at) || new Date(0).toISOString(), @@ -569,7 +569,7 @@ export class AuthService { if (!target.rowCount) throw new AuthError('account_not_found', '账号不存在。', 404); const row = target.rows[0] as Record; if (normalizeRole(row.role) === 'admin') { - throw new AuthError('admin_business_authorization_fixed', '管理员固定拥有全部业务权限,无需单独授权。', 409); + throw new AuthError('admin_business_authorization_fixed', '管理员不参与业务任务,任务权限固定为空。', 409); } const currentRevision = Math.max(0, Number(row.business_authorization_revision || 0)); if (currentRevision !== expectedRevision) { diff --git a/control-plane/src/db.ts b/control-plane/src/db.ts index da569f6..8d85eb4 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 = '020_leader_task_summary_notifications'; +export const REQUIRED_SCHEMA_VERSION = '021_admin_task_data_plane_isolation'; export interface DatabaseReadiness { ready: boolean; diff --git a/control-plane/src/server.ts b/control-plane/src/server.ts index a599256..75df6e9 100644 --- a/control-plane/src/server.ts +++ b/control-plane/src/server.ts @@ -22,6 +22,7 @@ import { LeaderNotificationService } from './leader-notification-service.js'; import { TaskError, TaskService, + canUseTaskDataPlane, canViewOperationsDashboard, type ParseDecisionInput, type ParseTaskClaim, @@ -336,6 +337,18 @@ function setAuthNoStore(reply: FastifyReply): void { reply.header('Pragma', 'no-cache'); } +export function isTaskDataPlaneRoute(routePath: unknown): boolean { + const path = String(routePath || ''); + return path === '/api/messages' + || path === '/api/connections/heartbeat' + || path === '/api/events' + || path === '/api/tasks' + || path.startsWith('/api/tasks/') + || path.startsWith('/api/parser-decisions/') + || path === '/api/operations-dashboard' + || path.startsWith('/api/operations-dashboard/'); +} + function attachmentContentDisposition(fileName: string): string { const safeName = String(fileName || 'team-file.bin') .replace(/[\\"\r\n\u0000-\u001f\u007f]/g, '_') @@ -545,7 +558,14 @@ export async function buildServer({ const requireLeadership = (session: ActiveSession): ActiveSession => { if (!canViewOperationsDashboard(session.user.role)) { - throw new AuthError('leadership_required', '需要组长或管理员权限。', 403); + throw new AuthError('leadership_required', '需要组长权限。', 403); + } + return session; + }; + + const requireTaskDataPlane = (session: ActiveSession): ActiveSession => { + if (!canUseTaskDataPlane(session.user.role)) { + throw new AuthError('task_access_forbidden', '管理员账号仅用于平台管理,不能访问业务任务。', 403); } return session; }; @@ -570,6 +590,14 @@ export async function buildServer({ const requireMutationSession = requireAuthenticatedMutationSession; + const requireTaskSession = async (request: FastifyRequest): Promise => ( + requireTaskDataPlane(await getSession(request)) + ); + + const requireTaskMutationSession = async (request: FastifyRequest): Promise => ( + requireTaskDataPlane(await requireMutationSession(request)) + ); + const requireAdminSession = async (request: FastifyRequest): Promise => ( requireAdmin(await getSession(request)) ); @@ -582,6 +610,11 @@ export async function buildServer({ requireLeadership(await getSession(request)) ); + app.addHook('preHandler', async (request) => { + if (!isTaskDataPlaneRoute(request.routeOptions.url)) return; + requireTaskDataPlane(await getSession(request)); + }); + async function persistParseOutcome( claim: ParseTaskClaim, result: unknown, @@ -1090,7 +1123,7 @@ export async function buildServer({ }); app.post('/api/tasks/:taskId/reparse', async (request) => { - const session = await requireAdminMutationSession(request); + const session = requireTaskDataPlane(await requireAdminMutationSession(request)); const taskId = String((request.params as { taskId?: string }).taskId || ''); const body = parserReparseSchema.parse(request.body); const task = await tasks.reparseTaskWithAi(contextFor(session, request), taskId, body.reason); @@ -1102,7 +1135,7 @@ export async function buildServer({ }); app.put('/api/parser-decisions/:decisionId/review', async (request) => { - const session = await requireAdminMutationSession(request); + const session = requireTaskDataPlane(await requireAdminMutationSession(request)); const decisionId = String((request.params as { decisionId?: string }).decisionId || ''); const body = parserDecisionReviewSchema.parse(request.body); return { @@ -1192,7 +1225,7 @@ export async function buildServer({ }); app.get('/api/tasks', async (request) => { - const session = await getSession(request); + const session = await requireTaskSession(request); const query = listTasksQuerySchema.parse(request.query || {}); const page = await tasks.listTasksPage(session.user.organizationId, { status: query.status || undefined, @@ -1217,7 +1250,7 @@ export async function buildServer({ }); app.post('/api/tasks', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = createTaskSchema.parse(request.body); const task = await tasks.createTask( contextFor(session, request), @@ -1231,7 +1264,7 @@ export async function buildServer({ }); app.post('/api/messages', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = taskMessageSchema.parse(request.body); const result = await tasks.ingestMessage(contextFor(session, request), { message: body.message, @@ -1245,19 +1278,19 @@ export async function buildServer({ }); app.get('/api/tasks/:taskId', async (request) => { - const session = await getSession(request); + const session = await requireTaskSession(request); const params = request.params as { taskId: string }; return { ok: true, task: await tasks.getTask(session.user.organizationId, params.taskId, contextFor(session, request)) }; }); app.get('/api/tasks/:taskId/input-history', async (request) => { - const session = await getSession(request); + const session = await requireTaskSession(request); const params = request.params as { taskId: string }; return { ok: true, ...(await tasks.getTaskInputHistory(contextFor(session, request), params.taskId)) }; }); app.get('/api/tasks/:taskId/artifacts/:artifactId', async (request, reply) => { - const session = await getSession(request); + const session = await requireTaskSession(request); const params = request.params as { taskId: string; artifactId: string }; const artifactId = z.string().uuid().safeParse(params.artifactId); if (!artifactId.success) throw new TaskError('artifact_not_found', '附件不存在或无权访问。', 404); @@ -1280,13 +1313,13 @@ export async function buildServer({ }); app.post('/api/tasks/:taskId/confirm', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const params = request.params as { taskId: string }; return { ok: true, task: await tasks.confirmTask(contextFor(session, request), params.taskId) }; }); app.post('/api/tasks/:taskId/claim', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = claimSchema.parse(request.body); const params = request.params as { taskId: string }; const claim = await tasks.claimForBrowser(contextFor(session, request), params.taskId, body.connection_id); @@ -1302,7 +1335,7 @@ export async function buildServer({ }); app.post('/api/tasks/:taskId/result', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = resultSchema.parse(request.body); const params = request.params as { taskId: string }; return { @@ -1318,13 +1351,13 @@ export async function buildServer({ }); app.post('/api/tasks/:taskId/cancel', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const params = request.params as { taskId: string }; return { ok: true, task: await tasks.cancelTask(contextFor(session, request), params.taskId) }; }); app.post('/api/tasks/bulk-delete', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = taskBulkDeleteSchema.parse(request.body); return { ok: true, @@ -1334,7 +1367,7 @@ export async function buildServer({ }); app.post('/api/tasks/bulk-archive', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = taskBulkArchiveSchema.parse(request.body); return { ok: true, @@ -1344,7 +1377,7 @@ export async function buildServer({ }); app.delete('/api/tasks/:taskId', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const params = request.params as { taskId: string }; return { ok: true, @@ -1353,7 +1386,7 @@ export async function buildServer({ }); app.post('/api/tasks/:taskId/archive', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const params = request.params as { taskId: string }; const body = taskArchiveSchema.parse(request.body || {}); return { @@ -1364,13 +1397,13 @@ export async function buildServer({ }); app.post('/api/tasks/:taskId/restore', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const params = request.params as { taskId: string }; return { ok: true, restored: true, task: await tasks.restoreTask(contextFor(session, request), params.taskId) }; }); app.post('/api/connections/heartbeat', async (request) => { - const session = await requireMutationSession(request); + const session = await requireTaskMutationSession(request); const body = heartbeatSchema.parse(request.body); const worker = await tasks.heartbeat( contextFor(session, request), @@ -1445,7 +1478,7 @@ export async function buildServer({ }); app.get('/api/events', async (request, reply) => { - const session = await getSession(request); + const session = await requireTaskSession(request); const query = taskEventsQuerySchema.parse(request.query || {}); const querySince = query.since; const reconnectSince = Number(request.headers['last-event-id'] || 0); @@ -1463,9 +1496,8 @@ export async function buildServer({ }); const send = (event: TaskEvent) => { if (event.organization_id !== session.user.organizationId) return; - // This is the executable wake-up feed, not the administrator's read - // model. Every role, including admin, receives only its own assigned - // task events so visibility can never turn into plugin dispatch. + // This is an employee execution feed. Administrators are rejected + // before the stream is opened, and workers receive only assigned work. if (event.owner_user_id !== session.user.id) return; const publicEvent = { id: event.id, diff --git a/control-plane/src/task-service.ts b/control-plane/src/task-service.ts index a9dc97d..400a9c8 100644 --- a/control-plane/src/task-service.ts +++ b/control-plane/src/task-service.ts @@ -551,14 +551,19 @@ export function isTaskOwnerRestricted(role: TaskRole | undefined): boolean { return role === 'team_lead' || role === 'user'; } +export function canUseTaskDataPlane(role: TaskRole | undefined): boolean { + return role === 'team_lead' || role === 'user'; +} + export function canViewOperationsDashboard(role: TaskRole | undefined): boolean { - return role === 'admin' || role === 'team_lead'; + return role === 'team_lead'; } export function canAccessTask( access: TaskAccessScope, task: { assignedUserId: string | null | undefined } ): boolean { + if (access.role === 'admin') return false; if (!isTaskOwnerRestricted(access.role)) return true; return Boolean(access.userId) && task.assignedUserId === access.userId; @@ -575,7 +580,7 @@ export function canExecuteBusinessRoute({ routeId: BusinessRouteId | null; authorizedRouteIds: readonly BusinessRouteId[]; }): boolean { - if (role === 'admin') return source === 'manual'; + if (role === 'admin') return false; if (!isTaskOwnerRestricted(role) || !routeId) return false; return authorizedRouteIds.includes(routeId); } @@ -1587,7 +1592,7 @@ export function failureSummary( ? erpFeedback || text(result.failure_message) : text(result.failure_message)) || (uncertainStatus - ? `ERP 写入结果不确定,请勿重复提交同一任务,等待管理员只读核验:${blockers.join(';') || text(result.message) || '未取得确定回执。'}` + ? `ERP 写入结果不确定,请勿重复提交同一任务,等待任务所属账号只读核验:${blockers.join(';') || text(result.message) || '未取得确定回执。'}` : contractFailure ? `Agent 已返回,但 operation 未通过标准契约校验${blockers[0] ? `:${blockers[0]}` : ''};未进入插件,也未写入 ERP。` : blockers.length @@ -2497,9 +2502,15 @@ export class TaskService { return isTaskOwnerRestricted(context.role); } + private requireTaskDataPlaneAccess(access: TaskAccessScope): void { + if (access.role === 'admin') { + throw new TaskError('task_access_forbidden', '管理员账号仅用于平台管理,不能访问业务任务。', 403); + } + } + private requireOperationsDashboardAccess(context: TaskContext): void { if (!canViewOperationsDashboard(context.role)) { - throw new TaskError('leadership_required', '需要组长或管理员权限。', 403); + throw new TaskError('leadership_required', '需要组长权限。', 403); } } @@ -2667,6 +2678,7 @@ export class TaskService { taskId: string, { allowArchived = false }: { allowArchived?: boolean } = {} ): Promise> { + this.requireTaskDataPlaneAccess(context); const result = await client.query( `SELECT * FROM tasks WHERE organization_id = $1 @@ -3281,6 +3293,7 @@ export class TaskService { verdict: ParserReviewStatus, note?: string ): Promise> { + this.requireTaskDataPlaneAccess(context); const allowed: ParserReviewStatus[] = ['equivalent', 'program_correct', 'ai_correct', 'both_wrong']; if (!allowed.includes(verdict)) throw new TaskError('parser_review_invalid', '人工判定值无效。', 400); const normalizedNote = String(note || '').trim().slice(0, 2_000) || null; @@ -3328,6 +3341,7 @@ export class TaskService { } async reparseTaskWithAi(context: TaskContext, taskId: string, reason: string): Promise { + this.requireTaskDataPlaneAccess(context); const normalizedReason = String(reason || '').trim().slice(0, 1_000); if (!normalizedReason) throw new TaskError('reason_required', 'AI 重解析必须填写原因。', 400); const event = await withTransaction(this.config, async (client) => { @@ -3663,6 +3677,7 @@ export class TaskService { attachment: TaskInputAttachmentInput, selection: PassengerRosterAttachmentSelection ): Promise { + this.requireTaskDataPlaneAccess(context); const target = await this.passengerRosterAttachmentTarget(context, selection); const digest = sha256Bytes(attachment.content); const requestHash = sha256Text(`passenger-workbook\0${digest}`); @@ -4030,6 +4045,7 @@ export class TaskService { conversationId?: string, attachments: TaskInputAttachmentInput[] = [] ): Promise { + this.requireTaskDataPlaneAccess(context); const normalized = String(rawText || '').trim(); if (!normalized) throw new TaskError('empty_input', '输入区域不能为空。', 400); if (Buffer.byteLength(normalized, 'utf8') > 200_000) throw new TaskError('input_too_large', '输入内容超过限制。', 413); @@ -4158,6 +4174,7 @@ export class TaskService { } async ingestMessage(context: TaskContext, input: TaskMessageInput): Promise { + this.requireTaskDataPlaneAccess(context); const normalizedMessage = String(input.message || '').trim(); const attachments = Array.isArray(input.attachments) ? input.attachments : []; if (attachments.length > 1) throw new TaskError('too_many_attachments', '名单业务每次只能发送一个 Excel 附件。', 400); @@ -4567,6 +4584,7 @@ export class TaskService { } async getTask(organizationId: string, taskId: string, access?: TaskAccessScope): Promise { + if (access) this.requireTaskDataPlaneAccess(access); const result = await getPool(this.config).query( `SELECT t.*, uc.display_name AS channel_name, s.conversation_id, creator.id AS creator_id, creator.username AS creator_username, @@ -4627,6 +4645,7 @@ export class TaskService { messages: PublicTaskInputHistoryEntry[]; attachments: PublicTaskInputAttachmentAudit[]; }> { + this.requireTaskDataPlaneAccess(context); const task = await getPool(this.config).query( `SELECT id FROM tasks WHERE organization_id = $1 AND task_id = $2 @@ -5315,6 +5334,7 @@ export class TaskService { artifactId: string, access?: TaskAccessScope ): Promise { + if (access) this.requireTaskDataPlaneAccess(access); const authorized = await getPool(this.config).query( `SELECT 1 FROM tasks WHERE organization_id = $1 AND task_id = $2 @@ -5369,6 +5389,7 @@ export class TaskService { access?: TaskAccessScope; } = {} ): Promise { + if (options.access) this.requireTaskDataPlaneAccess(options.access); const status = text(options.status).trim(); const search = text(options.search).trim(); const params: unknown[] = [organizationId]; @@ -6132,6 +6153,7 @@ export class TaskService { } async confirmTask(context: TaskContext, taskId: string): Promise { + this.requireTaskDataPlaneAccess(context); const outcome = await withTransaction(this.config, async (client) => { const row = await this.lockTaskForAccess(client, context, taskId); if (!context.userId || text(row.assigned_user_id) !== context.userId) { @@ -6193,6 +6215,7 @@ export class TaskService { } async claimForBrowser(context: TaskContext, taskId: string, connectionId: string): Promise { + this.requireTaskDataPlaneAccess(context); const leaseOwner = `browser:${connectionId}`; const outcome = await withTransaction(this.config, async (client) => { const organization = await client.query( @@ -6392,6 +6415,7 @@ export class TaskService { executionId: string, options: { skipArtifactPersistence?: boolean } = {} ): Promise { + this.requireTaskDataPlaneAccess(context); const suppliedResult = jsonObject(resultPayload); if (suppliedResult.execution_id && text(suppliedResult.execution_id) !== executionId) { throw new TaskError('execution_id_mismatch', '插件回执与已领取的执行记录不一致。'); @@ -6470,7 +6494,7 @@ export class TaskService { resultObject.status = status; const lifecycle = executionLifecycleFacts(resultObject, status, true); Object.assign(resultObject, lifecycle); - const message = text(resultObject.message || (uncertain ? 'ERP 结果不确定,请勿重复提交同一任务,等待管理员只读核验。' : '插件已返回状态。')).slice(0, 2_000); + const message = text(resultObject.message || (uncertain ? 'ERP 结果不确定,请勿重复提交同一任务,等待任务所属账号只读核验。' : '插件已返回状态。')).slice(0, 2_000); const executionFailure = failureSummary(resultObject, '', status, stage); const operation = decryptedJson(this.config, row.operation_ciphertext) || row.operation || null; const baseSuccessReceipt = status === 'completed' ? successReceiptFromResult(resultObject) : null; @@ -6626,8 +6650,9 @@ export class TaskService { message: string, details: Record = {} ): Promise { + this.requireTaskDataPlaneAccess(context); const normalizedReason = text(reasonCode || 'erp_result_uncertain').slice(0, 160); - const normalizedMessage = text(message || 'ERP 执行结果不确定,已停止自动执行。请勿重复提交同一任务,等待管理员只读核验。').slice(0, 2_000); + const normalizedMessage = text(message || 'ERP 执行结果不确定,已停止自动执行。请勿重复提交同一任务,等待任务所属账号只读核验。').slice(0, 2_000); const outcome = await withTransaction(this.config, async (client) => { const lookup = await client.query( 'SELECT * FROM tasks WHERE organization_id = $1 AND task_id = $2 FOR UPDATE', @@ -7025,6 +7050,7 @@ export class TaskService { taskIds: string[], reason = '' ): Promise<{ task_ids: string[]; archived_count: number }> { + this.requireTaskDataPlaneAccess(context); const normalizedTaskIds = [...new Set(taskIds.map((taskId) => text(taskId).trim()).filter(Boolean))]; if (!normalizedTaskIds.length) throw new TaskError('invalid_task_ids', '请至少选择一个待归档任务。', 400); const archiveReason = text(reason).trim().slice(0, 500) || null; @@ -7079,6 +7105,7 @@ export class TaskService { } async restoreTask(context: TaskContext, taskId: string): Promise { + this.requireTaskDataPlaneAccess(context); const outcome = await withTransaction(this.config, async (client) => { const row = await this.lockTaskForAccess(client, context, taskId, { allowArchived: true }); if (!row.archived_at) return { event: null }; @@ -7115,6 +7142,7 @@ export class TaskService { context: TaskContext, taskIds: string[] ): Promise<{ task_ids: string[]; deleted_count: number }> { + this.requireTaskDataPlaneAccess(context); const normalizedTaskIds = [...new Set(taskIds.map((taskId) => text(taskId).trim()).filter(Boolean))]; if (!normalizedTaskIds.length) throw new TaskError('invalid_task_ids', '请至少选择一个待删除任务。', 400); const outcome = await withTransaction(this.config, async (client) => { @@ -7225,6 +7253,7 @@ export class TaskService { } async cancelTask(context: TaskContext, taskId: string): Promise { + this.requireTaskDataPlaneAccess(context); const outcome = await withTransaction(this.config, async (client) => { const row = await this.lockTaskForAccess(client, context, taskId); if (['completed', 'reconciliation_pending'].includes(text(row.status))) { @@ -7274,6 +7303,7 @@ export class TaskService { metadata: Record = {}, routing: { erpAccount?: string; erpAccountMatched?: boolean } = {} ): Promise<{ execution_ready: boolean; erp_account_matched: boolean; worker_connection_id: string }> { + this.requireTaskDataPlaneAccess(context); return withTransaction(this.config, async (client) => { const account = await client.query( `SELECT id, role, is_active, erp_account @@ -7288,15 +7318,15 @@ export class TaskService { const accountRow = account.rows[0] as Record; const expectedErpAccount = text(accountRow.erp_account); const reportedErpAccount = text(routing.erpAccount); - const isAdmin = text(accountRow.role) === 'admin'; - const erpAccountMatched = isAdmin - ? true - : Boolean( - expectedErpAccount - && routing.erpAccountMatched === true - && reportedErpAccount.toLocaleLowerCase() === expectedErpAccount.toLocaleLowerCase() - ); - const executionReady = isAdmin || erpAccountMatched; + if (text(accountRow.role) === 'admin') { + throw new TaskError('task_access_forbidden', '管理员账号不能注册任务执行浏览器。', 403); + } + const erpAccountMatched = Boolean( + expectedErpAccount + && routing.erpAccountMatched === true + && reportedErpAccount.toLocaleLowerCase() === expectedErpAccount.toLocaleLowerCase() + ); + const executionReady = erpAccountMatched; if (executionReady) { const competing = await client.query( diff --git a/control-plane/test/account-authorization.test.ts b/control-plane/test/account-authorization.test.ts index c491c3f..bca5998 100644 --- a/control-plane/test/account-authorization.test.ts +++ b/control-plane/test/account-authorization.test.ts @@ -71,6 +71,30 @@ test('AgentBus account-worker migration adds fail-closed channel, task, browser, assert.match(sql, /owner_user_id IS NULL[\s\S]+enabled = true/); }); +test('administrator task-isolation migration removes legacy bindings and rejects future data-plane principals', async () => { + const sql = await source('../migrations/021_admin_task_data_plane_isolation.sql'); + assert.match(sql, /UPDATE tasks task[\s\S]+assigned_user_id = NULL[\s\S]+account\.role = 'admin'/); + assert.match(sql, /UPDATE user_channels channel[\s\S]+owner_user_id = NULL[\s\S]+account\.role = 'admin'/); + assert.match(sql, /UPDATE browser_connections connection[\s\S]+status = 'superseded'[\s\S]+account\.role = 'admin'/); + assert.match(sql, /DELETE FROM user_business_route_authorizations[\s\S]+account\.role = 'admin'/); + assert.match(sql, /UPDATE leader_task_summary_subscriptions subscription[\s\S]+enabled = false[\s\S]+account\.role = 'admin'/); + assert.match(sql, /UPDATE leader_task_summary_deliveries delivery[\s\S]+delivery_status = 'cancelled'[\s\S]+account\.role = 'admin'/); + assert.match(sql, /CREATE OR REPLACE FUNCTION reject_admin_task_principal/); + assert.match(sql, /FOR SHARE/); + for (const constraint of [ + 'tasks_admin_assignee_forbidden', + 'tasks_admin_creator_forbidden', + 'user_channels_admin_owner_forbidden', + 'browser_connections_admin_worker_forbidden', + 'user_business_routes_admin_grantee_forbidden', + 'leader_task_summary_admin_subscriber_forbidden', + 'leader_task_summary_admin_recipient_forbidden' + ]) assert.match(sql, new RegExp(constraint)); + assert.match(sql, /CREATE OR REPLACE FUNCTION isolate_administrator_from_task_runtime/); + assert.match(sql, /AFTER UPDATE OF role ON users/); + assert.doesNotMatch(sql, /DELETE FROM tasks/); +}); + test('AgentBus channel keys and owners are unique so one inbound identity cannot fan out to multiple employees', async () => { const channels = await source('../src/agentbus-channels.ts'); assert.match(channels, /requireAssignableOwner/); @@ -109,7 +133,7 @@ test('account lifecycle is administrator-gated and protects passwords, sessions, assert.doesNotMatch(publicUser, /organization/); }); -test('administrators manage task-type grants and manual intake enforces them before parsing or ERP dispatch', async () => { +test('administrators manage employee task-type grants while remaining outside manual intake', async () => { const [auth, tasks, server] = await Promise.all([ source('../src/auth.ts'), source('../src/task-service.ts'), @@ -120,6 +144,7 @@ test('administrators manage task-type grants and manual intake enforces them bef assert.match(auth, /business_authorization_revision_conflict/); assert.match(auth, /account\.business_authorizations_updated/); assert.match(auth, /admin_business_authorization_fixed/); + assert.match(auth, /authorized_business_route_ids: role === 'admin' \? \[\] : storedRouteIds/); assert.match(server, /task_types: BUSINESS_ROUTES\.map/); assert.match(server, /app\.put\('\/api\/accounts\/:userId\/business-authorizations'[\s\S]+requireAdminMutationSession\(request\)/); assert.match(tasks, /export function canExecuteBusinessRoute/); @@ -155,7 +180,7 @@ test('operations dashboard is leadership-gated, cross-source, business-facing, a source('../src/task-service.ts'), source('../src/server.ts') ]); - assert.match(tasks, /canViewOperationsDashboard[\s\S]+role === 'admin' \|\| role === 'team_lead'/); + assert.match(tasks, /canViewOperationsDashboard[\s\S]+return role === 'team_lead'/); assert.match(tasks, /isTaskOwnerRestricted[\s\S]+role === 'team_lead' \|\| role === 'user'/); assert.match(tasks, /async listOperationsDashboard[\s\S]+t\.assigned_user_id IS NOT NULL[\s\S]+t\.source IN \('manual', 'agentbus'\)/); assert.match(tasks, /actorUserId[\s\S]+t\.assigned_user_id = \$\$\{params\.length\}/); @@ -243,9 +268,11 @@ test('ordinary task access is enforced across reads, mutations, artifacts, event assert.match(server, /tasks\.listTasksPage[\s\S]+access: contextFor\(session, request\)/); assert.match(server, /tasks\.getTaskArtifact[\s\S]+contextFor\(session, request\)/); assert.match(server, /tasks\.eventsSince\(session\.user\.organizationId, session\.user\.id, since\)/); + assert.match(server, /app\.addHook\('preHandler'[\s\S]+isTaskDataPlaneRoute\(request\.routeOptions\.url\)[\s\S]+requireTaskDataPlane\(await getSession\(request\)\)/); + assert.match(server, /task_access_forbidden/); }); -test('administrator visibility is isolated from executable events and plugin result routing', async () => { +test('administrators are excluded from task events and plugin result routing', async () => { const [tasks, server, app] = await Promise.all([ source('../src/task-service.ts'), source('../src/server.ts'), @@ -264,6 +291,7 @@ test('administrator visibility is isolated from executable events and plugin res assert.match(eventRoute, /command\.assigned_user_id !== session\.user\.id/); assert.match(eventRoute, /event: browser-command/); assert.match(eventRoute, /tasks\.eventsSince\(session\.user\.organizationId, session\.user\.id, since\)/); + assert.match(eventRoute, /requireTaskSession\(request\)/); const eventStream = app.slice( app.indexOf('function startRemoteEventStream()'), @@ -377,6 +405,12 @@ test('operator UI exposes role-aware accounts, executive drill-through, archive, assert.match(app, /\/business-authorizations/); assert.match(app, /当前默认不能执行任何业务/); assert.match(app, /function canViewOperationsDashboard/); + assert.match(app, /function canUseTaskDataPlane/); + assert.match(app, /isAdministrator\(\) && \(IS_TASK_PAGE \|\| IS_OPERATIONS_DASHBOARD_PAGE\)[\s\S]+window\.location\.replace\('\/accounts'\)/); + assert.match(app, /browserConnectionId = canUseTaskDataPlane\(\) \? connectionIdForUser\(user\) : ''/); + assert.match(app, /if \(bridgeState\) bridgeState\.hidden = !canUseTaskDataPlane\(\)/); + assert.match(app, /if \(!canUseTaskDataPlane\(\)\) return false;[\s\S]+sendToExtension\('PING'/); + assert.match(app, /account\.role === 'admin' \? '不参与任务' : '任务权限'/); assert.match(app, /\/api\/operations-dashboard\?/); assert.match(app, /\/api\/operations-dashboard\/tasks\/\$\{encodeURIComponent\(taskId\)\}/); assert.match(app, /business_route_id/); @@ -457,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-leader-summary-2/); - assert.match(index, /app\.js\?v=20260907-leader-summary-2/); + assert.match(index, /styles\.css\?v=20260907-admin-task-isolation-1/); + assert.match(index, /app\.js\?v=20260907-admin-task-isolation-1/); }); diff --git a/control-plane/test/agentbus.test.ts b/control-plane/test/agentbus.test.ts index 70f049d..6040b92 100644 --- a/control-plane/test/agentbus.test.ts +++ b/control-plane/test/agentbus.test.ts @@ -734,7 +734,7 @@ test('AgentBus result text uses the unified important message and preserves the assert.equal(taskResultStatus(uncertainWrite), 'failed'); assert.equal( taskResultText(uncertainWrite), - 'ERP 写入已发起,但系统尚未确认最终结果。请勿重复提交同一任务,等待管理员只读核验。' + 'ERP 写入已发起,但系统尚未确认最终结果。请勿重复提交同一任务,等待任务所属账号只读核验。' ); }); diff --git a/control-plane/test/control-plane.test.ts b/control-plane/test/control-plane.test.ts index 5ff1ae2..26264a5 100644 --- a/control-plane/test/control-plane.test.ts +++ b/control-plane/test/control-plane.test.ts @@ -5,7 +5,7 @@ import { tmpdir } from 'node:os'; import test from 'node:test'; import { loadConfig } from '../src/config.js'; import { decryptBytes, decryptText, encryptBytes, encryptText, hashToken, sameTokenHash, sha256Bytes } from '../src/crypto.js'; -import { aiServiceConnected, buildServer } from '../src/server.js'; +import { aiServiceConnected, buildServer, isTaskDataPlaneRoute } from '../src/server.js'; import { REQUIRED_SCHEMA_VERSION } from '../src/db.js'; import { TaskService, @@ -15,6 +15,7 @@ import { classifyExecutionResult, canAccessTask, canExecuteBusinessRoute, + canUseTaskDataPlane, canViewOperationsDashboard, executionLifecycleFacts, failureSummary, @@ -45,9 +46,10 @@ test('message routing starts a new session for a business directive, not for a s assert.equal(isNewDirectiveMessage('数量改为 2'), false); }); -test('task access contract isolates users and team leads while preserving administrator and worker access', () => { +test('task access contract excludes administrators and isolates employee owners', () => { const cases = [ - { name: 'administrator can inspect another assigned task', access: { userId: 'admin', role: 'admin' as const }, assignedUserId: 'user-a', allowed: true }, + { name: 'administrator cannot inspect another assigned task', access: { userId: 'admin', role: 'admin' as const }, assignedUserId: 'user-a', allowed: false }, + { name: 'administrator cannot inspect a legacy admin-assigned task', access: { userId: 'admin', role: 'admin' as const }, assignedUserId: 'admin', allowed: false }, { name: 'trusted worker can inspect an unassigned task', access: { userId: '', role: undefined }, assignedUserId: null, allowed: true }, { name: 'team lead sees own manual task', access: { userId: 'lead-a', role: 'team_lead' as const }, assignedUserId: 'lead-a', allowed: true }, { name: 'team lead cannot use the normal task path for another task', access: { userId: 'lead-a', role: 'team_lead' as const }, assignedUserId: 'user-b', allowed: false }, @@ -63,16 +65,85 @@ test('task access contract isolates users and team leads while preserving admini assert.equal(isTaskOwnerRestricted('admin'), false); assert.equal(isTaskOwnerRestricted('team_lead'), true); assert.equal(isTaskOwnerRestricted('user'), true); - assert.equal(canViewOperationsDashboard('admin'), true); + assert.equal(canUseTaskDataPlane('admin'), false); + assert.equal(canUseTaskDataPlane('team_lead'), true); + assert.equal(canUseTaskDataPlane('user'), true); + assert.equal(canUseTaskDataPlane(undefined), false); + assert.equal(canViewOperationsDashboard('admin'), false); assert.equal(canViewOperationsDashboard('team_lead'), true); assert.equal(canViewOperationsDashboard('user'), false); }); +test('task data-plane HTTP classifier covers every task-bearing surface and excludes management APIs', () => { + for (const route of [ + '/api/tasks', + '/api/tasks/:taskId', + '/api/messages', + '/api/connections/heartbeat', + '/api/events', + '/api/parser-decisions/:decisionId/review', + '/api/operations-dashboard', + '/api/operations-dashboard/tasks/:taskId' + ]) assert.equal(isTaskDataPlaneRoute(route), true, route); + for (const route of [ + '/api/accounts', + '/api/channels', + '/api/settings/parser-routing', + '/api/settings/automation', + '/api/audit' + ]) assert.equal(isTaskDataPlaneRoute(route), false, route); +}); + +test('task service rejects administrator operations before touching task storage', async () => { + const config = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: Buffer.alloc(32, 19).toString('base64'), + DATABASE_URL: 'postgresql://invalid:invalid@127.0.0.1:1/invalid' + }); + const service = new TaskService(config); + const context = { + organizationId: '11111111-1111-4111-8111-111111111111', + userId: '22222222-2222-4222-8222-222222222222', + requestId: 'admin-task-isolation-test', + role: 'admin' as const, + source: 'manual' as const + }; + const operations: Array<[string, () => Promise]> = [ + ['create', () => service.createTask(context, '安排用车')], + ['message', () => service.ingestMessage(context, { message: '安排用车' })], + ['attachment', () => service.attachPassengerRosterAttachment(context, { + fileName: 'roster.xlsx', contentType: 'application/octet-stream', content: Buffer.from('x'), source: 'manual' + }, { taskId: 'TASK-1' })], + ['list', () => service.listTasksPage(context.organizationId, { access: context })], + ['read', () => service.getTask(context.organizationId, 'TASK-1', context)], + ['input history', () => service.getTaskInputHistory(context, 'TASK-1')], + ['artifact', () => service.getTaskArtifact(context.organizationId, 'TASK-1', '33333333-3333-4333-8333-333333333333', context)], + ['confirm', () => service.confirmTask(context, 'TASK-1')], + ['claim', () => service.claimForBrowser(context, 'TASK-1', 'platform-browser:test')], + ['result', () => service.recordExecutionResult(context, 'TASK-1', {}, 'platform-browser:test', '44444444-4444-4444-8444-444444444444')], + ['reconcile', () => service.markReconciliationRequired(context, 'TASK-1', 'uncertain', '')], + ['archive', () => service.archiveTasks(context, ['TASK-1'])], + ['restore', () => service.restoreTask(context, 'TASK-1')], + ['delete', () => service.hardDeleteTasks(context, ['TASK-1'])], + ['cancel', () => service.cancelTask(context, 'TASK-1')], + ['heartbeat', () => service.heartbeat(context, 'platform-browser:test', '0.0.0')], + ['AI reparse', () => service.reparseTaskWithAi(context, 'TASK-1', 'test')], + ['parser review', () => service.reviewParserDecision(context, '33333333-3333-4333-8333-333333333333', 'equivalent')] + ]; + for (const [name, operation] of operations) { + await assert.rejects(operation, (error: unknown) => ( + error instanceof Error + && 'code' in error + && error.code === 'task_access_forbidden' + ), name); + } +}); + test('business route authorization is an explicit allowlist for team leads and ordinary users', () => { const routeId = 'arrangement_hotel_create' as const; assert.equal(canExecuteBusinessRoute({ role: 'admin', source: 'manual', routeId: null, authorizedRouteIds: [] - }), true, 'administrators retain all registered and unclassified manual intake'); + }), false, 'administrators cannot create manual tasks'); assert.equal(canExecuteBusinessRoute({ role: 'team_lead', source: 'manual', routeId, authorizedRouteIds: [routeId] }), true, 'team lead can use a granted route'); @@ -293,6 +364,65 @@ test('control plane exposes a live health endpoint without a database connection await app.close(); }); +test('administrator sessions receive 403 before every task data-plane handler', async () => { + const config = loadConfig({ + NODE_ENV: 'test', + FIELD_ENCRYPTION_KEY: Buffer.alloc(32, 20).toString('base64'), + DATABASE_URL: 'postgresql://invalid:invalid@127.0.0.1:1/invalid' + }); + const { app, auth } = await buildServer({ + config, + startParserLoop: false, + parser: { + async parse() { return { blockers: ['test parser'] }; }, + async checkConnection() { return { ok: false, configured: false }; } + } + }); + (auth as unknown as { getActiveSession: () => Promise }).getActiveSession = async () => ({ + id: '33333333-3333-4333-8333-333333333333', + csrfTokenHash: Buffer.alloc(32), + user: { + id: '22222222-2222-4222-8222-222222222222', + organizationId: '11111111-1111-4111-8111-111111111111', + username: 'admin', + role: 'admin', + erpAccount: null + } + }); + const requests = [ + { method: 'GET', url: '/api/tasks' }, + { method: 'POST', url: '/api/tasks', payload: { raw_text: '安排用车' } }, + { method: 'POST', url: '/api/messages', payload: { message: '安排用车' } }, + { method: 'GET', url: '/api/tasks/TASK-1' }, + { method: 'GET', url: '/api/tasks/TASK-1/input-history' }, + { method: 'GET', url: '/api/tasks/TASK-1/artifacts/44444444-4444-4444-8444-444444444444' }, + { method: 'POST', url: '/api/tasks/TASK-1/confirm', payload: {} }, + { method: 'POST', url: '/api/tasks/TASK-1/claim', payload: { connection_id: 'platform-browser:1234567890abcdef' } }, + { method: 'POST', url: '/api/tasks/TASK-1/result', payload: { result: {}, connection_id: 'platform-browser:1234567890abcdef', execution_id: '44444444-4444-4444-8444-444444444444' } }, + { method: 'POST', url: '/api/tasks/TASK-1/cancel', payload: {} }, + { method: 'POST', url: '/api/tasks/bulk-delete', payload: { task_ids: ['TASK-1'] } }, + { method: 'POST', url: '/api/tasks/bulk-archive', payload: { task_ids: ['TASK-1'] } }, + { method: 'DELETE', url: '/api/tasks/TASK-1' }, + { method: 'POST', url: '/api/tasks/TASK-1/archive', payload: {} }, + { method: 'POST', url: '/api/tasks/TASK-1/restore', payload: {} }, + { method: 'POST', url: '/api/connections/heartbeat', payload: { connection_id: 'platform-browser:1234567890abcdef' } }, + { method: 'GET', url: '/api/events' }, + { method: 'GET', url: '/api/operations-dashboard' }, + { method: 'GET', url: '/api/operations-dashboard/tasks/TASK-1' }, + { method: 'POST', url: '/api/tasks/TASK-1/reparse', payload: { engine: 'ai', reason: 'test' } }, + { method: 'PUT', url: '/api/parser-decisions/44444444-4444-4444-8444-444444444444/review', payload: { verdict: 'equivalent' } } + ]; + try { + for (const request of requests) { + const response = await app.inject(request as any); + assert.equal(response.statusCode, 403, `${request.method} ${request.url}`); + assert.equal(response.json().error_code, 'task_access_forbidden', `${request.method} ${request.url}`); + } + } finally { + await app.close(); + } +}); + test('AI primary connection state follows service reachability, not historical authentication evidence', () => { assert.equal(aiServiceConnected(true, { configured: true, @@ -313,11 +443,11 @@ test('migration contains the durable state tables and safety fields', async () = } }); -test('control plane requires the latest durable task-outcome migration before readiness', async () => { +test('control plane requires the administrator task-isolation 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, '020_leader_task_summary_notifications'); + assert.equal(REQUIRED_SCHEMA_VERSION, '021_admin_task_data_plane_isolation'); assert.match(db, /schema_migrations/); assert.match(db, /databaseReadiness/); assert.match(db, /assertDatabaseSchema/); @@ -1231,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-leader-summary-2/); - assert.match(index, /app\.js\?v=20260907-leader-summary-2/); + 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, /id="statusDetailsPopover"/); assert.match(index, /id="statusDetailsRefresh"/); assert.match(app, /apiRequest\(`\/api\/tasks\?\$\{params\.toString\(\)\}`/); @@ -1312,7 +1442,7 @@ test('operator page has a login gate and uses the durable task API', async () => assert.match(app, /fetchWithTimeout\('\/api\/auth\/login'/); assert.match(app, /showAuthChecking\(\);[\s\S]*pingAi\(\)\.catch/); assert.match(app, /cache: options\.cache \|\| 'no-store'/); - assert.match(app, /showAuthenticatedApp\(me\.user\);[\s\S]*任务同步失败/); + assert.match(app, /if \(!showAuthenticatedApp\(me\.user\)\) return false;[\s\S]*任务同步失败/); assert.match(app, /async function toggleStatusDetails/); assert.match(app, /已连接,有告警/); assert.match(app, /历史验证失败(不影响链路状态)/);