From 3d8b5f1d2147c75be40e0f8b6d160bfaead924ff Mon Sep 17 00:00:00 2001 From: zhangzherui Date: Fri, 24 Jul 2026 15:49:08 +0800 Subject: [PATCH] fix(planner): handle concatenated plan responses --- .../ai/dataagent/util/PlanProcessUtil.java | 24 ++++++++++++++++++- .../dataagent/util/PlanProcessUtilTest.java | 12 ++++++++++ 2 files changed, 35 insertions(+), 1 deletion(-) diff --git a/data-agent-management/src/main/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtil.java b/data-agent-management/src/main/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtil.java index 8ab6a8553..e863e6a7b 100644 --- a/data-agent-management/src/main/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtil.java +++ b/data-agent-management/src/main/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtil.java @@ -18,9 +18,11 @@ import com.alibaba.cloud.ai.graph.OverAllState; import com.alibaba.cloud.ai.dataagent.dto.planner.ExecutionStep; import com.alibaba.cloud.ai.dataagent.dto.planner.Plan; +import com.fasterxml.jackson.databind.MappingIterator; import org.springframework.ai.converter.BeanOutputConverter; import org.springframework.core.ParameterizedTypeReference; +import java.io.IOException; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -101,13 +103,33 @@ public static ExecutionStep getCurrentExecutionStep(Plan plan, Integer currentSt public static Plan getPlan(OverAllState state) { String plannerNodeOutput = (String) state.value(PLANNER_NODE_OUTPUT) .orElseThrow(() -> new IllegalStateException("计划节点输出为空")); - Plan plan = converter.convert(plannerNodeOutput); + Plan plan = convertPlan(plannerNodeOutput); if (plan == null) { throw new IllegalStateException("计划解析失败"); } return plan; } + private static Plan convertPlan(String plannerNodeOutput) { + try (MappingIterator plans = JsonUtil.getObjectMapper() + .readerFor(Plan.class) + .readValues(plannerNodeOutput)) { + Plan lastPlan = null; + int planCount = 0; + while (plans.hasNextValue()) { + lastPlan = plans.nextValue(); + planCount++; + } + if (planCount > 1) { + return lastPlan; + } + } + catch (IOException ignored) { + // Let the existing converter handle its supported wrappers and report errors. + } + return converter.convert(plannerNodeOutput); + } + /** * Get the current step number from state * @param state the overall state diff --git a/data-agent-management/src/test/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtilTest.java b/data-agent-management/src/test/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtilTest.java index d6f2dcd76..dfbf960fb 100644 --- a/data-agent-management/src/test/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtilTest.java +++ b/data-agent-management/src/test/java/com/alibaba/cloud/ai/dataagent/util/PlanProcessUtilTest.java @@ -42,6 +42,18 @@ void getPlan_withValidJson_returnsParsedPlan() { assertEquals(SQL_GENERATE_NODE, plan.getExecutionPlan().get(0).getToolToUse()); } + @Test + void getPlan_withConcatenatedPlans_returnsLastCompletePlan() { + String plannerOutput = TestFixtures.createSingleSqlPlanJson() + TestFixtures.createMultiStepPlanJson(); + OverAllState state = TestFixtures.createStateWith(Map.of(PLANNER_NODE_OUTPUT, plannerOutput)); + + Plan plan = PlanProcessUtil.getPlan(state); + + assertEquals("Multi-step analysis", plan.getThoughtProcess()); + assertEquals(3, plan.getExecutionPlan().size()); + assertEquals(REPORT_GENERATOR_NODE, plan.getExecutionPlan().get(2).getToolToUse()); + } + @Test void getPlan_withEmptyState_throwsException() { OverAllState state = TestFixtures.createStateWith(PLANNER_NODE_OUTPUT);