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 395bacc..9ceaf0c 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 @@ -11,9 +11,6 @@ import java.util.Objects; import java.util.Locale; 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; @@ -99,7 +96,7 @@ public class FlowAssistantService { // 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( + static final Set PLAN_SENSITIVE_ERROR_CODES = EnumSet.of( ValidationErrorCode.FLOW_CONTAINS_NULL_NODE, ValidationErrorCode.NODE_ID_REQUIRED, ValidationErrorCode.UNSUPPORTED_NODE_TYPE, @@ -465,19 +462,19 @@ public class FlowAssistantService { catalogByType.put(descriptor.type(), descriptor); } - boolean targetedRepair = mode == OperationMode.FIX && isTargetedBlockRepairEligible(currentFlow, errors); + boolean targetedRepair = mode == OperationMode.FIX && PlanValidationSupport.isTargetedBlockRepairEligible(currentFlow, errors); progressListener.onProgress("planning", "Planning workflow blocks"); ParsedPlan parsedPlan; if (targetedRepair) { - parsedPlan = buildReusedPlanForTargetedRepair(currentFlow, errors); + parsedPlan = PlanValidationSupport.buildReusedPlanForTargetedRepair(currentFlow, errors); } else { String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog, mcpServerCatalogEntries()); parsedPlan = invokeStructuredAndValidate(provider, authorization, planningModelFor(mode, phaseModels), phaseModels.repairModel(), planPrompt, "plan", rawResponse -> { ParsedPlan plan = AssistantResponseParser.parsePlan(rawResponse); - AssistantFlowPlan normalizedPlan = validateAndNormalizePlan(plan.plan(), mode, userPrompt, currentFlow, + AssistantFlowPlan normalizedPlan = PlanValidationSupport.validateAndNormalizePlan(plan.plan(), mode, userPrompt, currentFlow, errors, catalogByType.keySet()); return new ParsedPlan(normalizedPlan, plan.rationale()); }); @@ -499,8 +496,8 @@ public class FlowAssistantService { List blockPlans = parsedPlan.plan().blocks(); for (int blockIndex = 0; blockIndex < blockPlans.size(); blockIndex++) { AssistantBlockPlan blockPlan = blockPlans.get(blockIndex); - PlanOperation operation = parsePlanOperation(blockPlan.operation()); - Block existingBlock = resolveExistingBlockForPlan(blockPlan, currentFlow, blockIndex, usedExistingBlockIds); + PlanOperation operation = PlanValidationSupport.parsePlanOperation(blockPlan.operation()); + Block existingBlock = PlanValidationSupport.resolveExistingBlockForPlan(blockPlan, currentFlow, blockIndex, usedExistingBlockIds); if (existingBlock != null) { usedExistingBlockIds.add(existingBlock.getId()); } @@ -569,8 +566,8 @@ public class FlowAssistantService { for (AssistantContainerPlan containerPlan : parsedPlan.plan().containers() == null ? List.of() : parsedPlan.plan().containers()) { - PlanOperation containerOperation = parsePlanOperation(containerPlan.operation()); - Container existingContainer = resolveExistingContainerForPlan(containerPlan, currentFlow, -1, + PlanOperation containerOperation = PlanValidationSupport.parsePlanOperation(containerPlan.operation()); + Container existingContainer = PlanValidationSupport.resolveExistingContainerForPlan(containerPlan, currentFlow, -1, usedExistingContainerIds); if (existingContainer != null) { usedExistingContainerIds.add(existingContainer.getId()); @@ -696,8 +693,8 @@ public class FlowAssistantService { List innerBlockPlans = containerPlan.blocks(); for (int blockIndex = 0; blockIndex < innerBlockPlans.size(); blockIndex++) { AssistantBlockPlan blockPlan = innerBlockPlans.get(blockIndex); - PlanOperation operation = parsePlanOperation(blockPlan.operation()); - Block existingInnerBlock = resolveExistingBlock(blockPlan, existingInnerBlocks, blockIndex, + PlanOperation operation = PlanValidationSupport.parsePlanOperation(blockPlan.operation()); + Block existingInnerBlock = PlanValidationSupport.resolveExistingBlock(blockPlan, existingInnerBlocks, blockIndex, usedExistingInnerIds); if (existingInnerBlock != null) { usedExistingInnerIds.add(existingInnerBlock.getId()); @@ -1502,156 +1499,6 @@ public class FlowAssistantService { return (Block) factory.create(configuration); } - private PlanOperation parsePlanOperation(String rawOperation) { - if (rawOperation == null || rawOperation.isBlank()) { - return PlanOperation.ADD; - } - try { - return PlanOperation.valueOf(rawOperation.trim().toUpperCase(Locale.ROOT)); - } catch (IllegalArgumentException e) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an unknown block operation: " + rawOperation); - } - } - - /** - * 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) { - // buildReusedPlanForTargetedRepair carries existing containers forward as KEEP (untouched), - // so this stays eligible even when the flow has containers - as long as every error is - // still scoped to an existing *top-level block*. An error scoped to a container (or - // anything else) can't be fixed by a KEEP-only container plan, so it still disqualifies - // targeted repair and falls back to a full replan. - 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(); - // isTargetedBlockRepairEligible only allows this path when every error is scoped to a - // top-level block, so any existing containers are always untouched here - just carried - // forward as KEEP so they aren't silently dropped from the rebuilt plan. - List containerPlans = hasCurrentFlowContainers(currentFlow) - ? currentFlow.flow().getContainers().stream() - .map(container -> new AssistantContainerPlan( - container.getId(), - container.getType().getName(), - container.getName(), - PlanOperation.KEEP.name(), - null, - null, - null, - null, - List.of())) - .toList() - : List.of(); - AssistantFlowPlan plan = new AssistantFlowPlan(currentFlow.name(), currentFlow.description(), blockPlans, - containerPlans); - 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 - || currentFlow.flow().getBlocks() == null - ? List.of() - : currentFlow.flow().getBlocks(); - return resolveExistingBlock(blockPlan, existingBlocks, planIndex, usedExistingBlockIds); - } - - /** - * Matches a block plan entry to one of an arbitrary list of existing blocks (the top-level - * flow's blocks, or a container subflow's inner blocks), by id, then normalized name, then - * purpose-as-name, then same-position-same-type, then unique-same-type. Shared by top-level - * and container-inner block resolution so both use identical matching rules. - */ - private Block resolveExistingBlock(AssistantBlockPlan blockPlan, List> existingBlocks, - int planIndex, Set usedExistingBlockIds) { - if (existingBlocks == null || existingBlocks.isEmpty()) { - return null; - } - - Block exact = findExistingBlock(existingBlocks, usedExistingBlockIds, - block -> Objects.equals(block.getId(), blockPlan.blockId()) - || AssistantTextSupport.normalizeBlockReference(blockPlan.blockId()) != null - && AssistantTextSupport.normalizeBlockReference(blockPlan.blockId()) - .equals(AssistantTextSupport.normalizeBlockReference(block.getName()))); - if (exact != null) { - return exact; - } - - Block byPurpose = findExistingBlock(existingBlocks, usedExistingBlockIds, - block -> AssistantTextSupport.normalizeBlockReference(blockPlan.purpose()) != null - && AssistantTextSupport.normalizeBlockReference(blockPlan.purpose()).equals(AssistantTextSupport.normalizeBlockReference(block.getName()))); - if (byPurpose != null) { - return byPurpose; - } - - if (planIndex >= 0 && planIndex < existingBlocks.size()) { - Block byPosition = existingBlocks.get(planIndex); - if (!usedExistingBlockIds.contains(byPosition.getId()) - && Objects.equals(blockPlan.blockType(), blockTypeName(byPosition))) { - return byPosition; - } - } - - List> sameType = existingBlocks.stream() - .filter(block -> !usedExistingBlockIds.contains(block.getId())) - .filter(block -> Objects.equals(blockPlan.blockType(), blockTypeName(block))) - .toList(); - return sameType.size() == 1 ? sameType.getFirst() : null; - } - - private Block findExistingBlock(List> existingBlocks, Set usedExistingBlockIds, - java.util.function.Predicate> predicate) { - for (Block block : existingBlocks) { - if (block == null || usedExistingBlockIds.contains(block.getId())) { - continue; - } - if (predicate.test(block)) { - return block; - } - } - return null; - } - - private String blockTypeName(Block block) { - return block == null || block.getType() == null ? null : block.getType().getName(); - } - private List preserveCurrentConnections(FlowCreateRequest currentFlow, Map oldNodeIdToAssembledNode, Set removedExistingNodeIds) { List sourceConnections = currentFlow == null || currentFlow.flow() == null @@ -1664,404 +1511,7 @@ public class FlowAssistantService { * MCP shared-session producer blocks by their session name (shareSession=true with a name). */ - private AssistantFlowPlan validateAndNormalizePlan(AssistantFlowPlan plan, OperationMode mode, String userPrompt, - FlowCreateRequest currentFlow, List errors, - Set availableBlockTypes) { - if (plan == null) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant returned an empty plan"); - } - boolean hasBlocks = plan.blocks() != null && !plan.blocks().isEmpty(); - boolean hasContainers = plan.containers() != null && !plan.containers().isEmpty(); - // Only fall back to the single-LLMBlock draft when the plan is genuinely empty. A plan - // that groups all its work inside containers (no flat top-level blocks) is legitimate - - // treating it as "empty" would silently discard every container. - if (!hasBlocks && !hasContainers) { - if (canUseMinimalDraftFallback(userPrompt, currentFlow, availableBlockTypes)) { - return new AssistantFlowPlan( - AssistantTextSupport.defaultIfBlank(plan.name(), "Assistant flow"), - AssistantTextSupport.defaultIfBlank(plan.description(), userPrompt.trim()), - List.of(new AssistantBlockPlan("b1", "LLMBlock", userPrompt.trim(), PlanOperation.ADD.name()))); - } - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant returned a plan with no blocks"); - } - Set ids = new LinkedHashSet<>(); - Set usedExistingBlockIds = new LinkedHashSet<>(); - List normalizedBlocks = new ArrayList<>(); - List rawBlocks = plan.blocks() == null ? List.of() : plan.blocks(); - for (int blockIndex = 0; blockIndex < rawBlocks.size(); blockIndex++) { - AssistantBlockPlan block = rawBlocks.get(blockIndex); - if (block.blockId() == null || block.blockId().isBlank() || !ids.add(block.blockId())) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned invalid or duplicate block ids in the plan"); - } - if (block.blockType() == null || block.blockType().isBlank()) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned a block without blockType"); - } - if (isContainerBlockType(block.blockType())) { - // A container type placed in the flat "blocks" list (instead of "containers") used - // to fail hard as an unrecoverable 502 - the block-assembly loop has no catalog - // descriptor for a container type and never gets a chance to retry. Thrown here, - // inside the plan's own retry-wrapped parser callback, it becomes a normal - // structured-repair retry instead. - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant placed a container type '" + block.blockType() + "' (blockId " + block.blockId() - + ") in the top-level \"blocks\" list - container types must be entries in the" - + " \"containers\" list instead, with their own inner \"blocks\""); - } - if (isSharedMemoryContext(userPrompt, currentFlow, plan) - && "LLMBlock".equals(block.blockType()) - && SharedMemoryIntentClassifier.isSharedStatePurpose(block.purpose())) { - if (!availableBlockTypes.contains("MCPAgent")) { - // MCPAgent not available in this deployment: skip the constraint rather than - // silently accepting an invalid plan — surface a clear error so the operator - // knows the catalog is incomplete for shared-memory workflows. - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned a plan that requires shared-memory semantics but MCPAgent is not available in the block catalog"); - } - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an inconsistent plan: workflows that require shared memory or reusable state between steps must use MCPAgent blocks with a shared MCP session instead of standalone LLMBlock nodes"); - } - PlanOperation operation = normalizePlanOperation(block, mode, currentFlow, errors, blockIndex, - usedExistingBlockIds); - normalizedBlocks.add(new AssistantBlockPlan( - block.blockId(), - block.blockType(), - block.purpose(), - operation.name())); - } - - List normalizedContainers = new ArrayList<>(); - Set usedExistingContainerIds = new LinkedHashSet<>(); - List rawContainers = plan.containers() == null ? List.of() : plan.containers(); - for (int containerIndex = 0; containerIndex < rawContainers.size(); containerIndex++) { - AssistantContainerPlan containerPlan = rawContainers.get(containerIndex); - if (containerPlan.containerId() == null || containerPlan.containerId().isBlank() - || !ids.add(containerPlan.containerId())) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned invalid or duplicate container ids in the plan"); - } - if (!GenericContainerType.TYPE.equals(containerPlan.containerType()) - && !IteratorContainerType.TYPE.equals(containerPlan.containerType()) - && !LoopContainerType.TYPE.equals(containerPlan.containerType())) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant selected an unsupported container type: " + containerPlan.containerType() - + " (only " + GenericContainerType.TYPE + ", " + IteratorContainerType.TYPE - + " and " + LoopContainerType.TYPE + " can currently be authored by the assistant)"); - } - Container matchedContainer = resolveExistingContainerForPlan(containerPlan, currentFlow, containerIndex, - usedExistingContainerIds); - PlanOperation containerOperation = resolveContainerOperation(containerPlan, matchedContainer, mode, errors); - boolean rebuildsBody = containerOperation == PlanOperation.ADD || containerOperation == PlanOperation.UPDATE; - if (LoopContainerType.TYPE.equals(containerPlan.containerType()) && rebuildsBody - && AssistantTextSupport.trimToNull(containerPlan.guardCondition()) == null) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned a LoopContainer without a guardCondition: " + containerPlan.containerId()); - } - Container existingContainer = containerOperation == PlanOperation.ADD ? null : matchedContainer; - if (existingContainer != null) { - usedExistingContainerIds.add(existingContainer.getId()); - } - List innerBlocks; - if (containerOperation == PlanOperation.KEEP || containerOperation == PlanOperation.REMOVE) { - innerBlocks = List.of(); - } else if (containerOperation == PlanOperation.UPDATE && existingContainer != null - && (mode == OperationMode.REFINE || mode == OperationMode.FIX)) { - // Incremental container diffing: resolve each inner block against the existing - // container's subflow so unchanged ones become KEEP and are reused as-is. - List incremental = - normalizeContainerInnerBlocksAgainstExisting(containerPlan, existingContainer, mode, errors); - // Container-internal validation errors are remapped to the container id (the inner - // block id is lost), so FIX can't auto-target an inner block - it relies on the - // model marking one. If a FIX round would keep every inner block unchanged, the - // error would never get fixed, so fall back to full regeneration in that case. - boolean anyChange = incremental.stream() - .anyMatch(block -> parsePlanOperation(block.operation()) != PlanOperation.KEEP); - innerBlocks = (mode == OperationMode.REFINE || anyChange) - ? incremental - : normalizeContainerInnerBlocks(containerPlan); - } else { - innerBlocks = normalizeContainerInnerBlocks(containerPlan); - } - // iterationInput only applies to IteratorContainer; drop it for GenericContainer so a - // stray value doesn't leak into the config. It may be blank for IteratorContainer too: - // the factory infers it when the subflow has exactly one open input, and any genuinely - // wrong value is caught by IteratorContainerConfiguration's bean-validation on the - // execution-validator pass (which then triggers a repair round). - String iterationInput = IteratorContainerType.TYPE.equals(containerPlan.containerType()) - ? AssistantTextSupport.trimToNull(containerPlan.iterationInput()) - : null; - boolean isLoop = LoopContainerType.TYPE.equals(containerPlan.containerType()); - String guardCondition = isLoop ? AssistantTextSupport.trimToNull(containerPlan.guardCondition()) : null; - Integer maxIterations = isLoop ? containerPlan.maxIterations() : null; - String feedbackInput = isLoop ? AssistantTextSupport.trimToNull(containerPlan.feedbackInput()) : null; - normalizedContainers.add(new AssistantContainerPlan( - containerPlan.containerId(), - containerPlan.containerType(), - containerPlan.purpose(), - containerOperation.name(), - iterationInput, - guardCondition, - maxIterations, - feedbackInput, - innerBlocks)); - } - - return new AssistantFlowPlan(plan.name(), plan.description(), normalizedBlocks, normalizedContainers); - } - - /** - * Container inner blocks fully (re)generated: every inner block becomes ADD. Used for a new - * (ADD) container and for FIX-mode updates, where rebuilding the whole subflow is the safe way - * to resolve a container-internal error. - */ - private List normalizeContainerInnerBlocks(AssistantContainerPlan containerPlan) { - List rawInnerBlocks = validateInnerBlockShapes(containerPlan); - List innerBlocks = new ArrayList<>(); - for (AssistantBlockPlan innerBlock : rawInnerBlocks) { - innerBlocks.add(new AssistantBlockPlan(innerBlock.blockId(), innerBlock.blockType(), innerBlock.purpose(), - PlanOperation.ADD.name())); - } - return innerBlocks; - } - - /** - * Incremental variant (REFINE only): resolves each inner block against the existing container's - * subflow so unchanged blocks become KEEP (reused verbatim, no LLM reconfiguration), changed - * ones UPDATE, new ones ADD, and dropped ones REMOVE - mirroring how top-level blocks are - * diffed, one level down. - */ - private List normalizeContainerInnerBlocksAgainstExisting(AssistantContainerPlan containerPlan, - Container existingContainer, OperationMode mode, List errors) { - List rawInnerBlocks = validateInnerBlockShapes(containerPlan); - List> existingInner = FlowAssemblySupport.subFlowBlocks(existingContainer); - Set usedInnerIds = new LinkedHashSet<>(); - List innerBlocks = new ArrayList<>(); - for (int index = 0; index < rawInnerBlocks.size(); index++) { - AssistantBlockPlan innerBlock = rawInnerBlocks.get(index); - PlanOperation operation = normalizeInnerBlockOperation(innerBlock, existingInner, mode, errors, index, - usedInnerIds); - innerBlocks.add(new AssistantBlockPlan(innerBlock.blockId(), innerBlock.blockType(), innerBlock.purpose(), - operation.name())); - } - return innerBlocks; - } - - /** - * Validates the shape of a container's inner block list (non-empty, unique non-blank ids, - * present blockType, no nested containers) and returns the raw entries. Shared by both the - * full-regeneration and incremental inner-block normalizers. - */ - private List validateInnerBlockShapes(AssistantContainerPlan containerPlan) { - List rawInnerBlocks = containerPlan.blocks() == null ? List.of() : containerPlan.blocks(); - if (rawInnerBlocks.isEmpty()) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned a container with no inner blocks: " + containerPlan.containerId()); - } - Set innerIds = new LinkedHashSet<>(); - for (AssistantBlockPlan innerBlock : rawInnerBlocks) { - if (innerBlock.blockId() == null || innerBlock.blockId().isBlank() || !innerIds.add(innerBlock.blockId())) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned invalid or duplicate inner block ids in container: " - + containerPlan.containerId()); - } - if (innerBlock.blockType() == null || innerBlock.blockType().isBlank()) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an inner block without blockType in container: " - + containerPlan.containerId()); - } - if (isContainerBlockType(innerBlock.blockType())) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Nested containers are not supported: container " + containerPlan.containerId() - + " cannot contain another container"); - } - } - return rawInnerBlocks; - } - - private PlanOperation normalizeInnerBlockOperation(AssistantBlockPlan block, List> existingInner, - OperationMode mode, List errors, int blockIndex, Set usedInnerIds) { - if (block.operation() != null && !block.operation().isBlank()) { - PlanOperation explicitOperation = parsePlanOperation(block.operation()); - if (explicitOperation != PlanOperation.ADD) { - Block matched = resolveExistingBlock(block, existingInner, blockIndex, usedInnerIds); - if (matched != null) { - usedInnerIds.add(matched.getId()); - } - } - return explicitOperation; - } - if (existingInner.isEmpty()) { - return PlanOperation.ADD; - } - Block matched = resolveExistingBlock(block, existingInner, blockIndex, usedInnerIds); - if (matched == null) { - return PlanOperation.ADD; - } - usedInnerIds.add(matched.getId()); - if (mode == OperationMode.FIX && hasValidationErrorForBlock(matched, errors)) { - return PlanOperation.UPDATE; - } - return PlanOperation.KEEP; - } - - - private boolean isContainerBlockType(String blockType) { - return GenericContainerType.TYPE.equals(blockType) - || "IteratorContainer".equals(blockType) - || "LoopContainer".equals(blockType); - } - - private PlanOperation normalizePlanOperation(AssistantBlockPlan block, OperationMode mode, - FlowCreateRequest currentFlow, List errors, int blockIndex, - Set usedExistingBlockIds) { - if (block.operation() != null && !block.operation().isBlank()) { - PlanOperation explicitOperation = parsePlanOperation(block.operation()); - if (explicitOperation != PlanOperation.ADD) { - Block matched = resolveExistingBlockForPlan(block, currentFlow, blockIndex, usedExistingBlockIds); - if (matched != null) { - usedExistingBlockIds.add(matched.getId()); - } - } - return explicitOperation; - } - if (mode == OperationMode.DRAFT || !hasCurrentFlowBlocks(currentFlow)) { - return PlanOperation.ADD; - } - - Block matched = resolveExistingBlockForPlan(block, currentFlow, blockIndex, usedExistingBlockIds); - if (matched == null) { - return PlanOperation.ADD; - } - usedExistingBlockIds.add(matched.getId()); - if (mode == OperationMode.FIX && hasValidationErrorForBlock(matched, errors)) { - return PlanOperation.UPDATE; - } - return PlanOperation.KEEP; - } - - private boolean hasValidationErrorForBlock(Block block, List errors) { - if (block == null || block.getId() == null || errors == null || errors.isEmpty()) { - return false; - } - for (ValidationError error : errors) { - if (Objects.equals(block.getId(), error.id()) - || error.relatedNodeIds() != null && error.relatedNodeIds().contains(block.getId())) { - return true; - } - } - return false; - } - - /** - * Decides a container plan entry's operation from its already-resolved existing match. Kept - * separate from resolution so the caller can reuse the same matched container for incremental - * inner-block diffing. - */ - private PlanOperation resolveContainerOperation(AssistantContainerPlan containerPlan, Container existingContainer, - OperationMode mode, List errors) { - if (containerPlan.operation() != null && !containerPlan.operation().isBlank()) { - return parsePlanOperation(containerPlan.operation()); - } - if (mode == OperationMode.DRAFT || existingContainer == null) { - return PlanOperation.ADD; - } - if (mode == OperationMode.FIX && hasValidationErrorForContainer(existingContainer, errors)) { - return PlanOperation.UPDATE; - } - return PlanOperation.KEEP; - } - - private boolean hasValidationErrorForContainer(Container container, List errors) { - if (container == null || container.getId() == null || errors == null || errors.isEmpty()) { - return false; - } - for (ValidationError error : errors) { - if (Objects.equals(container.getId(), error.id()) - || error.relatedNodeIds() != null && error.relatedNodeIds().contains(container.getId())) { - return true; - } - } - return false; - } - - private boolean hasCurrentFlowContainers(FlowCreateRequest currentFlow) { - return currentFlow != null - && currentFlow.flow() != null - && currentFlow.flow().getContainers() != null - && !currentFlow.flow().getContainers().isEmpty(); - } - - private Container resolveExistingContainerForPlan(AssistantContainerPlan containerPlan, - FlowCreateRequest currentFlow, int planIndex, Set usedExistingContainerIds) { - List> existingContainers = currentFlow == null || currentFlow.flow() == null - || currentFlow.flow().getContainers() == null - ? List.of() - : currentFlow.flow().getContainers(); - if (existingContainers.isEmpty()) { - return null; - } - - Container exact = findExistingContainer(existingContainers, usedExistingContainerIds, - container -> Objects.equals(container.getId(), containerPlan.containerId()) - || AssistantTextSupport.normalizeBlockReference(containerPlan.containerId()) != null - && AssistantTextSupport.normalizeBlockReference(containerPlan.containerId()) - .equals(AssistantTextSupport.normalizeBlockReference(container.getName()))); - if (exact != null) { - return exact; - } - - Container byPurpose = findExistingContainer(existingContainers, usedExistingContainerIds, - container -> AssistantTextSupport.normalizeBlockReference(containerPlan.purpose()) != null - && AssistantTextSupport.normalizeBlockReference(containerPlan.purpose()) - .equals(AssistantTextSupport.normalizeBlockReference(container.getName()))); - if (byPurpose != null) { - return byPurpose; - } - - if (planIndex >= 0 && planIndex < existingContainers.size()) { - Container byPosition = existingContainers.get(planIndex); - if (!usedExistingContainerIds.contains(byPosition.getId())) { - return byPosition; - } - } - - List> unused = existingContainers.stream() - .filter(container -> !usedExistingContainerIds.contains(container.getId())) - .toList(); - return unused.size() == 1 ? unused.getFirst() : null; - } - - private Container findExistingContainer(List> existingContainers, - Set usedExistingContainerIds, java.util.function.Predicate> predicate) { - for (Container container : existingContainers) { - if (container == null || usedExistingContainerIds.contains(container.getId())) { - continue; - } - if (predicate.test(container)) { - return container; - } - } - return null; - } - - private boolean canUseMinimalDraftFallback(String userPrompt, FlowCreateRequest currentFlow, - Set availableBlockTypes) { - return userPrompt != null - && !userPrompt.isBlank() - && !hasCurrentFlowBlocks(currentFlow) - && availableBlockTypes != null - && availableBlockTypes.contains("LLMBlock"); - } - - private boolean hasCurrentFlowBlocks(FlowCreateRequest currentFlow) { - return currentFlow != null - && currentFlow.flow() != null - && currentFlow.flow().getBlocks() != null - && !currentFlow.flow().getBlocks().isEmpty(); - } - - private boolean isSharedMemoryContext(String userPrompt, FlowCreateRequest currentFlow, AssistantFlowPlan plan) { + static boolean isSharedMemoryContext(String userPrompt, FlowCreateRequest currentFlow, AssistantFlowPlan plan) { if (SharedMemoryIntentClassifier.isSharedMemoryRequest(userPrompt)) { return true; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/PlanValidationSupport.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/PlanValidationSupport.java new file mode 100644 index 0000000..b7c3a9e --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/PlanValidationSupport.java @@ -0,0 +1,578 @@ +package it.cnr.isti.workflow.manager.assistant; + +import java.util.ArrayList; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Locale; +import java.util.Objects; +import java.util.Set; +import java.util.stream.Collectors; + +import org.springframework.http.HttpStatus; +import org.springframework.web.server.ResponseStatusException; + +import it.cnr.isti.workflow.manager.assistant.FlowAssistantPromptService.OperationMode; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantBlockPlan; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantContainerPlan; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantFlowPlan; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.ParsedPlan; +import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.PlanOperation; +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.containers.Container; +import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; +import it.cnr.isti.workflow.manager.containers.types.IteratorContainerType; +import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; +import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; +import it.cnr.isti.workflow.manager.flows.validation.ValidationError; + +final class PlanValidationSupport { + + private PlanValidationSupport() { + } + + static PlanOperation parsePlanOperation(String rawOperation) { + if (rawOperation == null || rawOperation.isBlank()) { + return PlanOperation.ADD; + } + try { + return PlanOperation.valueOf(rawOperation.trim().toUpperCase(Locale.ROOT)); + } catch (IllegalArgumentException e) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an unknown block operation: " + rawOperation); + } + } + + /** + * 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. + */ + static boolean isTargetedBlockRepairEligible(FlowCreateRequest currentFlow, List errors) { + // buildReusedPlanForTargetedRepair carries existing containers forward as KEEP (untouched), + // so this stays eligible even when the flow has containers - as long as every error is + // still scoped to an existing *top-level block*. An error scoped to a container (or + // anything else) can't be fixed by a KEEP-only container plan, so it still disqualifies + // targeted repair and falls back to a full replan. + 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 (FlowAssistantService.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. + */ + static 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(); + // isTargetedBlockRepairEligible only allows this path when every error is scoped to a + // top-level block, so any existing containers are always untouched here - just carried + // forward as KEEP so they aren't silently dropped from the rebuilt plan. + List containerPlans = hasCurrentFlowContainers(currentFlow) + ? currentFlow.flow().getContainers().stream() + .map(container -> new AssistantContainerPlan( + container.getId(), + container.getType().getName(), + container.getName(), + PlanOperation.KEEP.name(), + null, + null, + null, + null, + List.of())) + .toList() + : List.of(); + AssistantFlowPlan plan = new AssistantFlowPlan(currentFlow.name(), currentFlow.description(), blockPlans, + containerPlans); + return new ParsedPlan(plan, "Targeted repair: reconfiguring only the blocks flagged by validation errors."); + } + + static Block resolveExistingBlockForPlan(AssistantBlockPlan blockPlan, FlowCreateRequest currentFlow, + int planIndex, Set usedExistingBlockIds) { + List> existingBlocks = currentFlow == null || currentFlow.flow() == null + || currentFlow.flow().getBlocks() == null + ? List.of() + : currentFlow.flow().getBlocks(); + return resolveExistingBlock(blockPlan, existingBlocks, planIndex, usedExistingBlockIds); + } + + /** + * Matches a block plan entry to one of an arbitrary list of existing blocks (the top-level + * flow's blocks, or a container subflow's inner blocks), by id, then normalized name, then + * purpose-as-name, then same-position-same-type, then unique-same-type. Shared by top-level + * and container-inner block resolution so both use identical matching rules. + */ + static Block resolveExistingBlock(AssistantBlockPlan blockPlan, List> existingBlocks, + int planIndex, Set usedExistingBlockIds) { + if (existingBlocks == null || existingBlocks.isEmpty()) { + return null; + } + + Block exact = findExistingBlock(existingBlocks, usedExistingBlockIds, + block -> Objects.equals(block.getId(), blockPlan.blockId()) + || AssistantTextSupport.normalizeBlockReference(blockPlan.blockId()) != null + && AssistantTextSupport.normalizeBlockReference(blockPlan.blockId()) + .equals(AssistantTextSupport.normalizeBlockReference(block.getName()))); + if (exact != null) { + return exact; + } + + Block byPurpose = findExistingBlock(existingBlocks, usedExistingBlockIds, + block -> AssistantTextSupport.normalizeBlockReference(blockPlan.purpose()) != null + && AssistantTextSupport.normalizeBlockReference(blockPlan.purpose()).equals(AssistantTextSupport.normalizeBlockReference(block.getName()))); + if (byPurpose != null) { + return byPurpose; + } + + if (planIndex >= 0 && planIndex < existingBlocks.size()) { + Block byPosition = existingBlocks.get(planIndex); + if (!usedExistingBlockIds.contains(byPosition.getId()) + && Objects.equals(blockPlan.blockType(), blockTypeName(byPosition))) { + return byPosition; + } + } + + List> sameType = existingBlocks.stream() + .filter(block -> !usedExistingBlockIds.contains(block.getId())) + .filter(block -> Objects.equals(blockPlan.blockType(), blockTypeName(block))) + .toList(); + return sameType.size() == 1 ? sameType.getFirst() : null; + } + + private static Block findExistingBlock(List> existingBlocks, Set usedExistingBlockIds, + java.util.function.Predicate> predicate) { + for (Block block : existingBlocks) { + if (block == null || usedExistingBlockIds.contains(block.getId())) { + continue; + } + if (predicate.test(block)) { + return block; + } + } + return null; + } + + static String blockTypeName(Block block) { + return block == null || block.getType() == null ? null : block.getType().getName(); + } + + static AssistantFlowPlan validateAndNormalizePlan(AssistantFlowPlan plan, OperationMode mode, String userPrompt, + FlowCreateRequest currentFlow, List errors, + Set availableBlockTypes) { + if (plan == null) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant returned an empty plan"); + } + boolean hasBlocks = plan.blocks() != null && !plan.blocks().isEmpty(); + boolean hasContainers = plan.containers() != null && !plan.containers().isEmpty(); + // Only fall back to the single-LLMBlock draft when the plan is genuinely empty. A plan + // that groups all its work inside containers (no flat top-level blocks) is legitimate - + // treating it as "empty" would silently discard every container. + if (!hasBlocks && !hasContainers) { + if (canUseMinimalDraftFallback(userPrompt, currentFlow, availableBlockTypes)) { + return new AssistantFlowPlan( + AssistantTextSupport.defaultIfBlank(plan.name(), "Assistant flow"), + AssistantTextSupport.defaultIfBlank(plan.description(), userPrompt.trim()), + List.of(new AssistantBlockPlan("b1", "LLMBlock", userPrompt.trim(), PlanOperation.ADD.name()))); + } + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant returned a plan with no blocks"); + } + Set ids = new LinkedHashSet<>(); + Set usedExistingBlockIds = new LinkedHashSet<>(); + List normalizedBlocks = new ArrayList<>(); + List rawBlocks = plan.blocks() == null ? List.of() : plan.blocks(); + for (int blockIndex = 0; blockIndex < rawBlocks.size(); blockIndex++) { + AssistantBlockPlan block = rawBlocks.get(blockIndex); + if (block.blockId() == null || block.blockId().isBlank() || !ids.add(block.blockId())) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned invalid or duplicate block ids in the plan"); + } + if (block.blockType() == null || block.blockType().isBlank()) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned a block without blockType"); + } + if (isContainerBlockType(block.blockType())) { + // A container type placed in the flat "blocks" list (instead of "containers") used + // to fail hard as an unrecoverable 502 - the block-assembly loop has no catalog + // descriptor for a container type and never gets a chance to retry. Thrown here, + // inside the plan's own retry-wrapped parser callback, it becomes a normal + // structured-repair retry instead. + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant placed a container type '" + block.blockType() + "' (blockId " + block.blockId() + + ") in the top-level \"blocks\" list - container types must be entries in the" + + " \"containers\" list instead, with their own inner \"blocks\""); + } + if (FlowAssistantService.isSharedMemoryContext(userPrompt, currentFlow, plan) + && "LLMBlock".equals(block.blockType()) + && SharedMemoryIntentClassifier.isSharedStatePurpose(block.purpose())) { + if (!availableBlockTypes.contains("MCPAgent")) { + // MCPAgent not available in this deployment: skip the constraint rather than + // silently accepting an invalid plan — surface a clear error so the operator + // knows the catalog is incomplete for shared-memory workflows. + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned a plan that requires shared-memory semantics but MCPAgent is not available in the block catalog"); + } + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an inconsistent plan: workflows that require shared memory or reusable state between steps must use MCPAgent blocks with a shared MCP session instead of standalone LLMBlock nodes"); + } + PlanOperation operation = normalizePlanOperation(block, mode, currentFlow, errors, blockIndex, + usedExistingBlockIds); + normalizedBlocks.add(new AssistantBlockPlan( + block.blockId(), + block.blockType(), + block.purpose(), + operation.name())); + } + + List normalizedContainers = new ArrayList<>(); + Set usedExistingContainerIds = new LinkedHashSet<>(); + List rawContainers = plan.containers() == null ? List.of() : plan.containers(); + for (int containerIndex = 0; containerIndex < rawContainers.size(); containerIndex++) { + AssistantContainerPlan containerPlan = rawContainers.get(containerIndex); + if (containerPlan.containerId() == null || containerPlan.containerId().isBlank() + || !ids.add(containerPlan.containerId())) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned invalid or duplicate container ids in the plan"); + } + if (!GenericContainerType.TYPE.equals(containerPlan.containerType()) + && !IteratorContainerType.TYPE.equals(containerPlan.containerType()) + && !LoopContainerType.TYPE.equals(containerPlan.containerType())) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant selected an unsupported container type: " + containerPlan.containerType() + + " (only " + GenericContainerType.TYPE + ", " + IteratorContainerType.TYPE + + " and " + LoopContainerType.TYPE + " can currently be authored by the assistant)"); + } + Container matchedContainer = resolveExistingContainerForPlan(containerPlan, currentFlow, containerIndex, + usedExistingContainerIds); + PlanOperation containerOperation = resolveContainerOperation(containerPlan, matchedContainer, mode, errors); + boolean rebuildsBody = containerOperation == PlanOperation.ADD || containerOperation == PlanOperation.UPDATE; + if (LoopContainerType.TYPE.equals(containerPlan.containerType()) && rebuildsBody + && AssistantTextSupport.trimToNull(containerPlan.guardCondition()) == null) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned a LoopContainer without a guardCondition: " + containerPlan.containerId()); + } + Container existingContainer = containerOperation == PlanOperation.ADD ? null : matchedContainer; + if (existingContainer != null) { + usedExistingContainerIds.add(existingContainer.getId()); + } + List innerBlocks; + if (containerOperation == PlanOperation.KEEP || containerOperation == PlanOperation.REMOVE) { + innerBlocks = List.of(); + } else if (containerOperation == PlanOperation.UPDATE && existingContainer != null + && (mode == OperationMode.REFINE || mode == OperationMode.FIX)) { + // Incremental container diffing: resolve each inner block against the existing + // container's subflow so unchanged ones become KEEP and are reused as-is. + List incremental = + normalizeContainerInnerBlocksAgainstExisting(containerPlan, existingContainer, mode, errors); + // Container-internal validation errors are remapped to the container id (the inner + // block id is lost), so FIX can't auto-target an inner block - it relies on the + // model marking one. If a FIX round would keep every inner block unchanged, the + // error would never get fixed, so fall back to full regeneration in that case. + boolean anyChange = incremental.stream() + .anyMatch(block -> parsePlanOperation(block.operation()) != PlanOperation.KEEP); + innerBlocks = (mode == OperationMode.REFINE || anyChange) + ? incremental + : normalizeContainerInnerBlocks(containerPlan); + } else { + innerBlocks = normalizeContainerInnerBlocks(containerPlan); + } + // iterationInput only applies to IteratorContainer; drop it for GenericContainer so a + // stray value doesn't leak into the config. It may be blank for IteratorContainer too: + // the factory infers it when the subflow has exactly one open input, and any genuinely + // wrong value is caught by IteratorContainerConfiguration's bean-validation on the + // execution-validator pass (which then triggers a repair round). + String iterationInput = IteratorContainerType.TYPE.equals(containerPlan.containerType()) + ? AssistantTextSupport.trimToNull(containerPlan.iterationInput()) + : null; + boolean isLoop = LoopContainerType.TYPE.equals(containerPlan.containerType()); + String guardCondition = isLoop ? AssistantTextSupport.trimToNull(containerPlan.guardCondition()) : null; + Integer maxIterations = isLoop ? containerPlan.maxIterations() : null; + String feedbackInput = isLoop ? AssistantTextSupport.trimToNull(containerPlan.feedbackInput()) : null; + normalizedContainers.add(new AssistantContainerPlan( + containerPlan.containerId(), + containerPlan.containerType(), + containerPlan.purpose(), + containerOperation.name(), + iterationInput, + guardCondition, + maxIterations, + feedbackInput, + innerBlocks)); + } + + return new AssistantFlowPlan(plan.name(), plan.description(), normalizedBlocks, normalizedContainers); + } + + /** + * Container inner blocks fully (re)generated: every inner block becomes ADD. Used for a new + * (ADD) container and for FIX-mode updates, where rebuilding the whole subflow is the safe way + * to resolve a container-internal error. + */ + static List normalizeContainerInnerBlocks(AssistantContainerPlan containerPlan) { + List rawInnerBlocks = validateInnerBlockShapes(containerPlan); + List innerBlocks = new ArrayList<>(); + for (AssistantBlockPlan innerBlock : rawInnerBlocks) { + innerBlocks.add(new AssistantBlockPlan(innerBlock.blockId(), innerBlock.blockType(), innerBlock.purpose(), + PlanOperation.ADD.name())); + } + return innerBlocks; + } + + /** + * Incremental variant (REFINE only): resolves each inner block against the existing container's + * subflow so unchanged blocks become KEEP (reused verbatim, no LLM reconfiguration), changed + * ones UPDATE, new ones ADD, and dropped ones REMOVE - mirroring how top-level blocks are + * diffed, one level down. + */ + static List normalizeContainerInnerBlocksAgainstExisting(AssistantContainerPlan containerPlan, + Container existingContainer, OperationMode mode, List errors) { + List rawInnerBlocks = validateInnerBlockShapes(containerPlan); + List> existingInner = FlowAssemblySupport.subFlowBlocks(existingContainer); + Set usedInnerIds = new LinkedHashSet<>(); + List innerBlocks = new ArrayList<>(); + for (int index = 0; index < rawInnerBlocks.size(); index++) { + AssistantBlockPlan innerBlock = rawInnerBlocks.get(index); + PlanOperation operation = normalizeInnerBlockOperation(innerBlock, existingInner, mode, errors, index, + usedInnerIds); + innerBlocks.add(new AssistantBlockPlan(innerBlock.blockId(), innerBlock.blockType(), innerBlock.purpose(), + operation.name())); + } + return innerBlocks; + } + + /** + * Validates the shape of a container's inner block list (non-empty, unique non-blank ids, + * present blockType, no nested containers) and returns the raw entries. Shared by both the + * full-regeneration and incremental inner-block normalizers. + */ + static List validateInnerBlockShapes(AssistantContainerPlan containerPlan) { + List rawInnerBlocks = containerPlan.blocks() == null ? List.of() : containerPlan.blocks(); + if (rawInnerBlocks.isEmpty()) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned a container with no inner blocks: " + containerPlan.containerId()); + } + Set innerIds = new LinkedHashSet<>(); + for (AssistantBlockPlan innerBlock : rawInnerBlocks) { + if (innerBlock.blockId() == null || innerBlock.blockId().isBlank() || !innerIds.add(innerBlock.blockId())) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned invalid or duplicate inner block ids in container: " + + containerPlan.containerId()); + } + if (innerBlock.blockType() == null || innerBlock.blockType().isBlank()) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Assistant returned an inner block without blockType in container: " + + containerPlan.containerId()); + } + if (isContainerBlockType(innerBlock.blockType())) { + throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, + "Nested containers are not supported: container " + containerPlan.containerId() + + " cannot contain another container"); + } + } + return rawInnerBlocks; + } + + static PlanOperation normalizeInnerBlockOperation(AssistantBlockPlan block, List> existingInner, + OperationMode mode, List errors, int blockIndex, Set usedInnerIds) { + if (block.operation() != null && !block.operation().isBlank()) { + PlanOperation explicitOperation = parsePlanOperation(block.operation()); + if (explicitOperation != PlanOperation.ADD) { + Block matched = resolveExistingBlock(block, existingInner, blockIndex, usedInnerIds); + if (matched != null) { + usedInnerIds.add(matched.getId()); + } + } + return explicitOperation; + } + if (existingInner.isEmpty()) { + return PlanOperation.ADD; + } + Block matched = resolveExistingBlock(block, existingInner, blockIndex, usedInnerIds); + if (matched == null) { + return PlanOperation.ADD; + } + usedInnerIds.add(matched.getId()); + if (mode == OperationMode.FIX && hasValidationErrorForBlock(matched, errors)) { + return PlanOperation.UPDATE; + } + return PlanOperation.KEEP; + } + + static boolean isContainerBlockType(String blockType) { + return GenericContainerType.TYPE.equals(blockType) + || "IteratorContainer".equals(blockType) + || "LoopContainer".equals(blockType); + } + + static PlanOperation normalizePlanOperation(AssistantBlockPlan block, OperationMode mode, + FlowCreateRequest currentFlow, List errors, int blockIndex, + Set usedExistingBlockIds) { + if (block.operation() != null && !block.operation().isBlank()) { + PlanOperation explicitOperation = parsePlanOperation(block.operation()); + if (explicitOperation != PlanOperation.ADD) { + Block matched = resolveExistingBlockForPlan(block, currentFlow, blockIndex, usedExistingBlockIds); + if (matched != null) { + usedExistingBlockIds.add(matched.getId()); + } + } + return explicitOperation; + } + if (mode == OperationMode.DRAFT || !hasCurrentFlowBlocks(currentFlow)) { + return PlanOperation.ADD; + } + + Block matched = resolveExistingBlockForPlan(block, currentFlow, blockIndex, usedExistingBlockIds); + if (matched == null) { + return PlanOperation.ADD; + } + usedExistingBlockIds.add(matched.getId()); + if (mode == OperationMode.FIX && hasValidationErrorForBlock(matched, errors)) { + return PlanOperation.UPDATE; + } + return PlanOperation.KEEP; + } + + static boolean hasValidationErrorForBlock(Block block, List errors) { + if (block == null || block.getId() == null || errors == null || errors.isEmpty()) { + return false; + } + for (ValidationError error : errors) { + if (Objects.equals(block.getId(), error.id()) + || error.relatedNodeIds() != null && error.relatedNodeIds().contains(block.getId())) { + return true; + } + } + return false; + } + + /** + * Decides a container plan entry's operation from its already-resolved existing match. Kept + * separate from resolution so the caller can reuse the same matched container for incremental + * inner-block diffing. + */ + static PlanOperation resolveContainerOperation(AssistantContainerPlan containerPlan, Container existingContainer, + OperationMode mode, List errors) { + if (containerPlan.operation() != null && !containerPlan.operation().isBlank()) { + return parsePlanOperation(containerPlan.operation()); + } + if (mode == OperationMode.DRAFT || existingContainer == null) { + return PlanOperation.ADD; + } + if (mode == OperationMode.FIX && hasValidationErrorForContainer(existingContainer, errors)) { + return PlanOperation.UPDATE; + } + return PlanOperation.KEEP; + } + + static boolean hasValidationErrorForContainer(Container container, List errors) { + if (container == null || container.getId() == null || errors == null || errors.isEmpty()) { + return false; + } + for (ValidationError error : errors) { + if (Objects.equals(container.getId(), error.id()) + || error.relatedNodeIds() != null && error.relatedNodeIds().contains(container.getId())) { + return true; + } + } + return false; + } + + static boolean hasCurrentFlowContainers(FlowCreateRequest currentFlow) { + return currentFlow != null + && currentFlow.flow() != null + && currentFlow.flow().getContainers() != null + && !currentFlow.flow().getContainers().isEmpty(); + } + + static Container resolveExistingContainerForPlan(AssistantContainerPlan containerPlan, + FlowCreateRequest currentFlow, int planIndex, Set usedExistingContainerIds) { + List> existingContainers = currentFlow == null || currentFlow.flow() == null + || currentFlow.flow().getContainers() == null + ? List.of() + : currentFlow.flow().getContainers(); + if (existingContainers.isEmpty()) { + return null; + } + + Container exact = findExistingContainer(existingContainers, usedExistingContainerIds, + container -> Objects.equals(container.getId(), containerPlan.containerId()) + || AssistantTextSupport.normalizeBlockReference(containerPlan.containerId()) != null + && AssistantTextSupport.normalizeBlockReference(containerPlan.containerId()) + .equals(AssistantTextSupport.normalizeBlockReference(container.getName()))); + if (exact != null) { + return exact; + } + + Container byPurpose = findExistingContainer(existingContainers, usedExistingContainerIds, + container -> AssistantTextSupport.normalizeBlockReference(containerPlan.purpose()) != null + && AssistantTextSupport.normalizeBlockReference(containerPlan.purpose()) + .equals(AssistantTextSupport.normalizeBlockReference(container.getName()))); + if (byPurpose != null) { + return byPurpose; + } + + if (planIndex >= 0 && planIndex < existingContainers.size()) { + Container byPosition = existingContainers.get(planIndex); + if (!usedExistingContainerIds.contains(byPosition.getId())) { + return byPosition; + } + } + + List> unused = existingContainers.stream() + .filter(container -> !usedExistingContainerIds.contains(container.getId())) + .toList(); + return unused.size() == 1 ? unused.getFirst() : null; + } + + private static Container findExistingContainer(List> existingContainers, + Set usedExistingContainerIds, java.util.function.Predicate> predicate) { + for (Container container : existingContainers) { + if (container == null || usedExistingContainerIds.contains(container.getId())) { + continue; + } + if (predicate.test(container)) { + return container; + } + } + return null; + } + + static boolean canUseMinimalDraftFallback(String userPrompt, FlowCreateRequest currentFlow, + Set availableBlockTypes) { + return userPrompt != null + && !userPrompt.isBlank() + && !hasCurrentFlowBlocks(currentFlow) + && availableBlockTypes != null + && availableBlockTypes.contains("LLMBlock"); + } + + static boolean hasCurrentFlowBlocks(FlowCreateRequest currentFlow) { + return currentFlow != null + && currentFlow.flow() != null + && currentFlow.flow().getBlocks() != null + && !currentFlow.flow().getBlocks().isEmpty(); + } +}