gnodet-bot commented on code in PR #26768: URL: https://github.com/apache/camel/pull/26768#discussion_r4081785804
########## test-infra/camel-test-infra-typesafe-ai-mock/src/main/java/org/apache/camel/test/infra/typesafeai/mock/TypeSafeAiService.java: ########## @@ -0,0 +1,252 @@ +/* + * 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.test.infra.typesafeai.mock; + +import java.io.IOException; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.function.Function; + +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; +import org.apache.camel.util.json.DeserializationException; +import org.apache.camel.util.json.JsonObject; +import org.apache.camel.util.json.Jsoner; +import org.junit.jupiter.api.extension.AfterEachCallback; +import org.junit.jupiter.api.extension.BeforeEachCallback; +import org.junit.jupiter.api.extension.ExtensionContext; + +/** + * A local System One HTTP mock, or an existing compatible service when its base URL is supplied. Configure + * {@code TYPESAFE_AI_BASE_URL}, {@code TYPESAFE_AI_API_KEY}, and optionally {@code TYPESAFE_AI_MODEL} to run the same + * tests against a real API. The Laya equivalents {@code LAYA_BASE_URL} and {@code LAYA_API_KEY} are also accepted. + * + * @since 4.23 + */ +public class TypeSafeAiService implements BeforeEachCallback, AfterEachCallback { + private final String remoteBaseUrl; + private final String apiKey; + private final String model; + private final List<JsonObject> requests = new CopyOnWriteArrayList<>(); + private volatile Function<JsonObject, Map<String, ?>> responder = TypeSafeAiService::defaultAnswers; + private HttpServer server; + private ExecutorService executor; + + /** + * Creates a service using environment configuration, or a local mock when absent. + * + * @since 4.23 + */ + public TypeSafeAiService() { + this(firstConfigured("TYPESAFE_AI_BASE_URL", "LAYA_BASE_URL"), + firstConfigured("TYPESAFE_AI_API_KEY", "LAYA_API_KEY"), + firstConfigured("TYPESAFE_AI_MODEL", "LAYA_MODEL")); + } + + /** + * Explicit configuration is useful when a test starts its own compatible server. + * + * @since 4.23 + */ + public TypeSafeAiService(String baseUrl, String apiKey, String model) { + remoteBaseUrl = baseUrl; + if (baseUrl != null && !baseUrl.isBlank() && (apiKey == null || apiKey.isBlank())) { + throw new IllegalArgumentException("TYPESAFE_AI_API_KEY is required with TYPESAFE_AI_BASE_URL"); + } + this.apiKey = apiKey == null || apiKey.isBlank() ? "test-key" : apiKey; + this.model = model == null || model.isBlank() ? "jev-latest" : model; + } + + /** + * Returns the URL used by the component. + * + * @since 4.23 + */ + public String getBaseUrl() { + if (isRemote()) { + return remoteBaseUrl; + } + if (server == null) { + throw new IllegalStateException("TypeSafe AI mock has not started"); + } + return "http://127.0.0.1:" + server.getAddress().getPort(); + } + + /** + * Returns the configured API key or the local mock's test key. + * + * @since 4.23 + */ + public String getApiKey() { + return apiKey; + } + + /** + * Returns the configured model name. + * + * @since 4.23 + */ + public String getModel() { + return model; + } + + /** + * Whether an existing API is used instead of a local mock. + * + * @since 4.23 + */ + public boolean isRemote() { + return remoteBaseUrl != null && !remoteBaseUrl.isBlank(); + } + + /** + * Returns requests received by the local mock. A remote service is not recorded. + * + * @since 4.23 + */ + public List<JsonObject> getRequests() { + return List.copyOf(requests); + } + + /** + * Supplies answer objects keyed by the question names in the incoming request. + * + * @since 4.23 + */ + public void setResponder(Function<JsonObject, Map<String, ?>> responder) { + this.responder = Objects.requireNonNull(responder); + } + + @Override + public void beforeEach(ExtensionContext context) throws IOException { + requests.clear(); + if (isRemote()) { + return; + } + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/v1/systemone", this::handle); + executor = Executors.newCachedThreadPool(); + server.setExecutor(executor); + server.start(); + } + + @Override + public void afterEach(ExtensionContext context) { + if (server != null) { + server.stop(0); + server = null; + } + if (executor != null) { + executor.shutdownNow(); + executor = null; + } + } + + private void handle(HttpExchange exchange) throws IOException { + try { + if (!"POST".equals(exchange.getRequestMethod())) { + send(exchange, 405, "{}"); + return; + } + if (!("Bearer " + apiKey).equals(exchange.getRequestHeaders().getFirst("Authorization"))) { + send(exchange, 401, "{}"); + return; + } + Object parsed; + try { + parsed = Jsoner.deserialize(new String(exchange.getRequestBody().readAllBytes(), StandardCharsets.UTF_8)); + } catch (DeserializationException e) { + send(exchange, 400, "{}"); + return; + } + if (!(parsed instanceof JsonObject request) || !(request.get("questions") instanceof Map<?, ?> questions) + || questions.isEmpty() || !(request.get("model") instanceof String)) { + send(exchange, 400, "{}"); + return; + } + requests.add(request); + Map<String, ?> answers; + try { + answers = responder.apply(request); + } catch (RuntimeException e) { + send(exchange, 400, "{}"); Review Comment: ⚠️ **Silent exception swallow turns test bugs into cryptic 400s.** When a custom responder (set via `setResponder`) throws a `RuntimeException`, the exception is silently discarded and a 400 is returned. From the test's perspective this surfaces as a confusing assertion failure on the response body — not the actual bug in the responder. At minimum, log the exception before sending 400; better still, propagate it so the test fails with the real cause: ```suggestion } catch (RuntimeException e) { LOG.error("Responder threw: {}", e.getMessage(), e); send(exchange, 500, "{}"); return; } ``` Using 500 also distinguishes responder bugs from malformed requests (400) — easier to grep in test output. ########## test-infra/camel-test-infra-typesafe-ai-mock/src/main/java/org/apache/camel/test/infra/typesafeai/mock/TypeSafeAiService.java: ########## @@ -0,0 +1,252 @@ +/* + * 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.test.infra.typesafeai.mock; + +import java.io.IOException; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.function.Function; + +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; +import org.apache.camel.util.json.DeserializationException; +import org.apache.camel.util.json.JsonObject; +import org.apache.camel.util.json.Jsoner; +import org.junit.jupiter.api.extension.AfterEachCallback; +import org.junit.jupiter.api.extension.BeforeEachCallback; +import org.junit.jupiter.api.extension.ExtensionContext; + +/** + * A local System One HTTP mock, or an existing compatible service when its base URL is supplied. Configure + * {@code TYPESAFE_AI_BASE_URL}, {@code TYPESAFE_AI_API_KEY}, and optionally {@code TYPESAFE_AI_MODEL} to run the same + * tests against a real API. The Laya equivalents {@code LAYA_BASE_URL} and {@code LAYA_API_KEY} are also accepted. + * + * @since 4.23 + */ +public class TypeSafeAiService implements BeforeEachCallback, AfterEachCallback { + private final String remoteBaseUrl; + private final String apiKey; + private final String model; + private final List<JsonObject> requests = new CopyOnWriteArrayList<>(); + private volatile Function<JsonObject, Map<String, ?>> responder = TypeSafeAiService::defaultAnswers; + private HttpServer server; + private ExecutorService executor; + + /** + * Creates a service using environment configuration, or a local mock when absent. + * + * @since 4.23 + */ + public TypeSafeAiService() { + this(firstConfigured("TYPESAFE_AI_BASE_URL", "LAYA_BASE_URL"), + firstConfigured("TYPESAFE_AI_API_KEY", "LAYA_API_KEY"), + firstConfigured("TYPESAFE_AI_MODEL", "LAYA_MODEL")); + } + + /** + * Explicit configuration is useful when a test starts its own compatible server. + * + * @since 4.23 + */ + public TypeSafeAiService(String baseUrl, String apiKey, String model) { + remoteBaseUrl = baseUrl; + if (baseUrl != null && !baseUrl.isBlank() && (apiKey == null || apiKey.isBlank())) { + throw new IllegalArgumentException("TYPESAFE_AI_API_KEY is required with TYPESAFE_AI_BASE_URL"); + } + this.apiKey = apiKey == null || apiKey.isBlank() ? "test-key" : apiKey; + this.model = model == null || model.isBlank() ? "jev-latest" : model; + } + + /** + * Returns the URL used by the component. + * + * @since 4.23 + */ + public String getBaseUrl() { + if (isRemote()) { + return remoteBaseUrl; + } + if (server == null) { + throw new IllegalStateException("TypeSafe AI mock has not started"); + } + return "http://127.0.0.1:" + server.getAddress().getPort(); + } + + /** + * Returns the configured API key or the local mock's test key. + * + * @since 4.23 + */ + public String getApiKey() { + return apiKey; + } + + /** + * Returns the configured model name. + * + * @since 4.23 + */ + public String getModel() { + return model; + } + + /** + * Whether an existing API is used instead of a local mock. + * + * @since 4.23 + */ + public boolean isRemote() { + return remoteBaseUrl != null && !remoteBaseUrl.isBlank(); + } + + /** + * Returns requests received by the local mock. A remote service is not recorded. + * + * @since 4.23 + */ + public List<JsonObject> getRequests() { + return List.copyOf(requests); + } + + /** + * Supplies answer objects keyed by the question names in the incoming request. + * + * @since 4.23 + */ + public void setResponder(Function<JsonObject, Map<String, ?>> responder) { + this.responder = Objects.requireNonNull(responder); + } + + @Override + public void beforeEach(ExtensionContext context) throws IOException { + requests.clear(); + if (isRemote()) { + return; + } + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/v1/systemone", this::handle); + executor = Executors.newCachedThreadPool(); + server.setExecutor(executor); + server.start(); + } + + @Override + public void afterEach(ExtensionContext context) { + if (server != null) { + server.stop(0); + server = null; + } + if (executor != null) { + executor.shutdownNow(); + executor = null; + } + } Review Comment: ⚠️ **`afterEach` race: zero-grace server stop + no `awaitTermination`.** `server.stop(0)` tells the server to stop with 0 seconds of grace — any in-flight handler thread continues running after `stop` returns because `stop` only waits for the exchange handler, not for the executor's thread pool. Calling `shutdownNow()` without `awaitTermination` means the thread pool threads may still be active when the next test's `beforeEach` starts a new server on the same loopback port. In a fast test suite this can cause the stale thread to accidentally respond to the next test's requests. ```suggestion public void afterEach(ExtensionContext context) { if (server != null) { server.stop(1); server = null; } if (executor != null) { executor.shutdownNow(); try { executor.awaitTermination(5, TimeUnit.SECONDS); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } executor = null; } } ``` (Add `import java.util.concurrent.TimeUnit;`) ########## components/camel-ai/camel-typesafe-ai/src/main/java/org/apache/camel/component/typesafeai/TypeSafeAiEndpoint.java: ########## @@ -74,6 +78,14 @@ protected void doStart() throws Exception { if (configuration.getQuestions() != null) { configuredQuestions = TypeSafeAiJson.questions(configuration.getQuestions()); stateExpression = createStateExpression(); + } else if (configuration.getQuestionsResource() != null) { + String location = configuration.getQuestionsResource(); + try (InputStream input = ResourceHelper.resolveMandatoryResourceAsInputStream(getCamelContext(), location)) { + configuredQuestions = TypeSafeAiJson.questions(new String(input.readAllBytes(), StandardCharsets.UTF_8)); Review Comment: ⚠️ **`readAllBytes()` has no size cap — OOM risk on large or adversarial file resources.** A `file:` URI pointing at a large file (misconfiguration, accidental symlink to `/dev/urandom`, path traversal) will attempt to allocate the entire content into heap at endpoint startup with no guard. Add a bounded read before parsing: ```suggestion byte[] bytes = input.readNBytes(4 * 1024 * 1024); // 4 MB cap if (input.read() != -1) { throw new IOException("questionsResource exceeds 4 MB limit"); } configuredQuestions = TypeSafeAiJson.questions(new String(bytes, StandardCharsets.UTF_8)); ``` 4 MB is generous for a question set; adjust if the API supports genuinely larger files. ########## components/camel-ai/camel-typesafe-ai/src/main/java/org/apache/camel/component/typesafeai/TypeSafeAiConfiguration.java: ########## @@ -91,6 +93,23 @@ public void setQuestions(String questions) { this.questions = questions; } + /** + * @since 4.23 + */ + public String getQuestionsResource() { + return questionsResource; + } Review Comment: 🔹 **Missing `@return` on `getQuestionsResource()`.** The setter has a full Javadoc description; the getter only has `@since 4.23`. For a new public API method, add `@return`: ```suggestion /** * Returns the Camel resource URI for the JSON questions file, or {@code null} if not set. * * @return the questions resource URI, or {@code null} * @since 4.23 */ public String getQuestionsResource() { return questionsResource; } ``` -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
