refactor(assistant): extract connection resolution into ConnectionAssembler

Cluster J: connection drafting/preservation/merging/dangling-drop and
LoopContainer body-output chaining - 23 methods, field-free except a
logger, using only FlowNode/Block/Container/Connection/IODescriptor and
the now-shared AssistantConnectionDraft record.

- New ConnectionAssembler holds toValidConnections/toConnection,
  preserveConnections, mergeConnections, dropDanglingConnections,
  chainStrandedLoopBodyOutputs (+ their private helpers: resolveIoName,
  findIoByName, stripHandleNoise, resolveConnectionBlock, inferBlockByIo,
  isOpenBodyOutput, firstOpenDataInput, registerNodeAlias, etc.)
- preserveCurrentConnections stays in FlowAssistantService (thin wrapper
  over FlowCreateRequest) but now delegates to
  ConnectionAssembler.preserveConnections

2988 -> 2571 lines (-417). 3350 -> 2571 total so far (-779, ~23%).
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:16:47 +02:00
parent 76aa3fa71d
commit cf068ffe7f
2 changed files with 462 additions and 436 deletions

View File

@ -0,0 +1,443 @@
package it.cnr.isti.workflow.manager.assistant;
import java.util.ArrayList;
import java.util.Collection;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.http.HttpStatus;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.assistant.FlowAssistantService.AssistantConnectionDraft;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.containers.Container;
import it.cnr.isti.workflow.manager.containers.iresolvers.ContainerFlowInterfaceResolver;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.FlowNode;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
final class ConnectionAssembler {
private static final Logger log = LoggerFactory.getLogger(ConnectionAssembler.class);
private ConnectionAssembler() {
}
static List<Connection> preserveConnections(List<Connection> sourceConnections,
Map<String, FlowNode> oldNodeIdToAssembledNode, Set<String> removedExistingNodeIds) {
if (sourceConnections == null || sourceConnections.isEmpty()) {
return List.of();
}
List<Connection> preserved = new ArrayList<>();
for (Connection connection : sourceConnections) {
if (connection == null
|| removedExistingNodeIds.contains(connection.getSourceId())
|| removedExistingNodeIds.contains(connection.getTargetId())) {
continue;
}
FlowNode source = oldNodeIdToAssembledNode.get(connection.getSourceId());
FlowNode target = oldNodeIdToAssembledNode.get(connection.getTargetId());
if (source == null || target == null) {
continue;
}
// A reused node id may point at a freshly reconfigured block whose I/O shape changed
// (e.g. a different prompt placeholder), so the old connection's handle names can be
// stale. Carrying it forward unchanged would leave a dangling reference that only
// surfaces later as CONNECTION_TARGET_INPUT_NOT_FOUND; drop it here instead.
if (!hasOutputNamed(source, connection.getSourceName()) || !hasInputNamed(target, connection.getTargetName())) {
continue;
}
preserved.add(Connection.builder()
.sourceId(source.getId())
.sourceName(connection.getSourceName())
.targetId(target.getId())
.targetName(connection.getTargetName())
.build());
}
return preserved;
}
static boolean hasOutputNamed(FlowNode node, String name) {
return node != null && node.getOutputs() != null && node.getOutputs().stream()
.anyMatch(io -> io != null && Objects.equals(io.getName(), name));
}
static boolean hasInputNamed(FlowNode node, String name) {
return node != null && node.getInputs() != null && node.getInputs().stream()
.anyMatch(io -> io != null && Objects.equals(io.getName(), name));
}
static List<Connection> mergeConnections(List<Connection> preservedConnections, List<Connection> generatedConnections) {
Map<String, Connection> merged = new LinkedHashMap<>();
for (Connection connection : preservedConnections == null ? List.<Connection>of() : preservedConnections) {
merged.put(connectionKey(connection), connection);
}
for (Connection connection : generatedConnections == null ? List.<Connection>of() : generatedConnections) {
merged.putIfAbsent(connectionKey(connection), connection);
}
return List.copyOf(merged.values());
}
private static String connectionKey(Connection connection) {
if (connection == null) {
return "";
}
return connection.getSourceId() + "|" + connection.getSourceName() + "|"
+ connection.getTargetId() + "|" + connection.getTargetName();
}
/**
* Final safety net applied to the fully merged connection list (preserved + freshly generated),
* at both the top-level flow and every container subflow: drops any connection whose source or
* target node/handle doesn't resolve against the CURRENT node graph. A dangling connection like
* this is a hard, save-blocking bean-validation failure (FlowDataValidator.validateConnection),
* not merely a "not executable" one - without this, an assistant-produced flow could come back
* unable to be saved even as a draft. Dropping a connection whose handle name was never real to
* begin with does not change which handles a node exposes, so this is safe regardless of where
* the staleness came from (a known preserve-path bug, or a not-yet-discovered one).
*/
static List<Connection> dropDanglingConnections(List<Connection> connections, List<Block<?>> blocks,
List<Container<?>> containers) {
if (connections == null || connections.isEmpty()) {
return List.of();
}
Map<String, FlowNode> nodesById = new LinkedHashMap<>();
for (Block<?> block : blocks == null ? List.<Block<?>>of() : blocks) {
nodesById.put(block.getId(), block);
}
for (Container<?> container : containers == null ? List.<Container<?>>of() : containers) {
nodesById.put(container.getId(), container);
}
List<Connection> kept = new ArrayList<>();
for (Connection connection : connections) {
if (connection == null) {
continue;
}
FlowNode source = nodesById.get(connection.getSourceId());
FlowNode target = nodesById.get(connection.getTargetId());
if (source == null || target == null
|| !hasOutputNamed(source, connection.getSourceName())
|| !hasInputNamed(target, connection.getTargetName())) {
log.warn("Dropping dangling assistant connection {} -> {} ({} -> {}): endpoint no longer resolves",
connection.getSourceId(), connection.getTargetId(),
connection.getSourceName(), connection.getTargetName());
continue;
}
kept.add(connection);
}
return List.copyOf(kept);
}
/**
* Collapses a LoopContainer body to a single open output by forwarding any stranded
* intermediate producer output into a later block's first open data input, in plan order.
* Only adds connections (never removes), and only forwards a block that has exactly ONE output
* (a clear primary producer like LLMBlock/MCPAgent "response") - branch blocks (HumanDecision
* yes/no, Conditional true/false, Switch cases) are left untouched, so exclusive routing can
* never be mis-wired. Forward-only (source earlier than target in plan order), so no cycle can
* be introduced. When the model already produced a clean chain this is a no-op.
*/
static List<Connection> chainStrandedLoopBodyOutputs(List<Block<?>> blocks, List<Connection> connections) {
if (blocks == null || blocks.size() < 2) {
return connections;
}
List<Connection> result = new ArrayList<>(connections == null ? List.of() : connections);
for (int i = 0; i < blocks.size() - 1; i++) {
Block<?> source = blocks.get(i);
if (source.getOutputs() == null || source.getOutputs().size() != 1) {
continue;
}
String outputName = source.getOutputs().getFirst().getName();
if (!isOpenBodyOutput(blocks, result, source.getId(), outputName)) {
continue;
}
for (int j = i + 1; j < blocks.size(); j++) {
Block<?> target = blocks.get(j);
String inputName = firstOpenDataInput(blocks, result, target);
if (inputName == null) {
continue;
}
result.add(Connection.builder()
.sourceId(source.getId())
.sourceName(outputName)
.targetId(target.getId())
.targetName(inputName)
.build());
break;
}
}
return List.copyOf(result);
}
/** True when {@code outputName} on {@code blockId} is still open (not sourced by any connection). */
private static boolean isOpenBodyOutput(List<Block<?>> blocks, List<Connection> connections, String blockId,
String outputName) {
FlowData probe = FlowData.builder().blocks(blocks).connections(connections).build();
return ContainerFlowInterfaceResolver.getOpenOutputs(probe).stream()
.anyMatch(handle -> handle.blockId().equals(blockId) && handle.io().getName().equals(outputName));
}
/** First still-open, non-technical (not {@code model}) input name of {@code target}, or null. */
private static String firstOpenDataInput(List<Block<?>> blocks, List<Connection> connections, Block<?> target) {
FlowData probe = FlowData.builder().blocks(blocks).connections(connections).build();
return ContainerFlowInterfaceResolver.getOpenInputs(probe).stream()
.filter(handle -> handle.blockId().equals(target.getId()))
.map(handle -> handle.io().getName())
.filter(name -> !"model".equals(name))
.findFirst()
.orElse(null);
}
static List<Connection> toValidConnections(List<AssistantConnectionDraft> draftedConnections,
Map<String, FlowNode> nodesByPlanId, Map<String, FlowNode> nodesByAlias) {
if (draftedConnections == null || draftedConnections.isEmpty()) {
return List.of();
}
List<Connection> validConnections = new ArrayList<>();
for (AssistantConnectionDraft draftedConnection : draftedConnections) {
try {
validConnections.add(toConnection(draftedConnection, nodesByPlanId, nodesByAlias));
} catch (ResponseStatusException e) {
if (!HttpStatus.BAD_GATEWAY.equals(e.getStatusCode())) {
throw e;
}
log.warn("Skipping invalid assistant connection draft: {}. Reason: {}",
draftedConnection,
e.getReason());
}
}
return List.copyOf(validConnections);
}
private static Connection toConnection(AssistantConnectionDraft connection, Map<String, FlowNode> nodesByPlanId,
Map<String, FlowNode> nodesByAlias) {
FlowNode source = resolveConnectionBlock(connection.fromBlockId(), nodesByPlanId, nodesByAlias);
FlowNode target = resolveConnectionBlock(connection.toBlockId(), nodesByPlanId, nodesByAlias);
if (source == null) {
source = inferBlockByIo(connection.fromOutput(), nodesByPlanId.values(), true);
}
if (target == null) {
target = inferBlockByIo(connection.toInput(), nodesByPlanId.values(), false);
}
if (source == null || target == null) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection with unknown block ids"
+ " (fromBlockId=" + connection.fromBlockId()
+ ", toBlockId=" + connection.toBlockId()
+ ", allowedBlockIds=" + nodesByPlanId.keySet() + ")");
}
String sourceName = resolveConnectionOutputName(source, connection.fromOutput());
String targetName = resolveConnectionInputName(target, connection.toInput());
if (sourceName == null) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection with unresolved output '" + connection.fromOutput()
+ "' on block '" + source.getName() + "' (available: "
+ (source.getOutputs() == null ? "none" : source.getOutputs().stream()
.map(io -> io.getName()).toList()) + ")");
}
if (targetName == null) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection with unresolved input '" + connection.toInput()
+ "' on block '" + target.getName() + "' (available: "
+ (target.getInputs() == null ? "none" : target.getInputs().stream()
.map(io -> io.getName()).toList()) + ")");
}
if (isModelInputConnection(sourceName, targetName)) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection to technical input 'model' on block '" + target.getName()
+ "'");
}
return Connection.builder()
.sourceId(source.getId())
.sourceName(sourceName)
.targetId(target.getId())
.targetName(targetName)
.build();
}
private static boolean isModelInputConnection(String sourceName, String targetName) {
return "model".equals(AssistantTextSupport.normalizeBlockReference(targetName))
&& !"model".equals(AssistantTextSupport.normalizeBlockReference(sourceName));
}
static List<AssistantConnectionDraft> inferSequentialConnections(List<Block<?>> blocks) {
if (blocks == null || blocks.size() < 2) {
return List.of();
}
List<AssistantConnectionDraft> connections = new ArrayList<>();
for (int i = 0; i < blocks.size() - 1; i++) {
Block<?> source = blocks.get(i);
Block<?> target = blocks.get(i + 1);
String sourceOutput = resolveConnectionOutputName(source, "response");
String targetInput = resolveSequentialInputName(target);
if (sourceOutput == null || targetInput == null) {
continue;
}
connections.add(new AssistantConnectionDraft(
source.getName(),
sourceOutput,
target.getName(),
targetInput));
}
return connections;
}
private static String resolveSequentialInputName(Block<?> block) {
List<IODescriptor> inputs = block.getInputs();
if (inputs == null || inputs.isEmpty()) {
return null;
}
if (inputs.size() == 1) {
return inputs.getFirst().getName();
}
return inputs.stream()
.map(IODescriptor::getName)
.filter(name -> !isLikelyUserProvidedInput(name))
.findFirst()
.orElse(inputs.getFirst().getName());
}
private static boolean isLikelyUserProvidedInput(String name) {
String normalized = AssistantTextSupport.normalizeBlockReference(name);
return normalized != null
&& Set.of("query", "question", "user_query", "userquery", "url", "file_url", "fileurl")
.contains(normalized);
}
private static FlowNode resolveConnectionBlock(String rawReference, Map<String, FlowNode> nodesByPlanId,
Map<String, FlowNode> nodesByAlias) {
if (rawReference == null || rawReference.isBlank()) {
return null;
}
FlowNode direct = nodesByPlanId.get(rawReference);
if (direct != null) {
return direct;
}
FlowNode byAlias = nodesByAlias.get(AssistantTextSupport.normalizeBlockReference(rawReference));
if (byAlias != null) {
return byAlias;
}
// Same syntactic-noise fallback as findIoByName: a brace-mangled or ${{...}}-wrapped block
// reference should still resolve to its node instead of dropping the whole connection.
String sanitized = AssistantTextSupport.normalizeBlockReference(stripHandleNoise(rawReference));
return sanitized == null ? null : nodesByAlias.get(sanitized);
}
private static FlowNode inferBlockByIo(String ioName, Collection<FlowNode> nodes, boolean output) {
String normalizedIo = AssistantTextSupport.normalizeBlockReference(ioName);
if (normalizedIo == null) {
return null;
}
FlowNode match = null;
for (FlowNode node : nodes) {
List<IODescriptor> ioDescriptors = output ? node.getOutputs() : node.getInputs();
if (ioDescriptors == null) {
continue;
}
boolean hasMatch = ioDescriptors.stream()
.anyMatch(io -> normalizedIo.equals(AssistantTextSupport.normalizeBlockReference(io.getName())));
if (!hasMatch) {
continue;
}
if (match != null) {
return null;
}
match = node;
}
return match;
}
private static String resolveConnectionOutputName(FlowNode node, String requestedOutput) {
return resolveIoName(node.getOutputs(), requestedOutput, List.of("response", "output", "true", "false"));
}
private static String resolveConnectionInputName(FlowNode node, String requestedInput) {
return resolveIoName(node.getInputs(), requestedInput, List.of("input", "prompt"));
}
private static String resolveIoName(List<IODescriptor> descriptors, String requestedName, List<String> preferredNames) {
if (descriptors == null || descriptors.isEmpty()) {
return null;
}
if (requestedName != null && !requestedName.isBlank()) {
IODescriptor exact = findIoByName(descriptors, requestedName);
if (exact != null) {
return exact.getName();
}
}
if (descriptors.size() == 1) {
return descriptors.getFirst().getName();
}
for (String preferredName : preferredNames) {
IODescriptor preferred = findIoByName(descriptors, preferredName);
if (preferred != null) {
return preferred.getName();
}
}
return null;
}
private static IODescriptor findIoByName(List<IODescriptor> descriptors, String requestedName) {
String normalized = AssistantTextSupport.normalizeBlockReference(requestedName);
if (normalized == null) {
return null;
}
IODescriptor exact = descriptors.stream()
.filter(io -> normalized.equals(AssistantTextSupport.normalizeBlockReference(io.getName())))
.findFirst()
.orElse(null);
if (exact != null) {
return exact;
}
// Fallback for syntactic noise the model sometimes emits in a handle name - a stray or
// unbalanced brace ("{userReview"), an accidental ${{...}} wrapper ("${{userReview}}"), or a
// block-qualified reference ("SomeBlock.response"). Strip that noise and retry; this only
// runs after the exact match already failed, so it can never change a currently-resolving
// connection - it only rescues one that would otherwise be silently dropped.
String sanitized = AssistantTextSupport.normalizeBlockReference(stripHandleNoise(requestedName));
if (sanitized == null || sanitized.equals(normalized)) {
return null;
}
return descriptors.stream()
.filter(io -> sanitized.equals(AssistantTextSupport.normalizeBlockReference(io.getName())))
.findFirst()
.orElse(null);
}
/**
* Removes purely-syntactic noise from a handle/block reference the model produced: placeholder
* wrappers ({@code ${{...}}} / {@code {{...}}}), stray braces, dollar signs and quotes, and a
* leading block-name qualifier (keeps the last dotted segment, so {@code "Review.response"}
* becomes {@code "response"}). Deterministic clean-up only - it never invents a name, so a
* reference that was genuinely wrong (not just mangled) still fails to resolve and is dropped.
*/
private static String stripHandleNoise(String raw) {
if (raw == null) {
return null;
}
String cleaned = raw.replace("${{", "").replace("{{", "").replace("}}", "");
cleaned = cleaned.replaceAll("[{}$\"']", "").trim();
int lastDot = cleaned.lastIndexOf('.');
if (lastDot >= 0 && lastDot < cleaned.length() - 1) {
cleaned = cleaned.substring(lastDot + 1).trim();
}
return cleaned.isBlank() ? null : cleaned;
}
static void registerNodeAlias(Map<String, FlowNode> nodesByAlias, String reference, FlowNode node) {
String normalized = AssistantTextSupport.normalizeBlockReference(reference);
if (normalized != null) {
nodesByAlias.putIfAbsent(normalized, node);
}
}
}

View File

@ -557,9 +557,9 @@ public class FlowAssistantService {
if (existingBlock != null && operation != PlanOperation.ADD) {
oldNodeIdToAssembledNode.put(existingBlock.getId(), assembledBlock);
}
registerNodeAlias(nodesByAlias, blockPlan.blockId(), assembledBlock);
registerNodeAlias(nodesByAlias, assembledBlock.getName(), assembledBlock);
registerNodeAlias(nodesByAlias, blockPlan.purpose(), assembledBlock);
ConnectionAssembler.registerNodeAlias(nodesByAlias, blockPlan.blockId(), assembledBlock);
ConnectionAssembler.registerNodeAlias(nodesByAlias, assembledBlock.getName(), assembledBlock);
ConnectionAssembler.registerNodeAlias(nodesByAlias, blockPlan.purpose(), assembledBlock);
configuredNodes.add(new ConfiguredBlockSummary(
blockPlan.blockId(),
blockPlan.blockType(),
@ -607,9 +607,9 @@ public class FlowAssistantService {
if (existingContainer != null && containerOperation != PlanOperation.ADD) {
oldNodeIdToAssembledNode.put(existingContainer.getId(), assembledContainer);
}
registerNodeAlias(nodesByAlias, containerPlan.containerId(), assembledContainer);
registerNodeAlias(nodesByAlias, assembledContainer.getName(), assembledContainer);
registerNodeAlias(nodesByAlias, containerPlan.purpose(), assembledContainer);
ConnectionAssembler.registerNodeAlias(nodesByAlias, containerPlan.containerId(), assembledContainer);
ConnectionAssembler.registerNodeAlias(nodesByAlias, assembledContainer.getName(), assembledContainer);
ConnectionAssembler.registerNodeAlias(nodesByAlias, containerPlan.purpose(), assembledContainer);
configuredNodes.add(new ConfiguredBlockSummary(
containerPlan.containerId(),
containerPlan.containerType(),
@ -638,15 +638,15 @@ public class FlowAssistantService {
// MCP shared-session ordering is expressed as a Dependency (see
// buildTopLevelSharedMemoryDependencies), not a data connection, so no
// sequential-connection completion is needed here.
toValidConnections(parsed.connections(), nodesByPlanId, nodesByAlias);
ConnectionAssembler.toValidConnections(parsed.connections(), nodesByPlanId, nodesByAlias);
return parsed;
});
}
appendRationale(rationaleParts, parsedConnections.rationale());
List<Connection> generatedConnections = toValidConnections(parsedConnections.connections(), nodesByPlanId,
List<Connection> generatedConnections = ConnectionAssembler.toValidConnections(parsedConnections.connections(), nodesByPlanId,
nodesByAlias);
List<Connection> connections = dropDanglingConnections(mergeConnections(
List<Connection> connections = ConnectionAssembler.dropDanglingConnections(ConnectionAssembler.mergeConnections(
preserveCurrentConnections(currentFlow, oldNodeIdToAssembledNode, removedExistingNodeIds),
generatedConnections), assembledBlocks, assembledContainers);
@ -761,9 +761,9 @@ public class FlowAssistantService {
if (existingInnerBlock != null && operation != PlanOperation.ADD) {
oldInnerIdToAssembled.put(existingInnerBlock.getId(), assembledBlock);
}
registerNodeAlias(innerNodesByAlias, blockPlan.blockId(), assembledBlock);
registerNodeAlias(innerNodesByAlias, assembledBlock.getName(), assembledBlock);
registerNodeAlias(innerNodesByAlias, blockPlan.purpose(), assembledBlock);
ConnectionAssembler.registerNodeAlias(innerNodesByAlias, blockPlan.blockId(), assembledBlock);
ConnectionAssembler.registerNodeAlias(innerNodesByAlias, assembledBlock.getName(), assembledBlock);
ConnectionAssembler.registerNodeAlias(innerNodesByAlias, blockPlan.purpose(), assembledBlock);
innerConfiguredNodes.add(new ConfiguredBlockSummary(
blockPlan.blockId(),
blockPlan.blockType(),
@ -787,20 +787,20 @@ public class FlowAssistantService {
ParsedConnections parsed = parseConnectionsOrInferSequential(rawResponse, innerBlocks);
// Within-container MCP ordering is a Dependency in the subflow (see below),
// not a data connection.
toValidConnections(parsed.connections(), innerNodesByPlanId, innerNodesByAlias);
ConnectionAssembler.toValidConnections(parsed.connections(), innerNodesByPlanId, innerNodesByAlias);
return parsed;
});
}
appendRationale(rationaleParts, innerConnections.rationale());
List<Connection> generatedConnections = toValidConnections(innerConnections.connections(), innerNodesByPlanId,
List<Connection> generatedConnections = ConnectionAssembler.toValidConnections(innerConnections.connections(), innerNodesByPlanId,
innerNodesByAlias);
List<Connection> connections = dropDanglingConnections(mergeConnections(
preserveConnections(subFlowConnections(existingContainer), oldInnerIdToAssembled, removedInnerIds),
List<Connection> connections = ConnectionAssembler.dropDanglingConnections(ConnectionAssembler.mergeConnections(
ConnectionAssembler.preserveConnections(subFlowConnections(existingContainer), oldInnerIdToAssembled, removedInnerIds),
generatedConnections), innerBlocks, List.of());
if (LoopContainerType.TYPE.equals(containerPlan.containerType())) {
// Keep the loop body a single chain so it exposes exactly one output for the guard,
// instead of leaking several qualified/dotted outputs when the model under-connects it.
connections = chainStrandedLoopBodyOutputs(innerBlocks, connections);
connections = ConnectionAssembler.chainStrandedLoopBodyOutputs(innerBlocks, connections);
}
FlowData subFlow = FlowData.builder()
@ -1039,7 +1039,7 @@ public class FlowAssistantService {
if (e.getReason() == null || !e.getReason().contains("invalid connections payload")) {
throw e;
}
List<AssistantConnectionDraft> inferred = inferSequentialConnections(assembledBlocks);
List<AssistantConnectionDraft> inferred = ConnectionAssembler.inferSequentialConnections(assembledBlocks);
if (inferred.isEmpty()) {
throw e;
}
@ -1741,275 +1741,7 @@ public class FlowAssistantService {
List<Connection> sourceConnections = currentFlow == null || currentFlow.flow() == null
? List.of()
: currentFlow.flow().getConnections();
return preserveConnections(sourceConnections, oldNodeIdToAssembledNode, removedExistingNodeIds);
}
/**
* Carries an existing connection forward (rebound to the assembled nodes' ids) when both its
* endpoints survived - i.e. were reused or reconfigured rather than removed. Shared by the
* top-level flow and a container's inner subflow so both preserve connections the same way.
*/
private List<Connection> preserveConnections(List<Connection> sourceConnections,
Map<String, FlowNode> oldNodeIdToAssembledNode, Set<String> removedExistingNodeIds) {
if (sourceConnections == null || sourceConnections.isEmpty()) {
return List.of();
}
List<Connection> preserved = new ArrayList<>();
for (Connection connection : sourceConnections) {
if (connection == null
|| removedExistingNodeIds.contains(connection.getSourceId())
|| removedExistingNodeIds.contains(connection.getTargetId())) {
continue;
}
FlowNode source = oldNodeIdToAssembledNode.get(connection.getSourceId());
FlowNode target = oldNodeIdToAssembledNode.get(connection.getTargetId());
if (source == null || target == null) {
continue;
}
// A reused node id may point at a freshly reconfigured block whose I/O shape changed
// (e.g. a different prompt placeholder), so the old connection's handle names can be
// stale. Carrying it forward unchanged would leave a dangling reference that only
// surfaces later as CONNECTION_TARGET_INPUT_NOT_FOUND; drop it here instead.
if (!hasOutputNamed(source, connection.getSourceName()) || !hasInputNamed(target, connection.getTargetName())) {
continue;
}
preserved.add(Connection.builder()
.sourceId(source.getId())
.sourceName(connection.getSourceName())
.targetId(target.getId())
.targetName(connection.getTargetName())
.build());
}
return preserved;
}
private boolean hasOutputNamed(FlowNode node, String name) {
return node != null && node.getOutputs() != null && node.getOutputs().stream()
.anyMatch(io -> io != null && Objects.equals(io.getName(), name));
}
private boolean hasInputNamed(FlowNode node, String name) {
return node != null && node.getInputs() != null && node.getInputs().stream()
.anyMatch(io -> io != null && Objects.equals(io.getName(), name));
}
private List<Connection> mergeConnections(List<Connection> preservedConnections, List<Connection> generatedConnections) {
Map<String, Connection> merged = new LinkedHashMap<>();
for (Connection connection : preservedConnections == null ? List.<Connection>of() : preservedConnections) {
merged.put(connectionKey(connection), connection);
}
for (Connection connection : generatedConnections == null ? List.<Connection>of() : generatedConnections) {
merged.putIfAbsent(connectionKey(connection), connection);
}
return List.copyOf(merged.values());
}
private String connectionKey(Connection connection) {
if (connection == null) {
return "";
}
return connection.getSourceId() + "|" + connection.getSourceName() + "|"
+ connection.getTargetId() + "|" + connection.getTargetName();
}
/**
* Final safety net applied to the fully merged connection list (preserved + freshly generated),
* at both the top-level flow and every container subflow: drops any connection whose source or
* target node/handle doesn't resolve against the CURRENT node graph. A dangling connection like
* this is a hard, save-blocking bean-validation failure (FlowDataValidator.validateConnection),
* not merely a "not executable" one - without this, an assistant-produced flow could come back
* unable to be saved even as a draft. Dropping a connection whose handle name was never real to
* begin with does not change which handles a node exposes, so this is safe regardless of where
* the staleness came from (a known preserve-path bug, or a not-yet-discovered one).
*/
private List<Connection> dropDanglingConnections(List<Connection> connections, List<Block<?>> blocks,
List<Container<?>> containers) {
if (connections == null || connections.isEmpty()) {
return List.of();
}
Map<String, FlowNode> nodesById = new LinkedHashMap<>();
for (Block<?> block : blocks == null ? List.<Block<?>>of() : blocks) {
nodesById.put(block.getId(), block);
}
for (Container<?> container : containers == null ? List.<Container<?>>of() : containers) {
nodesById.put(container.getId(), container);
}
List<Connection> kept = new ArrayList<>();
for (Connection connection : connections) {
if (connection == null) {
continue;
}
FlowNode source = nodesById.get(connection.getSourceId());
FlowNode target = nodesById.get(connection.getTargetId());
if (source == null || target == null
|| !hasOutputNamed(source, connection.getSourceName())
|| !hasInputNamed(target, connection.getTargetName())) {
log.warn("Dropping dangling assistant connection {} -> {} ({} -> {}): endpoint no longer resolves",
connection.getSourceId(), connection.getTargetId(),
connection.getSourceName(), connection.getTargetName());
continue;
}
kept.add(connection);
}
return List.copyOf(kept);
}
/**
* Collapses a LoopContainer body to a single open output by forwarding any stranded
* intermediate producer output into a later block's first open data input, in plan order.
* Only adds connections (never removes), and only forwards a block that has exactly ONE output
* (a clear primary producer like LLMBlock/MCPAgent "response") - branch blocks (HumanDecision
* yes/no, Conditional true/false, Switch cases) are left untouched, so exclusive routing can
* never be mis-wired. Forward-only (source earlier than target in plan order), so no cycle can
* be introduced. When the model already produced a clean chain this is a no-op.
*/
private List<Connection> chainStrandedLoopBodyOutputs(List<Block<?>> blocks, List<Connection> connections) {
if (blocks == null || blocks.size() < 2) {
return connections;
}
List<Connection> result = new ArrayList<>(connections == null ? List.of() : connections);
for (int i = 0; i < blocks.size() - 1; i++) {
Block<?> source = blocks.get(i);
if (source.getOutputs() == null || source.getOutputs().size() != 1) {
continue;
}
String outputName = source.getOutputs().getFirst().getName();
if (!isOpenBodyOutput(blocks, result, source.getId(), outputName)) {
continue;
}
for (int j = i + 1; j < blocks.size(); j++) {
Block<?> target = blocks.get(j);
String inputName = firstOpenDataInput(blocks, result, target);
if (inputName == null) {
continue;
}
result.add(Connection.builder()
.sourceId(source.getId())
.sourceName(outputName)
.targetId(target.getId())
.targetName(inputName)
.build());
break;
}
}
return List.copyOf(result);
}
/** True when {@code outputName} on {@code blockId} is still open (not sourced by any connection). */
private boolean isOpenBodyOutput(List<Block<?>> blocks, List<Connection> connections, String blockId,
String outputName) {
FlowData probe = FlowData.builder().blocks(blocks).connections(connections).build();
return ContainerFlowInterfaceResolver.getOpenOutputs(probe).stream()
.anyMatch(handle -> handle.blockId().equals(blockId) && handle.io().getName().equals(outputName));
}
/** First still-open, non-technical (not {@code model}) input name of {@code target}, or null. */
private String firstOpenDataInput(List<Block<?>> blocks, List<Connection> connections, Block<?> target) {
FlowData probe = FlowData.builder().blocks(blocks).connections(connections).build();
return ContainerFlowInterfaceResolver.getOpenInputs(probe).stream()
.filter(handle -> handle.blockId().equals(target.getId()))
.map(handle -> handle.io().getName())
.filter(name -> !"model".equals(name))
.findFirst()
.orElse(null);
}
private List<Connection> toValidConnections(List<AssistantConnectionDraft> draftedConnections,
Map<String, FlowNode> nodesByPlanId, Map<String, FlowNode> nodesByAlias) {
if (draftedConnections == null || draftedConnections.isEmpty()) {
return List.of();
}
List<Connection> validConnections = new ArrayList<>();
for (AssistantConnectionDraft draftedConnection : draftedConnections) {
try {
validConnections.add(toConnection(draftedConnection, nodesByPlanId, nodesByAlias));
} catch (ResponseStatusException e) {
if (!HttpStatus.BAD_GATEWAY.equals(e.getStatusCode())) {
throw e;
}
log.warn("Skipping invalid assistant connection draft: {}. Reason: {}",
draftedConnection,
e.getReason());
}
}
return List.copyOf(validConnections);
}
private Connection toConnection(AssistantConnectionDraft connection, Map<String, FlowNode> nodesByPlanId,
Map<String, FlowNode> nodesByAlias) {
FlowNode source = resolveConnectionBlock(connection.fromBlockId(), nodesByPlanId, nodesByAlias);
FlowNode target = resolveConnectionBlock(connection.toBlockId(), nodesByPlanId, nodesByAlias);
if (source == null) {
source = inferBlockByIo(connection.fromOutput(), nodesByPlanId.values(), true);
}
if (target == null) {
target = inferBlockByIo(connection.toInput(), nodesByPlanId.values(), false);
}
if (source == null || target == null) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection with unknown block ids"
+ " (fromBlockId=" + connection.fromBlockId()
+ ", toBlockId=" + connection.toBlockId()
+ ", allowedBlockIds=" + nodesByPlanId.keySet() + ")");
}
String sourceName = resolveConnectionOutputName(source, connection.fromOutput());
String targetName = resolveConnectionInputName(target, connection.toInput());
if (sourceName == null) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection with unresolved output '" + connection.fromOutput()
+ "' on block '" + source.getName() + "' (available: "
+ (source.getOutputs() == null ? "none" : source.getOutputs().stream()
.map(io -> io.getName()).toList()) + ")");
}
if (targetName == null) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection with unresolved input '" + connection.toInput()
+ "' on block '" + target.getName() + "' (available: "
+ (target.getInputs() == null ? "none" : target.getInputs().stream()
.map(io -> io.getName()).toList()) + ")");
}
if (isModelInputConnection(sourceName, targetName)) {
throw new ResponseStatusException(HttpStatus.BAD_GATEWAY,
"Assistant returned a connection to technical input 'model' on block '" + target.getName()
+ "'");
}
return Connection.builder()
.sourceId(source.getId())
.sourceName(sourceName)
.targetId(target.getId())
.targetName(targetName)
.build();
}
private boolean isModelInputConnection(String sourceName, String targetName) {
return "model".equals(AssistantTextSupport.normalizeBlockReference(targetName))
&& !"model".equals(AssistantTextSupport.normalizeBlockReference(sourceName));
}
private List<AssistantConnectionDraft> inferSequentialConnections(List<Block<?>> blocks) {
if (blocks == null || blocks.size() < 2) {
return List.of();
}
List<AssistantConnectionDraft> connections = new ArrayList<>();
for (int i = 0; i < blocks.size() - 1; i++) {
Block<?> source = blocks.get(i);
Block<?> target = blocks.get(i + 1);
String sourceOutput = resolveConnectionOutputName(source, "response");
String targetInput = resolveSequentialInputName(target);
if (sourceOutput == null || targetInput == null) {
continue;
}
connections.add(new AssistantConnectionDraft(
source.getName(),
sourceOutput,
target.getName(),
targetInput));
}
return connections;
return ConnectionAssembler.preserveConnections(sourceConnections, oldNodeIdToAssembledNode, removedExistingNodeIds);
}
/**
@ -2133,155 +1865,6 @@ public class FlowAssistantService {
return List.copyOf(merged.values());
}
private String resolveSequentialInputName(Block<?> block) {
List<IODescriptor> inputs = block.getInputs();
if (inputs == null || inputs.isEmpty()) {
return null;
}
if (inputs.size() == 1) {
return inputs.getFirst().getName();
}
return inputs.stream()
.map(IODescriptor::getName)
.filter(name -> !isLikelyUserProvidedInput(name))
.findFirst()
.orElse(inputs.getFirst().getName());
}
private boolean isLikelyUserProvidedInput(String name) {
String normalized = AssistantTextSupport.normalizeBlockReference(name);
return normalized != null
&& Set.of("query", "question", "user_query", "userquery", "url", "file_url", "fileurl")
.contains(normalized);
}
private FlowNode resolveConnectionBlock(String rawReference, Map<String, FlowNode> nodesByPlanId,
Map<String, FlowNode> nodesByAlias) {
if (rawReference == null || rawReference.isBlank()) {
return null;
}
FlowNode direct = nodesByPlanId.get(rawReference);
if (direct != null) {
return direct;
}
FlowNode byAlias = nodesByAlias.get(AssistantTextSupport.normalizeBlockReference(rawReference));
if (byAlias != null) {
return byAlias;
}
// Same syntactic-noise fallback as findIoByName: a brace-mangled or ${{...}}-wrapped block
// reference should still resolve to its node instead of dropping the whole connection.
String sanitized = AssistantTextSupport.normalizeBlockReference(stripHandleNoise(rawReference));
return sanitized == null ? null : nodesByAlias.get(sanitized);
}
private FlowNode inferBlockByIo(String ioName, Collection<FlowNode> nodes, boolean output) {
String normalizedIo = AssistantTextSupport.normalizeBlockReference(ioName);
if (normalizedIo == null) {
return null;
}
FlowNode match = null;
for (FlowNode node : nodes) {
List<IODescriptor> ioDescriptors = output ? node.getOutputs() : node.getInputs();
if (ioDescriptors == null) {
continue;
}
boolean hasMatch = ioDescriptors.stream()
.anyMatch(io -> normalizedIo.equals(AssistantTextSupport.normalizeBlockReference(io.getName())));
if (!hasMatch) {
continue;
}
if (match != null) {
return null;
}
match = node;
}
return match;
}
private String resolveConnectionOutputName(FlowNode node, String requestedOutput) {
return resolveIoName(node.getOutputs(), requestedOutput, List.of("response", "output", "true", "false"));
}
private String resolveConnectionInputName(FlowNode node, String requestedInput) {
return resolveIoName(node.getInputs(), requestedInput, List.of("input", "prompt"));
}
private String resolveIoName(List<IODescriptor> descriptors, String requestedName, List<String> preferredNames) {
if (descriptors == null || descriptors.isEmpty()) {
return null;
}
if (requestedName != null && !requestedName.isBlank()) {
IODescriptor exact = findIoByName(descriptors, requestedName);
if (exact != null) {
return exact.getName();
}
}
if (descriptors.size() == 1) {
return descriptors.getFirst().getName();
}
for (String preferredName : preferredNames) {
IODescriptor preferred = findIoByName(descriptors, preferredName);
if (preferred != null) {
return preferred.getName();
}
}
return null;
}
private IODescriptor findIoByName(List<IODescriptor> descriptors, String requestedName) {
String normalized = AssistantTextSupport.normalizeBlockReference(requestedName);
if (normalized == null) {
return null;
}
IODescriptor exact = descriptors.stream()
.filter(io -> normalized.equals(AssistantTextSupport.normalizeBlockReference(io.getName())))
.findFirst()
.orElse(null);
if (exact != null) {
return exact;
}
// Fallback for syntactic noise the model sometimes emits in a handle name - a stray or
// unbalanced brace ("{userReview"), an accidental ${{...}} wrapper ("${{userReview}}"), or a
// block-qualified reference ("SomeBlock.response"). Strip that noise and retry; this only
// runs after the exact match already failed, so it can never change a currently-resolving
// connection - it only rescues one that would otherwise be silently dropped.
String sanitized = AssistantTextSupport.normalizeBlockReference(stripHandleNoise(requestedName));
if (sanitized == null || sanitized.equals(normalized)) {
return null;
}
return descriptors.stream()
.filter(io -> sanitized.equals(AssistantTextSupport.normalizeBlockReference(io.getName())))
.findFirst()
.orElse(null);
}
/**
* Removes purely-syntactic noise from a handle/block reference the model produced: placeholder
* wrappers ({@code ${{...}}} / {@code {{...}}}), stray braces, dollar signs and quotes, and a
* leading block-name qualifier (keeps the last dotted segment, so {@code "Review.response"}
* becomes {@code "response"}). Deterministic clean-up only - it never invents a name, so a
* reference that was genuinely wrong (not just mangled) still fails to resolve and is dropped.
*/
private String stripHandleNoise(String raw) {
if (raw == null) {
return null;
}
String cleaned = raw.replace("${{", "").replace("{{", "").replace("}}", "");
cleaned = cleaned.replaceAll("[{}$\"']", "").trim();
int lastDot = cleaned.lastIndexOf('.');
if (lastDot >= 0 && lastDot < cleaned.length() - 1) {
cleaned = cleaned.substring(lastDot + 1).trim();
}
return cleaned.isBlank() ? null : cleaned;
}
private void registerNodeAlias(Map<String, FlowNode> nodesByAlias, String reference, FlowNode node) {
String normalized = AssistantTextSupport.normalizeBlockReference(reference);
if (normalized != null) {
nodesByAlias.putIfAbsent(normalized, node);
}
}
private AssistantFlowPlan validateAndNormalizePlan(AssistantFlowPlan plan, OperationMode mode, String userPrompt,
FlowCreateRequest currentFlow, List<ValidationError> errors,