diff --git a/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/executor/CriterionEvaluationAction.java b/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/executor/CriterionEvaluationAction.java index 3c2c0845..fddf8189 100644 --- a/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/executor/CriterionEvaluationAction.java +++ b/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/executor/CriterionEvaluationAction.java @@ -28,6 +28,7 @@ import com.alibaba.assistant.agent.evaluation.model.EvaluationCriterion; import com.alibaba.assistant.agent.evaluation.model.EvaluationSuite; import com.alibaba.assistant.agent.evaluation.model.ExecutionContextFactory; +import com.alibaba.assistant.agent.evaluation.util.EvaluationLogContextHelper; import com.alibaba.cloud.ai.graph.OverAllState; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -92,7 +93,8 @@ public Map apply(OverAllState state) { throw new IllegalStateException("Required components not found in OverAllState"); } - logger.info("Executing criterion: {}", criterion.getName()); + String sessionId = EvaluationLogContextHelper.getSessionId(evaluationContext); + logger.info("Executing criterion: {}, sessionId={}", criterion.getName(), sessionId); // Build dependency results map from individual state keys Map dependencyResults = buildDependencyResults(state); @@ -103,10 +105,10 @@ public Map apply(OverAllState state) { CriterionResult skipResult = checkConditionalExecution(conditionalConfig, dependencyResults); if (skipResult != null) { // Condition not met, return skipped result - logger.info("Criterion '{}' skipped: {}", criterion.getName(), conditionalConfig.getSkipReason()); + logger.info("Criterion '{}' skipped, sessionId={}: {}", criterion.getName(), sessionId, conditionalConfig.getSkipReason()); return buildSkippedResultUpdates(skipResult); } - logger.debug("Conditional execution check passed for criterion: {}", criterion.getName()); + logger.debug("Conditional execution check passed for criterion: {}, sessionId={}", criterion.getName(), sessionId); } // Determine timeout: criterion-specific > suite default @@ -139,18 +141,18 @@ public Map apply(OverAllState state) { } else { // No executor or timeout disabled, execute directly if (batchingEnabled) { - logger.debug("Batching enabled for criterion: {}", criterion.getName()); + logger.debug("Batching enabled for criterion: {}, sessionId={}", criterion.getName(), sessionId); result = executeWithBatching(evaluationContext, dependencyResults, batchingConfig); } else { - logger.debug("Batching not enabled for criterion: {}", criterion.getName()); + logger.debug("Batching not enabled for criterion: {}, sessionId={}", criterion.getName(), sessionId); result = executeWithoutBatching(evaluationContext, dependencyResults); } } } catch (TimeoutException te) { // Criterion timed out - return timeout result with default value long elapsedMs = System.currentTimeMillis() - startTime; - logger.warn("Criterion '{}' timed out after {}ms (timeout={}ms), using default value: {}", - criterion.getName(), elapsedMs, timeoutMs, criterion.getDefaultValue()); + logger.warn("Criterion '{}' timed out after {}ms (timeout={}ms), sessionId={}, using default value: {}", + criterion.getName(), elapsedMs, timeoutMs, sessionId, criterion.getDefaultValue()); result = buildTimeoutResult(startTime, timeoutMs); } @@ -171,11 +173,11 @@ public Map apply(OverAllState state) { com.alibaba.assistant.agent.evaluation.observation.EvaluationObservationLifecycleListener .registerCriterionResult(criterion.getName(), result); - logger.info("Criterion {} completed with status: {}, value: {}", - criterion.getName(), result.getStatus(), result.getValue()); + logger.info("Criterion {} completed, sessionId={}, status={}, result={}", + criterion.getName(), sessionId, result.getStatus(), result.getValue()); } catch (Exception e) { - logger.error("Error executing criterion {}: {}", criterion.getName(), e.getMessage(), e); + logger.error("Error executing criterion {}, sessionId={}: {}", criterion.getName(), extractSessionIdFromState(state), e.getMessage(), e); // Create error result with default value if available CriterionResult errorResult = buildErrorResult(startTime, e.getMessage()); @@ -259,6 +261,14 @@ private Map buildDependencyResults(OverAllState state) return dependencyResults; } + private String extractSessionIdFromState(OverAllState state) { + Object evaluationContext = state != null ? state.data().get("evaluationContext") : null; + if (evaluationContext instanceof EvaluationContext context) { + return EvaluationLogContextHelper.getSessionId(context); + } + return null; + } + /** * Execute criterion without batching (original logic). */ @@ -298,23 +308,27 @@ private CriterionResult executeWithBatching(EvaluationContext evaluationContext, // Check if source is a collection if (!SourcePathResolver.isCollection(sourceObject)) { - logger.warn("Source path '{}' did not resolve to a collection for criterion '{}', falling back to non-batching execution", - batchingConfig.getSourcePath(), criterion.getName()); + logger.warn("Source path '{}' did not resolve to a collection for criterion '{}', sessionId={}, falling back to non-batching execution", + batchingConfig.getSourcePath(), criterion.getName(), EvaluationLogContextHelper.getSessionId(evaluationContext)); return executeWithoutBatching(evaluationContext, dependencyResults); } Collection sourceCollection = SourcePathResolver.toCollection(sourceObject); if (sourceCollection == null || sourceCollection.isEmpty()) { - logger.debug("Source collection is empty for criterion '{}', returning empty result", criterion.getName()); + logger.debug("Source collection is empty for criterion '{}', sessionId={}, returning empty result", + criterion.getName(), EvaluationLogContextHelper.getSessionId(evaluationContext)); return createEmptyCollectionResult(); } - logger.info("Criterion '{}': processing {} items with batchSize={}, maxConcurrentBatches={}", - criterion.getName(), sourceCollection.size(), batchingConfig.getBatchSize(), batchingConfig.getMaxConcurrentBatches()); + logger.info("Criterion '{}': processing {} items with batchSize={}, maxConcurrentBatches={}, sessionId={}", + criterion.getName(), sourceCollection.size(), batchingConfig.getBatchSize(), batchingConfig.getMaxConcurrentBatches(), + EvaluationLogContextHelper.getSessionId(evaluationContext)); // Split into batches List> batches = splitIntoBatches(sourceCollection, batchingConfig.getBatchSize()); - logger.debug("Split {} items into {} batches", sourceCollection.size(), batches.size()); + logger.debug("Split {} items into {} batches, criterion={}, sessionId={}", + sourceCollection.size(), batches.size(), criterion.getName(), + EvaluationLogContextHelper.getSessionId(evaluationContext)); // Process batches List allBatchResults = processBatches( @@ -328,7 +342,8 @@ private CriterionResult executeWithBatching(EvaluationContext evaluationContext, return aggregateResults(evaluationContext, dependencyResults, batchingConfig, allBatchResults); } catch (Exception e) { - logger.error("Error during batching execution for criterion '{}': {}", criterion.getName(), e.getMessage(), e); + logger.error("Error during batching execution for criterion '{}', sessionId={}: {}", + criterion.getName(), EvaluationLogContextHelper.getSessionId(evaluationContext), e.getMessage(), e); CriterionResult errorResult = new CriterionResult(); errorResult.setCriterionName(criterion.getName()); errorResult.setStatus(CriterionStatus.ERROR); diff --git a/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/observation/EvaluationObservationLifecycleListener.java b/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/observation/EvaluationObservationLifecycleListener.java index 7417bf4f..95206d00 100644 --- a/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/observation/EvaluationObservationLifecycleListener.java +++ b/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/observation/EvaluationObservationLifecycleListener.java @@ -16,6 +16,7 @@ package com.alibaba.assistant.agent.evaluation.observation; import com.alibaba.assistant.agent.evaluation.model.CriterionResult; +import com.alibaba.assistant.agent.evaluation.util.EvaluationLogContextHelper; import com.alibaba.cloud.ai.graph.GraphLifecycleListener; import com.alibaba.cloud.ai.graph.RunnableConfig; import io.opentelemetry.api.trace.Span; @@ -164,7 +165,8 @@ public EvaluationObservationLifecycleListener(Tracer tracer, Span parentSpan) { public void onStart(String nodeId, Map state, RunnableConfig config) { // 评估 Graph 开始时不需要特殊处理 if (nodeId != null && nodeId.equalsIgnoreCase("__start__")) { - log.debug("EvaluationObservationLifecycleListener#onStart - reason=评估Graph开始"); + log.debug("EvaluationObservationLifecycleListener#onStart - reason=评估Graph开始, sessionId={}", + EvaluationLogContextHelper.getSessionId(config)); } } @@ -172,13 +174,15 @@ public void onStart(String nodeId, Map state, RunnableConfig con public void onComplete(String nodeId, Map state, RunnableConfig config) { // 评估 Graph 完成时不需要特殊处理 if (nodeId != null && nodeId.equalsIgnoreCase("__end__")) { - log.debug("EvaluationObservationLifecycleListener#onComplete - reason=评估Graph完成"); + log.debug("EvaluationObservationLifecycleListener#onComplete - reason=评估Graph完成, sessionId={}", + EvaluationLogContextHelper.getSessionId(config)); } } @Override public void onError(String nodeId, Map state, Throwable ex, RunnableConfig config) { - log.error("EvaluationObservationLifecycleListener#onError - reason=评估执行出错, nodeId={}", nodeId, ex); + log.error("EvaluationObservationLifecycleListener#onError - reason=评估执行出错, nodeId={}, sessionId={}", + nodeId, EvaluationLogContextHelper.getSessionId(config), ex); // 停止对应节点的 Span String nodeKey = getNodeKey(nodeId, config); @@ -223,7 +227,8 @@ public void before(String nodeId, Map state, RunnableConfig conf nodeSpans.put(nodeKey, span); nodeScopes.put(nodeKey, scope); - log.debug("EvaluationObservationLifecycleListener#before - reason=评估项开始执行, criterionName={}", nodeId); + log.debug("EvaluationObservationLifecycleListener#before - reason=评估项开始执行, criterionName={}, sessionId={}", + nodeId, EvaluationLogContextHelper.getSessionId(config)); } /** @@ -309,10 +314,11 @@ public void after(String nodeId, Map state, RunnableConfig confi } log.info("EvaluationObservationLifecycleListener#after - reason=评估项执行完成, " + - "criterionName={}, status={}, durationMs={}", - nodeId, result.getStatus(), durationMs); + "criterionName={}, sessionId={}, status={}, durationMs={}", + nodeId, EvaluationLogContextHelper.getSessionId(config), result.getStatus(), durationMs); } else { - log.warn("EvaluationObservationLifecycleListener#after - reason=评估结果未找到, criterionName={}", nodeId); + log.warn("EvaluationObservationLifecycleListener#after - reason=评估结果未找到, criterionName={}, sessionId={}", + nodeId, EvaluationLogContextHelper.getSessionId(config)); } span.end(); @@ -360,4 +366,3 @@ private String truncate(String str, int maxLength) { return str.substring(0, maxLength) + "...[truncated]"; } } - diff --git a/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/util/EvaluationLogContextHelper.java b/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/util/EvaluationLogContextHelper.java new file mode 100644 index 00000000..b6abadfc --- /dev/null +++ b/assistant-agent-evaluation/src/main/java/com/alibaba/assistant/agent/evaluation/util/EvaluationLogContextHelper.java @@ -0,0 +1,31 @@ +package com.alibaba.assistant.agent.evaluation.util; + +import com.alibaba.assistant.agent.evaluation.model.EvaluationContext; +import com.alibaba.cloud.ai.graph.RunnableConfig; + +/** + * 评估链路日志上下文提取工具。 + */ +public final class EvaluationLogContextHelper { + + private EvaluationLogContextHelper() { + } + + public static String getSessionId(EvaluationContext context) { + if (context == null) { + return null; + } + Object sessionId = context.getEnvironmentValue("sessionId"); + if (sessionId == null) { + sessionId = context.getEnvironmentValue("threadId"); + } + return sessionId != null ? String.valueOf(sessionId) : null; + } + + public static String getSessionId(RunnableConfig config) { + if (config == null) { + return null; + } + return config.threadId().orElse(null); + } +} diff --git a/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/experience/ExperienceRetrievalEvaluatorFactory.java b/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/experience/ExperienceRetrievalEvaluatorFactory.java index 1edecd75..4132051a 100644 --- a/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/experience/ExperienceRetrievalEvaluatorFactory.java +++ b/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/experience/ExperienceRetrievalEvaluatorFactory.java @@ -19,6 +19,7 @@ import com.alibaba.assistant.agent.evaluation.model.CriterionExecutionContext; import com.alibaba.assistant.agent.evaluation.model.CriterionResult; import com.alibaba.assistant.agent.evaluation.model.CriterionStatus; +import com.alibaba.assistant.agent.evaluation.util.EvaluationLogContextHelper; import com.alibaba.assistant.agent.extension.experience.model.Experience; import com.alibaba.assistant.agent.extension.experience.model.ExperienceQuery; import com.alibaba.assistant.agent.extension.experience.model.ExperienceQueryContext; @@ -114,23 +115,29 @@ private static RuleBasedEvaluator createEvaluator( try { // 优先使用 enhanced_user_input,否则使用原始 userInput String queryText = extractEnhancedOrOriginalInput(ctx); + String sessionId = EvaluationLogContextHelper.getSessionId(ctx.getInputContext()); + log.info("ExperienceRetrievalEvaluatorFactory#{} - reason=开始检索经验, sessionId={}, queryText={}, experienceTypes={}", + evaluatorId, sessionId, queryText, experienceTypes); List experiences = queryExperiencesByStringIntersection( experienceProvider, + ctx, queryText, experienceTypes, maxExperiencesPerType ); if (experiences.isEmpty()) { - log.info("ExperienceRetrievalEvaluatorFactory#{} - reason=未检索到经验", evaluatorId); + log.info("ExperienceRetrievalEvaluatorFactory#{} - reason=未检索到经验, sessionId={}", evaluatorId, sessionId); result.setStatus(CriterionStatus.SUCCESS); result.setValue(""); result.getMetadata().put("is_empty", true); return result; } - log.info("ExperienceRetrievalEvaluatorFactory#{} - reason=检索到经验, count={}", evaluatorId, experiences.size()); + log.info("ExperienceRetrievalEvaluatorFactory#{} - reason=检索到经验, sessionId={}, count={}, experienceIds={}", + evaluatorId, sessionId, experiences.size(), + experiences.stream().map(Experience::getId).collect(Collectors.toList())); // 构建 ref_entries,每个经验作为一个独立的条目,experience ID 作为 ref-id List> refEntries = buildRefEntries(experiences); @@ -142,7 +149,8 @@ private static RuleBasedEvaluator createEvaluator( result.getMetadata().put("experience_count", experiences.size()); } catch (Exception e) { - log.error("ExperienceRetrievalEvaluatorFactory#{} - reason=经验检索失败", evaluatorId, e); + log.error("ExperienceRetrievalEvaluatorFactory#{} - reason=经验检索失败, sessionId={}", + evaluatorId, EvaluationLogContextHelper.getSessionId(ctx.getInputContext()), e); result.setStatus(CriterionStatus.ERROR); result.setErrorMessage(e.getMessage()); } @@ -181,6 +189,7 @@ private static String extractEnhancedOrOriginalInput(CriterionExecutionContext c */ private static List queryExperiencesByStringIntersection( ExperienceProvider experienceProvider, + CriterionExecutionContext executionContext, String queryText, List types, int maxExperiencesPerType) { @@ -191,6 +200,20 @@ private static List queryExperiencesByStringIntersection( ExperienceQueryContext queryContext = new ExperienceQueryContext(); queryContext.setUserQuery(queryText); + if (executionContext != null && executionContext.getInputContext() != null) { + Object sessionId = executionContext.getInputContext().getEnvironmentValue("sessionId"); + if (sessionId != null) { + queryContext.setSessionId(String.valueOf(sessionId)); + } + Object tenantId = executionContext.getInputContext().getEnvironmentValue("tenantId"); + if (tenantId != null) { + queryContext.setTenantId(String.valueOf(tenantId)); + } + Object userId = executionContext.getInputContext().getEnvironmentValue("userId"); + if (userId != null) { + queryContext.setUserId(String.valueOf(userId)); + } + } List allExperiences = new ArrayList<>(); @@ -318,4 +341,3 @@ private static String formatExperiences(List experiences, String pha return sb.toString(); } } - diff --git a/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/hook/BeforeAgentEvaluationHook.java b/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/hook/BeforeAgentEvaluationHook.java index 73edc434..bb06f6ac 100644 --- a/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/hook/BeforeAgentEvaluationHook.java +++ b/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/evaluation/hook/BeforeAgentEvaluationHook.java @@ -19,6 +19,7 @@ import com.alibaba.assistant.agent.evaluation.model.EvaluationContext; import com.alibaba.assistant.agent.evaluation.model.EvaluationResult; import com.alibaba.assistant.agent.evaluation.model.EvaluationSuite; +import com.alibaba.assistant.agent.evaluation.util.EvaluationLogContextHelper; import com.alibaba.assistant.agent.extension.evaluation.store.OverAllStateEvaluationResultStore; import com.alibaba.cloud.ai.graph.OverAllState; import com.alibaba.cloud.ai.graph.RunnableConfig; @@ -152,11 +153,16 @@ public CompletableFuture> beforeAgent(OverAllState state, Ru return CompletableFuture.completedFuture(Map.of()); } - log.info("BeforeAgentEvaluationHook#beforeAgent - reason=开始执行评估, suiteId={}, async={}", suiteId, async); + String runnableSessionId = EvaluationLogContextHelper.getSessionId(config); + log.info("BeforeAgentEvaluationHook#beforeAgent - reason=开始执行评估, suiteId={}, sessionId={}, async={}", + suiteId, runnableSessionId, async); try { // 1. 构造评估上下文(支持 config) EvaluationContext context = contextBuilder.apply(state, config); + String sessionId = EvaluationLogContextHelper.getSessionId(context); + log.info("BeforeAgentEvaluationHook#beforeAgent - reason=评估上下文已构建, suiteId={}, sessionId={}", + suiteId, sessionId); // 2. 加载评估套件 EvaluationSuite suite = loadSuite(); @@ -172,7 +178,8 @@ public CompletableFuture> beforeAgent(OverAllState state, Ru } } catch (Exception e) { - log.error("BeforeAgentEvaluationHook#beforeAgent - reason=评估执行失败, suiteId=" + suiteId, e); + log.error("BeforeAgentEvaluationHook#beforeAgent - reason=评估执行失败, suiteId={}, sessionId={}", + suiteId, runnableSessionId, e); return CompletableFuture.completedFuture(Map.of()); } } @@ -195,10 +202,10 @@ private EvaluationSuite loadSuite() { * 同步执行评估 */ private CompletableFuture> executeSync(OverAllState state, - EvaluationSuite suite, - EvaluationContext context) { + EvaluationSuite suite, + EvaluationContext context) { EvaluationResult result = evaluationService.evaluate(suite, context); - return CompletableFuture.completedFuture(buildResultMap(state, result)); + return CompletableFuture.completedFuture(buildResultMap(state, result, context)); } /** @@ -209,13 +216,14 @@ private CompletableFuture> executeAsync(OverAllState state, EvaluationContext context) { return evaluationService.evaluateAsync(suite, context) .orTimeout(timeoutMs, TimeUnit.MILLISECONDS) - .thenApply(result -> buildResultMap(state, result)) + .thenApply(result -> buildResultMap(state, result, context)) .exceptionally(e -> { if (e.getCause() instanceof TimeoutException) { - log.warn("BeforeAgentEvaluationHook#beforeAgent - reason=评估超时, suiteId={}, timeoutMs={}", - suiteId, timeoutMs); + log.warn("BeforeAgentEvaluationHook#beforeAgent - reason=评估超时, suiteId={}, sessionId={}, timeoutMs={}", + suiteId, EvaluationLogContextHelper.getSessionId(context), timeoutMs); } else { - log.error("BeforeAgentEvaluationHook#beforeAgent - reason=异步评估失败, suiteId=" + suiteId, e); + log.error("BeforeAgentEvaluationHook#beforeAgent - reason=异步评估失败, suiteId={}, sessionId={}", + suiteId, EvaluationLogContextHelper.getSessionId(context), e); } return Map.of(); }); @@ -224,9 +232,9 @@ private CompletableFuture> executeAsync(OverAllState state, /** * 构建结果 Map 并记录日志 */ - private Map buildResultMap(OverAllState state, EvaluationResult result) { - log.info("BeforeAgentEvaluationHook#beforeAgent - reason=评估完成, suiteId={}, statistics={}", - suiteId, result.getStatistics()); + private Map buildResultMap(OverAllState state, EvaluationResult result, EvaluationContext context) { + log.info("BeforeAgentEvaluationHook#beforeAgent - reason=评估完成, suiteId={}, sessionId={}, statistics={}", + suiteId, EvaluationLogContextHelper.getSessionId(context), result.getStatistics()); OverAllStateEvaluationResultStore store = new OverAllStateEvaluationResultStore(state); return store.createUpdateMap(suiteId, result); diff --git a/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/experience/model/ExperienceQueryContext.java b/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/experience/model/ExperienceQueryContext.java index ae34e780..050983ef 100644 --- a/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/experience/model/ExperienceQueryContext.java +++ b/assistant-agent-extensions/src/main/java/com/alibaba/assistant/agent/extension/experience/model/ExperienceQueryContext.java @@ -15,11 +15,21 @@ public class ExperienceQueryContext { */ private String userQuery; + /** + * 会话ID + */ + private String sessionId; + /** * 租户ID */ private String tenantId; + /** + * 用户ID + */ + private String userId; + public ExperienceQueryContext() { } @@ -31,6 +41,14 @@ public void setUserQuery(String userQuery) { this.userQuery = userQuery; } + public String getSessionId() { + return sessionId; + } + + public void setSessionId(String sessionId) { + this.sessionId = StringUtils.hasText(sessionId) ? sessionId.trim() : null; + } + public String getTenantId() { return tenantId; } @@ -39,4 +57,11 @@ public void setTenantId(String tenantId) { this.tenantId = StringUtils.hasText(tenantId) ? tenantId.trim() : null; } + public String getUserId() { + return userId; + } + + public void setUserId(String userId) { + this.userId = StringUtils.hasText(userId) ? userId.trim() : null; + } }