This is an automated email from the ASF dual-hosted git repository. luigidemasi pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/camel.git
commit 1cad5ffff9d5c3f6ab03ecfbbf84d2060a47af77 Author: Luigi De Masi <[email protected]> AuthorDate: Mon Sep 28 11:26:23 2026 +0200 CAMEL-25049: Support batching named semantic questions Evaluate explicit question batches against one state and definition snapshot, returning named decisions and validated per-question diagnostics. Preserve single-question adapters through a sequential default and use one HTTP request for native TypeSafe AI batches. Cover batch parsing, policies, reload, failures and Java/YAML result reuse, and document provider fallback and shared-state restrictions. Co-authored-by: Codex <[email protected]> Signed-off-by: Luigi De Masi <[email protected]> --- .../camel/catalog/docs/semantic-language.adoc | 84 ++++++ .../src/main/docs/semantic-language.adoc | 84 ++++++ .../camel/language/semantic/SemanticLanguage.java | 145 +++++++--- .../org/apache/camel/semantic/SemanticAdapter.java | 19 ++ .../apache/camel/semantic/SemanticQuestions.java | 17 ++ .../apache/camel/semantic/SemanticBatchTest.java | 322 +++++++++++++++++++++ .../org/apache/camel/semantic/SemanticEipTest.java | 18 ++ .../typesafeai/TypeSafeAiSemanticAdapter.java | 39 ++- .../typesafeai/TypeSafeAiSemanticAdapterTest.java | 88 +++++- .../camel/dsl/yaml/SemanticQuestionTest.java | 34 +++ 10 files changed, 800 insertions(+), 50 deletions(-) diff --git a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/semantic-language.adoc b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/semantic-language.adoc index 0e516b1015e8..a077f0b645b4 100644 --- a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/semantic-language.adoc +++ b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/semantic-language.adoc @@ -130,6 +130,90 @@ do not imply equivalent quality or calibration when switching providers. Timeout answers and unsupported capabilities are errors, distinct from valid negative decisions or unmatched Choice results. +== Batching named questions + +Use `refs:name1,name2` to evaluate independent questions against the same selected state. +The expression returns a map of normalized decisions keyed by question name. For example, +using the declarations above: + +[source,java] +---- +from("direct:tickets") + .setProperty("decision").language("semantic", "refs:actionable,department,urgency") + .choice() + .when(simple("${exchangeProperty.decision[actionable]} == true && ${exchangeProperty.decision[department]} == 'billing'")) + .to("direct:billing") + .otherwise() + .to("direct:review"); +---- + +The same expression works through the generic language integration in YAML and XML. +For example, replace the route above with: + +[source,yaml] +---- +- route: + from: + uri: direct:tickets + steps: + - setProperty: + name: decision + expression: + language: + language: semantic + expression: refs:actionable,department,urgency + - setHeader: + name: urgency + expression: + simple: + expression: "${exchangeProperty.decision[urgency]}" + - choice: + when: + - expression: + simple: + expression: "${exchangeProperty.decision[department]} == 'billing'" + steps: + - to: direct:billing + otherwise: + steps: + - to: direct:review +---- + +A result might contain `actionable -> true`, `department -> "billing"` and `urgency -> 1.2`. +Reading the stored map does not invoke the provider again. Each question retains its own +threshold and uncertainty policy. There is no implicit AND/OR: batch expressions, including +a batch containing one boolean question, cannot be used as predicates. + +`CamelSemanticResults` holds an immutable map of the detailed `SemanticResult` objects, +keyed by the same names. These retain available probabilities, confidence and provider +metadata. `CamelSemanticResult` remains the single-question diagnostic property for +`ref:name`. Both diagnostic properties are cleared before every semantic evaluation; +a batch publishes its decisions and diagnostics only after every answer passes validation. +Missing, extra or invalid answers, operational failures and an uncertainty policy of `fail` +fail the entire evaluation through Camel error handling. A valid negative decision or +uncertainty handled by `non-match` remains a normal `false` result. A route property assigned +by Set Property is owned by the route; clear it explicitly before reevaluation if an error +handler could otherwise reuse a decision from an earlier attempt. + +Batch references are comma-separated names with surrounding whitespace removed. Empty, +duplicate and unknown references are rejected. There is no quoting or escaping in this +syntax: names containing commas or leading/trailing whitespace require a single `ref:name` +evaluation. Existing single references are unchanged. + +All selected questions must use the same effective Simple state selector, after applying +`camel.language.semantic.default-state` and resolving property placeholders. Selectors +must match exactly; different expressions that happen to return equal values are not +interchangeable. The selector is evaluated once per batch. Missing or unsupported state +fails before inference. One snapshot of the question definitions is used throughout each +batch, including validation and decision policies; the next evaluation sees replacements +from resource reload. Questions that need a previous answer require separate evaluations. + +`SemanticAdapter.evaluateBatch(Map<String, SemanticQuestion>, Object)` has a default +implementation that invokes the existing single-question method sequentially, preserving +compatibility with existing adapters. TypeSafe AI overrides it to send all questions in +one HTTP request using the component's existing timeout and concurrency limits. No results +are cached or combined across exchanges. + == EIP integration Use `language("semantic", "ref:name")` wherever an expression or boolean predicate is accepted. diff --git a/components/camel-ai/camel-semantic/src/main/docs/semantic-language.adoc b/components/camel-ai/camel-semantic/src/main/docs/semantic-language.adoc index 0e516b1015e8..a077f0b645b4 100644 --- a/components/camel-ai/camel-semantic/src/main/docs/semantic-language.adoc +++ b/components/camel-ai/camel-semantic/src/main/docs/semantic-language.adoc @@ -130,6 +130,90 @@ do not imply equivalent quality or calibration when switching providers. Timeout answers and unsupported capabilities are errors, distinct from valid negative decisions or unmatched Choice results. +== Batching named questions + +Use `refs:name1,name2` to evaluate independent questions against the same selected state. +The expression returns a map of normalized decisions keyed by question name. For example, +using the declarations above: + +[source,java] +---- +from("direct:tickets") + .setProperty("decision").language("semantic", "refs:actionable,department,urgency") + .choice() + .when(simple("${exchangeProperty.decision[actionable]} == true && ${exchangeProperty.decision[department]} == 'billing'")) + .to("direct:billing") + .otherwise() + .to("direct:review"); +---- + +The same expression works through the generic language integration in YAML and XML. +For example, replace the route above with: + +[source,yaml] +---- +- route: + from: + uri: direct:tickets + steps: + - setProperty: + name: decision + expression: + language: + language: semantic + expression: refs:actionable,department,urgency + - setHeader: + name: urgency + expression: + simple: + expression: "${exchangeProperty.decision[urgency]}" + - choice: + when: + - expression: + simple: + expression: "${exchangeProperty.decision[department]} == 'billing'" + steps: + - to: direct:billing + otherwise: + steps: + - to: direct:review +---- + +A result might contain `actionable -> true`, `department -> "billing"` and `urgency -> 1.2`. +Reading the stored map does not invoke the provider again. Each question retains its own +threshold and uncertainty policy. There is no implicit AND/OR: batch expressions, including +a batch containing one boolean question, cannot be used as predicates. + +`CamelSemanticResults` holds an immutable map of the detailed `SemanticResult` objects, +keyed by the same names. These retain available probabilities, confidence and provider +metadata. `CamelSemanticResult` remains the single-question diagnostic property for +`ref:name`. Both diagnostic properties are cleared before every semantic evaluation; +a batch publishes its decisions and diagnostics only after every answer passes validation. +Missing, extra or invalid answers, operational failures and an uncertainty policy of `fail` +fail the entire evaluation through Camel error handling. A valid negative decision or +uncertainty handled by `non-match` remains a normal `false` result. A route property assigned +by Set Property is owned by the route; clear it explicitly before reevaluation if an error +handler could otherwise reuse a decision from an earlier attempt. + +Batch references are comma-separated names with surrounding whitespace removed. Empty, +duplicate and unknown references are rejected. There is no quoting or escaping in this +syntax: names containing commas or leading/trailing whitespace require a single `ref:name` +evaluation. Existing single references are unchanged. + +All selected questions must use the same effective Simple state selector, after applying +`camel.language.semantic.default-state` and resolving property placeholders. Selectors +must match exactly; different expressions that happen to return equal values are not +interchangeable. The selector is evaluated once per batch. Missing or unsupported state +fails before inference. One snapshot of the question definitions is used throughout each +batch, including validation and decision policies; the next evaluation sees replacements +from resource reload. Questions that need a previous answer require separate evaluations. + +`SemanticAdapter.evaluateBatch(Map<String, SemanticQuestion>, Object)` has a default +implementation that invokes the existing single-question method sequentially, preserving +compatibility with existing adapters. TypeSafe AI overrides it to send all questions in +one HTTP request using the component's existing timeout and concurrency limits. No results +are cached or combined across exchanges. + == EIP integration Use `language("semantic", "ref:name")` wherever an expression or boolean predicate is accepted. diff --git a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/language/semantic/SemanticLanguage.java b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/language/semantic/SemanticLanguage.java index 89257cb6ee4d..376677a50a06 100644 --- a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/language/semantic/SemanticLanguage.java +++ b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/language/semantic/SemanticLanguage.java @@ -19,7 +19,11 @@ package org.apache.camel.language.semantic; import java.io.IOException; import java.io.InputStream; import java.net.URL; +import java.util.Arrays; +import java.util.Collections; import java.util.Enumeration; +import java.util.HashSet; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Properties; @@ -52,6 +56,7 @@ import org.apache.camel.util.IOHelper; label = "language,ai", firstVersion = "4.23.0") public class SemanticLanguage extends LanguageSupport { public static final String RESULT = "CamelSemanticResult"; + public static final String RESULTS = "CamelSemanticResults"; public static final String ADAPTER_NAME = "camelSemanticAdapter"; public static final String ADAPTER_FACTORY = "semantic-adapter"; public static final String ADAPTER_RESOURCE = FactoryFinder.DEFAULT_PATH + ADAPTER_FACTORY; @@ -92,19 +97,41 @@ public class SemanticLanguage extends LanguageSupport { } public boolean validateExpression(String expression) { - if (expression == null || !expression.startsWith("ref:") || expression.substring(4).isBlank()) { - throw new IllegalArgumentException("Semantic expression must reference a named question using ref:name"); - } + references(expression); return true; } public boolean validatePredicate(String expression) { - return validateExpression(expression); + validateExpression(expression); + if (expression.startsWith("refs:")) { + throw new IllegalArgumentException("Semantic batch expressions cannot be predicates"); + } + return true; + } + + private List<String> references(String expression) { + if (expression != null && expression.startsWith("refs:")) { + List<String> names = Arrays.stream(expression.substring(5).split(",", -1)).map(String::strip).toList(); + if (names.stream().anyMatch(String::isEmpty)) { + throw new IllegalArgumentException("Semantic batch requires nonblank question names separated by commas"); + } + if (new HashSet<>(names).size() != names.size()) { + throw new IllegalArgumentException("Duplicate semantic question reference in batch"); + } + return names; + } + if (expression == null || !expression.startsWith("ref:") || expression.substring(4).isBlank()) { + throw new IllegalArgumentException("Semantic expression must use ref:name or refs:name1,name2"); + } + return List.of(expression.substring(4)); } private Evaluation createEvaluation(String expression, boolean predicate) { - validateExpression(expression); - Evaluation evaluation = new Evaluation(expression.substring(4), predicate); + List<String> names = references(expression); + if (predicate) { + validatePredicate(expression); + } + Evaluation evaluation = new Evaluation(names, expression.startsWith("refs:"), predicate); if (getCamelContext() != null) { evaluation.init(getCamelContext()); } @@ -252,14 +279,16 @@ public class SemanticLanguage extends LanguageSupport { } private final class Evaluation extends ExpressionAdapter { - private final String name; + private final List<String> names; + private final boolean batch; private final boolean predicate; private volatile Compiled compiled; private volatile SemanticQuestions questions; private volatile SemanticAdapter provider; - private Evaluation(String name, boolean predicate) { - this.name = name; + private Evaluation(List<String> names, boolean batch, boolean predicate) { + this.names = names; + this.batch = batch; this.predicate = predicate; } @@ -267,53 +296,101 @@ public class SemanticLanguage extends LanguageSupport { public void init(CamelContext context) { super.init(context); questions = SemanticQuestions.get(context); - SemanticQuestion question = questions.get(name); - if (predicate && question.getType() != SemanticQuestion.Type.BOOLEAN) { - throw new IllegalArgumentException("Semantic predicate requires a boolean question: " + name); + Map<String, SemanticQuestion> selected = questions.get(names); + if (predicate) { + requireBoolean(selected); } provider = adapter(); - compile(question); + compile(selected); } - private synchronized Compiled compile(SemanticQuestion question) { - if (compiled == null || compiled.question != question) { - if (predicate && question.getType() != SemanticQuestion.Type.BOOLEAN) { - throw new IllegalArgumentException("Semantic predicate requires a boolean question: " + name); + private void requireBoolean(Map<String, SemanticQuestion> selected) { + if (batch) { + throw new IllegalArgumentException("Semantic batch expressions cannot be predicates"); + } + if (selected.get(names.get(0)).getType() != SemanticQuestion.Type.BOOLEAN) { + throw new IllegalArgumentException("Semantic predicate requires a boolean question: " + names.get(0)); + } + } + + private synchronized Compiled compile(Map<String, SemanticQuestion> selected) { + if (compiled == null || !compiled.questions.equals(selected)) { + if (predicate) { + requireBoolean(selected); } - provider.validate(question); - String selector = question.getState() != null ? question.getState() : defaultState; - if (selector == null || selector.isBlank()) { - throw new IllegalArgumentException("Semantic state selector must not be blank for question: " + name); + String selector = null; + for (var entry : selected.entrySet()) { + SemanticQuestion question = entry.getValue(); + provider.validate(question); + String effective = question.getState() != null ? question.getState() : defaultState; + if (effective != null) { + effective = getCamelContext().resolvePropertyPlaceholders(effective); + } + if (effective == null || effective.isBlank()) { + throw new IllegalArgumentException( + "Semantic state selector must not be blank for question: " + entry.getKey()); + } + if (selector != null && !selector.equals(effective)) { + throw new IllegalArgumentException( + "Semantic batch questions must use the same effective state selector"); + } + selector = effective; } - selector = getCamelContext().resolvePropertyPlaceholders(selector); Expression state = getCamelContext().resolveLanguage("simple").createExpression(selector); state.init(getCamelContext()); - compiled = new Compiled(question, state); + compiled = new Compiled(selected, state); } return compiled; } @Override public Object evaluate(Exchange exchange) { + return evaluate(exchange, predicate); + } + + private Object evaluate(Exchange exchange, boolean asPredicate) { exchange.removeProperty(RESULT); + exchange.removeProperty(RESULTS); try { - SemanticQuestion question = questions.get(name); + Map<String, SemanticQuestion> selected = questions.get(names); + if (asPredicate) { + requireBoolean(selected); + } Compiled current = compiled; - if (current == null || current.question != question) { - current = compile(question); + if (current == null || !current.questions.equals(selected)) { + current = compile(selected); } Object state = current.state.evaluate(exchange, Object.class); if (state == null) { - throw new IllegalArgumentException("Missing selected state for semantic question: " + name); + throw new IllegalArgumentException("Missing selected state for semantic questions: " + names); } if (!(state instanceof String || state instanceof Map<?, ?> || state instanceof List<?>)) { throw new IllegalArgumentException( - "Unsupported state type for semantic question: " + name + "Unsupported state type for semantic questions: " + names + ". Select strings, maps or lists explicitly"); } + if (batch) { + Map<String, SemanticResult> results = provider.evaluateBatch(current.questions, state); + if (results == null || !results.keySet().equals(current.questions.keySet())) { + throw new IllegalArgumentException("Semantic batch result names must match question names"); + } + Map<String, Object> decisions = new LinkedHashMap<>(); + Map<String, SemanticResult> details = new LinkedHashMap<>(); + for (var entry : current.questions.entrySet()) { + SemanticResult result = results.get(entry.getKey()); + if (result == null) { + throw new IllegalArgumentException("Missing result for semantic question: " + entry.getKey()); + } + decisions.put(entry.getKey(), result.decision(entry.getValue())); + details.put(entry.getKey(), result); + } + exchange.setProperty(RESULTS, Collections.unmodifiableMap(details)); + return Collections.unmodifiableMap(decisions); + } + SemanticQuestion question = current.questions.get(names.get(0)); SemanticResult result = provider.evaluate(question, state); if (result == null) { - throw new IllegalArgumentException("Missing result for semantic question: " + name); + throw new IllegalArgumentException("Missing result for semantic question: " + names.get(0)); } Object decision = result.decision(question); exchange.setProperty(RESULT, result); @@ -328,19 +405,15 @@ public class SemanticLanguage extends LanguageSupport { @Override public boolean matches(Exchange exchange) { - exchange.removeProperty(RESULT); - if (questions.get(name).getType() != SemanticQuestion.Type.BOOLEAN) { - throw new IllegalArgumentException("Semantic predicate requires a boolean question: " + name); - } - return (Boolean) evaluate(exchange); + return (Boolean) evaluate(exchange, true); } @Override public String toString() { - return "semantic[ref:" + name + "]"; + return "semantic[" + (batch ? "refs:" : "ref:") + String.join(",", names) + "]"; } } - private record Compiled(SemanticQuestion question, Expression state) { + private record Compiled(Map<String, SemanticQuestion> questions, Expression state) { } } diff --git a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticAdapter.java b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticAdapter.java index 16d9f1c99e66..74d0d37f951a 100644 --- a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticAdapter.java +++ b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticAdapter.java @@ -16,6 +16,9 @@ */ package org.apache.camel.semantic; +import java.util.LinkedHashMap; +import java.util.Map; + /** * Provider-independent semantic evaluation. Implementations must support concurrent calls, bound evaluation time and * resource use, honor interruption, and release outstanding work on shutdown. Exceptions must not expose state or @@ -30,4 +33,20 @@ public interface SemanticAdapter { /** Synchronous, potentially blocking evaluation. Operational errors must be thrown, never returned as decisions. */ SemanticResult evaluate(SemanticQuestion question, Object state) throws Exception; + + /** + * Evaluate named questions against the same selected state, returning exactly one result per name. The default + * implementation calls the single-question method sequentially; providers may override it to use one request. + * Operational errors must be thrown, never returned as partial results. The language applies each decision policy. + */ + default Map<String, SemanticResult> evaluateBatch(Map<String, SemanticQuestion> questions, Object state) throws Exception { + Map<String, SemanticResult> results = new LinkedHashMap<>(); + for (var entry : questions.entrySet()) { + if (Thread.currentThread().isInterrupted()) { + throw new InterruptedException("Semantic batch evaluation interrupted"); + } + results.put(entry.getKey(), evaluate(entry.getValue(), state)); + } + return results; + } } diff --git a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticQuestions.java b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticQuestions.java index b71dffa90220..a326b3f766fb 100644 --- a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticQuestions.java +++ b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/SemanticQuestions.java @@ -16,7 +16,10 @@ */ package org.apache.camel.semantic; +import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import org.apache.camel.CamelContext; @@ -86,4 +89,18 @@ public final class SemanticQuestions { } return question; } + + /** Resolve all requested names from one immutable snapshot, preserving reference order. */ + public Map<String, SemanticQuestion> get(List<String> names) { + Map<String, SemanticQuestion> snapshot = questions; + Map<String, SemanticQuestion> selected = new LinkedHashMap<>(); + for (String name : names) { + SemanticQuestion question = snapshot.get(name); + if (question == null) { + throw new IllegalArgumentException("Unknown semantic question: " + name); + } + selected.put(name, question); + } + return Collections.unmodifiableMap(selected); + } } diff --git a/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticBatchTest.java b/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticBatchTest.java new file mode 100644 index 000000000000..c434a4b58176 --- /dev/null +++ b/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticBatchTest.java @@ -0,0 +1,322 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.camel.semantic; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.Expression; +import org.apache.camel.Predicate; +import org.apache.camel.impl.DefaultCamelContext; +import org.apache.camel.language.semantic.SemanticLanguage; +import org.apache.camel.support.DefaultExchange; +import org.apache.camel.support.DefaultMessage; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class SemanticBatchTest { + private DefaultCamelContext context; + private SemanticLanguage language; + private DefaultExchange exchange; + private RecordingAdapter adapter; + + @BeforeEach + void setup() throws Exception { + context = new DefaultCamelContext(); + language = new SemanticLanguage(); + language.setCamelContext(context); + language.setAdapter("adapter"); + adapter = new RecordingAdapter(); + context.getRegistry().bind("adapter", adapter); + context.start(); + exchange = new DefaultExchange(context); + exchange.getMessage().setBody("original"); + SemanticQuestions.get(context).replace("test", Map.of( + "urgent", question(SemanticQuestion.Type.BOOLEAN, null, 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL), + "department", question(SemanticQuestion.Type.CHOICE, null, 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL), + "priority", question(SemanticQuestion.Type.SCORE, null, 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL))); + } + + @AfterEach + void stop() throws Exception { + context.stop(); + } + + private static SemanticQuestion question( + SemanticQuestion.Type type, String state, double threshold, + double uncertainty, SemanticQuestion.UncertaintyPolicy policy) { + return new SemanticQuestion( + type, "Classify", state, + type == SemanticQuestion.Type.CHOICE ? Map.of("billing", "Payments", "technical", "Bugs") : Map.of(), + type == SemanticQuestion.Type.SCORE ? List.of("low", "medium", "high") : List.of(), + threshold, uncertainty, policy); + } + + @Test + void existingAdapterEvaluatesMixedQuestionsSequentiallyAndRetainsDetails() { + Expression expression = language.createExpression("refs: urgent, department, priority "); + assertThat(adapter.calls).isEmpty(); + assertThat(expression.evaluate(exchange, Map.class)) + .containsAllEntriesOf(Map.of("urgent", true, "department", "billing", "priority", 1.2)); + Map<String, Object> decisions = expression.evaluate(exchange, Map.class); + assertThat(new ArrayList<>(decisions.keySet())).containsExactly("urgent", "department", "priority"); + assertThat(adapter.calls).containsExactly(SemanticQuestion.Type.BOOLEAN, SemanticQuestion.Type.CHOICE, + SemanticQuestion.Type.SCORE, SemanticQuestion.Type.BOOLEAN, SemanticQuestion.Type.CHOICE, + SemanticQuestion.Type.SCORE); + assertThat(adapter.states).containsOnly("original"); + assertThat(exchange.getMessage().getBody()).isEqualTo("original"); + Map<?, ?> details = exchange.getProperty(SemanticLanguage.RESULTS, Map.class); + SemanticResult department = (SemanticResult) details.get("department"); + assertThat(department.getProbabilities()).containsEntry("billing", 0.9); + assertThat(department.getConfidence()).isEqualTo(0.8); + assertThat(department.getMetadata()).containsEntry("provider", "fixture"); + assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isNull(); + language.createExpression("ref:urgent").evaluate(exchange, Boolean.class); + assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isInstanceOf(SemanticResult.class); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + } + + @ParameterizedTest + @ValueSource(strings = { + "refs:", "refs: ", "refs:,urgent", "refs:urgent,", "refs:urgent,,priority", + "refs:urgent,urgent", "refs:urgent, urgent", "refs:urgent,unknown" }) + void rejectsInvalidReferencesBeforeInference(String expression) { + assertThatThrownBy(() -> language.createExpression(expression)).isInstanceOf(IllegalArgumentException.class); + assertThat(adapter.calls).isEmpty(); + } + + @Test + void singleReferenceStillAcceptsNamesContainingCommas() { + SemanticQuestions.get(context).replace("comma", Map.of("a,b", SemanticQuestions.get(context).get("urgent"))); + assertThat(language.createExpression("ref:a,b").evaluate(exchange, Boolean.class)).isTrue(); + assertThatThrownBy(() -> language.createExpression("refs:a,b")).hasMessageContaining("Unknown"); + } + + @Test + void batchesCannotBePredicatesEvenWithOneBooleanQuestion() { + assertThatThrownBy(() -> language.validatePredicate("refs:urgent")).hasMessageContaining("cannot be predicates"); + assertThatThrownBy(() -> language.createPredicate("refs:urgent")).hasMessageContaining("cannot be predicates"); + Expression expression = language.createExpression("refs:urgent"); + expression.evaluate(exchange, Object.class); + assertThatThrownBy(() -> ((Predicate) expression).matches(exchange)).hasMessageContaining("cannot be predicates"); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + assertThat(adapter.calls).hasSize(1); + } + + @Test + void sharedStateSelectorIsEvaluatedOnlyOnce() { + AtomicInteger reads = new AtomicInteger(); + exchange.setIn(new DefaultMessage(context) { + @Override + public Object getHeader(String name) { + if (name.equals("selected")) { + reads.incrementAndGet(); + return "selected"; + } + return super.getHeader(name); + } + }); + language.setDefaultState("${header.selected}"); + language.createExpression("refs:urgent,department,priority").evaluate(exchange, Map.class); + assertThat(reads).hasValue(1); + assertThat(adapter.states).containsExactly("selected", "selected", "selected"); + } + + @Test + void resolvesSharedDefaultAndExplicitSelectorsBeforeComparison() { + Properties properties = new Properties(); + properties.setProperty("selected", "${header.selected}"); + context.getPropertiesComponent().setInitialProperties(properties); + language.setDefaultState("{{selected}}"); + SemanticQuestions.get(context).replace("explicit", Map.of("other", + question(SemanticQuestion.Type.BOOLEAN, "${header.selected}", 0.5, 0, + SemanticQuestion.UncertaintyPolicy.FAIL))); + exchange.getMessage().setHeader("selected", Map.of("text", "invoice")); + Expression expression = language.createExpression("refs:urgent,other"); + expression.evaluate(exchange, Map.class); + assertThat(adapter.states).containsExactly(Map.of("text", "invoice"), Map.of("text", "invoice")); + exchange.getMessage().removeHeader("selected"); + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)).hasMessageContaining("Missing selected state"); + assertThat(adapter.calls).hasSize(2); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + } + + @Test + void incompatibleSelectorsAndUnsupportedCapabilitiesFailBeforeInference() { + SemanticQuestions.get(context).replace("other", Map.of("other", + question(SemanticQuestion.Type.BOOLEAN, "${header.other}", 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL))); + assertThatThrownBy(() -> language.createExpression("refs:urgent,other")) + .hasMessageContaining("same effective state selector"); + context.getRegistry().unbind("adapter"); + context.getRegistry().bind("adapter", new RecordingAdapter() { + @Override + public void validate(SemanticQuestion question) { + if (question.getType() == SemanticQuestion.Type.SCORE) { + throw new IllegalArgumentException("Unsupported score"); + } + } + }); + // Use a fresh language because adapter selection is cached for the context lifetime. + SemanticLanguage other = new SemanticLanguage(); + other.setCamelContext(context); + other.setAdapter("adapter"); + assertThatThrownBy(() -> other.createExpression("refs:urgent,priority")).hasMessageContaining("Unsupported score"); + assertThat(adapter.calls).isEmpty(); + } + + @Test + void eachQuestionKeepsItsDecisionPolicyAndFailureClearsAllDiagnostics() { + SemanticQuestion negative + = question(SemanticQuestion.Type.BOOLEAN, null, 0.95, 0, SemanticQuestion.UncertaintyPolicy.FAIL); + SemanticQuestion uncertain + = question(SemanticQuestion.Type.BOOLEAN, null, 0.9, 0.05, SemanticQuestion.UncertaintyPolicy.NON_MATCH); + SemanticQuestions.get(context).replace("policy", Map.of("negative", negative, "uncertain", uncertain)); + Expression expression = language.createExpression("refs:urgent,negative,uncertain"); + assertThat(expression.evaluate(exchange, Map.class)).containsEntry("urgent", true) + .containsEntry("negative", false).containsEntry("uncertain", false); + SemanticQuestions.get(context).replace("policy", Map.of("negative", negative, "uncertain", + question(SemanticQuestion.Type.BOOLEAN, null, 0.9, 0.05, SemanticQuestion.UncertaintyPolicy.FAIL))); + exchange.setProperty(SemanticLanguage.RESULT, "stale"); + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)).hasMessageContaining("uncertain"); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isNull(); + } + + @ParameterizedTest + @ValueSource(strings = { "missing", "extra", "null", "invalid", "failure" }) + void rejectsIncompleteOrInvalidBatchResultsWithoutPublishingPartialDiagnostics(String failure) { + context.getRegistry().unbind("adapter"); + context.getRegistry().bind("adapter", new RecordingAdapter() { + @Override + public Map<String, SemanticResult> evaluateBatch(Map<String, SemanticQuestion> questions, Object state) { + if (failure.equals("failure")) { + throw new IllegalStateException("provider failed"); + } + Map<String, SemanticResult> results = new LinkedHashMap<>(); + results.put("urgent", result(SemanticQuestion.Type.BOOLEAN)); + if (!failure.equals("missing")) { + results.put("department", failure.equals("null") ? null + : failure.equals("invalid") ? new SemanticResult("undeclared", null, null, null, null) + : result(SemanticQuestion.Type.CHOICE)); + } + if (failure.equals("extra")) { + results.put("extra", result(SemanticQuestion.Type.BOOLEAN)); + } + return results; + } + }); + Expression expression = language.createExpression("refs:urgent,department"); + exchange.setProperty(SemanticLanguage.RESULT, "stale"); + exchange.setProperty(SemanticLanguage.RESULTS, "stale"); + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)).isInstanceOf(RuntimeException.class); + assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isNull(); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + } + + @Test + void reloadDuringValidationOrEvaluationCannotMixQuestionDefinitions() { + SemanticQuestion old = SemanticQuestions.get(context).get("urgent"); + SemanticQuestion updated + = question(SemanticQuestion.Type.BOOLEAN, null, 0.95, 0, SemanticQuestion.UncertaintyPolicy.FAIL); + SemanticQuestions questions = SemanticQuestions.get(context); + questions.replace("test", Map.of("first", old, "second", old)); + AtomicInteger validations = new AtomicInteger(); + context.getRegistry().unbind("adapter"); + context.getRegistry().bind("adapter", new RecordingAdapter() { + @Override + public void validate(SemanticQuestion question) { + if (validations.incrementAndGet() == 1) { + questions.replace("test", Map.of("first", updated, "second", updated)); + } + } + + @Override + public SemanticResult evaluate(SemanticQuestion question, Object state) throws Exception { + questions.replace("test", Map.of("first", old, "second", old)); + return super.evaluate(question, state); + } + }); + Expression expression = language.createExpression("refs:first,second"); + assertThat(expression.evaluate(exchange, Map.class)).containsEntry("first", false).containsEntry("second", false); + assertThat(expression.evaluate(exchange, Map.class)).containsEntry("first", true).containsEntry("second", true); + questions.replace("test", Map.of("first", old)); + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)) + .hasMessageContaining("Unknown semantic question: second"); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + questions.replace("test", Map.of("first", old, "second", + question(SemanticQuestion.Type.BOOLEAN, "${header.changed}", 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL))); + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)) + .hasMessageContaining("same effective state selector"); + } + + @Test + void interruptionStopsSequentialFallbackAndPreservesInterruptFlag() { + AtomicInteger invocations = new AtomicInteger(); + context.getRegistry().unbind("adapter"); + context.getRegistry().bind("adapter", new RecordingAdapter() { + @Override + public SemanticResult evaluate(SemanticQuestion question, Object state) throws Exception { + invocations.incrementAndGet(); + throw new InterruptedException("interrupted"); + } + }); + Expression expression = language.createExpression("refs:urgent,department"); + try { + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)).hasCauseInstanceOf(InterruptedException.class); + assertThat(Thread.currentThread().isInterrupted()).isTrue(); + assertThat(invocations).hasValue(1); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + } finally { + Thread.interrupted(); + } + } + + private static SemanticResult result(SemanticQuestion.Type type) { + return switch (type) { + case BOOLEAN -> new SemanticResult(null, 0.9, null, null, Map.of("provider", "fixture")); + case CHOICE -> new SemanticResult( + "billing", null, Map.of("billing", 0.9, "technical", 0.1), 0.8, Map.of("provider", "fixture")); + case SCORE -> new SemanticResult(1.2, null, null, 0.7, Map.of("provider", "fixture")); + }; + } + + private static class RecordingAdapter implements SemanticAdapter { + final List<SemanticQuestion.Type> calls = new ArrayList<>(); + final List<Object> states = new ArrayList<>(); + + @Override + public void validate(SemanticQuestion question) { + } + + @Override + public SemanticResult evaluate(SemanticQuestion question, Object state) throws Exception { + calls.add(question.getType()); + states.add(state); + return result(question.getType()); + } + } +} diff --git a/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticEipTest.java b/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticEipTest.java index c6fc488cc534..868f5e68d3b6 100644 --- a/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticEipTest.java +++ b/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticEipTest.java @@ -74,6 +74,13 @@ class SemanticEipTest extends CamelTestSupport { "unfinished", question("unfinished", SemanticQuestion.Type.BOOLEAN), "urgency", question("urgency", SemanticQuestion.Type.SCORE))); + from("direct:batch").setProperty("decision").language("semantic", "refs:actionable,department,urgency") + .setHeader("urgency").simple("${exchangeProperty.decision[urgency]}") + .choice() + .when(simple( + "${exchangeProperty.decision[actionable]} == true && ${exchangeProperty.decision[department]} == 'billing'")) + .to("mock:batchBilling") + .otherwise().to("mock:batchReview"); from("direct:choice").setProperty("department").language("semantic", "ref:department") .choice() .when(exchangeProperty("department").isEqualTo("billing")).to("mock:billing") @@ -151,6 +158,17 @@ class SemanticEipTest extends CamelTestSupport { 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL); } + @Test + void batchDecisionsAreReusedByOrdinaryEips() throws Exception { + getMockEndpoint("mock:batchBilling").expectedBodiesReceived("please fix this urgent invoice"); + getMockEndpoint("mock:batchBilling").expectedHeaderReceived("urgency", 2.0); + template.sendBody("direct:batch", "please fix this urgent invoice"); + MockEndpoint.assertIsSatisfied(context); + assertThat(calls.get("actionable")).hasValue(1); + assertThat(calls.get("department")).hasValue(1); + assertThat(calls.get("urgency")).hasValue(1); + } + @Test void choiceFilterValidationAndStoredMetadata() throws Exception { getMockEndpoint("mock:technical").expectedMessageCount(1); diff --git a/components/camel-ai/camel-typesafe-ai/src/main/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapter.java b/components/camel-ai/camel-typesafe-ai/src/main/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapter.java index f486e7e4729c..3dd1ed114fef 100644 --- a/components/camel-ai/camel-typesafe-ai/src/main/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapter.java +++ b/components/camel-ai/camel-typesafe-ai/src/main/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapter.java @@ -17,6 +17,7 @@ package org.apache.camel.component.typesafeai; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.Map; import org.apache.camel.CamelContext; @@ -67,28 +68,48 @@ public class TypeSafeAiSemanticAdapter implements SemanticAdapter, CamelContextA @Override public SemanticResult evaluate(SemanticQuestion question, Object state) throws Exception { + return evaluateBatch(Map.of("question", question), state).get("question"); + } + + @Override + public Map<String, SemanticResult> evaluateBatch(Map<String, SemanticQuestion> questions, Object state) throws Exception { + Map<String, Object> definitions = new LinkedHashMap<>(); + questions.forEach((name, question) -> definitions.put(name, definition(question))); + JsonObject response = endpoint().evaluate(Map.of("state", state, "questions", definitions)); + Map<String, SemanticResult> results = new LinkedHashMap<>(); + questions.forEach((name, question) -> results.put(name, + result(question, response.getJsonObject("answers").getJsonObject(name), response))); + return results; + } + + private Map<String, Object> definition(SemanticQuestion question) { validate(question); Map<String, Object> definition = new HashMap<>(); definition.put("instructions", question.getInstructions()); - String type = switch (question.getType()) { - case BOOLEAN -> "noul"; - case CHOICE -> "choice"; - case SCORE -> "score"; - }; - definition.put("type", type); + definition.put("type", type(question)); if (question.getType() == SemanticQuestion.Type.SCORE) { definition.put("criteria", question.getLevels()); } else if (!question.getCriteria().isEmpty()) { definition.put("criteria", question.getCriteria()); } - JsonObject response = endpoint().evaluate(Map.of("state", state, "questions", Map.of("question", definition))); - JsonObject answer = response.getJsonObject("answers").getJsonObject("question"); + return definition; + } + + private String type(SemanticQuestion question) { + return switch (question.getType()) { + case BOOLEAN -> "noul"; + case CHOICE -> "choice"; + case SCORE -> "score"; + }; + } + + private SemanticResult result(SemanticQuestion question, JsonObject answer, JsonObject response) { Map<String, Double> probabilities = new HashMap<>(); if (answer.get("probabilities") instanceof Map<?, ?> values) { values.forEach((key, value) -> probabilities.put((String) key, ((Number) value).doubleValue())); } return new SemanticResult( - question.getType() == SemanticQuestion.Type.BOOLEAN ? null : answer.get(type), + question.getType() == SemanticQuestion.Type.BOOLEAN ? null : answer.get(type(question)), question.getType() == SemanticQuestion.Type.BOOLEAN ? answer.getDouble("noul") : null, probabilities, answer.get("confidence") == null ? null : answer.getDouble("confidence"), Map.of("provider", "typesafe-ai", "model", response.get("model"), "usage", response.get("usage"))); diff --git a/components/camel-ai/camel-typesafe-ai/src/test/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapterTest.java b/components/camel-ai/camel-typesafe-ai/src/test/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapterTest.java index ac90da5eb6e9..3215add1d79b 100644 --- a/components/camel-ai/camel-typesafe-ai/src/test/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapterTest.java +++ b/components/camel-ai/camel-typesafe-ai/src/test/java/org/apache/camel/component/typesafeai/TypeSafeAiSemanticAdapterTest.java @@ -27,7 +27,10 @@ import org.apache.camel.semantic.SemanticQuestion; import org.apache.camel.semantic.SemanticQuestions; import org.apache.camel.semantic.SemanticResult; import org.apache.camel.support.DefaultExchange; +import org.apache.camel.util.json.JsonObject; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -108,9 +111,83 @@ class TypeSafeAiSemanticAdapterTest extends TypeSafeAiTestSupport { assertThat(expression(SemanticQuestion.Type.SCORE).evaluate(exchange, Double.class)).isEqualTo(0.7); } + private Expression batchExpression() { + SemanticQuestions.get(context).replace("test", Map.of( + "refund", new SemanticQuestion( + SemanticQuestion.Type.BOOLEAN, "Refund requested?", null, + Map.of(), List.of(), 0.95, 0, SemanticQuestion.UncertaintyPolicy.FAIL), + "department", new SemanticQuestion( + SemanticQuestion.Type.CHOICE, "Which department?", null, + Map.of("billing", "Refunds", "technical", "Faults", "other", "Anything else"), List.of(), + 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL), + "urgency", new SemanticQuestion( + SemanticQuestion.Type.SCORE, "How urgent?", null, + Map.of(), List.of("Routine", "Urgent", "Critical"), 0.5, 0, SemanticQuestion.UncertaintyPolicy.FAIL))); + return context.resolveLanguage("semantic").createExpression("refs:refund,department,urgency"); + } + @Test - void providerErrorsClearPreviousResult() { - Expression expression = expression(SemanticQuestion.Type.BOOLEAN); + void mixedBatchUsesOneRequestAndPreservesEveryResult() { + respond = request -> mixedResponse(); + Expression expression = batchExpression(); + assertThat(requests).isEmpty(); + var exchange = new DefaultExchange(context); + exchange.getMessage().setBody(Map.of("ticket", "refund requested")); + assertThat(expression.evaluate(exchange, Map.class)) + .containsEntry("refund", false).containsEntry("department", "billing").containsEntry("urgency", 1.2); + assertThat(requests).hasSize(1); + JsonObject request = requests.peek(); + assertThat(request.getJsonObject("questions").keySet()).containsExactlyInAnyOrder("refund", "department", "urgency"); + assertThat(request.get("state")).isEqualTo(exchange.getMessage().getBody()); + assertThat(authorization).containsExactly("Bearer test-key"); + Map<?, ?> results = exchange.getProperty(SemanticLanguage.RESULTS, Map.class); + SemanticResult refund = (SemanticResult) results.get("refund"); + SemanticResult department = (SemanticResult) results.get("department"); + SemanticResult urgency = (SemanticResult) results.get("urgency"); + assertThat(refund.getProbability()).isEqualTo(0.9); + assertThat(refund.getConfidence()).isNull(); + assertThat(department.getProbabilities()).containsEntry("billing", 0.9); + assertThat(department.getConfidence()).isEqualTo(0.8); + assertThat(urgency.getProbabilities()).containsEntry("1", 0.8); + assertThat(urgency.getConfidence()).isEqualTo(0.7); + assertThat(urgency.getMetadata()).containsEntry("provider", "typesafe-ai").containsEntry("model", "jev-1.13.0") + .containsKey("usage"); + assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isNull(); + } + + @ParameterizedTest + @ValueSource(strings = { "missing", "extra", "type", "probability", "choice", "score", "confidence" }) + void malformedBatchResponseClearsPreviousDiagnostics(String failure) throws Exception { + JsonObject response = TypeSafeAiJson.parse(mixedResponse()); + JsonObject answers = response.getJsonObject("answers"); + switch (failure) { + case "missing" -> answers.remove("urgency"); + case "extra" -> answers.put("extra", answers.get("refund")); + case "type" -> answers.getJsonObject("refund").put("type", "choice"); + case "probability" -> answers.getJsonObject("refund").put("noul", 1.1); + case "choice" -> answers.getJsonObject("department").put("choice", "undeclared"); + case "score" -> answers.getJsonObject("urgency").put("score", 3); + case "confidence" -> answers.getJsonObject("urgency").remove("confidence"); + default -> throw new IllegalArgumentException(failure); + } + respond = request -> response.toJson(); + Expression expression = batchExpression(); + var exchange = new DefaultExchange(context); + exchange.getMessage().setBody("private-input"); + exchange.setProperty(SemanticLanguage.RESULT, "old"); + exchange.setProperty(SemanticLanguage.RESULTS, "old"); + assertThatThrownBy(() -> expression.evaluate(exchange, Map.class)) + .hasStackTraceContaining("Invalid TypeSafe AI response") + .hasMessageNotContaining("private-input"); + assertThat(requests).hasSize(1); + assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isNull(); + assertThat(exchange.getProperty(SemanticLanguage.RESULTS)).isNull(); + } + + @ParameterizedTest + @ValueSource(booleans = { false, true }) + void providerErrorsClearPreviousResult(boolean batch) { + Expression expression = batch ? batchExpression() : expression(SemanticQuestion.Type.BOOLEAN); var exchange = new DefaultExchange(context); exchange.getMessage().setBody("private-input"); exchange.setProperty(SemanticLanguage.RESULT, "old"); @@ -120,11 +197,12 @@ class TypeSafeAiSemanticAdapterTest extends TypeSafeAiTestSupport { assertThat(exchange.getProperty(SemanticLanguage.RESULT)).isNull(); } - @Test - void componentTimeoutBoundsSemanticEvaluation() { + @ParameterizedTest + @ValueSource(booleans = { false, true }) + void componentTimeoutBoundsSemanticEvaluation(boolean batch) { context.getComponent("typesafe-ai", TypeSafeAiComponent.class).getConfiguration().setRequestTimeout(100); holdHeaders = true; - Expression expression = expression(SemanticQuestion.Type.BOOLEAN); + Expression expression = batch ? batchExpression() : expression(SemanticQuestion.Type.BOOLEAN); var exchange = new DefaultExchange(context); exchange.getMessage().setBody("private-input"); assertThatThrownBy(() -> expression.evaluate(exchange, Object.class)) diff --git a/dsl/camel-yaml-dsl/camel-yaml-dsl/src/test/java/org/apache/camel/dsl/yaml/SemanticQuestionTest.java b/dsl/camel-yaml-dsl/camel-yaml-dsl/src/test/java/org/apache/camel/dsl/yaml/SemanticQuestionTest.java index 639141d0dc06..7008fde6d70e 100644 --- a/dsl/camel-yaml-dsl/camel-yaml-dsl/src/test/java/org/apache/camel/dsl/yaml/SemanticQuestionTest.java +++ b/dsl/camel-yaml-dsl/camel-yaml-dsl/src/test/java/org/apache/camel/dsl/yaml/SemanticQuestionTest.java @@ -19,6 +19,7 @@ package org.apache.camel.dsl.yaml; import java.nio.file.Files; import java.nio.file.Path; import java.util.List; +import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Stream; @@ -123,6 +124,39 @@ class SemanticQuestionTest extends YamlTestSupport { """; } + @Test + void batchExpressionReusesNamedQuestionsAndResults() throws Exception { + loadRoutes(declarations("${body}") + declarations("${body}").replace("department:", "second:") + """ + - route: + from: + uri: direct:batch + steps: + - setProperty: + name: decision + expression: + language: + language: semantic + expression: refs:department,second + - setHeader: + name: selectedDepartment + expression: + simple: + expression: "${exchangeProperty.decision[department]}" + - to: mock:batch + """); + context.start(); + MockEndpoint mock = context.getEndpoint("mock:batch", MockEndpoint.class); + mock.expectedBodiesReceived("invoice"); + mock.expectedHeaderReceived("selectedDepartment", "billing"); + try (var template = context.createProducerTemplate()) { + template.sendBody("direct:batch", "invoice"); + } + mock.assertIsSatisfied(); + assertThat(calls).hasValue(2); + assertThat(mock.getExchanges().get(0).getProperty(SemanticLanguage.RESULTS, Map.class)) + .containsKeys("department", "second"); + } + @Test void declarationAfterRouteEvaluatesOnceAndPreservesMessage() throws Exception { loadRoutes(route() + declarations("${header.selected}"));
