refactor: order MCP shared sessions with a dependency, drop the state_ready connection hack
MCP shared-session ordering was enforced with a fake data connection: the
consumer's prompt carried a magic ${{state_ready}} placeholder to synthesize
an input, and the producer's output was wired into it purely to force
execution order - the consumer never actually used that data (the real
shared state is the MCP session, accessed by name). This replaces that hack
with a Dependency, the mechanism meant exactly for ordering-without-data.
No engine change: the executor already gates readiness on dependencies and
the validator already counts them for reachability; MCPAgent uses the
default activity() capabilities so it can be a dependency source/target.
- FlowAssistantService: after assembly, generate ordering dependencies
deterministically for all three cases:
- top-level chain: Dependency(producer, consumer)
- within-container chain: Dependency(producer, consumer) in the subflow
- cross-boundary (top-level producer -> consumer inside a container):
Dependency(producer, container) - the container runs after the producer
and inherits its session.
Existing dependencies are preserved on REFINE (preserveCurrentDependencies).
Removed completeRequiredSequentialConnections, validateSharedMemorySemantics
and their now-unused helpers (isReachable, summarizeMcpSessionBlocks); the
ordering no longer flows through connections.
- Prompt: removed the ${{state_ready}} instructions; the model is told the
backend orders the consumer after the producer via a dependency and must
not add a placeholder or connection for it.
- Tests: updated all MCP mocks to stop emitting state_ready (prompt and
connections) and assert the dependency instead (top-level, within-container,
cross-boundary). Zero state_ready references remain.
Full suite green (431).
This commit is contained in:
parent
68a86564fa
commit
9cc110464d
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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<Connection> 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<Connection> 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<Block<?>> assembledBlocks, Map<String, FlowNode> nodesByPlanId, Map<String, FlowNode> nodesByAlias) {
|
||||
if (!required) {
|
||||
return parsed;
|
||||
}
|
||||
|
||||
List<AssistantConnectionDraft> 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<String, Block<?>> mcpProducersByName(Collection<Block<?>> blocks) {
|
||||
Map<String, Block<?>> 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<Dependency> buildInnerSharedMemoryDependencies(boolean required, List<Block<?>> blocks) {
|
||||
if (!required) {
|
||||
return List.of();
|
||||
}
|
||||
Map<String, Block<?>> producers = mcpProducersByName(blocks);
|
||||
if (producers.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<Dependency> 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<Dependency> buildTopLevelSharedMemoryDependencies(boolean required, List<Block<?>> topBlocks,
|
||||
List<Container<?>> containers) {
|
||||
if (!required) {
|
||||
return List.of();
|
||||
}
|
||||
Map<String, Block<?>> producers = mcpProducersByName(topBlocks);
|
||||
if (producers.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<Dependency> 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<Dependency> preserveCurrentDependencies(FlowCreateRequest currentFlow,
|
||||
Map<String, FlowNode> oldNodeIdToAssembledNode, Set<String> removedExistingNodeIds) {
|
||||
if (currentFlow == null || currentFlow.flow() == null || currentFlow.flow().getDependencies() == null
|
||||
|| currentFlow.flow().getDependencies().isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<Dependency> 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<Dependency> mergeDependencies(List<Dependency> preserved, List<Dependency> generated) {
|
||||
Map<String, Dependency> merged = new LinkedHashMap<>();
|
||||
for (Dependency dependency : preserved == null ? List.<Dependency>of() : preserved) {
|
||||
merged.put(dependency.getSourceId() + "|" + dependency.getTargetId(), dependency);
|
||||
}
|
||||
for (Dependency dependency : generated == null ? List.<Dependency>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<Block<?>> blocks,
|
||||
List<Connection> 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<Connection> connections, Set<String> 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<Block<?>> 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) {
|
||||
|
|
|
|||
|
|
@ -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<Dependency> 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
|
||||
|
|
|
|||
Loading…
Reference in New Issue