feat: targeted per-error repair instead of full replan on FIX rounds
Every FIX round previously re-ran the whole pipeline - PLAN, every non-KEEP block's BLOCK_CONFIG, and CONNECTIONS - even when the only problems were plain bean-validation errors scoped to one existing block's own configuration (e.g. an overlong EndBlock outcomeLabel, or the placeholder-naming/BranchRejoin-size violations added earlier this session). That's wasted LLM calls and unnecessary churn risk on parts of the flow that were already correct. FlowAssistantService.assembleFlow now checks, per FIX round, whether every reported error is cleanly scoped to an existing block and none of them is in PLAN_SENSITIVE_ERROR_CODES (anything that could imply the flow's shape needs to change - connections, dependencies, block/ container I/O, lanes, BranchRejoin fan-in, shared-session wiring, ...). When that holds: - PLAN is skipped entirely; the plan is rebuilt locally from currentFlow's existing blocks (UPDATE for the flagged ones, KEEP for the rest) with no LLM call. - CONNECTIONS is skipped entirely; existing connections are carried forward unchanged. Only BLOCK_CONFIG still runs, and only for the flagged blocks - the same targeting that already happened via KEEP/UPDATE, now extended to skip the two other phases too. This is deliberately conservative: any error outside the safe set, or any error not attributable to one existing block, falls back to the existing full-repair behavior unchanged. If a "safe" fix unexpectedly changes a block's derived I/O anyway (e.g. a factory that derives ports from placeholder text), the next round's validation will surface a sensitive code and force a normal full repair - never worse than before, just self-correcting one round later. Added a regression test that only mocks BLOCK_CONFIG (not PLAN/ CONNECTIONS) for a block-scoped-only error, so it fails loudly if the skip logic regresses.
This commit is contained in:
parent
8299c50c65
commit
94ed6b245d
|
|
@ -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<ValidationErrorCode> 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<Block<?>> 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<ValidationError> errors) {
|
||||
if (errors == null || errors.isEmpty() || !hasCurrentFlowBlocks(currentFlow)) {
|
||||
return false;
|
||||
}
|
||||
Set<String> 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<ValidationError> errors) {
|
||||
Set<String> erroringBlockIds = errors.stream()
|
||||
.map(ValidationError::id)
|
||||
.filter(Objects::nonNull)
|
||||
.collect(Collectors.toSet());
|
||||
List<AssistantBlockPlan> 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<String> usedExistingBlockIds) {
|
||||
List<Block<?>> existingBlocks = currentFlow == null || currentFlow.flow() == null
|
||||
|
|
|
|||
|
|
@ -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<String> 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<ValidationError> 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<LLMBlockType> classifierBlock = Block.<LLMBlockType>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<EndBlockType> endBlock = Block.<EndBlockType>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)
|
||||
|
|
|
|||
Loading…
Reference in New Issue