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 fcf09afc4800187e598f5b4078f9fa31d16a66a2 Author: Luigi De Masi <[email protected]> AuthorDate: Thu Sep 24 11:59:44 2026 +0200 CAMEL-24977: Locate semantic YAML errors and isolate initialization locks Report invalid enums, duplicate declarations, unknown fields and malformed questions with YAML source marks and preserve the underlying validation cause. Use private locks for question registration and adapter ownership. Keep adapter construction isolated per context and coordinate identity-checked cleanup with registration without locking on the public context or registry objects. Cover source locations, concurrent initialization, ownership collisions and independent context progress with focused regression tests. Co-authored-by: Codex <[email protected]> Signed-off-by: Luigi De Masi <[email protected]> --- .../camel/language/semantic/SemanticLanguage.java | 31 +++- .../apache/camel/semantic/SemanticQuestions.java | 3 +- .../yaml/SemanticDefinitionDeserializer.java | 101 +++++++---- .../camel/semantic/SemanticInitializationTest.java | 194 +++++++++++++++++++++ .../camel/dsl/yaml/SemanticQuestionTest.java | 75 ++++++++ 5 files changed, 363 insertions(+), 41 deletions(-) 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 ce790c86dda4..0d97c29d8a8e 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 @@ -153,13 +153,14 @@ public class SemanticLanguage extends LanguageSupport { throw new IllegalArgumentException("No semantic adapter bean or class found: " + className); } Class<? extends SemanticAdapter> type = resolved.asSubclass(SemanticAdapter.class); - synchronized (context.getRegistry()) { + AdapterLock lock = AdapterLock.get(context); + synchronized (lock.monitor) { if (context.getRegistry().lookupByName(ADAPTER_NAME) != null) { throw new IllegalArgumentException("Semantic adapter registry name is already bound: " + ADAPTER_NAME); } SemanticAdapter instance = context.getInjector().newInstance(type); CamelContextAware.trySetCamelContext(instance, context); - owned = new ManagedAdapter(context, instance); + owned = new ManagedAdapter(context, instance, lock); context.getRegistry().bind(ADAPTER_NAME, instance); } context.addService(owned, true, true); @@ -178,13 +179,31 @@ public class SemanticLanguage extends LanguageSupport { } } + private static final class AdapterLock { + private static final Object CREATION_LOCK = new Object(); + private final Object monitor = new Object(); + + private static AdapterLock get(CamelContext context) { + synchronized (CREATION_LOCK) { + AdapterLock lock = context.getCamelContextExtension().getContextPlugin(AdapterLock.class); + if (lock == null) { + lock = new AdapterLock(); + context.getCamelContextExtension().addContextPlugin(AdapterLock.class, lock); + } + return lock; + } + } + } + private static final class ManagedAdapter extends ServiceSupport { private final CamelContext context; private final SemanticAdapter instance; + private final AdapterLock lock; - private ManagedAdapter(CamelContext context, SemanticAdapter instance) { + private ManagedAdapter(CamelContext context, SemanticAdapter instance, AdapterLock lock) { this.context = context; this.instance = instance; + this.lock = lock; } @Override @@ -207,8 +226,10 @@ public class SemanticLanguage extends LanguageSupport { try { ServiceHelper.stopAndShutdownService(instance); } finally { - if (context.getRegistry().lookupByName(ADAPTER_NAME) == instance) { - context.getRegistry().unbind(ADAPTER_NAME); + synchronized (lock.monitor) { + if (context.getRegistry().lookupByName(ADAPTER_NAME) == instance) { + context.getRegistry().unbind(ADAPTER_NAME); + } } } } 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 5a67f7b63e15..b71dffa90220 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 @@ -24,12 +24,13 @@ import org.apache.camel.spi.Resource; /** Context-local named questions, replaced atomically per source when a route resource is reloaded. */ public final class SemanticQuestions { + private static final Object CREATION_LOCK = new Object(); private final Map<String, Map<String, SemanticQuestion>> sources = new HashMap<>(); private final Map<String, Resource> resources = new HashMap<>(); private volatile Map<String, SemanticQuestion> questions = Map.of(); public static SemanticQuestions get(CamelContext context) { - synchronized (context) { + synchronized (CREATION_LOCK) { SemanticQuestions answer = context.getCamelContextExtension().getContextPlugin(SemanticQuestions.class); if (answer == null) { answer = new SemanticQuestions(); diff --git a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/yaml/SemanticDefinitionDeserializer.java b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/yaml/SemanticDefinitionDeserializer.java index 355a05a2e521..a849031c2f2b 100644 --- a/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/yaml/SemanticDefinitionDeserializer.java +++ b/components/camel-ai/camel-semantic/src/main/java/org/apache/camel/semantic/yaml/SemanticDefinitionDeserializer.java @@ -73,7 +73,8 @@ public class SemanticDefinitionDeserializer extends YamlDeserializerSupport impl if ("semantic".equals(asText(tuple.getKeyNode()))) { read(tuple.getValueNode()).forEach((name, question) -> { if (definitions.putIfAbsent(name, question) != null) { - throw new IllegalArgumentException("Duplicate semantic question: " + name); + throw new YamlDeserializationException( + tuple.getValueNode(), "Duplicate semantic question: " + name); } }); } @@ -84,51 +85,81 @@ public class SemanticDefinitionDeserializer extends YamlDeserializerSupport impl ? context.getCamelContextExtension().getContextPlugin(SemanticQuestions.class) : SemanticQuestions.get(context); if (questions != null) { - questions.replace(dc.getResource(), definitions); + try { + questions.replace(dc.getResource(), definitions); + } catch (IllegalArgumentException e) { + throw new YamlDeserializationException(root, e.getMessage(), e); + } } } private static Map<String, SemanticQuestion> read(Node node) { - Map<String, Node> semantic = fields(node); + Map<String, Node> semantic = fields(node, "semantic declaration"); if (!semantic.keySet().equals(Set.of("question"))) { - throw new IllegalArgumentException("Semantic declaration requires only question"); + throw new YamlDeserializationException(node, "Semantic declaration requires only question"); } Map<String, SemanticQuestion> result = new LinkedHashMap<>(); - fields(semantic.get("question")).forEach((name, definition) -> { - Map<String, Node> values = fields(definition); - if (!FIELDS.containsAll(values.keySet())) { - throw new IllegalArgumentException("Unknown property in semantic question: " + name); - } - String typeName = asText(values.get("type")); - if (typeName == null) { - throw new IllegalArgumentException("Semantic question type is required: " + name); + fields(semantic.get("question"), "semantic questions").forEach((name, definition) -> { + if (name.isBlank()) { + throw new YamlDeserializationException(definition, "Semantic question requires a nonblank name"); } - SemanticQuestion.Type type = SemanticQuestion.Type.valueOf(typeName.toUpperCase(Locale.ROOT)); - Map<String, String> criteria = new LinkedHashMap<>(); - List<String> levels = List.of(); - if (values.containsKey("criteria")) { - if (type == SemanticQuestion.Type.SCORE) { - levels = asSequenceNode(values.get("criteria")).getValue().stream().map(YamlDeserializerSupport::asText) - .toList(); - } else { - fields(values.get("criteria")).forEach((key, value) -> criteria.put(key, asText(value))); - } + try { + result.put(name, readQuestion(name, definition)); + } catch (IllegalArgumentException e) { + throw new YamlDeserializationException( + definition, "Invalid semantic question '" + name + "': " + e.getMessage(), e); } - if (type != SemanticQuestion.Type.BOOLEAN && (values.containsKey("threshold") || values.containsKey("uncertainty") - || values.containsKey("uncertaintyPolicy"))) { - throw new IllegalArgumentException("Threshold and uncertainty policy require a boolean question: " + name); - } - SemanticQuestion.UncertaintyPolicy policy = values.containsKey("uncertaintyPolicy") - ? SemanticQuestion.UncertaintyPolicy - .valueOf(asText(values.get("uncertaintyPolicy")).replace('-', '_').toUpperCase(Locale.ROOT)) - : SemanticQuestion.UncertaintyPolicy.FAIL; - result.put(name, new SemanticQuestion( - type, asText(values.get("instructions")), asText(values.get("state")), - criteria, levels, number(values, name, "threshold", 0.5), number(values, name, "uncertainty", 0), policy)); }); return result; } + private static SemanticQuestion readQuestion(String name, Node definition) { + Map<String, Node> values = fields(definition, "semantic question '" + name + "'"); + values.forEach((field, value) -> { + if (!FIELDS.contains(field)) { + throw new YamlDeserializationException( + value, "Unknown property '" + field + "' in semantic question '" + name + "'"); + } + }); + if (!values.containsKey("type")) { + throw new YamlDeserializationException(definition, "Semantic question type is required: " + name); + } + SemanticQuestion.Type type = enumeration(values.get("type"), name, "type", SemanticQuestion.Type.class); + Map<String, String> criteria = new LinkedHashMap<>(); + List<String> levels = List.of(); + if (values.containsKey("criteria")) { + if (type == SemanticQuestion.Type.SCORE) { + levels = asSequenceNode(values.get("criteria")).getValue().stream().map(YamlDeserializerSupport::asText) + .toList(); + } else { + fields(values.get("criteria"), "criteria for semantic question '" + name + "'") + .forEach((key, value) -> criteria.put(key, asText(value))); + } + } + if (type != SemanticQuestion.Type.BOOLEAN && (values.containsKey("threshold") || values.containsKey("uncertainty") + || values.containsKey("uncertaintyPolicy"))) { + throw new YamlDeserializationException( + definition, "Threshold and uncertainty policy require a boolean question: " + name); + } + SemanticQuestion.UncertaintyPolicy policy = values.containsKey("uncertaintyPolicy") + ? enumeration(values.get("uncertaintyPolicy"), name, "uncertaintyPolicy", + SemanticQuestion.UncertaintyPolicy.class) + : SemanticQuestion.UncertaintyPolicy.FAIL; + return new SemanticQuestion( + type, asText(values.get("instructions")), asText(values.get("state")), + criteria, levels, number(values, name, "threshold", 0.5), number(values, name, "uncertainty", 0), policy); + } + + private static <T extends Enum<T>> T enumeration(Node node, String question, String field, Class<T> type) { + String raw = asText(node); + try { + return Enum.valueOf(type, raw.replace('-', '_').toUpperCase(Locale.ROOT)); + } catch (IllegalArgumentException e) { + throw new YamlDeserializationException( + node, "Invalid value for '" + field + "' in semantic question '" + question + "': " + raw, e); + } + } + private static double number(Map<String, Node> values, String question, String name, double fallback) { if (!values.containsKey(name)) { return fallback; @@ -144,12 +175,12 @@ public class SemanticDefinitionDeserializer extends YamlDeserializerSupport impl } } - private static Map<String, Node> fields(Node node) { + private static Map<String, Node> fields(Node node, String description) { Map<String, Node> result = new LinkedHashMap<>(); for (NodeTuple tuple : asMappingNode(node).getValue()) { String name = asText(tuple.getKeyNode()); if (result.putIfAbsent(name, tuple.getValueNode()) != null) { - throw new IllegalArgumentException("Duplicate semantic declaration key: " + name); + throw new YamlDeserializationException(tuple.getKeyNode(), "Duplicate key '" + name + "' in " + description); } } return result; diff --git a/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticInitializationTest.java b/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticInitializationTest.java new file mode 100644 index 000000000000..8c9745ae5324 --- /dev/null +++ b/components/camel-ai/camel-semantic/src/test/java/org/apache/camel/semantic/SemanticInitializationTest.java @@ -0,0 +1,194 @@ +/* + * 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.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.RuntimeCamelException; +import org.apache.camel.impl.DefaultCamelContext; +import org.apache.camel.language.semantic.SemanticLanguage; +import org.apache.camel.support.service.ServiceSupport; +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +class SemanticInitializationTest { + @Test + void initializationDoesNotAcquirePublicContextOrRegistryMonitors() throws Exception { + ExecutorService callers = Executors.newSingleThreadExecutor(); + try (var context = new DefaultCamelContext()) { + context.start(); + synchronized (context) { + var questions = callers.submit(() -> SemanticQuestions.get(context)); + assertThat(questions.get(10, TimeUnit.SECONDS)).isSameAs(SemanticQuestions.get(context)); + } + var language = language(context, Adapter.class); + synchronized (context.getRegistry()) { + var expression = callers.submit(() -> language.createExpression("ref:q")); + assertThat(expression.get(10, TimeUnit.SECONDS)).isNotNull(); + } + } finally { + callers.shutdownNow(); + } + } + + @Test + void concurrentQuestionRegistryCreationReturnsOneInstance() throws Exception { + ExecutorService callers = Executors.newFixedThreadPool(4); + CountDownLatch start = new CountDownLatch(1); + try (var context = new DefaultCamelContext()) { + List<Future<SemanticQuestions>> results = new ArrayList<>(); + for (int i = 0; i < 8; i++) { + results.add(callers.submit(() -> { + assertThat(start.await(10, TimeUnit.SECONDS)).isTrue(); + return SemanticQuestions.get(context); + })); + } + start.countDown(); + var expected = results.get(0).get(10, TimeUnit.SECONDS); + for (var result : results) { + assertThat(result.get(10, TimeUnit.SECONDS)).isSameAs(expected); + } + } finally { + start.countDown(); + callers.shutdownNow(); + } + } + + @Test + void concurrentLanguagesCannotOverwriteTheAdapterOwner() throws Exception { + ExecutorService callers = Executors.newFixedThreadPool(2); + CyclicBarrier start = new CyclicBarrier(2); + CountingAdapter.constructed.set(0); + CountingAdapter.started.set(0); + CountingAdapter.stopped.set(0); + try (var context = new DefaultCamelContext()) { + context.start(); + var first = language(context, CountingAdapter.class); + var second = language(context, CountingAdapter.class); + List<Future<String>> results = new ArrayList<>(); + for (var language : List.of(first, second)) { + results.add(callers.submit(() -> { + start.await(10, TimeUnit.SECONDS); + try { + language.createExpression("ref:q"); + return "created"; + } catch (RuntimeCamelException e) { + assertThat(e).hasRootCauseInstanceOf(IllegalArgumentException.class) + .hasRootCauseMessage( + "Semantic adapter registry name is already bound: " + SemanticLanguage.ADAPTER_NAME); + return "collision"; + } + })); + } + assertThat(List.of(results.get(0).get(10, TimeUnit.SECONDS), results.get(1).get(10, TimeUnit.SECONDS))) + .containsExactlyInAnyOrder("created", "collision"); + assertThat(CountingAdapter.constructed).hasValue(1); + assertThat(CountingAdapter.started).hasValue(1); + assertThat(context.getRegistry().lookupByName(SemanticLanguage.ADAPTER_NAME)).isInstanceOf(CountingAdapter.class); + } finally { + callers.shutdownNow(); + } + assertThat(CountingAdapter.stopped).hasValue(1); + } + + @Test + void blockedConstructorDoesNotBlockAnotherContext() throws Exception { + ExecutorService callers = Executors.newFixedThreadPool(2); + BlockingAdapter.entered = new CountDownLatch(1); + BlockingAdapter.release = new CountDownLatch(1); + try (var first = new DefaultCamelContext(); var second = new DefaultCamelContext()) { + first.start(); + second.start(); + var blocked = language(first, BlockingAdapter.class); + var independent = language(second, Adapter.class); + Future<?> creation = callers.submit(() -> blocked.createExpression("ref:q")); + try { + assertThat(BlockingAdapter.entered.await(10, TimeUnit.SECONDS)).isTrue(); + var other = callers.submit(() -> independent.createExpression("ref:q")); + assertThat(other.get(10, TimeUnit.SECONDS)).isNotNull(); + } finally { + BlockingAdapter.release.countDown(); + creation.get(10, TimeUnit.SECONDS); + } + } finally { + BlockingAdapter.release.countDown(); + callers.shutdownNow(); + } + } + + private static SemanticLanguage language(DefaultCamelContext context, Class<? extends SemanticAdapter> type) { + SemanticQuestions.get(context).replace("test", Map.of("q", + new SemanticQuestion( + SemanticQuestion.Type.BOOLEAN, "Is this valid?", null, null, null, 0.5, 0, + SemanticQuestion.UncertaintyPolicy.FAIL))); + SemanticLanguage language = new SemanticLanguage(); + language.setCamelContext(context); + language.setAdapter(type.getName()); + return language; + } + + public static class Adapter extends ServiceSupport implements SemanticAdapter { + @Override + public void validate(SemanticQuestion question) { + } + + @Override + public SemanticResult evaluate(SemanticQuestion question, Object state) { + return new SemanticResult(true, null, null, null, null); + } + } + + public static class CountingAdapter extends Adapter { + static final AtomicInteger constructed = new AtomicInteger(); + static final AtomicInteger started = new AtomicInteger(); + static final AtomicInteger stopped = new AtomicInteger(); + + public CountingAdapter() { + constructed.incrementAndGet(); + } + + @Override + protected void doStart() { + started.incrementAndGet(); + } + + @Override + protected void doStop() { + stopped.incrementAndGet(); + } + } + + public static class BlockingAdapter extends Adapter { + static CountDownLatch entered; + static CountDownLatch release; + + public BlockingAdapter() throws InterruptedException { + entered.countDown(); + assertThat(release.await(30, TimeUnit.SECONDS)).isTrue(); + } + } +} 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 677931c3dbc5..97031b07a3f9 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 @@ -20,6 +20,7 @@ import java.nio.file.Files; import java.nio.file.Path; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Stream; import org.apache.camel.component.mock.MockEndpoint; import org.apache.camel.dsl.yaml.common.exception.YamlDeserializationException; @@ -36,6 +37,8 @@ import org.apache.camel.support.RouteWatcherReloadStrategy; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; import org.junit.jupiter.params.provider.ValueSource; import static org.assertj.core.api.Assertions.assertThat; @@ -250,4 +253,76 @@ class SemanticQuestionTest extends YamlTestSupport { }); }); } + + @ParameterizedTest + @MethodSource("invalidEnums") + void invalidEnumValuesIdentifyQuestionFieldAndLocation(String field, String value, int line) { + String yaml = """ + - semantic: + question: + spam: + type: boolean + instructions: Is this spam? + """; + yaml = field.equals("type") + ? yaml.replace("type: boolean", "type: " + value) + : yaml + " uncertaintyPolicy: " + value + "\n"; + String source = yaml; + assertThatThrownBy(() -> loadRoutesNoValidate(source)) + .hasMessageContaining("route-0.yaml") + .hasRootCauseInstanceOf(IllegalArgumentException.class) + .cause().isInstanceOfSatisfying(YamlDeserializationException.class, error -> { + assertThat(error).hasMessageContaining("Invalid value for '" + field + "' in semantic question 'spam'"); + assertThat(error.getProblemMark()).hasValueSatisfying(mark -> { + assertThat(mark.getLine()).isEqualTo(line); + assertThat(mark.getColumn()).isEqualTo(8 + field.length() + 2); + }); + }); + } + + static Stream<Arguments> invalidEnums() { + return Stream.of( + Arguments.of("type", "unsupported", 3), + Arguments.of("type", "''", 3), + Arguments.of("uncertaintyPolicy", "unsupported", 5), + Arguments.of("uncertaintyPolicy", "''", 5), + Arguments.of("uncertaintyPolicy", "null", 5)); + } + + @ParameterizedTest + @MethodSource("invalidStructures") + void invalidStructuresIdentifySourceAndOffendingNode(String yaml, String message, int line, int column) { + assertThatThrownBy(() -> loadRoutesNoValidate(yaml)) + .hasMessageContaining("route-0.yaml") + .cause().isInstanceOfSatisfying(YamlDeserializationException.class, error -> { + assertThat(error).hasMessageContaining(message); + assertThat(error.getProblemMark()).hasValueSatisfying(mark -> { + assertThat(mark.getLine()).isEqualTo(line); + assertThat(mark.getColumn()).isEqualTo(column); + }); + }); + } + + static Stream<Arguments> invalidStructures() { + String declaration = declarations("${body}"); + String question = """ + - semantic: + question: + q: {type: boolean, instructions: Is it valid?} + """; + return Stream.of( + Arguments.of(declaration.replace("instructions:", "typo:"), + "Unknown property 'typo' in semantic question 'department'", 5, 14), + Arguments.of(declaration.replace(" type: choice\n", ""), + "Semantic question type is required: department", 3, 8), + Arguments.of(declaration.replace("type: choice", "type: choice\n type: choice"), + "Duplicate key 'type' in semantic question 'department'", 4, 8), + Arguments.of("- semantic: {other: {}}", "Semantic declaration requires only question", 0, 12), + Arguments.of(question + question, "Duplicate semantic question: q", 4, 4), + Arguments.of(question + " q: {type: boolean, instructions: Is it valid?}\n", + "Duplicate key 'q' in semantic questions", 3, 6), + Arguments.of(question.replace("Is it valid?", "''"), + "Invalid semantic question 'q': Question instructions must not be blank", 2, 9)); + } + }
