refactor(assistant): extract plan validation/normalization into PlanValidationSupport

Clusters E (plan validation, normalization, KEEP/ADD/UPDATE/REMOVE
operation-diffing - including the 145-line validateAndNormalizePlan
"god method") and F (existing block/container identity matching) from
the structural analysis. Field-free except the PLAN_SENSITIVE_ERROR_CODES
constant.

- New PlanValidationSupport holds validateAndNormalizePlan and its full
  dependency tree: normalizeContainerInnerBlocks(AgainstExisting),
  validateInnerBlockShapes, normalize{Inner}BlockOperation,
  resolveContainerOperation, hasValidationErrorFor{Block,Container},
  resolveExisting{Block,Container}(ForPlan), findExisting{Block,Container},
  parsePlanOperation, isTargetedBlockRepairEligible,
  buildReusedPlanForTargetedRepair, canUseMinimalDraftFallback,
  hasCurrentFlow{Blocks,Containers}, isContainerBlockType, blockTypeName
- PLAN_SENSITIVE_ERROR_CODES promoted from private to package-private so
  the new class can reference it without duplicating the 48-entry set
- isSharedMemoryContext made static + package-private (it already didn't
  touch any instance field) so PlanValidationSupport can call it while it
  stays in FlowAssistantService, next to the AssistantFlowPlan/
  AssistantBlockPlan records it needs

2352 -> 1805 lines (-547, now under 1900 and past the halfway point of
the original file). 3350 -> 1805 total so far (-1545, ~46%).
Behavior-preserving: pure extraction, no logic changes.

Verified with `mvn test`: 462 tests, 0 failures, 0 errors.

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
This commit is contained in:
Lucio Lelii 2026-09-02 10:26:25 +02:00
parent 848c703619
commit 46a1fb2d3e
2 changed files with 589 additions and 561 deletions

View File

@ -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<ValidationErrorCode> PLAN_SENSITIVE_ERROR_CODES = EnumSet.of(
static final Set<ValidationErrorCode> 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<AssistantBlockPlan> 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.<AssistantContainerPlan>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<AssistantBlockPlan> 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<ValidationError> 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<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();
// 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<AssistantContainerPlan> 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<String> usedExistingBlockIds) {
List<Block<?>> 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<Block<?>> existingBlocks,
int planIndex, Set<String> 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<Block<?>> 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<Block<?>> existingBlocks, Set<String> usedExistingBlockIds,
java.util.function.Predicate<Block<?>> 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<Connection> preserveCurrentConnections(FlowCreateRequest currentFlow,
Map<String, FlowNode> oldNodeIdToAssembledNode, Set<String> removedExistingNodeIds) {
List<Connection> 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<ValidationError> errors,
Set<String> 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<String> ids = new LinkedHashSet<>();
Set<String> usedExistingBlockIds = new LinkedHashSet<>();
List<AssistantBlockPlan> normalizedBlocks = new ArrayList<>();
List<AssistantBlockPlan> 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<AssistantContainerPlan> normalizedContainers = new ArrayList<>();
Set<String> usedExistingContainerIds = new LinkedHashSet<>();
List<AssistantContainerPlan> 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<AssistantBlockPlan> 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<AssistantBlockPlan> 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<AssistantBlockPlan> normalizeContainerInnerBlocks(AssistantContainerPlan containerPlan) {
List<AssistantBlockPlan> rawInnerBlocks = validateInnerBlockShapes(containerPlan);
List<AssistantBlockPlan> 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<AssistantBlockPlan> normalizeContainerInnerBlocksAgainstExisting(AssistantContainerPlan containerPlan,
Container<?> existingContainer, OperationMode mode, List<ValidationError> errors) {
List<AssistantBlockPlan> rawInnerBlocks = validateInnerBlockShapes(containerPlan);
List<Block<?>> existingInner = FlowAssemblySupport.subFlowBlocks(existingContainer);
Set<String> usedInnerIds = new LinkedHashSet<>();
List<AssistantBlockPlan> 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<AssistantBlockPlan> validateInnerBlockShapes(AssistantContainerPlan containerPlan) {
List<AssistantBlockPlan> 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<String> 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<Block<?>> existingInner,
OperationMode mode, List<ValidationError> errors, int blockIndex, Set<String> 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<ValidationError> errors, int blockIndex,
Set<String> 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<ValidationError> 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<ValidationError> 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<ValidationError> 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<String> usedExistingContainerIds) {
List<Container<?>> 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<Container<?>> unused = existingContainers.stream()
.filter(container -> !usedExistingContainerIds.contains(container.getId()))
.toList();
return unused.size() == 1 ? unused.getFirst() : null;
}
private Container<?> findExistingContainer(List<Container<?>> existingContainers,
Set<String> usedExistingContainerIds, java.util.function.Predicate<Container<?>> 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<String> 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;
}

View File

@ -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<ValidationError> 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<String> 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<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();
// 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<AssistantContainerPlan> 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<String> usedExistingBlockIds) {
List<Block<?>> 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<Block<?>> existingBlocks,
int planIndex, Set<String> 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<Block<?>> 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<Block<?>> existingBlocks, Set<String> usedExistingBlockIds,
java.util.function.Predicate<Block<?>> 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<ValidationError> errors,
Set<String> 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<String> ids = new LinkedHashSet<>();
Set<String> usedExistingBlockIds = new LinkedHashSet<>();
List<AssistantBlockPlan> normalizedBlocks = new ArrayList<>();
List<AssistantBlockPlan> 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<AssistantContainerPlan> normalizedContainers = new ArrayList<>();
Set<String> usedExistingContainerIds = new LinkedHashSet<>();
List<AssistantContainerPlan> 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<AssistantBlockPlan> 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<AssistantBlockPlan> 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<AssistantBlockPlan> normalizeContainerInnerBlocks(AssistantContainerPlan containerPlan) {
List<AssistantBlockPlan> rawInnerBlocks = validateInnerBlockShapes(containerPlan);
List<AssistantBlockPlan> 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<AssistantBlockPlan> normalizeContainerInnerBlocksAgainstExisting(AssistantContainerPlan containerPlan,
Container<?> existingContainer, OperationMode mode, List<ValidationError> errors) {
List<AssistantBlockPlan> rawInnerBlocks = validateInnerBlockShapes(containerPlan);
List<Block<?>> existingInner = FlowAssemblySupport.subFlowBlocks(existingContainer);
Set<String> usedInnerIds = new LinkedHashSet<>();
List<AssistantBlockPlan> 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<AssistantBlockPlan> validateInnerBlockShapes(AssistantContainerPlan containerPlan) {
List<AssistantBlockPlan> 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<String> 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<Block<?>> existingInner,
OperationMode mode, List<ValidationError> errors, int blockIndex, Set<String> 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<ValidationError> errors, int blockIndex,
Set<String> 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<ValidationError> 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<ValidationError> 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<ValidationError> 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<String> usedExistingContainerIds) {
List<Container<?>> 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<Container<?>> unused = existingContainers.stream()
.filter(container -> !usedExistingContainerIds.contains(container.getId()))
.toList();
return unused.size() == 1 ? unused.getFirst() : null;
}
private static Container<?> findExistingContainer(List<Container<?>> existingContainers,
Set<String> usedExistingContainerIds, java.util.function.Predicate<Container<?>> 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<String> 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();
}
}