diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java index 1ceea41..686cc16 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantPromptService.java @@ -210,12 +210,12 @@ public class FlowAssistantPromptService { - The config object must contain only task-specific fields. - For enum fields, use only the allowedValues listed in the block descriptor. - For HTTPServerCall without authentication, set requiresAuthorization to false and omit authorization details. - - MCPAgent shared session rule: when multiple blocks need to work with the SAME MCP tools AND share the same execution context or state, model them as a producer-consumer chain. The FIRST block that opens the tool session is the producer: set shareSession to true and sharedSessionName to "sharedMemorySession". Subsequent blocks that must reuse the same tools and the same live context (same session, same accumulated state) are consumers: set useSharedSession to true and sharedSessionRef to "sharedMemorySession", and include an upstream ordering placeholder such as ${{state_ready}} in their prompt. A later block that does NOT need the same tools or the same context should NOT be made a consumer — use a separate independent MCPAgent or a different block type instead. + - MCPAgent shared session rule: when multiple blocks need to work with the SAME MCP tools AND share the same execution context or state, model them as a producer-consumer chain. The FIRST block that opens the tool session is the producer: set shareSession to true and sharedSessionName to "sharedMemorySession". Subsequent blocks that must reuse the same tools and the same live context (same session, same accumulated state) are consumers: set useSharedSession to true and sharedSessionRef to "sharedMemorySession". The backend orders the consumer after the producer automatically (via a dependency) - do NOT add any ordering placeholder or connection between them. A later block that does NOT need the same tools or the same context should NOT be made a consumer — use a separate independent MCPAgent or a different block type instead. - The producer is always the first block in the chain that initialises the shared tool session. Never mark it as a consumer. - A shared-session MCP producer/consumer chain may live all at the top level, all inside one container, or span the boundary from a TOP-LEVEL producer to a consumer inside a container - but in the last case the producer must be wired (with a connection) to run before that container, so its session is ready when the container runs. Do NOT put the producer inside a container and the consumer at the top level, and do NOT split a chain across two different containers. - Do not include system-managed fields like provider, model, llmDescriptor, ids, skills. The exception is BranchRejoinBlock's "inputs" field, which is a structural config array you must set (see below), not a derived IO list. - Use placeholders like ${{variable}} when needed. - - If this block depends on data, context, or completion from an earlier workflow step, include a placeholder such as ${{input}}, ${{previous_response}}, or ${{state_ready}} in its prompt so the backend can create a real input for the connection. + - If this block consumes DATA produced by an earlier workflow step, include a placeholder such as ${{input}} or ${{previous_response}} in its prompt so the backend can create a real input for the connection. (Pure execution ordering that carries no data - e.g. an MCP shared-session consumer waiting for its producer - is handled by the backend as a dependency and needs no placeholder.) - Do not use placeholders for literal instructions that do not depend on another block. - Placeholder names must match ^[A-Za-z][A-Za-z0-9_.-]*$ (start with a letter; letters/digits/underscore/dot/hyphen only). - If a placeholder should collect several values as one array-valued input instead of a single value, suffix its name with []: ${{name[]}}. Keep the same name and the same [] usage everywhere it appears in this block's text. @@ -319,7 +319,7 @@ public class FlowAssistantPromptService { - Do not connect workflow outputs to technical configuration inputs such as model. - Do not leave processing blocks disconnected when one step depends on another step's output or execution order. - Leaving a block's OUTPUT unconnected is intentional and correct for the result-producing block(s): an open output is what surfaces as the flow's readable result. Do not connect the final block's output into an EndBlock (or anything else) just to "terminate" it - only connect an output when a downstream block actually consumes it. - - For MCPAgent blocks that reuse a shared session, still connect the producer response to the consumer's upstream ordering placeholder, for example state_ready, so execution order is explicit. + - Do NOT add a connection between an MCPAgent shared-session producer and its consumer just to order them: the backend enforces that order with a dependency. Only connect them if the consumer genuinely consumes the producer's output as data. - In REFINE/FIX, existing connections between KEEP/UPDATE blocks are preserved by the backend; return only connections that are new or intentionally changed. - Allowed block ids are listed in the flow plan and configured blocks. Use only those exact values. - Return the minimal set of connections required by the user request. diff --git a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java index 0b07c61..18b7b8c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/assistant/FlowAssistantService.java @@ -63,6 +63,7 @@ import it.cnr.isti.workflow.manager.containers.types.LoopContainerType; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import it.cnr.isti.workflow.manager.flows.model.FlowNode; import it.cnr.isti.workflow.manager.flows.model.Connection; +import it.cnr.isti.workflow.manager.flows.model.Dependency; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; @@ -569,12 +570,10 @@ public class FlowAssistantService { parsedConnections = invokeStructuredAndValidate(provider, jsonModelFor(mode, phaseModels), phaseModels.repairModel(), connectionsPrompt, "connections", rawResponse -> { ParsedConnections parsed = parseConnectionsOrInferSequential(rawResponse, assembledBlocks); - parsed = completeRequiredSequentialConnections(requireSharedMemorySemantics, parsed, - assembledBlocks, nodesByPlanId, nodesByAlias); - List candidateConnections = toValidConnections(parsed.connections(), nodesByPlanId, - nodesByAlias); - validateSharedMemorySemantics(requireSharedMemorySemantics, assembledBlocks, - candidateConnections); + // 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); return parsed; }); } @@ -593,6 +592,10 @@ public class FlowAssistantService { .blocks(assembledBlocks) .containers(assembledContainers) .connections(connections) + .dependencies(mergeDependencies( + preserveCurrentDependencies(currentFlow, oldNodeIdToAssembledNode, removedExistingNodeIds), + buildTopLevelSharedMemoryDependencies(requireSharedMemorySemantics, assembledBlocks, + assembledContainers))) .globalInputs(collectGlobalInputs(assembledBlocks, assembledContainers, currentFlow)) .build()); @@ -708,11 +711,9 @@ public class FlowAssistantService { phaseModels.repairModel(), connectionsPrompt, "connections for container " + containerPlan.containerId(), rawResponse -> { ParsedConnections parsed = parseConnectionsOrInferSequential(rawResponse, innerBlocks); - parsed = completeRequiredSequentialConnections(innerRequiresSharedMemory, parsed, innerBlocks, - innerNodesByPlanId, innerNodesByAlias); - List candidateConnections = toValidConnections(parsed.connections(), - innerNodesByPlanId, innerNodesByAlias); - validateSharedMemorySemantics(innerRequiresSharedMemory, innerBlocks, candidateConnections); + // Within-container MCP ordering is a Dependency in the subflow (see below), + // not a data connection. + toValidConnections(parsed.connections(), innerNodesByPlanId, innerNodesByAlias); return parsed; }); } @@ -726,6 +727,8 @@ public class FlowAssistantService { FlowData subFlow = FlowData.builder() .blocks(innerBlocks) .connections(connections) + // Within-container MCP producer/consumer ordering is a Dependency inside the subflow. + .dependencies(buildInnerSharedMemoryDependencies(innerRequiresSharedMemory, innerBlocks)) // A ${{global.x}} referenced inside the subflow must be declared in the subflow's // own globalInputs: the container subflow is validated (and, at runtime, fed its // required globals from the parent) as its own scope. The same global is also @@ -1646,39 +1649,125 @@ public class FlowAssistantService { return connections; } - private ParsedConnections completeRequiredSequentialConnections(boolean required, ParsedConnections parsed, - List> assembledBlocks, Map nodesByPlanId, Map nodesByAlias) { - if (!required) { - return parsed; - } - - List merged = new ArrayList<>( - parsed == null || parsed.connections() == null ? List.of() : parsed.connections()); - for (AssistantConnectionDraft inferred : inferSequentialConnections(assembledBlocks)) { - FlowNode inferredSource = resolveConnectionBlock(inferred.fromBlockId(), nodesByPlanId, nodesByAlias); - FlowNode inferredTarget = resolveConnectionBlock(inferred.toBlockId(), nodesByPlanId, nodesByAlias); - if (inferredSource == null || inferredTarget == null) { - continue; + /** + * MCP shared-session producer blocks by their session name (shareSession=true with a name). + */ + private Map> mcpProducersByName(Collection> blocks) { + Map> producers = new LinkedHashMap<>(); + for (Block block : blocks) { + if (block.getSpecificConfiguration() instanceof MCPAgentBlockConfiguration configuration + && Boolean.TRUE.equals(configuration.getShareSession()) + && configuration.getSharedSessionName() != null + && !configuration.getSharedSessionName().isBlank()) { + producers.putIfAbsent(configuration.getSharedSessionName().trim(), block); } - boolean alreadyConnected = false; - for (AssistantConnectionDraft existing : merged) { - FlowNode existingSource = resolveConnectionBlock(existing.fromBlockId(), nodesByPlanId, nodesByAlias); - FlowNode existingTarget = resolveConnectionBlock(existing.toBlockId(), nodesByPlanId, nodesByAlias); - if (existingSource != null && existingTarget != null - && Objects.equals(existingSource.getId(), inferredSource.getId()) - && Objects.equals(existingTarget.getId(), inferredTarget.getId())) { - alreadyConnected = true; + } + return producers; + } + + /** + * The shared-session name a block consumes (useSharedSession=true with a ref), or null. + */ + private String mcpConsumerSessionRef(Block block) { + if (block.getSpecificConfiguration() instanceof MCPAgentBlockConfiguration configuration + && Boolean.TRUE.equals(configuration.getUseSharedSession()) + && configuration.getSharedSessionRef() != null + && !configuration.getSharedSessionRef().isBlank()) { + return configuration.getSharedSessionRef().trim(); + } + return null; + } + + /** + * Dependencies enforcing MCP shared-session order inside a single subflow: each consumer block + * depends on the producer of the session it references, so the producer opens the session + * before the consumer runs - a pure ordering constraint, no data connection. + */ + private List buildInnerSharedMemoryDependencies(boolean required, List> blocks) { + if (!required) { + return List.of(); + } + Map> producers = mcpProducersByName(blocks); + if (producers.isEmpty()) { + return List.of(); + } + List dependencies = new ArrayList<>(); + for (Block block : blocks) { + String ref = mcpConsumerSessionRef(block); + Block producer = ref == null ? null : producers.get(ref); + if (producer != null && !producer.getId().equals(block.getId())) { + dependencies.add(Dependency.builder().sourceId(producer.getId()).targetId(block.getId()).build()); + } + } + return dependencies; + } + + /** + * Top-level MCP shared-session ordering dependencies: a top-level consumer depends on the + * top-level producer of its session; a consumer inside a container makes the whole container + * depend on the producer (the container inherits the producer's session and must run after it). + */ + private List buildTopLevelSharedMemoryDependencies(boolean required, List> topBlocks, + List> containers) { + if (!required) { + return List.of(); + } + Map> producers = mcpProducersByName(topBlocks); + if (producers.isEmpty()) { + return List.of(); + } + List dependencies = new ArrayList<>(); + for (Block block : topBlocks) { + String ref = mcpConsumerSessionRef(block); + Block producer = ref == null ? null : producers.get(ref); + if (producer != null && !producer.getId().equals(block.getId())) { + dependencies.add(Dependency.builder().sourceId(producer.getId()).targetId(block.getId()).build()); + } + } + for (Container container : containers) { + for (Block inner : subFlowBlocks(container)) { + String ref = mcpConsumerSessionRef(inner); + Block producer = ref == null ? null : producers.get(ref); + if (producer != null) { + dependencies.add(Dependency.builder().sourceId(producer.getId()).targetId(container.getId()).build()); break; } } - if (!alreadyConnected) { - merged.add(inferred); - } } + return dependencies; + } - String rationale = parsed == null ? "" : parsed.rationale(); - return new ParsedConnections(merged, - defaultIfBlank(rationale, "Completed required sequential shared-memory connections.")); + private List preserveCurrentDependencies(FlowCreateRequest currentFlow, + Map oldNodeIdToAssembledNode, Set removedExistingNodeIds) { + if (currentFlow == null || currentFlow.flow() == null || currentFlow.flow().getDependencies() == null + || currentFlow.flow().getDependencies().isEmpty()) { + return List.of(); + } + List preserved = new ArrayList<>(); + for (Dependency dependency : currentFlow.flow().getDependencies()) { + if (dependency == null || removedExistingNodeIds.contains(dependency.getSourceId()) + || removedExistingNodeIds.contains(dependency.getTargetId())) { + continue; + } + FlowNode source = oldNodeIdToAssembledNode.get(dependency.getSourceId()); + FlowNode target = oldNodeIdToAssembledNode.get(dependency.getTargetId()); + if (source == null || target == null) { + continue; + } + preserved.add(Dependency.builder().sourceId(source.getId()).targetId(target.getId()).build()); + } + return preserved; + } + + private List mergeDependencies(List preserved, List generated) { + Map merged = new LinkedHashMap<>(); + for (Dependency dependency : preserved == null ? List.of() : preserved) { + merged.put(dependency.getSourceId() + "|" + dependency.getTargetId(), dependency); + } + for (Dependency dependency : generated == null ? List.of() : generated) { + merged.putIfAbsent(dependency.getSourceId() + "|" + dependency.getTargetId(), dependency); + } + return List.copyOf(merged.values()); } private String resolveSequentialInputName(Block block) { @@ -2298,78 +2387,6 @@ public class FlowAssistantService { && !currentFlow.flow().getBlocks().isEmpty(); } - private void validateSharedMemorySemantics(boolean required, Collection> blocks, - List connections) { - if (!required) { - return; - } - - Block sessionProducer = null; - Block sessionConsumer = null; - for (Block block : blocks) { - if (block.getSpecificConfiguration() instanceof MCPAgentBlockConfiguration configuration) { - if (Boolean.TRUE.equals(configuration.getShareSession()) - && SHARED_MEMORY_SESSION_NAME.equals(configuration.getSharedSessionName())) { - sessionProducer = block; - } - if (Boolean.TRUE.equals(configuration.getUseSharedSession()) - && SHARED_MEMORY_SESSION_NAME.equals(configuration.getSharedSessionRef())) { - sessionConsumer = block; - } - } - } - - if (sessionProducer == null || sessionConsumer == null) { - log.warn( - "Assistant returned shared-memory MCP flags without a complete producer/consumer pair for {}. Continuing because shared-session flags are advisory. Blocks: {}", - SHARED_MEMORY_SESSION_NAME, - summarizeMcpSessionBlocks(blocks)); - return; - } - if (!isReachable(sessionProducer.getId(), sessionConsumer.getId(), connections, new LinkedHashSet<>())) { - log.warn( - "Assistant returned shared-memory MCP flags where consumer is not reachable from producer. Continuing because shared-session flags are advisory. Producer={}, consumer={}, connections={}", - sessionProducer.getName(), - sessionConsumer.getName(), - connections); - } - } - - private boolean isReachable(String sourceId, String targetId, List connections, Set visited) { - if (!visited.add(sourceId)) { - return false; - } - for (Connection connection : connections) { - if (!Objects.equals(sourceId, connection.getSourceId())) { - continue; - } - if (Objects.equals(targetId, connection.getTargetId()) - || isReachable(connection.getTargetId(), targetId, connections, visited)) { - return true; - } - } - return false; - } - - private String summarizeMcpSessionBlocks(Collection> blocks) { - if (blocks == null || blocks.isEmpty()) { - return "(none)"; - } - return blocks.stream() - .map(block -> { - if (block.getSpecificConfiguration() instanceof MCPAgentBlockConfiguration configuration) { - return block.getName() - + "[MCPAgent shareSession=" + configuration.getShareSession() - + ", sharedSessionName=" + configuration.getSharedSessionName() - + ", useSharedSession=" + configuration.getUseSharedSession() - + ", sharedSessionRef=" + configuration.getSharedSessionRef() + "]"; - } - return block.getName() + "[" + (block.getType() == null ? "unknown" : block.getType().getName()) + "]"; - }) - .toList() - .toString(); - } - private boolean isSharedMemoryRequest(String text) { String normalized = normalizeBlockReference(text); if (normalized == null) { diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java index a81ce6a..ebf9c2b 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/AssistantControllerTest.java @@ -54,6 +54,7 @@ import it.cnr.isti.workflow.manager.containers.factories.GenericContainerFactory import it.cnr.isti.workflow.manager.containers.types.GenericContainerType; import it.cnr.isti.workflow.manager.blocks.types.MCPAgentBlockType; import it.cnr.isti.workflow.manager.flows.model.Connection; +import it.cnr.isti.workflow.manager.flows.model.Dependency; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; import it.cnr.isti.workflow.manager.flows.model.FlowData; import it.cnr.isti.workflow.manager.flows.validation.ValidationError; @@ -807,10 +808,10 @@ public class AssistantControllerTest { assertTrue(queryConfiguration.getUseSharedSession()); assertEquals("sharedMemorySession", queryConfiguration.getSharedSessionRef()); - assertTrue(flow.flow().getConnections().stream() - .anyMatch(connection -> indexBlock.getId().equals(connection.getSourceId()) - && queryBlock.getId().equals(connection.getTargetId()) - && "state_ready".equals(connection.getTargetName()))); + // Ordering is a dependency (producer -> consumer), not a data connection. + assertTrue(flow.flow().getDependencies().stream() + .anyMatch(dependency -> indexBlock.getId().equals(dependency.getSourceId()) + && queryBlock.getId().equals(dependency.getTargetId()))); } private void mockAssistantResponses() { @@ -1071,24 +1072,20 @@ public class AssistantControllerTest { "blockId", "b3", "name", "Query indexed data", "config", java.util.Map.of( - "prompt", "The shared state is ready: ${{state_ready}}. Query it: ${{query}}", + "prompt", "Query the indexed data: ${{query}}", "useSharedSession", true, "sharedSessionRef", "sharedMemorySession")))); } if (prompt.contains("TASK: CONNECTIONS")) { + // b2 -> b3 ordering is a dependency (added by the backend), not a data connection. return TestAssistantResponses.wrap(java.util.Map.of( - "rationale", "Connected download, indexing, and query.", + "rationale", "Connected download to indexing.", "connections", java.util.List.of( java.util.Map.of( "fromBlockId", "b1", "fromOutput", "response", "toBlockId", "b2", - "toInput", "file_content"), - java.util.Map.of( - "fromBlockId", "b2", - "fromOutput", "response", - "toBlockId", "b3", - "toInput", "state_ready")))); + "toInput", "file_content")))); } throw new IllegalStateException("Unexpected assistant prompt:\n" + prompt); }; @@ -1138,19 +1135,15 @@ public class AssistantControllerTest { "blockId", "b2", "name", "Use shared context", "config", java.util.Map.of( - "prompt", "Use prepared context: ${{state_ready}}", + "prompt", "Use the prepared shared context to answer.", "useSharedSession", true, "sharedSessionRef", "sharedMemorySession")))); } if (prompt.contains("TASK: CONNECTIONS")) { + // Both blocks are consumer-only (no producer), so there is nothing to order. return TestAssistantResponses.wrap(java.util.Map.of( - "rationale", "Connected the two advisory shared-memory steps.", - "connections", java.util.List.of( - java.util.Map.of( - "fromBlockId", "b1", - "fromOutput", "response", - "toBlockId", "b2", - "toInput", "state_ready")))); + "rationale", "No producer, so no ordering.", + "connections", java.util.List.of())); } throw new IllegalStateException("Unexpected assistant prompt:\n" + prompt); }; @@ -1635,7 +1628,7 @@ public class AssistantControllerTest { "blockId", "b3", "name", "Query index", "config", java.util.Map.of( - "prompt", "The shared state is ready: ${{state_ready}}. Answer the user query from the shared memory session: ${{query}}")))); + "prompt", "Answer the user query from the shared memory session: ${{query}}")))); } if (prompt.contains("TASK: BLOCK_CONFIG") && prompt.contains("Current block to configure:\n{\n \"blockId\" : \"b4\"")) { @@ -1654,19 +1647,15 @@ public class AssistantControllerTest { return "Connect the query output to the validator input."; } if (prompt.contains("TASK: CONNECTIONS")) { + // b2 -> b3 ordering is a dependency (added by the backend), not a data connection. return TestAssistantResponses.wrap(java.util.Map.of( - "rationale", "Connected dependent steps and shared-memory ordering.", + "rationale", "Connected download to indexing.", "connections", java.util.List.of( java.util.Map.of( "fromBlockId", "b1", "fromOutput", "response", "toBlockId", "b2", - "toInput", "file_content"), - java.util.Map.of( - "fromBlockId", "b2", - "fromOutput", "response", - "toBlockId", "b3", - "toInput", "state_ready")))); + "toInput", "file_content")))); } throw new IllegalStateException("Unexpected assistant prompt:\n" + prompt); }; @@ -2376,6 +2365,16 @@ public class AssistantControllerTest { assertEquals("sharedMemorySession", producer.getSharedSessionName()); assertTrue(consumer.getUseSharedSession()); assertEquals("sharedMemorySession", consumer.getSharedSessionRef()); + + // Ordering is a dependency inside the subflow (producer -> consumer), not a data connection. + Block producerBlock = inner.stream().filter(b -> "index-and-create-session".equals(b.getName())) + .findFirst().orElseThrow(); + Block consumerBlock = inner.stream().filter(b -> "query-indexed-data".equals(b.getName())) + .findFirst().orElseThrow(); + java.util.List subFlowDeps = container.getSpecificConfiguration().getSubFlow().getDependencies(); + assertTrue(subFlowDeps.stream().anyMatch(d -> producerBlock.getId().equals(d.getSourceId()) + && consumerBlock.getId().equals(d.getTargetId()))); + assertTrue(container.getSpecificConfiguration().getSubFlow().getConnections().isEmpty()); } private void mockMcpSharedSessionChainInsideContainer() { @@ -2411,17 +2410,16 @@ public class AssistantControllerTest { if (prompt.contains("TASK: BLOCK_CONFIG") && prompt.contains("Current block to configure:\n{\n \"blockId\" : \"c1-b2\"")) { return TestAssistantResponses.wrap(java.util.Map.of( - "rationale", "Query consumer.", + "rationale", "Query consumer; the backend orders it after the producer via a dependency.", "block", java.util.Map.of("blockId", "c1-b2", "name", "query-indexed-data", "config", java.util.Map.of( - "prompt", "The shared state is ready: ${{state_ready}}. Query the indexed data.")))); + "prompt", "Query the indexed data from the shared knowledge base.")))); } if (prompt.contains("TASK: CONNECTIONS")) { + // No data connection between producer and consumer - ordering is a dependency. return TestAssistantResponses.wrap(java.util.Map.of( - "rationale", "Order the producer before the consumer.", - "connections", java.util.List.of( - java.util.Map.of("fromBlockId", "c1-b1", "fromOutput", "response", - "toBlockId", "c1-b2", "toInput", "state_ready")))); + "rationale", "No data connections needed.", + "connections", java.util.List.of())); } throw new IllegalStateException("Unexpected assistant prompt:\n" + prompt); }; @@ -2592,15 +2590,13 @@ public class AssistantControllerTest { && prompt.contains("Current block to configure:\n{\n \"blockId\" : \"c1-b1\"")) { return TestAssistantResponses.wrap(java.util.Map.of("rationale", "consumer", "block", java.util.Map.of("blockId", "c1-b1", "name", "query", - "config", java.util.Map.of("prompt", - "The shared state is ready: ${{state_ready}}. Query the indexed data.")))); + "config", java.util.Map.of("prompt", "Query the indexed data.")))); } if (prompt.contains("TASK: CONNECTIONS")) { - // Wire the top-level producer to the container's exposed state_ready input, so the - // producer runs before the container and its session is available inside. - return TestAssistantResponses.wrap(java.util.Map.of("rationale", "order producer before container", - "connections", java.util.List.of(java.util.Map.of("fromBlockId", "b1", "fromOutput", "response", - "toBlockId", "c1", "toInput", "state_ready")))); + // No data connection - the backend orders the container after the producer via a + // dependency to the container. + return TestAssistantResponses.wrap(java.util.Map.of("rationale", "no data connection needed", + "connections", java.util.List.of())); } throw new IllegalStateException("Unexpected assistant prompt:\n" + prompt); }; @@ -2616,12 +2612,18 @@ public class AssistantControllerTest { // the container - the cross-boundary chain now validates instead of being rejected. assertTrue(response.valid(), () -> "Unexpected validation errors: " + response.validationErrors()); assertEquals(1, response.flow().flow().getBlocks().size()); - MCPAgentBlockConfiguration producer = (MCPAgentBlockConfiguration) response.flow().flow().getBlocks() - .getFirst().getSpecificConfiguration(); - assertTrue(producer.getShareSession()); - MCPAgentBlockConfiguration consumer = (MCPAgentBlockConfiguration) response.flow().flow().getContainers() - .getFirst().getSpecificConfiguration().getSubFlow().getBlocks().getFirst().getSpecificConfiguration(); - assertTrue(consumer.getUseSharedSession()); + Block producerBlock = response.flow().flow().getBlocks().getFirst(); + assertTrue(((MCPAgentBlockConfiguration) producerBlock.getSpecificConfiguration()).getShareSession()); + Container container = response.flow().flow().getContainers().getFirst(); + assertTrue(((MCPAgentBlockConfiguration) container.getSpecificConfiguration().getSubFlow().getBlocks() + .getFirst().getSpecificConfiguration()).getUseSharedSession()); + + // Cross-boundary ordering is a top-level dependency from the producer to the CONTAINER + // (the container runs after the producer and inherits its session), not a data connection. + assertTrue(response.flow().flow().getConnections().isEmpty()); + assertTrue(response.flow().flow().getDependencies().stream() + .anyMatch(d -> producerBlock.getId().equals(d.getSourceId()) + && container.getId().equals(d.getTargetId()))); } @Test