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 d67d65c..b68b43d 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 @@ -2,6 +2,7 @@ package it.cnr.isti.workflow.manager.assistant; import java.util.ArrayList; import java.util.Collection; +import java.util.EnumSet; import java.util.LinkedHashMap; import java.util.LinkedHashSet; import java.util.List; @@ -12,6 +13,7 @@ import java.util.Set; import java.util.UUID; import java.util.regex.Matcher; import java.util.regex.Pattern; +import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -47,6 +49,7 @@ import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; import it.cnr.isti.workflow.manager.flows.validation.ValidationError; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode; import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCodec; import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.ios.IOType; @@ -69,6 +72,63 @@ public class FlowAssistantService { // GLOBAL_INPUT_NOT_DECLARED. The assistant has no way to satisfy that on its own, so this // is applied once per assembled flow to auto-declare whatever the configured blocks reference. private static final Pattern GLOBAL_PLACEHOLDER_PATTERN = Pattern.compile("\\$\\{\\{global\\.([^}]+)}}"); + // Error codes that imply the flow's shape itself (block set, block I/O, connections, + // dependencies, lanes, shared-session wiring) may need to change. Any of these - or any + // error not cleanly scoped to one existing block - forces a full replan+reconnect. Every + // other block-scoped error (mainly plain VALIDATION_ERROR bean-validation violations on a + // block's own specificConfiguration) can be fixed by reconfiguring just that block, so the + // PLAN and CONNECTIONS phases are skipped entirely for those FIX rounds. If that guess turns + // out wrong (e.g. fixing a block incidentally changes its I/O), the next round's validation + // will surface a "sensitive" code and fall back to a full repair - never worse than before. + private static final Set PLAN_SENSITIVE_ERROR_CODES = EnumSet.of( + ValidationErrorCode.FLOW_CONTAINS_NULL_NODE, + ValidationErrorCode.NODE_ID_REQUIRED, + ValidationErrorCode.UNSUPPORTED_NODE_TYPE, + ValidationErrorCode.BLOCK_CONFIGURATION_MISSING, + ValidationErrorCode.BLOCK_TYPE_MISMATCH, + ValidationErrorCode.BLOCK_INPUTS_MISMATCH, + ValidationErrorCode.BLOCK_OUTPUTS_MISMATCH, + ValidationErrorCode.CONTAINER_CONFIGURATION_MISSING, + ValidationErrorCode.CONTAINER_SUBFLOW_INVALID, + ValidationErrorCode.CONTAINER_SUBFLOW_EMPTY, + ValidationErrorCode.NESTED_CONTAINERS_NOT_SUPPORTED, + ValidationErrorCode.INTERACTIVE_NODES_IN_CONTAINER_NOT_SUPPORTED, + ValidationErrorCode.CONTAINER_TYPE_MISMATCH, + ValidationErrorCode.CONTAINER_INPUTS_MISMATCH, + ValidationErrorCode.CONTAINER_OUTPUTS_MISMATCH, + ValidationErrorCode.DUPLICATE_NODE_ID, + ValidationErrorCode.NULL_CONNECTION, + ValidationErrorCode.CONNECTION_SOURCE_NODE_NOT_FOUND, + ValidationErrorCode.CONNECTION_TARGET_NODE_NOT_FOUND, + ValidationErrorCode.CONNECTION_SOURCE_OUTPUT_NOT_FOUND, + ValidationErrorCode.CONNECTION_TARGET_INPUT_NOT_FOUND, + ValidationErrorCode.NODE_INCOMING_CONNECTION_NOT_ALLOWED, + ValidationErrorCode.NODE_OUTGOING_CONNECTION_NOT_ALLOWED, + ValidationErrorCode.NULL_DEPENDENCY, + ValidationErrorCode.DEPENDENCY_SOURCE_REQUIRED, + ValidationErrorCode.DEPENDENCY_TARGET_REQUIRED, + ValidationErrorCode.DEPENDENCY_SOURCE_NODE_NOT_FOUND, + ValidationErrorCode.DEPENDENCY_TARGET_NODE_NOT_FOUND, + ValidationErrorCode.DEPENDENCY_SELF_REFERENCE, + ValidationErrorCode.DUPLICATE_DEPENDENCY, + ValidationErrorCode.NODE_CANNOT_DEPEND_ON_OTHER_NODES, + ValidationErrorCode.NODE_CANNOT_HAVE_DEPENDENT_NODES, + ValidationErrorCode.EXCLUSIVE_BRANCH_MERGE, + ValidationErrorCode.BRANCH_REJOIN_MIN_INPUTS, + ValidationErrorCode.BRANCH_REJOIN_INPUT_NOT_CONNECTED, + ValidationErrorCode.BRANCH_REJOIN_INPUT_MULTIPLE_CONNECTIONS, + ValidationErrorCode.BRANCH_REJOIN_MULTIPLE_VALUES, + ValidationErrorCode.BRANCH_REJOIN_INPUT_UNAVAILABLE, + ValidationErrorCode.UNSTRUCTURED_BRANCH_REJOIN, + ValidationErrorCode.END_NODE_OUTGOING_CONNECTION, + ValidationErrorCode.END_NODE_MULTIPLE_INCOMING_CONNECTIONS, + ValidationErrorCode.DUPLICATE_LANE_ID, + ValidationErrorCode.LANE_NAME_REQUIRED, + ValidationErrorCode.NODE_LANE_NOT_FOUND, + ValidationErrorCode.LANE_ORDER_INVALID, + ValidationErrorCode.SHARED_SESSION_NAME_NOT_UNIQUE, + ValidationErrorCode.SHARED_SESSION_PRODUCER_UNREACHABLE, + 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() @@ -290,15 +350,22 @@ public class FlowAssistantService { catalogByType.put(descriptor.type(), descriptor); } + boolean targetedRepair = mode == OperationMode.FIX && isTargetedBlockRepairEligible(currentFlow, errors); + progressListener.onProgress("planning", "Planning workflow blocks"); - String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog); - ParsedPlan parsedPlan = invokeStructuredAndValidate(provider, planningModelFor(mode, phaseModels), - phaseModels.repairModel(), planPrompt, "plan", rawResponse -> { - ParsedPlan plan = parsePlan(rawResponse); - AssistantFlowPlan normalizedPlan = validateAndNormalizePlan(plan.plan(), mode, userPrompt, currentFlow, - errors, catalogByType.keySet()); - return new ParsedPlan(normalizedPlan, plan.rationale()); - }); + ParsedPlan parsedPlan; + if (targetedRepair) { + parsedPlan = buildReusedPlanForTargetedRepair(currentFlow, errors); + } else { + String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog); + parsedPlan = invokeStructuredAndValidate(provider, planningModelFor(mode, phaseModels), + phaseModels.repairModel(), planPrompt, "plan", rawResponse -> { + ParsedPlan plan = parsePlan(rawResponse); + AssistantFlowPlan normalizedPlan = validateAndNormalizePlan(plan.plan(), mode, userPrompt, currentFlow, + errors, catalogByType.keySet()); + return new ParsedPlan(normalizedPlan, plan.rationale()); + }); + } boolean requireSharedMemorySemantics = isSharedMemoryContext(userPrompt, currentFlow, parsedPlan.plan()); List> assembledBlocks = new ArrayList<>(); @@ -378,7 +445,10 @@ public class FlowAssistantService { progressListener.onProgress("connecting_blocks", "Connecting configured blocks"); ParsedConnections parsedConnections; - if (assembledBlocks.size() < 2) { + if (targetedRepair) { + parsedConnections = new ParsedConnections(List.of(), + "Targeted repair: only the flagged blocks were reconfigured, connections are unchanged."); + } else if (assembledBlocks.size() < 2) { parsedConnections = new ParsedConnections(List.of(), "No connections needed for a single-block flow."); } else { String connectionsPrompt = promptService.buildConnectionsPrompt(mode, userPrompt, parsedPlan.plan(), @@ -889,6 +959,52 @@ public class FlowAssistantService { } } + /** + * True when every reported error is cleanly scoped to one already-existing block (entity + * "block", id matching a block currentFlow already has) and none of them is a + * PLAN_SENSITIVE_ERROR_CODES code. In that case the flow's shape (block set, I/O, + * connections) doesn't need to change - only the flagged blocks' own configuration does - + * so the PLAN and CONNECTIONS phases can be skipped for this FIX round. + */ + private boolean isTargetedBlockRepairEligible(FlowCreateRequest currentFlow, List errors) { + if (errors == null || errors.isEmpty() || !hasCurrentFlowBlocks(currentFlow)) { + return false; + } + Set existingBlockIds = currentFlow.flow().getBlocks().stream() + .map(Block::getId) + .collect(Collectors.toSet()); + for (ValidationError error : errors) { + if (PLAN_SENSITIVE_ERROR_CODES.contains(error.code())) { + return false; + } + if (!"block".equals(error.entity()) || error.id() == null || !existingBlockIds.contains(error.id())) { + return false; + } + } + return true; + } + + /** + * Reuses currentFlow's existing block list as the plan for a targeted FIX round, marking + * only the blocks referenced by an error as UPDATE (so they get reconfigured) and every + * other block as KEEP - without any PLAN LLM call. + */ + private ParsedPlan buildReusedPlanForTargetedRepair(FlowCreateRequest currentFlow, List errors) { + Set erroringBlockIds = errors.stream() + .map(ValidationError::id) + .filter(Objects::nonNull) + .collect(Collectors.toSet()); + List blockPlans = currentFlow.flow().getBlocks().stream() + .map(block -> new AssistantBlockPlan( + block.getId(), + blockTypeName(block), + block.getName(), + erroringBlockIds.contains(block.getId()) ? PlanOperation.UPDATE.name() : PlanOperation.KEEP.name())) + .toList(); + AssistantFlowPlan plan = new AssistantFlowPlan(currentFlow.name(), currentFlow.description(), blockPlans); + return new ParsedPlan(plan, "Targeted repair: reconfiguring only the blocks flagged by validation errors."); + } + private Block resolveExistingBlockForPlan(AssistantBlockPlan blockPlan, FlowCreateRequest currentFlow, int planIndex, Set usedExistingBlockIds) { List> existingBlocks = currentFlow == null || currentFlow.flow() == null diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java index 28da023..f4c08bd 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java @@ -39,10 +39,12 @@ import it.cnr.isti.workflow.manager.assistant.model.AssistantSessionView; import it.cnr.isti.workflow.manager.auth.config.JwtUtil; import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.EndBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractiveBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.types.EndBlockType; import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType; import it.cnr.isti.workflow.manager.blocks.types.HumanInteractionBlockType; import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; @@ -50,6 +52,8 @@ import it.cnr.isti.workflow.manager.blocks.types.MCPAgentBlockType; import it.cnr.isti.workflow.manager.flows.model.Connection; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.validation.ValidationError; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode; import it.cnr.isti.workflow.manager.ios.IODescriptor; import it.cnr.isti.workflow.manager.ios.IOType; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; @@ -373,6 +377,68 @@ public class AssistantControllerTest { assertEquals("response", response.flow().flow().getBlocks().getFirst().getOutputs().getFirst().getName()); } + @Test + public void fixSkipsPlanAndConnectionsForBlockScopedError() { + FlowCreateRequest brokenFlow = TestAssistantResponses.twoBlockFlowWithOverlongEndOutcomeLabel(); + Block endBlock = brokenFlow.flow().getBlocks().stream() + .filter(block -> "Done".equals(block.getName())) + .findFirst() + .orElseThrow(); + + // Only mocking BLOCK_CONFIG for the flagged block: if the targeted-repair path failed to + // skip PLAN/CONNECTIONS, the mock's fallback would throw and fail this test. + Answer answer = invocation -> { + String prompt = invocation.getArgument(1, String.class); + if (prompt.contains("TASK: BLOCK_CONFIG")) { + return TestAssistantResponses.wrap(java.util.Map.of( + "rationale", "Shortened the overlong outcome label.", + "block", java.util.Map.of( + "blockId", endBlock.getId(), + "name", "Done", + "config", java.util.Map.of( + "outcomeCode", "DONE", + "outcomeLabel", "Done")))); + } + throw new IllegalStateException("Unexpected assistant prompt (expected only BLOCK_CONFIG):\n" + prompt); + }; + Mockito.when(internalOllamaLLMProvider.generate(Mockito.eq(MODEL), Mockito.anyString())).thenAnswer(answer); + Mockito.when(internalOllamaLLMProvider.generateJson(Mockito.eq(MODEL), Mockito.anyString())).thenAnswer(answer); + + java.util.List validationErrors = java.util.List.of(new ValidationError( + ValidationErrorCode.VALIDATION_ERROR, "block", endBlock.getId(), + "specificConfiguration.outcomeLabel", "size must be between 0 and 255")); + + AssistantFlowResponse response = assistantController.fix( + new AssistantFixRequest( + "Fix the overlong outcome label", + brokenFlow, + validationErrors, + MODEL, + 1)); + + assertTrue(response.valid()); + assertTrue(response.validationErrors().isEmpty()); + // fix() already starts in FIX mode, so resolving everything on that very first + // assembleFlow pass costs zero *extra* repair rounds. + assertEquals(0, response.repairAttempts()); + assertEquals(2, response.flow().flow().getBlocks().size()); + assertEquals(1, response.flow().flow().getConnections().size()); + + Block classifierBlock = response.flow().flow().getBlocks().stream() + .filter(block -> "Ticket classifier".equals(block.getName())) + .findFirst() + .orElseThrow(); + LLMBlockConfiguration classifierConfiguration = (LLMBlockConfiguration) classifierBlock.getSpecificConfiguration(); + assertEquals("Classify the ticket: ${{ticket}}", classifierConfiguration.getPrompt()); + + Block fixedEndBlock = response.flow().flow().getBlocks().stream() + .filter(block -> "Done".equals(block.getName())) + .findFirst() + .orElseThrow(); + EndBlockConfiguration fixedConfiguration = (EndBlockConfiguration) fixedEndBlock.getSpecificConfiguration(); + assertEquals("Done", fixedConfiguration.getOutcomeLabel()); + } + @Test public void explainReturnsNarrativeText() { mockAssistantResponses(); @@ -1788,6 +1854,60 @@ public class AssistantControllerTest { .build()); } + /** + * Two connected blocks where only the end block's outcomeLabel is too long - a bean- + * validation-only, block-scoped defect that cannot change the flow's shape (EndBlock + * always has a fixed "input" port regardless of its configuration, so fixing this never + * affects I/O or connections). + */ + static FlowCreateRequest twoBlockFlowWithOverlongEndOutcomeLabel() { + LLMDescriptor llmDescriptor = LLMDescriptor.builder() + .provider(PROVIDER) + .model(MODEL) + .build(); + + LLMBlockConfiguration classifierConfiguration = LLMBlockConfiguration.builder() + .name("Ticket classifier") + .llmDescriptor(llmDescriptor) + .prompt("Classify the ticket: ${{ticket}}") + .build(); + + Block classifierBlock = Block.builder() + .specificConfiguration(classifierConfiguration) + .input(IODescriptor.of("ticket", IOType.TEXT)) + .output(IODescriptor.of("response", IOType.TEXT)) + .type(new LLMBlockType()) + .build(); + + EndBlockConfiguration endConfiguration = EndBlockConfiguration.builder() + .name("Done") + .outcomeCode("DONE") + .outcomeLabel("x".repeat(300)) + .build(); + + Block endBlock = Block.builder() + .specificConfiguration(endConfiguration) + .input(IODescriptor.of("input", IOType.ANY)) + .type(new EndBlockType()) + .build(); + + Connection connection = Connection.builder() + .sourceId(classifierBlock.getId()) + .sourceName("response") + .targetId(endBlock.getId()) + .targetName("input") + .build(); + + return new FlowCreateRequest( + "Ticket classification with terminal outcome", + "Classify an incoming ticket and record a terminal outcome.", + FlowData.builder() + .block(classifierBlock) + .block(endBlock) + .connection(connection) + .build()); + } + static FlowCreateRequest invalidSingleBlockFlow() { LLMDescriptor llmDescriptor = LLMDescriptor.builder() .provider(PROVIDER)