diff --git a/shared/src/main/kotlin/io/askimo/core/providers/OpenAiCompatibleChatModelFactory.kt b/shared/src/main/kotlin/io/askimo/core/providers/OpenAiCompatibleChatModelFactory.kt index b498f2d19..713d327be 100644 --- a/shared/src/main/kotlin/io/askimo/core/providers/OpenAiCompatibleChatModelFactory.kt +++ b/shared/src/main/kotlin/io/askimo/core/providers/OpenAiCompatibleChatModelFactory.kt @@ -23,6 +23,7 @@ import io.askimo.core.context.AppContext import io.askimo.core.context.ExecutionMode import io.askimo.core.telemetry.TelemetryChatModelListener import io.askimo.core.util.ProxyUtil +import io.askimo.core.util.withLoggingIfDebug import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob @@ -155,7 +156,7 @@ abstract class OpenAiCompatibleChatModelFactory : ChatModelFactory ProxyUtil.configureProxy( HttpClient.newBuilder().version(httpVersion()), baseUrl, - ), + ).withLoggingIfDebug(), ).readTimeout(Duration.ofSeconds(AppConfig.models.timeouts.defaultModelTimeoutSeconds)) .connectTimeout(Duration.ofSeconds(AppConfig.models.timeouts.defaultModelTimeoutSeconds)) @@ -247,6 +248,7 @@ abstract class OpenAiCompatibleChatModelFactory : ChatModelFactory val reasoningLevel = ModelCapabilitiesCache.getReasoningLevel(getProvider(), settings.defaultModel) if (supportsThinking && reasoningLevel.isEnabled) { reasoningEffort(reasoningLevel.value) + reasoningSummary("detailed") } } .strictTools(true) @@ -290,9 +292,6 @@ abstract class OpenAiCompatibleChatModelFactory : ChatModelFactory .baseUrl(settings.baseUrl) .apiKey(resolveApiKey(settings)) .modelName(settings.imageModel.ifBlank { AppConfig.models[getProvider()].imageModel }) - .logger(log) - .logRequests(log.isDebugEnabled) - .logResponses(log.isTraceEnabled) .build() override fun createUtilityClient(settings: T): ChatClient = AiServices.builder(ChatClient::class.java) diff --git a/shared/src/main/kotlin/io/askimo/core/providers/anthropic/AnthropicModelFactory.kt b/shared/src/main/kotlin/io/askimo/core/providers/anthropic/AnthropicModelFactory.kt index e023c4697..e546e8ad5 100644 --- a/shared/src/main/kotlin/io/askimo/core/providers/anthropic/AnthropicModelFactory.kt +++ b/shared/src/main/kotlin/io/askimo/core/providers/anthropic/AnthropicModelFactory.kt @@ -34,6 +34,7 @@ import io.askimo.core.util.ApiKeyUtils.safeApiKey import io.askimo.core.util.ProxyUtil import io.askimo.core.util.appJson import io.askimo.core.util.httpGet +import io.askimo.core.util.withLoggingIfDebug import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob @@ -71,7 +72,7 @@ class AnthropicModelFactory : ChatModelFactory { chatMemory: ChatMemory?, ): ChatClient { // Configure HTTP client for thinking probe (probe needs its own builder) - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) // Probe thinking support once — result is persisted in ModelCapabilitiesCache @@ -165,7 +166,7 @@ class AnthropicModelFactory : ChatModelFactory { } override fun createStreamingModel(settings: AnthropicSettings): StreamingChatModel { - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) val telemetry = AppContext.getInstance().telemetry @@ -210,7 +211,7 @@ class AnthropicModelFactory : ChatModelFactory { } override fun createSecondaryModel(settings: AnthropicSettings): ChatModel { - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) return AnthropicChatModel.builder() .httpClientBuilder(jdkHttpClientBuilder) @@ -226,7 +227,7 @@ class AnthropicModelFactory : ChatModelFactory { } override fun createModel(settings: AnthropicSettings): ChatModel { - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) return AnthropicChatModel.builder() diff --git a/shared/src/main/kotlin/io/askimo/core/providers/gemini/GeminiModelFactory.kt b/shared/src/main/kotlin/io/askimo/core/providers/gemini/GeminiModelFactory.kt index b928ec491..add124569 100644 --- a/shared/src/main/kotlin/io/askimo/core/providers/gemini/GeminiModelFactory.kt +++ b/shared/src/main/kotlin/io/askimo/core/providers/gemini/GeminiModelFactory.kt @@ -39,6 +39,7 @@ import io.askimo.core.providers.sendStreamingMessageWithCallback import io.askimo.core.telemetry.TelemetryChatModelListener import io.askimo.core.util.ApiKeyUtils.safeApiKey import io.askimo.core.util.ProxyUtil +import io.askimo.core.util.withLoggingIfDebug import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob @@ -69,8 +70,7 @@ class GeminiModelFactory : ChatModelFactory { executionMode: ExecutionMode, chatMemory: ChatMemory?, ): ChatClient { - // Configure HTTP client for thinking probe (probe needs its own builder) - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) // Probe thinking support once — result is persisted in ModelCapabilitiesCache @@ -165,13 +165,10 @@ class GeminiModelFactory : ChatModelFactory { .apiKey(safeApiKey(settings.apiKey)) .baseUrl(settings.baseUrl) .modelName(settings.imageModel.ifBlank { AppConfig.models[GEMINI].imageModel }) - .logger(log) - .logRequests(log.isDebugEnabled) - .logResponses(log.isTraceEnabled) .build() override fun createStreamingModel(settings: GeminiSettings): StreamingChatModel { - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) val telemetry = AppContext.getInstance().telemetry @@ -194,6 +191,7 @@ class GeminiModelFactory : ChatModelFactory { thinkingConfig( GeminiThinkingConfig.builder() .thinkingLevel(geminiLevel) + .includeThoughts(true) .build(), ) sendThinking(true) @@ -204,7 +202,7 @@ class GeminiModelFactory : ChatModelFactory { } override fun createSecondaryModel(settings: GeminiSettings): ChatModel { - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) return GoogleAiGeminiChatModel.builder() .httpClientBuilder(jdkHttpClientBuilder) @@ -220,7 +218,7 @@ class GeminiModelFactory : ChatModelFactory { } override fun createModel(settings: GeminiSettings): ChatModel { - val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()) + val httpClientBuilder = ProxyUtil.configureProxy(HttpClient.newBuilder()).withLoggingIfDebug() val jdkHttpClientBuilder = JdkHttpClient.builder().httpClientBuilder(httpClientBuilder) return GoogleAiGeminiChatModel.builder() diff --git a/shared/src/main/kotlin/io/askimo/core/telemetry/TelemetryChatModelListener.kt b/shared/src/main/kotlin/io/askimo/core/telemetry/TelemetryChatModelListener.kt index 782632fe4..78948f9f9 100644 --- a/shared/src/main/kotlin/io/askimo/core/telemetry/TelemetryChatModelListener.kt +++ b/shared/src/main/kotlin/io/askimo/core/telemetry/TelemetryChatModelListener.kt @@ -4,7 +4,6 @@ */ package io.askimo.core.telemetry -import dev.langchain4j.data.message.SystemMessage import dev.langchain4j.model.chat.listener.ChatModelErrorContext import dev.langchain4j.model.chat.listener.ChatModelListener import dev.langchain4j.model.chat.listener.ChatModelRequestContext @@ -12,7 +11,6 @@ import dev.langchain4j.model.chat.listener.ChatModelResponseContext import io.askimo.core.analytics.Analytics import io.askimo.core.analytics.AnalyticsEvent import io.askimo.core.logging.logger -import io.askimo.core.memory.getTextContent import io.askimo.core.util.MachineId import java.net.URI import java.net.http.HttpClient @@ -97,57 +95,6 @@ class TelemetryChatModelListener( } attrs[ATTR_START_TIME] = System.currentTimeMillis() - - val request = context.chatRequest() - if (log.isDebugEnabled) { - val messages = request.messages() - val systemMessages = messages.filterIsInstance() - val nonSystemMessages = messages.filter { it !is SystemMessage } - val lastMessages = nonSystemMessages.takeLast(3) - val params = request.parameters() - - val debugMsg = buildString { - append("LLM request to $provider: ${messages.size} messages, ") - append("model=${request.modelName() ?: "default"}, ") - appendLine("clientId=$clientId, correlationId=$newCorrelationId") - - if (systemMessages.isNotEmpty()) { - appendLine(" [System messages (${systemMessages.size})]") - systemMessages.forEach { msg -> - val preview = msg.getTextContent().take(300).replace('\n', ' ') - appendLine(" SYSTEM: $preview") - } - } - - if (lastMessages.isNotEmpty()) { - appendLine(" [Last ${lastMessages.size} message(s)]") - lastMessages.forEach { msg -> - val role = msg.type().name - val preview = msg.getTextContent().take(300).replace('\n', ' ') - appendLine(" $role: $preview") - } - } - - val paramParts = buildList { - params.modelName()?.let { add("modelName=$it") } - params.temperature()?.let { add("temperature=$it") } - params.topP()?.let { add("topP=$it") } - params.topK()?.let { add("topK=$it") } - params.frequencyPenalty()?.let { add("frequencyPenalty=$it") } - params.presencePenalty()?.let { add("presencePenalty=$it") } - params.maxOutputTokens()?.let { add("maxOutputTokens=$it") } - params.stopSequences()?.takeIf { it.isNotEmpty() }?.let { add("stopSequences=$it") } - params.responseFormat()?.let { add("responseFormat=$it") } - params.toolChoice()?.let { add("toolChoice=$it") } - params.toolSpecifications()?.takeIf { it.isNotEmpty() }?.let { tools -> add("tools(${tools.size})=${tools.map { it.name() }}") } - } - if (paramParts.isNotEmpty()) { - append(" [Parameters] ") - append(paramParts.joinToString(", ")) - } - } - log.debug(debugMsg.trimEnd()) - } } override fun onResponse(context: ChatModelResponseContext) { diff --git a/shared/src/main/kotlin/io/askimo/core/util/LoggingHttpClientBuilder.kt b/shared/src/main/kotlin/io/askimo/core/util/LoggingHttpClientBuilder.kt new file mode 100644 index 000000000..023c5f5eb --- /dev/null +++ b/shared/src/main/kotlin/io/askimo/core/util/LoggingHttpClientBuilder.kt @@ -0,0 +1,189 @@ +/* SPDX-License-Identifier: AGPLv3 + * + * Copyright (c) 2026 Askimo + */ +package io.askimo.core.util + +import io.askimo.core.logging.currentFileLogger +import java.io.ByteArrayOutputStream +import java.net.Authenticator +import java.net.CookieHandler +import java.net.ProxySelector +import java.net.http.HttpClient +import java.net.http.HttpRequest +import java.net.http.HttpResponse +import java.net.http.WebSocket +import java.nio.ByteBuffer +import java.time.Duration +import java.util.Optional +import java.util.concurrent.CompletableFuture +import java.util.concurrent.Executor +import java.util.concurrent.Flow +import java.util.concurrent.TimeUnit +import javax.net.ssl.SSLContext +import javax.net.ssl.SSLParameters + +private val log = currentFileLogger() + +/** Header names whose values must be masked in log output. */ +private val SENSITIVE_HEADERS = setOf("authorization", "x-api-key", "api-key", "x-goog-api-key") + +/** + * Wraps this builder in a [LoggingHttpClientBuilder] only when DEBUG logging is enabled. + * When DEBUG is off the original builder is returned unchanged — zero overhead, no wrapper allocated. + */ +fun HttpClient.Builder.withLoggingIfDebug(): HttpClient.Builder = if (log.isDebugEnabled) LoggingHttpClientBuilder(this) else this + +/** + * A [HttpClient.Builder] decorator that wraps the built [HttpClient] in a [LoggingHttpClient]. + * Prefer constructing via [withLoggingIfDebug] so the wrapper is skipped entirely when DEBUG is off. + * + * Enable DEBUG logging for `io.askimo.core.util.LoggingHttpClientBuilder` to activate: + * ```xml + * + * ``` + */ +@Suppress("JAVA_DEFAULT_METHODS_NOT_OVERRIDDEN_BY_DELEGATION") +class LoggingHttpClientBuilder( + private val delegate: HttpClient.Builder, +) : HttpClient.Builder by delegate { + override fun build(): HttpClient = LoggingHttpClient(delegate.build()) +} + +/** + * A [HttpClient] wrapper that logs every outgoing HTTP request (URI, method, headers, body) + * at DEBUG / TRACE level before delegating to the real client. + * + * - **Headers**: sensitive values (`Authorization`, `x-api-key`, etc.) are masked via [Masking]. + * - **Body**: the `BodyPublisher` is drained into a buffer, logged, then rebuilt as a fresh + * [HttpRequest.BodyPublishers.ofByteArray] publisher so the actual HTTP call still carries + * its full payload intact. Body text is only printed at TRACE level; DEBUG shows byte-count. + */ +class LoggingHttpClient( + private val delegate: HttpClient, +) : HttpClient() { + + override fun send( + request: HttpRequest, + responseBodyHandler: HttpResponse.BodyHandler, + ): HttpResponse { + val (loggable, body) = interceptRequest(request) + logRequest(loggable, body) + return delegate.send(loggable, responseBodyHandler) + } + + override fun sendAsync( + request: HttpRequest, + responseBodyHandler: HttpResponse.BodyHandler, + ): CompletableFuture> { + val (loggable, body) = interceptRequest(request) + logRequest(loggable, body) + return delegate.sendAsync(loggable, responseBodyHandler) + } + + override fun sendAsync( + request: HttpRequest, + responseBodyHandler: HttpResponse.BodyHandler, + pushPromiseHandler: HttpResponse.PushPromiseHandler?, + ): CompletableFuture> { + val (loggable, body) = interceptRequest(request) + logRequest(loggable, body) + return delegate.sendAsync(loggable, responseBodyHandler, pushPromiseHandler) + } + + // ── HttpClient delegation boilerplate ────────────────────────────────────── + + override fun cookieHandler(): Optional = delegate.cookieHandler() + override fun connectTimeout(): Optional = delegate.connectTimeout() + override fun followRedirects(): Redirect = delegate.followRedirects() + override fun proxy(): Optional = delegate.proxy() + override fun sslContext(): SSLContext = delegate.sslContext() + override fun sslParameters(): SSLParameters = delegate.sslParameters() + override fun authenticator(): Optional = delegate.authenticator() + override fun version(): Version = delegate.version() + override fun executor(): Optional = delegate.executor() + override fun newWebSocketBuilder(): WebSocket.Builder = delegate.newWebSocketBuilder() + + // ── Private helpers ──────────────────────────────────────────────────────── + + /** + * Short-circuits when DEBUG is off (zero overhead). + * Otherwise drains the [HttpRequest.BodyPublisher], captures the bytes, and returns a + * rebuilt request with a fresh [HttpRequest.BodyPublishers.ofByteArray] publisher so + * the actual HTTP send still carries the full payload. + */ + private fun interceptRequest(request: HttpRequest): Pair { + // No isDebugEnabled guard needed — this class is only instantiated when debug is on. + val publisher = request.bodyPublisher().orElse(null) + ?: return request to null + if (publisher.contentLength() == 0L) return request to null + + val bodyBytes = drainPublisher(publisher) + ?: return request to "[unreadable body — ${publisher.contentLength()} bytes]" + + // Rebuild with a fresh publisher so the real HTTP send still carries its body. + val rebuilt = HttpRequest.newBuilder(request.uri()) + .method(request.method(), HttpRequest.BodyPublishers.ofByteArray(bodyBytes)) + .apply { + request.timeout().ifPresent { timeout(it) } + request.version().ifPresent { version(it) } + request.headers().map().forEach { (name, values) -> + values.forEach { value -> header(name, value) } + } + } + .build() + + val text = bodyBytes.toString(Charsets.UTF_8) + val bodyText = if (text.length > 4096) text.take(4096) + "\n…[truncated at 4 096 chars]" else text + + return rebuilt to bodyText + } + + private fun logRequest(request: HttpRequest, body: String?) { + // No isDebugEnabled guard needed — this class is only instantiated when debug is on. + val sb = StringBuilder() + sb.appendLine("──── Outgoing HTTP Request ────────────────────────────────────────────") + sb.appendLine("${request.method()} ${request.uri()}") + sb.appendLine() + sb.appendLine("Headers:") + request.headers().map().forEach { (name, values) -> + val display = if (name.lowercase() in SENSITIVE_HEADERS) { + values.map { Masking.maskSecret(it) } + } else { + values + } + sb.appendLine(" $name: ${display.joinToString(", ")}") + } + if (body != null) { + sb.appendLine() + sb.appendLine("Body:") + sb.append(" ") + sb.appendLine(body.replace("\n", "\n ")) + } + sb.append("───────────────────────────────────────────────────────────────────────") + log.debug(sb.toString()) + } + + /** + * Synchronously drains a [HttpRequest.BodyPublisher] into a [ByteArray]. + * Returns null if completion does not arrive within 2 s or an error is signalled. + */ + private fun drainPublisher(publisher: HttpRequest.BodyPublisher): ByteArray? = try { + val future = CompletableFuture() + val baos = ByteArrayOutputStream() + publisher.subscribe(object : Flow.Subscriber { + override fun onSubscribe(subscription: Flow.Subscription) = subscription.request(Long.MAX_VALUE) + override fun onNext(item: ByteBuffer) { + val bytes = ByteArray(item.remaining()) + item.get(bytes) + baos.write(bytes) + } + override fun onError(throwable: Throwable) = future.completeExceptionally(throwable).let {} + override fun onComplete() = future.complete(baos.toByteArray()).let {} + }) + future.get(2, TimeUnit.SECONDS) + } catch (e: Exception) { + log.trace("Could not drain body publisher for logging: {}", e.message) + null + } +}