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 4169ed0..60254a6 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 @@ -232,8 +232,9 @@ public class FlowAssistantService { String planPrompt = promptService.buildPlanPrompt(mode, userPrompt, currentFlow, errors, catalog); ParsedPlan parsedPlan = invokeStructuredAndValidate(provider, model, planPrompt, "plan", rawResponse -> { ParsedPlan plan = parsePlan(rawResponse); - validatePlan(plan.plan(), userPrompt, currentFlow, catalogByType.keySet()); - return plan; + AssistantFlowPlan normalizedPlan = validateAndNormalizePlan(plan.plan(), userPrompt, currentFlow, + catalogByType.keySet()); + return new ParsedPlan(normalizedPlan, plan.rationale()); }); boolean requireSharedMemorySemantics = isSharedMemoryContext(userPrompt, currentFlow, parsedPlan.plan()); @@ -499,44 +500,41 @@ public class FlowAssistantService { return; } - // Producer check must run before consumer: terms like "index" appear in both - // producer and consumer keyword sets. A block that creates/indexes shared state - // must be wired as producer (shareSession=true) even when its purpose also - // contains retrieval language. - if (isSharedStateProducerPurpose(blockPlan.purpose()) && !isSharedStateConsumerPurpose(blockPlan.purpose())) { - config.put("shareSession", true); - config.put("useSharedSession", false); - if (!hasTextValue(config.get("sharedSessionName"))) { - config.put("sharedSessionName", SHARED_MEMORY_SESSION_NAME); - } - if (!hasTextValue(config.get("model"))) { - config.put("model", model); - } - ensureDefaultRagServer(config); + boolean producerPurpose = isSharedStateProducerPurpose(blockPlan.purpose()); + boolean consumerPurpose = isSharedStateConsumerPurpose(blockPlan.purpose()); + + if (producerPurpose && (!consumerPurpose || isSharedStateProducerDominantPurpose(blockPlan.purpose()))) { + configureMcpSharedMemoryProducer(config, model); return; } - if (isSharedStateConsumerPurpose(blockPlan.purpose())) { - config.put("shareSession", false); - config.put("useSharedSession", true); - if (!hasTextValue(config.get("sharedSessionRef"))) { - config.put("sharedSessionRef", SHARED_MEMORY_SESSION_NAME); - } + if (consumerPurpose) { + configureMcpSharedMemoryConsumer(config); return; } - // Ambiguous purpose: matches both producer and consumer signals (e.g. "index for - // later retrieval"). Treat as producer so the shared session is actually created. - if (isSharedStateProducerPurpose(blockPlan.purpose())) { - config.put("shareSession", true); - config.put("useSharedSession", false); - if (!hasTextValue(config.get("sharedSessionName"))) { - config.put("sharedSessionName", SHARED_MEMORY_SESSION_NAME); - } - if (!hasTextValue(config.get("model"))) { - config.put("model", model); - } - ensureDefaultRagServer(config); + if (producerPurpose) { + configureMcpSharedMemoryProducer(config, model); + } + } + + private void configureMcpSharedMemoryProducer(ObjectNode config, String model) { + config.put("shareSession", true); + config.put("useSharedSession", false); + if (!hasTextValue(config.get("sharedSessionName"))) { + config.put("sharedSessionName", SHARED_MEMORY_SESSION_NAME); + } + if (!hasTextValue(config.get("model"))) { + config.put("model", model); + } + ensureDefaultRagServer(config); + } + + private void configureMcpSharedMemoryConsumer(ObjectNode config) { + config.put("shareSession", false); + config.put("useSharedSession", true); + if (!hasTextValue(config.get("sharedSessionRef"))) { + config.put("sharedSessionRef", SHARED_MEMORY_SESSION_NAME); } } @@ -927,12 +925,19 @@ public class FlowAssistantService { } } - private void validatePlan(AssistantFlowPlan plan, String userPrompt, FlowCreateRequest currentFlow, + private AssistantFlowPlan validateAndNormalizePlan(AssistantFlowPlan plan, String userPrompt, + FlowCreateRequest currentFlow, Set availableBlockTypes) { if (plan == null) { throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant returned an empty plan"); } if (plan.blocks() == null || plan.blocks().isEmpty()) { + if (canUseMinimalDraftFallback(userPrompt, currentFlow, availableBlockTypes)) { + return new AssistantFlowPlan( + defaultIfBlank(plan.name(), "Assistant flow"), + defaultIfBlank(plan.description(), userPrompt.trim()), + List.of(new AssistantBlockPlan("b1", "LLMBlock", userPrompt.trim()))); + } throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant returned a plan with no blocks"); } Set ids = new LinkedHashSet<>(); @@ -959,6 +964,23 @@ public class FlowAssistantService { "Assistant returned an inconsistent plan: workflows that require shared memory or reusable state between steps must use MCPAgent blocks with a shared MCP session instead of standalone LLMBlock nodes"); } } + return plan; + } + + private boolean canUseMinimalDraftFallback(String userPrompt, FlowCreateRequest currentFlow, + Set availableBlockTypes) { + return userPrompt != null + && !userPrompt.isBlank() + && !hasCurrentFlowBlocks(currentFlow) + && availableBlockTypes != null + && availableBlockTypes.contains("LLMBlock"); + } + + private boolean hasCurrentFlowBlocks(FlowCreateRequest currentFlow) { + return currentFlow != null + && currentFlow.flow() != null + && currentFlow.flow().getBlocks() != null + && !currentFlow.flow().getBlocks().isEmpty(); } private void validateSharedMemorySemantics(boolean required, Collection> blocks, @@ -983,13 +1005,18 @@ public class FlowAssistantService { } if (sessionProducer == null || sessionConsumer == null) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an inconsistent flow: shared-memory workflows require an MCPAgent producer and consumer sharing " - + SHARED_MEMORY_SESSION_NAME + ". Blocks: " + summarizeMcpSessionBlocks(blocks)); + 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<>())) { - throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, - "Assistant returned an inconsistent connections payload: the MCPAgent that consumes shared memory must execute after the producer MCPAgent"); + 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); } } @@ -1115,6 +1142,32 @@ public class FlowAssistantService { return normalized != null && containsSharedStateConsumerTerm(normalized); } + private boolean isSharedStateProducerDominantPurpose(String purpose) { + String normalized = normalizeBlockReference(purpose); + if (normalized == null) { + return false; + } + return normalized.contains("create session") + || normalized.contains("createsession") + || normalized.contains("open session") + || normalized.contains("opensession") + || normalized.contains("create shared") + || normalized.contains("createshared") + || normalized.contains("create index") + || normalized.contains("createindex") + || normalized.contains("build index") + || normalized.contains("buildindex") + || normalized.contains("index file") + || normalized.contains("indexfile") + || normalized.contains("index downloaded") + || normalized.contains("indexdownloaded") + || normalized.contains("ingest") + || containsWord(normalized, "store") + || normalized.contains("save") + || normalized.contains("embedding") + || normalized.contains("vector"); + } + private boolean containsSharedStateProducerTerm(String normalized) { return normalized.contains("index") || normalized.contains("indicizz") @@ -1348,6 +1401,7 @@ public class FlowAssistantService { } catch (ResponseStatusException e) { throw e; } catch (RuntimeException e) { + logAssistantProviderFailure(provider, model, "text", prompt, e); throw new ResponseStatusException(HttpStatus.BAD_GATEWAY, "Assistant provider call failed: " + e.getMessage(), e); } @@ -1365,6 +1419,7 @@ public class FlowAssistantService { } catch (ResponseStatusException e) { throw e; } catch (RuntimeException e) { + logAssistantProviderRetry(provider, model, "json", prompt, e); structuredFailure = e; } @@ -1390,6 +1445,7 @@ public class FlowAssistantService { } catch (ResponseStatusException e) { throw e; } catch (RuntimeException e) { + logAssistantProviderFailure(provider, model, "text-fallback", prompt, e); if (structuredFailure != null) { e.addSuppressed(structuredFailure); } @@ -1405,6 +1461,32 @@ public class FlowAssistantService { "Assistant provider call failed: Empty assistant response"); } + private void logAssistantProviderFailure(LLMProvider provider, String model, String mode, String prompt, + RuntimeException error) { + log.error( + "Assistant LLM provider call failed. provider={}, model={}, mode={}, promptLength={}, errorType={}, message={}", + provider == null ? "(none)" : provider.getName(), + model, + mode, + prompt == null ? 0 : prompt.length(), + error.getClass().getSimpleName(), + error.getMessage(), + error); + } + + private void logAssistantProviderRetry(LLMProvider provider, String model, String mode, String prompt, + RuntimeException error) { + log.warn( + "Assistant LLM provider call failed, trying fallback. provider={}, model={}, mode={}, promptLength={}, errorType={}, message={}", + provider == null ? "(none)" : provider.getName(), + model, + mode, + prompt == null ? 0 : prompt.length(), + error.getClass().getSimpleName(), + error.getMessage(), + error); + } + private boolean looksLikeStructuredJson(String response) { if (response == null) { return false;