From aecde46ff651a717a5a774d6bc5b9fe66617db5c Mon Sep 17 00:00:00 2001 From: john Date: Thu, 27 Aug 2026 12:15:15 +0800 Subject: [PATCH] fix(agent): restore active-task routing and active-mode validation observation Keep correction follow-ups on Agent using session transcript and active task context, and let Page Data validation reach Orchestrator when mode is active instead of requiring shadowEngine. Co-authored-by: Cursor --- agent-run-gateway.mjs | 121 ++++++++- agent-run-gateway.test.mjs | 57 ++++ chat-intent-router.mjs | 249 ++++++++++++++++-- chat-intent-router.test.mjs | 138 ++++++++++ services/orchestrator/shadow-observer.mjs | 44 +++- .../orchestrator/shadow-observer.test.mjs | 73 ++++- 6 files changed, 662 insertions(+), 20 deletions(-) diff --git a/agent-run-gateway.mjs b/agent-run-gateway.mjs index c2b9889..fd819cb 100644 --- a/agent-run-gateway.mjs +++ b/agent-run-gateway.mjs @@ -3,7 +3,7 @@ import fs from 'node:fs/promises'; import path from 'node:path'; import { isRunStreamReplayEnabled } from './agent-run-stream.mjs'; import { isDirectChatSessionId } from './direct-chat-service.mjs'; -import { CHAT_INTENT_ROUTE, resolveGatewayAgentSessionId, resolveLegacyRouteFromClassification, logRouterDecisionShadow } from './chat-intent-router.mjs'; +import { CHAT_INTENT_ROUTE, formatRouterTranscript, resolveGatewayAgentSessionId, resolveLegacyRouteFromClassification, logRouterDecisionShadow } from './chat-intent-router.mjs'; import { resolveSessionAccess } from './session-broker.mjs'; import { loadSnapshotMessages, @@ -617,6 +617,103 @@ function resolveEffectiveToolMode(runOptions) { return runOptions?.pageDataAiderWorkflow ? 'chat' : (runOptions?.toolMode ?? 'chat'); } +const ACTIVE_TASK_CONTEXT_MAX_AGE_MS = 4 * 60 * 60 * 1000; + +function parseIntentRoutedEvent(dataJson) { + if (!dataJson) return null; + const payload = typeof dataJson === 'string' + ? (() => { + try { + return JSON.parse(dataJson); + } catch { + return null; + } + })() + : dataJson; + if (!payload || typeof payload !== 'object') return null; + return { + route: payload.route ?? payload.decision?.route ?? null, + suggestedSkill: payload.suggestedSkill ?? payload.suggested_skill ?? null, + reason: payload.reason ?? null, + }; +} + +export async function resolveActiveTaskContext(pool, { + userId, + sessionId, + excludeRunId = null, + nowMs: nowMsFn = Date.now, +} = {}) { + if (!pool || !userId || !sessionId || isDirectChatSessionId(sessionId)) { + return null; + } + const params = [userId, sessionId]; + let excludeClause = ''; + if (excludeRunId) { + excludeClause = ' AND r.id <> ?'; + params.push(excludeRunId); + } + const [rows] = await pool.query( + `SELECT r.id, r.status, r.error_message, r.created_at, r.updated_at, + ( + SELECT e.data_json + FROM h5_agent_run_events e + WHERE e.run_id = r.id AND e.event_type = 'intent_routed' + ORDER BY e.created_at ASC + LIMIT 1 + ) AS intent_json + FROM h5_agent_runs r + WHERE r.user_id = ? AND r.agent_session_id = ?${excludeClause} + ORDER BY r.created_at DESC + LIMIT 1`, + params, + ); + const row = rows?.[0]; + if (!row) return null; + + const intent = parseIntentRoutedEvent(row.intent_json); + const legacyRoute = String(intent?.route ?? '').trim(); + const agentRoute = legacyRoute === 'agent' + || legacyRoute === CHAT_INTENT_ROUTE.AGENT + || legacyRoute === 'agent_orchestration'; + if (!agentRoute) return null; + + const referenceMs = Math.max(Number(row.updated_at) || 0, Number(row.created_at) || 0); + if (referenceMs > 0 && nowMsFn() - referenceMs > ACTIVE_TASK_CONTEXT_MAX_AGE_MS) { + return null; + } + + return { + hasRecentAgentTask: true, + lastRunId: row.id, + lastRunStatus: String(row.status ?? '').trim() || null, + lastRunFailed: String(row.status ?? '').trim() === 'failed', + lastSuggestedSkill: intent?.suggestedSkill ?? null, + lastRoute: legacyRoute, + lastIntentReason: intent?.reason ?? null, + lastErrorMessage: row.error_message + ? String(row.error_message).slice(0, 500) + : null, + }; +} + +export async function resolveRouterTranscript(sessionSnapshotService, { + sessionId, + userId = null, + excludeLatestUserText = null, +} = {}) { + if (!sessionSnapshotService?.get || !sessionId || isDirectChatSessionId(sessionId)) { + return ''; + } + const snapshot = await sessionSnapshotService.get(sessionId).catch(() => null); + if (!snapshot) return ''; + if (userId && snapshot?.session?.user_id && snapshot.session.user_id !== userId) { + return ''; + } + const messages = Array.isArray(snapshot?.messages) ? snapshot.messages : []; + return formatRouterTranscript(messages, { excludeLatestUserText }); +} + function restoreCursorFallbackUserMessage(userMessage) { const message = (userMessage && typeof userMessage === 'object' && !Array.isArray(userMessage)) ? { ...userMessage } @@ -1547,13 +1644,29 @@ export function createAgentRunGateway({ } } - async function resolveRunRouting(row, userMessage, runOptions) { + async function resolveRunRouting(row, userMessage, runOptions, { runId = null } = {}) { if (!chatIntentRouter?.classify) return null; const enabled = chatIntentRouter.isEnabled ? await Promise.resolve(chatIntentRouter.isEnabled()).catch(() => false) : true; if (!enabled) return null; const grantedSkills = await resolveGrantedSkills(row.user_id); + const activeTaskContext = await resolveActiveTaskContext(pool, { + userId: row.user_id, + sessionId: row.agent_session_id ?? null, + excludeRunId: runId, + }).catch(() => null); + const displayText = userMessage?.metadata?.displayText + ?? deriveUserFacingText( + typeof userMessage?.content?.find === 'function' + ? (userMessage.content.find((item) => item?.type === 'text')?.text ?? '') + : '', + ); + const recentTranscript = await resolveRouterTranscript(sessionSnapshotService, { + sessionId: row.agent_session_id ?? null, + userId: row.user_id, + excludeLatestUserText: displayText, + }).catch(() => ''); return chatIntentRouter.classify({ userId: row.user_id, userMessage, @@ -1562,6 +1675,8 @@ export function createAgentRunGateway({ toolMode: runOptions.toolMode, forceDeepReasoning: runOptions.forceDeepReasoning, grantedSkills, + activeTaskContext, + recentTranscript, }); } @@ -1624,7 +1739,7 @@ export function createAgentRunGateway({ policyBlocked: true, }; } - const routing = await resolveRunRouting(row, userMessage, runOptions); + const routing = await resolveRunRouting(row, userMessage, runOptions, { runId }); const routingDecision = resolveLegacyRouteFromClassification(routing) ?? routing?.route ?? null; const cursorFirstApplied = applyCursorFirstAgentExecution(userMessage, runOptions, { routingDecision, diff --git a/agent-run-gateway.test.mjs b/agent-run-gateway.test.mjs index 21478de..f96faaf 100644 --- a/agent-run-gateway.test.mjs +++ b/agent-run-gateway.test.mjs @@ -9,9 +9,66 @@ import { assertRequiredImageGenerationCompleted, createAgentRunGateway, normalizeAgentRunWorkerIdentity, + resolveActiveTaskContext, + resolveRouterTranscript, resolveRequiredCodeExecutor, } from './agent-run-gateway.mjs'; +test('resolveActiveTaskContext returns failed agent task metadata for same session', async () => { + const now = Date.now(); + const pool = { + async query() { + return [[{ + id: 'run-failed-1', + status: 'failed', + error_message: 'delivery failed', + created_at: now - 60_000, + updated_at: now - 30_000, + intent_json: { + route: 'agent_orchestration', + suggestedSkill: 'page-data-collect', + reason: '页面需要数据交互与持久化', + }, + }]]; + }, + }; + const context = await resolveActiveTaskContext(pool, { + userId: 'user-1', + sessionId: '20260827_3', + excludeRunId: 'run-current', + nowMs: () => now, + }); + assert.equal(context?.hasRecentAgentTask, true); + assert.equal(context?.lastRunFailed, true); + assert.equal(context?.lastSuggestedSkill, 'page-data-collect'); +}); + +test('resolveRouterTranscript formats recent portal snapshot messages', async () => { + const transcript = await resolveRouterTranscript({ + async get(sessionId) { + assert.equal(sessionId, '20260827_3'); + return { + messages: [ + { + role: 'user', + metadata: { displayText: '帮我做台账' }, + content: [{ type: 'text', text: '帮我做台账' }], + }, + { + role: 'assistant', + content: [{ type: 'text', text: '正在创建页面。' }], + }, + ], + }; + }, + }, { + sessionId: '20260827_3', + excludeLatestUserText: '继续完善', + }); + assert.match(transcript, /用户:帮我做台账/); + assert.match(transcript, /助手:正在创建页面/); +}); + test('required code executor is read from run metadata', () => { assert.equal(resolveRequiredCodeExecutor({ metadata: { memindRun: { executor: 'AIDER' } }, diff --git a/chat-intent-router.mjs b/chat-intent-router.mjs index 6983d38..4b1dda3 100644 --- a/chat-intent-router.mjs +++ b/chat-intent-router.mjs @@ -16,6 +16,10 @@ import { } from './memory-intervention.mjs'; import { filterMemoriesByQuery } from './memory-legacy-fallback.mjs'; import { matchDirectChatFaqRule, isExplicitTextOnlyRequest } from './chat-intent-router-rules.mjs'; +import { + deriveAssistantFacingText, + deriveUserFacingText, +} from './conversation-display.mjs'; import { isGoalRunIntent } from './goal-run-intent.mjs'; import { pgvectorMemoryBackendInternals } from './memory-v2-pgvector.mjs'; @@ -84,6 +88,17 @@ const AGENT_SESSION_CONTINUE_PATTERNS = [ /^(?:开始吧|按默认做|默认方案|继续|就这样|就这样吧)[!!。.\s]*$/iu, ]; +/** Corrections and refinements during an active Agent task must stay on Agent. */ +const AGENT_TASK_FOLLOWUP_PATTERNS = [ + /(?:还是|仍然|依然|并没|没有|并未).{0,16}(?:按照|按|理解|改|做|对|听)/u, + /(?:不对|错了|有误|有问题|不行|不好|不太对)/u, + /(?:更正|修正|修改|改一下|调整|优化|重新|再来|重做|再试)/u, + /(?:不听|没听|不按|没按|不照).{0,16}(?:指令|要求|说的|逻辑|意思)/u, + /(?:继续|接着).{0,16}(?:做|改|完善|优化|调整|处理|执行)/u, + /(?:按我(?:的|说)|照我(?:的|说)|我说的|我讲的|我上面)/u, + /(?:没(?:有)?按照|没有按).{0,20}(?:逻辑|要求|方案|口径|框架)/u, +]; + /** Pure text chat/creative prompts — fast-path even when LLM router is enabled. */ const OBVIOUS_DIRECT_CHAT_PATTERNS = [ /(?:讲|说|来|编).{0,10}(?:个|一段|一首|一个)?(?:睡前故事|故事|笑话|段子)/u, @@ -310,6 +325,84 @@ export function isAgentSessionContinueText(text) { return AGENT_SESSION_CONTINUE_PATTERNS.some((pattern) => pattern.test(normalized)); } +export function isAgentTaskFollowUpText(text) { + const normalized = String(text ?? '').trim(); + if (!normalized) return false; + return AGENT_TASK_FOLLOWUP_PATTERNS.some((pattern) => pattern.test(normalized)); +} + +function isExplicitDirectChatOnlyText(text) { + const normalized = String(text ?? '').trim(); + if (!normalized) return false; + if (isMemoryRecallQuestion(normalized)) return true; + if (isExplicitTextOnlyRequest(normalized)) return true; + if (OBVIOUS_DIRECT_PATTERNS.some((pattern) => pattern.test(normalized))) return true; + if (OBVIOUS_DIRECT_CHAT_PATTERNS.some((pattern) => pattern.test(normalized))) return true; + const faqMatch = matchDirectChatFaqRule(normalized); + return Boolean(faqMatch); +} + +export function shouldForceActiveAgentTaskContinuation({ + activeTaskContext = null, + text = '', + sessionId = null, + sessionMessageCount = null, +} = {}) { + if (!activeTaskContext?.hasRecentAgentTask) return false; + if (!hasPriorAgentConversation(sessionId, sessionMessageCount)) return false; + const normalized = String(text ?? '').trim(); + if (!normalized || isExplicitDirectChatOnlyText(normalized)) return false; + if (activeTaskContext.lastRunFailed === true) return true; + return isAgentTaskFollowUpText(normalized); +} + +export function shouldDeferActiveTaskRoutingToLlm({ + activeTaskContext = null, + text = '', + sessionId = null, + sessionMessageCount = null, +} = {}) { + if (!activeTaskContext?.hasRecentAgentTask) return false; + if (!hasPriorAgentConversation(sessionId, sessionMessageCount)) return false; + const normalized = String(text ?? '').trim(); + if (!normalized || isExplicitDirectChatOnlyText(normalized)) return false; + if (shouldForceActiveAgentTaskContinuation({ + activeTaskContext, + text: normalized, + sessionId, + sessionMessageCount, + })) { + return false; + } + return activeTaskContext.lastRoute !== 'direct_chat' + && activeTaskContext.lastRoute !== 'chat'; +} + +/** @deprecated use shouldForceActiveAgentTaskContinuation */ +export function shouldContinueActiveAgentTask(input = {}) { + return shouldForceActiveAgentTaskContinuation(input) + || shouldDeferActiveTaskRoutingToLlm(input); +} + +export function buildActiveTaskFollowUpClassification(activeTaskContext, { grantedSkills = [] } = {}) { + let suggestedSkill = String(activeTaskContext?.lastSuggestedSkill ?? '').trim() || null; + if (suggestedSkill && grantedSkills.length > 0 && !grantedSkills.includes(suggestedSkill)) { + suggestedSkill = null; + } + const agentBrief = suggestedSkill + ? `继续完成未完成的 ${suggestedSkill} 任务;用户正在纠正或补充要求,必须调用工具执行,不要只回复文字。` + : '用户正在纠正或补充上一轮 Agent 任务,必须继续执行并产出结果,不要只回复文字。'; + return normalizeClassification({ + route: CHAT_INTENT_ROUTE.AGENT, + confidence: 0.92, + reason: activeTaskContext?.lastRunFailed + ? '上一轮任务未完成,延续 Agent 执行' + : '活跃任务会话中的纠正或补充', + suggested_skill: suggestedSkill, + agent_brief: agentBrief, + }, { source: 'rule' }); +} + export function isRealtimeInfoQuestion(text) { const normalized = String(text ?? '').trim(); if (!normalized) return false; @@ -541,6 +634,8 @@ function buildRouterSystemPrompt(grantedSkills = []) { '- 用户要产出可访问页面、文件、链接,或需要工具/skills → agent_orchestration', '- 用户询问实时赛况、新闻、天气、行情等需要联网查询的问题 → agent_orchestration,suggested_skill 填 web(不要填 search)', '- 不确定时优先 agent_orchestration,避免漏执行', + '- 若 [Active Task] 显示上一轮 Agent 任务失败或进行中,用户当前消息是在纠正、补充、追问进度或要求按前文继续 → agent_orchestration,不要判成 direct_chat', + '- [Recent Conversation] 是对话摘要,当前消息可能引用上文(如「不对」「按我说的」「还是没有改」);结合上下文判断,不要孤立看当前一句', '- 记忆线索只用于辅助判断本轮意图,不能替用户扩写新需求', '', skills.length ? `当前用户已授权 skills:${skills.join(', ')}` : '当前用户未授权额外 skills。', @@ -626,15 +721,83 @@ export function buildRouterContext(resolveResult, { }; } -function buildRouterUserPrompt({ text, routerContext }) { - const context = String(routerContext ?? '').trim() || '无'; - return [ - '[Router Context]', - context, - '', - '[User]', - text || '(empty)', - ].join('\n'); +export function formatRouterTranscript(messages = [], { + maxTurns = 8, + maxChars = 4_000, + excludeLatestUserText = null, +} = {}) { + const normalizedExclude = String(excludeLatestUserText ?? '').trim(); + const lines = []; + let usedChars = 0; + const visible = []; + for (const message of Array.isArray(messages) ? messages : []) { + const role = String(message?.role ?? '').trim(); + if (role !== 'user' && role !== 'assistant') continue; + const raw = Array.isArray(message?.content) + ? message.content + .map((item) => (typeof item === 'string' ? item : item?.text ?? '')) + .join('\n') + : String(message?.content ?? message?.text ?? ''); + const displayText = String(message?.metadata?.displayText ?? '').trim(); + const text = role === 'user' + ? (displayText || deriveUserFacingText(raw)) + : deriveAssistantFacingText(raw); + const trimmed = String(text ?? '').trim(); + if (!trimmed) continue; + if (role === 'user' && normalizedExclude && trimmed === normalizedExclude) continue; + visible.push({ role, text: trimmed }); + } + for (let index = visible.length - 1; index >= 0 && lines.length < maxTurns; index -= 1) { + const item = visible[index]; + const line = `${item.role === 'user' ? '用户' : '助手'}:${item.text}`; + if (lines.length > 0 && usedChars + line.length > maxChars) break; + lines.unshift(line); + usedChars += line.length; + } + return lines.join('\n'); +} + +function formatActiveTaskContextForRouter(activeTaskContext) { + if (!activeTaskContext?.hasRecentAgentTask) return ''; + const parts = [ + `- 上一轮 Agent 任务状态:${activeTaskContext.lastRunStatus ?? 'unknown'}`, + ]; + if (activeTaskContext.lastRunFailed) { + parts.push('- 上一轮任务失败,用户可能在纠正、补充或要求重试'); + } + if (activeTaskContext.lastSuggestedSkill) { + parts.push(`- 上一轮建议 skill:${activeTaskContext.lastSuggestedSkill}`); + } + if (activeTaskContext.lastIntentReason) { + parts.push(`- 上一轮路由原因:${activeTaskContext.lastIntentReason}`); + } + if (activeTaskContext.lastErrorMessage) { + parts.push(`- 失败摘要:${String(activeTaskContext.lastErrorMessage).slice(0, 200)}`); + } + return parts.join('\n'); +} + +function buildRouterUserPrompt({ + text, + routerContext, + recentTranscript = '', + activeTaskContext = null, +} = {}) { + const sections = []; + const memoryContext = String(routerContext ?? '').trim(); + if (memoryContext) { + sections.push('[Memory Context]', memoryContext, ''); + } + const taskContext = formatActiveTaskContextForRouter(activeTaskContext); + if (taskContext) { + sections.push('[Active Task]', taskContext, ''); + } + const transcript = String(recentTranscript ?? '').trim(); + if (transcript) { + sections.push('[Recent Conversation]', transcript, ''); + } + sections.push('[Current User Message]', text || '(empty)'); + return sections.join('\n'); } function parseRouterJson(reply) { @@ -1043,6 +1206,8 @@ export function classifyWithRules({ sessionMessageCount = null, userMessage = null, includeIntentPatterns = true, + activeTaskContext = null, + grantedSkills = [], } = {}) { const decisionContext = { text, @@ -1123,6 +1288,21 @@ export function classifyWithRules({ reason: 'Agent 会话确认/续聊', }, { source: 'rule' }), decisionContext); } + if ( + includeIntentPatterns + && normalized + && shouldForceActiveAgentTaskContinuation({ + activeTaskContext, + text: normalized, + sessionId, + sessionMessageCount, + }) + ) { + return finalizeRouterClassification( + buildActiveTaskFollowUpClassification(activeTaskContext, { grantedSkills }), + decisionContext, + ); + } if (includeIntentPatterns && normalized && isPageDataDevIntent(normalized)) { return finalizeRouterClassification(normalizeClassification({ route: CHAT_INTENT_ROUTE.AGENT, @@ -1173,6 +1353,16 @@ export function classifyWithRules({ return finalizeRouterClassification(buildRealtimeInfoClassification(), decisionContext); } if (includeIntentPatterns && normalized) { + if ( + shouldDeferActiveTaskRoutingToLlm({ + activeTaskContext, + text: normalized, + sessionId, + sessionMessageCount, + }) + ) { + return null; + } return finalizeRouterClassification(normalizeClassification({ route: CHAT_INTENT_ROUTE.DIRECT_CHAT, confidence: 0.72, @@ -1229,6 +1419,8 @@ export function createChatIntentRouter(options = {}) { text, grantedSkills = [], routerContext = null, + recentTranscript = '', + activeTaskContext = null, } = {}) { if (typeof llmProviderService?.createChatCompletion !== 'function') return null; try { @@ -1244,6 +1436,8 @@ export function createChatIntentRouter(options = {}) { content: buildRouterUserPrompt({ text, routerContext: routerContext?.content ?? null, + recentTranscript, + activeTaskContext, }), }, ], @@ -1496,6 +1690,8 @@ export function createChatIntentRouter(options = {}) { toolMode = 'chat', forceDeepReasoning = false, grantedSkills = [], + activeTaskContext = null, + recentTranscript = '', } = {}) { const text = messageDisplayText(userMessage); const decisionContext = { @@ -1535,10 +1731,28 @@ export function createChatIntentRouter(options = {}) { sessionMessageCount, userMessage, includeIntentPatterns: true, + activeTaskContext, + grantedSkills, }); if (ruleResult) return finalizeWithCoercion(ruleResult); - const baseline = buildFallbackClassification(); + const activeTaskFallback = ( + shouldForceActiveAgentTaskContinuation({ + activeTaskContext, + text, + sessionId, + sessionMessageCount, + }) + || shouldDeferActiveTaskRoutingToLlm({ + activeTaskContext, + text, + sessionId, + sessionMessageCount, + }) + ) + ? buildActiveTaskFollowUpClassification(activeTaskContext, { grantedSkills }) + : null; + const baseline = activeTaskFallback ?? buildFallbackClassification(); if (!llmRouterEligible) { return finalizeWithCoercion(baseline); } @@ -1554,14 +1768,19 @@ export function createChatIntentRouter(options = {}) { .catch(() => buildRouterContext(null)) : Promise.resolve(buildRouterContext(null)); - // Memory is optional routing hint; run LLM in parallel to cut serial latency. - const [routerContext, llmResult] = await Promise.all([ - memoryPromise, + const llmPromise = memoryPromise.then((routerContext) => classifyWithLlm({ text, grantedSkills, - routerContext: null, + routerContext, + recentTranscript, + activeTaskContext, }), + ); + + const [routerContext, llmResult] = await Promise.all([ + memoryPromise, + llmPromise, ]); if (policy.shadowMode) { @@ -1592,7 +1811,7 @@ export function createChatIntentRouter(options = {}) { if (!llmResult || llmResult.confidence < policy.minConfidence) { return finalizeWithCoercion({ - ...baseline, + ...(activeTaskFallback ?? baseline), llmSuggestion: llmResult ? { route: llmResult.route, diff --git a/chat-intent-router.test.mjs b/chat-intent-router.test.mjs index 1050831..6cff922 100644 --- a/chat-intent-router.test.mjs +++ b/chat-intent-router.test.mjs @@ -11,6 +11,9 @@ import { coercePageGenerationSkill, createChatIntentRouter, createManagedChatIntentRouter, + formatRouterTranscript, + shouldDeferActiveTaskRoutingToLlm, + shouldForceActiveAgentTaskContinuation, isNormalizedRouterDecisionEnabled, isNormalizedRouterDecisionShadow, resolveGatewayAgentSessionId, @@ -469,6 +472,141 @@ test('classifyWithRules keeps agent session confirmation on agent path', () => { assert.match(fresh.reason, /记忆/); }); +test('classifyWithRules continues agent task after failed run in same session', () => { + const activeTaskContext = { + hasRecentAgentTask: true, + lastRunFailed: true, + lastSuggestedSkill: 'page-data-collect', + lastRunStatus: 'failed', + lastRoute: 'agent_orchestration', + }; + const result = classifyWithRules({ + text: '还是没有按照我的逻辑更正', + sessionId: '20260827_3', + sessionMessageCount: 50, + activeTaskContext, + grantedSkills: ['page-data-collect'], + }); + assert.equal(result.route, CHAT_INTENT_ROUTE.AGENT); + assert.equal(result.suggestedSkill, 'page-data-collect'); + assert.match(result.reason, /上一轮任务未完成/); +}); + +test('classifyWithRules routes correction follow-up to agent when recent agent task exists', () => { + const activeTaskContext = { + hasRecentAgentTask: true, + lastRunFailed: false, + lastSuggestedSkill: 'page-data-collect', + lastRunStatus: 'succeeded', + lastRoute: 'agent_orchestration', + }; + const result = classifyWithRules({ + text: '不听指令了吗', + sessionId: '20260827_3', + sessionMessageCount: 51, + activeTaskContext, + grantedSkills: ['page-data-collect'], + }); + assert.equal(result.route, CHAT_INTENT_ROUTE.AGENT); + assert.equal(result.suggestedSkill, 'page-data-collect'); + assert.match(result.reason, /纠正或补充/); +}); + +test('classifyWithRules defers ambiguous active-task follow-up to LLM router', () => { + const activeTaskContext = { + hasRecentAgentTask: true, + lastRunFailed: false, + lastSuggestedSkill: 'page-data-collect', + lastRunStatus: 'succeeded', + lastRoute: 'agent_orchestration', + }; + const result = classifyWithRules({ + text: '招商管理端可以展开招商漏斗,明确准入门槛', + sessionId: '20260827_3', + sessionMessageCount: 52, + activeTaskContext, + grantedSkills: ['page-data-collect'], + }); + assert.equal(result, null); +}); + +test('formatRouterTranscript keeps recent turns and strips skill prefixes from user text', () => { + const transcript = formatRouterTranscript([ + { + role: 'user', + metadata: { displayText: '帮我做商户沟通台账' }, + content: [{ type: 'text', text: '请使用 page-data-collect 技能:帮我做商户沟通台账' }], + }, + { + role: 'assistant', + content: [{ type: 'text', text: '好的,我先创建问卷页面。' }], + }, + ], { + excludeLatestUserText: '招商管理端展开漏斗', + }); + assert.match(transcript, /用户:帮我做商户沟通台账/); + assert.match(transcript, /助手:好的,我先创建问卷页面/); + assert.doesNotMatch(transcript, /page-data-collect 技能/); +}); + +test('createChatIntentRouter uses active-task fallback when LLM router is unavailable', async () => { + const router = createChatIntentRouter({ + llmProviderService: { + async createChatCompletion() { + return null; + }, + }, + env: { + ...process.env, + MEMIND_CHAT_LLM_ROUTER_ENABLED: '1', + MEMIND_CHAT_LLM_ROUTER_SHADOW: '0', + MEMIND_CHAT_ROUTER_CANARY_USER_IDS: '', + }, + }); + const result = await router.classify({ + userId: 'user-1', + sessionId: '20260827_3', + sessionMessageCount: 40, + grantedSkills: ['page-data-collect'], + activeTaskContext: { + hasRecentAgentTask: true, + lastRunFailed: true, + lastSuggestedSkill: 'page-data-collect', + lastRunStatus: 'failed', + lastRoute: 'agent_orchestration', + }, + userMessage: { + role: 'user', + content: [{ type: 'text', text: '还是没有按照我的逻辑更正' }], + metadata: { displayText: '还是没有按照我的逻辑更正' }, + }, + }); + assert.equal(result.route, CHAT_INTENT_ROUTE.AGENT); + assert.equal(result.suggestedSkill, 'page-data-collect'); +}); + +test('classifyWithRules keeps memory recall on direct chat even with failed agent task context', () => { + const activeTaskContext = { + hasRecentAgentTask: true, + lastRunFailed: true, + lastSuggestedSkill: 'page-data-collect', + lastRunStatus: 'failed', + }; + const result = classifyWithRules({ + text: '你记得我说想去哪儿吗', + sessionId: '20260827_3', + sessionMessageCount: 50, + activeTaskContext, + userMessage: { + role: 'user', + content: [{ type: 'text', text: '你记得我说想去哪儿吗' }], + metadata: { displayText: '你记得我说想去哪儿吗' }, + }, + }); + assert.equal(result.route, CHAT_INTENT_ROUTE.DIRECT_CHAT); + assert.match(result.reason, /记忆|历史对话/); +}); + test('classifyWithRules keeps memory recall on direct chat even when session already active', () => { const result = classifyWithRules({ text: '你记得我说想去哪儿吗', diff --git a/services/orchestrator/shadow-observer.mjs b/services/orchestrator/shadow-observer.mjs index 26fdb0f..b27ea69 100644 --- a/services/orchestrator/shadow-observer.mjs +++ b/services/orchestrator/shadow-observer.mjs @@ -28,6 +28,19 @@ async function recordValidationWithRetry(engine, runId, observation) { } } +export function resolveValidationEngineSelection(selection) { + if (selection.shadowEngine === WORKFLOW_ENGINE.LANGGRAPH) { + return selection; + } + if ( + selection.engine === WORKFLOW_ENGINE.LANGGRAPH + && ['active', 'canary', 'shadow'].includes(String(selection.mode ?? '')) + ) { + return selection; + } + return null; +} + export function createWorkflowShadowObserver({ configService, serviceToken = process.env.MEMIND_ORCHESTRATOR_SERVICE_TOKEN, @@ -66,6 +79,34 @@ export function createWorkflowShadowObserver({ }; } + async function selectValidationEngine({ + runId, + requestId, + userId, + workflowName, + }) { + const selection = await configService.selectEngine({ + runId, + requestId, + userId, + workflowName, + }); + if (!resolveValidationEngineSelection(selection)) { + return { selection, engine: null }; + } + const state = await configService.getRuntimeState(); + return { + selection, + engine: createRemoteWorkflowEngine({ + id: WORKFLOW_ENGINE.LANGGRAPH, + baseUrl: state.config.serviceUrl, + serviceToken, + timeoutMs: state.config.requestTimeoutMs, + fetchImpl, + }), + }; + } + async function observeWorkflowRun({ runId, requestId, @@ -152,7 +193,7 @@ export function createWorkflowShadowObserver({ workflowName = 'code-run-v1', observation, } = {}) { - const { selection, engine } = await selectShadowEngine({ + const { selection, engine } = await selectValidationEngine({ runId, requestId, userId, @@ -187,5 +228,6 @@ export function createWorkflowShadowObserver({ export const workflowShadowObserverInternals = { recordValidationWithRetry, + resolveValidationEngineSelection, safeError, }; diff --git a/services/orchestrator/shadow-observer.test.mjs b/services/orchestrator/shadow-observer.test.mjs index 40eaae2..a3a074b 100644 --- a/services/orchestrator/shadow-observer.test.mjs +++ b/services/orchestrator/shadow-observer.test.mjs @@ -3,7 +3,7 @@ import test from 'node:test'; import { MemorySaver } from '@langchain/langgraph'; import { createOrchestratorApp } from './app.mjs'; import { createLangGraphOrchestratorRuntime } from './runtime.mjs'; -import { createWorkflowShadowObserver } from './shadow-observer.mjs'; +import { createWorkflowShadowObserver, resolveValidationEngineSelection } from './shadow-observer.mjs'; function jsonResponse(body, status = 200) { return new Response(JSON.stringify(body), { @@ -25,6 +25,77 @@ async function listen(app) { }; } +test('resolveValidationEngineSelection accepts active-mode LangGraph engine', () => { + assert.ok(resolveValidationEngineSelection({ + engine: 'langgraph', + shadowEngine: null, + mode: 'active', + reason: 'active', + })); + assert.equal(resolveValidationEngineSelection({ + engine: 'native', + shadowEngine: null, + mode: 'active', + reason: 'active', + }), null); +}); + +test('validation observer records observations in active mode without shadowEngine', async () => { + let capturedUrl = null; + const observer = createWorkflowShadowObserver({ + configService: { + async selectEngine() { + return { + engine: 'langgraph', + candidateEngine: 'langgraph', + shadowEngine: null, + fallbackEngine: 'native', + reason: 'active', + candidateReason: 'active', + mode: 'active', + configVersion: 6, + dryRun: false, + }; + }, + async getRuntimeState() { + return { + config: { + serviceUrl: 'http://orchestrator.internal:8093', + requestTimeoutMs: 1200, + }, + }; + }, + }, + serviceToken: 'internal-token', + fetchImpl: async (url) => { + capturedUrl = url; + return jsonResponse({ + runId: 'run-active-1', + validation: { verdict: 'passed', kind: 'page-data-delivery' }, + }); + }, + }); + + const result = await observer.observeValidation({ + runId: 'run-active-1', + requestId: 'request-active-1', + userId: 'user-1', + observation: { + idempotencyKey: 'run-active-1:page-data-delivery:v1', + taskType: 'page_data_dev', + required: true, + checks: [{ id: 'page_data_binding', status: 'passed' }], + source: 'portal-agent-run', + observedAt: 123, + }, + }); + + assert.equal(result.observed, true); + assert.equal(result.mode, 'active'); + assert.equal(result.validation.verdict, 'passed'); + assert.match(String(capturedUrl), /\/v1\/runs\/run-active-1\/validation-observations$/); +}); + test('shadow observer skips without creating a remote client when mode does not select shadow', async () => { let fetchCalls = 0; const observer = createWorkflowShadowObserver({