diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/AssistantResponseParser.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/AssistantResponseParser.java new file mode 100644 index 0000000..689c7ab --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/AssistantResponseParser.java @@ -0,0 +1,215 @@ +package it.cnr.isti.workflow.manager.assistant; + +import java.util.ArrayList; +import java.util.List; +import java.util.Locale; + +import org.springframework.http.HttpStatus; +import org.springframework.web.server.ResponseStatusException; + +import tools.jackson.core.json.JsonReadFeature; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.json.JsonMapper; + +import it.cnr.isti.workflow.manager.app.ObjectMapperHolder; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantConfiguredBlockDraft; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantConnectionDraft; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantFlowPlan; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.ParsedBlockDraft; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.ParsedConnections; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.ParsedPlan; + +final class AssistantResponseParser { + + private static final ObjectMapper LENIENT_ASSISTANT_MAPPER = JsonMapper.builder() + .enable(JsonReadFeature.ALLOW_UNQUOTED_PROPERTY_NAMES) + .enable(JsonReadFeature.ALLOW_SINGLE_QUOTES) + .enable(JsonReadFeature.ALLOW_TRAILING_COMMA) + .build(); + + private AssistantResponseParser() { + } + + static ParsedPlan parsePlan(String rawResponse) { + try { + JsonNode root = readJsonObject(rawResponse); + JsonNode planNode = root.has("plan") ? root.get("plan") : root; + AssistantFlowPlan plan = ObjectMapperHolder.mapper.treeToValue(planNode, AssistantFlowPlan.class); + return new ParsedPlan(plan, AssistantTextSupport.textOrEmpty(root.path("rationale"))); + } catch (Exception e) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an invalid plan payload: " + e.getMessage()); + } + } + + static ParsedBlockDraft parseBlockDraft(String rawResponse) { + try { + JsonNode root = readJsonObject(rawResponse); + JsonNode blockNode = root.has("block") ? root.get("block") : root; + AssistantConfiguredBlockDraft block = new AssistantConfiguredBlockDraft( + AssistantTextSupport.textOrNull(blockNode.path("blockId")), + AssistantTextSupport.textOrNull(blockNode.path("name")), + blockNode.path("config")); + if (block.blockId() == null || block.config().isMissingNode()) { + throw new IllegalArgumentException("Missing blockId or config"); + } + return new ParsedBlockDraft(block, AssistantTextSupport.textOrEmpty(root.path("rationale"))); + } catch (Exception e) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an invalid block configuration payload: " + e.getMessage()); + } + } + + static ParsedConnections parseConnections(String rawResponse) { + try { + JsonNode root = readJsonObjectOrArray(rawResponse); + if (root.isArray()) { + List arrayConnections = new ArrayList<>(); + for (JsonNode node : root) { + arrayConnections.add(ObjectMapperHolder.mapper.treeToValue(node, AssistantConnectionDraft.class)); + } + return new ParsedConnections(arrayConnections, ""); + } + JsonNode connectionsNode = root.has("connections") ? root.get("connections") : root.path("connections"); + List connections = new ArrayList<>(); + if (connectionsNode.isArray()) { + for (JsonNode node : connectionsNode) { + connections.add(ObjectMapperHolder.mapper.treeToValue(node, AssistantConnectionDraft.class)); + } + } + return new ParsedConnections(connections, AssistantTextSupport.textOrEmpty(root.path("rationale"))); + } catch (IllegalArgumentException e) { + if (isLikelyNoConnectionsText(rawResponse)) { + return new ParsedConnections(List.of(), rawResponse == null ? "" : rawResponse.trim()); + } + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an invalid connections payload: " + e.getMessage()); + } catch (Exception e) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an invalid connections payload: " + e.getMessage()); + } + } + + static JsonNode readJsonObject(String rawResponse) { + String json = extractJsonObject(rawResponse); + try { + return ObjectMapperHolder.mapper.readTree(json); + } catch (Exception e) { + try { + return LENIENT_ASSISTANT_MAPPER.readTree(json); + } catch (Exception ignored) { + throw new IllegalArgumentException(e.getMessage(), e); + } + } + } + + static JsonNode readJsonObjectOrArray(String rawResponse) { + String json = extractJsonObjectOrArray(rawResponse); + try { + return ObjectMapperHolder.mapper.readTree(json); + } catch (Exception e) { + try { + return LENIENT_ASSISTANT_MAPPER.readTree(json); + } catch (Exception ignored) { + throw new IllegalArgumentException(e.getMessage(), e); + } + } + } + + static String extractJsonObject(String rawResponse) { + if (rawResponse == null || rawResponse.isBlank()) { + throw new IllegalArgumentException("Empty assistant response"); + } + + String trimmed = rawResponse.trim(); + boolean sawObjectStart = false; + for (int i = 0; i < trimmed.length(); i++) { + if (trimmed.charAt(i) != '{') { + continue; + } + sawObjectStart = true; + String candidate = tryExtractBalancedJson(trimmed, i); + if (candidate != null) { + return candidate; + } + } + if (!sawObjectStart) { + throw new IllegalArgumentException("No JSON object found in assistant response"); + } + throw new IllegalArgumentException("Incomplete JSON object found in assistant response"); + } + + static String extractJsonObjectOrArray(String rawResponse) { + if (rawResponse == null || rawResponse.isBlank()) { + throw new IllegalArgumentException("Empty assistant response"); + } + + String trimmed = rawResponse.trim(); + boolean sawJsonStart = false; + for (int i = 0; i < trimmed.length(); i++) { + char current = trimmed.charAt(i); + if (current != '{' && current != '[') { + continue; + } + sawJsonStart = true; + String candidate = tryExtractBalancedJson(trimmed, i); + if (candidate != null) { + return candidate; + } + } + if (!sawJsonStart) { + throw new IllegalArgumentException("No JSON object found in assistant response"); + } + throw new IllegalArgumentException("Incomplete JSON object found in assistant response"); + } + + static String tryExtractBalancedJson(String text, int start) { + int objectDepth = 0; + int arrayDepth = 0; + boolean inString = false; + boolean escaped = false; + for (int i = start; i < text.length(); i++) { + char current = text.charAt(i); + if (escaped) { + escaped = false; + continue; + } + if (current == '\\' && inString) { + escaped = true; + continue; + } + if (current == '"') { + inString = !inString; + continue; + } + if (inString) { + continue; + } + if (current == '{') { + objectDepth++; + } else if (current == '}') { + objectDepth--; + } else if (current == '[') { + arrayDepth++; + } else if (current == ']') { + arrayDepth--; + } + if (objectDepth == 0 && arrayDepth == 0) { + return text.substring(start, i + 1); + } + } + return null; + } + + static boolean isLikelyNoConnectionsText(String rawResponse) { + if (rawResponse == null || rawResponse.isBlank()) { + return false; + } + String normalized = rawResponse.trim().toLowerCase(Locale.ROOT); + return normalized.contains("no connection") + || normalized.contains("no connections") + || normalized.contains("none needed") + || normalized.contains("nessuna connessione"); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java index 4e2a30b..646cdb0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java @@ -24,10 +24,8 @@ import org.springframework.stereotype.Service; import org.springframework.beans.factory.annotation.Value; import org.springframework.web.server.ResponseStatusException; -import tools.jackson.core.json.JsonReadFeature; import tools.jackson.databind.JsonNode; import tools.jackson.databind.ObjectMapper; -import tools.jackson.databind.json.JsonMapper; import tools.jackson.databind.node.ArrayNode; import tools.jackson.databind.node.ObjectNode; @@ -157,11 +155,6 @@ public class FlowAssistantService { ValidationErrorCode.EXECUTION_DEADLOCK); private static final long DEFAULT_RETRY_BASE_DELAY_MILLIS = 120L; private static final long DEFAULT_RETRY_MAX_DELAY_MILLIS = 800L; - private static final ObjectMapper LENIENT_ASSISTANT_MAPPER = JsonMapper.builder() - .enable(JsonReadFeature.ALLOW_UNQUOTED_PROPERTY_NAMES) - .enable(JsonReadFeature.ALLOW_SINGLE_QUOTES) - .enable(JsonReadFeature.ALLOW_TRAILING_COMMA) - .build(); @FunctionalInterface public interface ProgressListener { @@ -176,14 +169,14 @@ public class FlowAssistantService { private static final ProgressListener NOOP_PROGRESS = (phase, message) -> { }; - private record AssistantFlowPlan(String name, String description, List blocks, + record AssistantFlowPlan(String name, String description, List blocks, List containers) { - private AssistantFlowPlan(String name, String description, List blocks) { + AssistantFlowPlan(String name, String description, List blocks) { this(name, description, blocks, List.of()); } } - private record AssistantBlockPlan(String blockId, String blockType, String purpose, String operation) { + record AssistantBlockPlan(String blockId, String blockType, String purpose, String operation) { } /** @@ -204,34 +197,34 @@ public class FlowAssistantService { * names the open non-multiple body input that receives the guard's feedback each iteration - * required when the body exposes more than one open input, inferred when it exposes exactly one. */ - private record AssistantContainerPlan(String containerId, String containerType, String purpose, String operation, + record AssistantContainerPlan(String containerId, String containerType, String purpose, String operation, String iterationInput, String guardCondition, Integer maxIterations, String feedbackInput, List blocks) { } - private enum PlanOperation { + enum PlanOperation { KEEP, ADD, UPDATE, REMOVE } - private record AssistantConfiguredBlockDraft(String blockId, String name, JsonNode config) { + record AssistantConfiguredBlockDraft(String blockId, String name, JsonNode config) { } - private record AssistantConnectionDraft(String fromBlockId, String fromOutput, String toBlockId, String toInput) { + record AssistantConnectionDraft(String fromBlockId, String fromOutput, String toBlockId, String toInput) { } - private record ParsedPlan(AssistantFlowPlan plan, String rationale) { + record ParsedPlan(AssistantFlowPlan plan, String rationale) { } - private record ParsedBlockDraft(AssistantConfiguredBlockDraft block, String rationale) { + record ParsedBlockDraft(AssistantConfiguredBlockDraft block, String rationale) { } private record ConfiguredBlockResult(ParsedBlockDraft parsedBlock, Block block) { } - private record ParsedConnections(List connections, String rationale) { + record ParsedConnections(List connections, String rationale) { } private record ConfiguredBlockSummary(String blockId, String blockType, String name, String purpose, List inputs, @@ -488,7 +481,7 @@ public class FlowAssistantService { mcpServerCatalogEntries()); parsedPlan = invokeStructuredAndValidate(provider, authorization, planningModelFor(mode, phaseModels), phaseModels.repairModel(), planPrompt, "plan", rawResponse -> { - ParsedPlan plan = parsePlan(rawResponse); + ParsedPlan plan = AssistantResponseParser.parsePlan(rawResponse); AssistantFlowPlan normalizedPlan = validateAndNormalizePlan(plan.plan(), mode, userPrompt, currentFlow, errors, catalogByType.keySet()); return new ParsedPlan(normalizedPlan, plan.rationale()); @@ -543,7 +536,7 @@ public class FlowAssistantService { ConfiguredBlockResult configuredBlock = invokeStructuredAndValidate(provider, authorization, jsonModelFor(mode, phaseModels), phaseModels.repairModel(), blockPrompt, "block configuration for " + blockPlan.blockId(), rawResponse -> { - ParsedBlockDraft parsedBlock = parseBlockDraft(rawResponse); + ParsedBlockDraft parsedBlock = AssistantResponseParser.parseBlockDraft(rawResponse); AssistantConfiguredBlockDraft normalizedDraft = normalizeBlockDraft(blockPlan, parsedBlock.block()); Block newBlock = buildBlock(descriptor, blockPlan, normalizedDraft, flowProvider, flowModel, @@ -745,7 +738,7 @@ public class FlowAssistantService { jsonModelFor(mode, phaseModels), phaseModels.repairModel(), blockPrompt, "block configuration for " + blockPlan.blockId() + " in container " + containerPlan.containerId(), rawResponse -> { - ParsedBlockDraft parsedBlock = parseBlockDraft(rawResponse); + ParsedBlockDraft parsedBlock = AssistantResponseParser.parseBlockDraft(rawResponse); AssistantConfiguredBlockDraft normalizedDraft = normalizeBlockDraft(blockPlan, parsedBlock.block()); Block newBlock = buildBlock(descriptor, blockPlan, normalizedDraft, flowProvider, flowModel, innerRequiresSharedMemory, currentBlockIndex, blockPlanCount); @@ -1041,7 +1034,7 @@ public class FlowAssistantService { private ParsedConnections parseConnectionsOrInferSequential(String rawResponse, List> assembledBlocks) { try { - return parseConnections(rawResponse); + return AssistantResponseParser.parseConnections(rawResponse); } catch (ResponseStatusException e) { if (e.getReason() == null || !e.getReason().contains("invalid connections payload")) { throw e; @@ -2289,91 +2282,6 @@ public class FlowAssistantService { } } - private ParsedPlan parsePlan(String rawResponse) { - try { - JsonNode root = readJsonObject(rawResponse); - JsonNode planNode = root.has("plan") ? root.get("plan") : root; - AssistantFlowPlan plan = ObjectMapperHolder.mapper.treeToValue(planNode, AssistantFlowPlan.class); - return new ParsedPlan(plan, AssistantTextSupport.textOrEmpty(root.path("rationale"))); - } catch (Exception e) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an invalid plan payload: " + e.getMessage()); - } - } - - private ParsedBlockDraft parseBlockDraft(String rawResponse) { - try { - JsonNode root = readJsonObject(rawResponse); - JsonNode blockNode = root.has("block") ? root.get("block") : root; - AssistantConfiguredBlockDraft block = new AssistantConfiguredBlockDraft( - AssistantTextSupport.textOrNull(blockNode.path("blockId")), - AssistantTextSupport.textOrNull(blockNode.path("name")), - blockNode.path("config")); - if (block.blockId() == null || block.config().isMissingNode()) { - throw new IllegalArgumentException("Missing blockId or config"); - } - return new ParsedBlockDraft(block, AssistantTextSupport.textOrEmpty(root.path("rationale"))); - } catch (Exception e) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an invalid block configuration payload: " + e.getMessage()); - } - } - - private ParsedConnections parseConnections(String rawResponse) { - try { - JsonNode root = readJsonObjectOrArray(rawResponse); - if (root.isArray()) { - List arrayConnections = new ArrayList<>(); - for (JsonNode node : root) { - arrayConnections.add(ObjectMapperHolder.mapper.treeToValue(node, AssistantConnectionDraft.class)); - } - return new ParsedConnections(arrayConnections, ""); - } - JsonNode connectionsNode = root.has("connections") ? root.get("connections") : root.path("connections"); - List connections = new ArrayList<>(); - if (connectionsNode.isArray()) { - for (JsonNode node : connectionsNode) { - connections.add(ObjectMapperHolder.mapper.treeToValue(node, AssistantConnectionDraft.class)); - } - } - return new ParsedConnections(connections, AssistantTextSupport.textOrEmpty(root.path("rationale"))); - } catch (IllegalArgumentException e) { - if (isLikelyNoConnectionsText(rawResponse)) { - return new ParsedConnections(List.of(), rawResponse == null ? "" : rawResponse.trim()); - } - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an invalid connections payload: " + e.getMessage()); - } catch (Exception e) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an invalid connections payload: " + e.getMessage()); - } - } - - private JsonNode readJsonObject(String rawResponse) { - String json = extractJsonObject(rawResponse); - try { - return ObjectMapperHolder.mapper.readTree(json); - } catch (Exception e) { - try { - return LENIENT_ASSISTANT_MAPPER.readTree(json); - } catch (Exception ignored) { - throw new IllegalArgumentException(e.getMessage(), e); - } - } - } - - private JsonNode readJsonObjectOrArray(String rawResponse) { - String json = extractJsonObjectOrArray(rawResponse); - try { - return ObjectMapperHolder.mapper.readTree(json); - } catch (Exception e) { - try { - return LENIENT_ASSISTANT_MAPPER.readTree(json); - } catch (Exception ignored) { - throw new IllegalArgumentException(e.getMessage(), e); - } - } - } private AssistantFlowPlan validateAndNormalizePlan(AssistantFlowPlan plan, OperationMode mode, String userPrompt, FlowCreateRequest currentFlow, List errors, @@ -2856,101 +2764,6 @@ public class FlowAssistantService { return new ValidationError("flow", null, violation.getPropertyPath().toString(), violation.getMessage()); } - private String extractJsonObject(String rawResponse) { - if (rawResponse == null || rawResponse.isBlank()) { - throw new IllegalArgumentException("Empty assistant response"); - } - - String trimmed = rawResponse.trim(); - boolean sawObjectStart = false; - for (int i = 0; i < trimmed.length(); i++) { - if (trimmed.charAt(i) != '{') { - continue; - } - sawObjectStart = true; - String candidate = tryExtractBalancedJson(trimmed, i); - if (candidate != null) { - return candidate; - } - } - if (!sawObjectStart) { - throw new IllegalArgumentException("No JSON object found in assistant response"); - } - throw new IllegalArgumentException("Incomplete JSON object found in assistant response"); - } - - private String extractJsonObjectOrArray(String rawResponse) { - if (rawResponse == null || rawResponse.isBlank()) { - throw new IllegalArgumentException("Empty assistant response"); - } - - String trimmed = rawResponse.trim(); - boolean sawJsonStart = false; - for (int i = 0; i < trimmed.length(); i++) { - char current = trimmed.charAt(i); - if (current != '{' && current != '[') { - continue; - } - sawJsonStart = true; - String candidate = tryExtractBalancedJson(trimmed, i); - if (candidate != null) { - return candidate; - } - } - if (!sawJsonStart) { - throw new IllegalArgumentException("No JSON object found in assistant response"); - } - throw new IllegalArgumentException("Incomplete JSON object found in assistant response"); - } - - private String tryExtractBalancedJson(String text, int start) { - int objectDepth = 0; - int arrayDepth = 0; - boolean inString = false; - boolean escaped = false; - for (int i = start; i < text.length(); i++) { - char current = text.charAt(i); - if (escaped) { - escaped = false; - continue; - } - if (current == '\\' && inString) { - escaped = true; - continue; - } - if (current == '"') { - inString = !inString; - continue; - } - if (inString) { - continue; - } - if (current == '{') { - objectDepth++; - } else if (current == '}') { - objectDepth--; - } else if (current == '[') { - arrayDepth++; - } else if (current == ']') { - arrayDepth--; - } - if (objectDepth == 0 && arrayDepth == 0) { - return text.substring(start, i + 1); - } - } - return null; - } - - private boolean isLikelyNoConnectionsText(String rawResponse) { - if (rawResponse == null || rawResponse.isBlank()) { - return false; - } - String normalized = rawResponse.trim().toLowerCase(Locale.ROOT); - return normalized.contains("no connection") - || normalized.contains("no connections") - || normalized.contains("none needed") - || normalized.contains("nessuna connessione"); - } private ObjectNode llmDescriptorNode(String provider, String model) { ObjectNode llmDescriptor = ObjectMapperHolder.mapper.createObjectNode();