diff --git a/src/main/java/com/xly/config/TracingChatModelListener.java b/src/main/java/com/xly/config/TracingChatModelListener.java
index 436fcc0..5c6c869 100644
--- a/src/main/java/com/xly/config/TracingChatModelListener.java
+++ b/src/main/java/com/xly/config/TracingChatModelListener.java
@@ -1,5 +1,6 @@
package com.xly.config;
+import com.fasterxml.jackson.databind.ObjectMapper;
import dev.langchain4j.model.chat.listener.ChatModelErrorContext;
import dev.langchain4j.model.chat.listener.ChatModelListener;
import dev.langchain4j.model.chat.listener.ChatModelRequestContext;
@@ -7,44 +8,147 @@ import dev.langchain4j.model.chat.listener.ChatModelResponseContext;
import dev.langchain4j.model.output.TokenUsage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
+import java.net.URI;
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.time.Instant;
+import java.util.Base64;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+
/**
- * LLM 调用观测(轻量 tracing)—— 每次模型调用记录耗时 / token 用量 / 错误到 `com.xly.trace.llm` 日志。
+ * LLM 调用观测(tracing)—— 每次模型调用记录耗时 / token 用量 / 错误。
*
- *
这是 **Langfuse 的接入点**:Langfuse 本体是需自托管的观测服务(需要实例 + 密钥),此处的
- * onResponse/onError 就是把 span 转发给 Langfuse 的挂钩位;在没有实例的环境下先落到日志,保证
- * 「LLM 可观测」这一能力有实现、可随时对接 Langfuse。业务审计另见 `ai_audit_log`(与 LLM tracing 分离)。
+ *
默认落 {@code com.xly.trace.llm} 日志(保证「LLM 可观测」这一能力始终有实现)。当配置
+ * {@code langfuse.enabled=true} 且给了 host + 公私钥时,额外把一条 generation span **转发到自托管 Langfuse**
+ * 的 ingestion API(不引第三方依赖,直接用 JDK HttpClient,best-effort、异步、失败不影响主流程)。
+ * 自托管方式见 {@code docker-compose.langfuse.yml}。业务审计另见 {@code ai_audit_log}(与 LLM tracing 分离)。
*/
@Component
public class TracingChatModelListener implements ChatModelListener {
private static final Logger log = LoggerFactory.getLogger("com.xly.trace.llm");
+ private final ObjectMapper mapper;
+ private final HttpClient http = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(5)).build();
+
+ @Value("${langfuse.enabled:false}")
+ private boolean langfuseEnabled;
+ @Value("${langfuse.host:http://localhost:3000}")
+ private String langfuseHost;
+ @Value("${langfuse.public-key:}")
+ private String publicKey;
+ @Value("${langfuse.secret-key:}")
+ private String secretKey;
+ @Value("${langchain4j.ollama.chat-model-name:unknown}")
+ private String modelName;
+
+ public TracingChatModelListener(ObjectMapper mapper) {
+ this.mapper = mapper;
+ }
+
@Override
public void onRequest(ChatModelRequestContext ctx) {
ctx.attributes().put("t0", System.nanoTime());
+ ctx.attributes().put("startTs", Instant.now().toString());
}
@Override
public void onResponse(ChatModelResponseContext ctx) {
long ms = elapsedMs(ctx.attributes().get("t0"));
- String tok = "?";
+ Integer in = null;
+ Integer out = null;
try {
TokenUsage u = ctx.chatResponse() == null ? null : ctx.chatResponse().tokenUsage();
if (u != null) {
- tok = u.inputTokenCount() + "/" + u.outputTokenCount();
+ in = u.inputTokenCount();
+ out = u.outputTokenCount();
}
} catch (Exception ignore) {
}
- log.info("LLM ok {}ms tokens(in/out)={}", ms, tok);
- // Langfuse 接入点:此处可 forward 一个 span(model, prompt, completion, latency, tokens)。
+ log.info("LLM ok {}ms tokens(in/out)={}/{}", ms, in, out);
+ exportToLangfuse(String.valueOf(ctx.attributes().get("startTs")), in, out, null);
}
@Override
public void onError(ChatModelErrorContext ctx) {
Throwable e = ctx.error();
- log.warn("LLM error: {}", e == null ? "?" : e.getMessage());
+ String msg = e == null ? "?" : e.getMessage();
+ log.warn("LLM error: {}", msg);
+ exportToLangfuse(String.valueOf(ctx.attributes().get("startTs")), null, null, msg);
+ }
+
+ /** 把一条 generation span 转发到 Langfuse(best-effort,异步,失败仅告警)。未启用则直接返回。 */
+ private void exportToLangfuse(String startTs, Integer in, Integer out, String error) {
+ if (!langfuseEnabled || publicKey == null || publicKey.isBlank() || secretKey == null || secretKey.isBlank()) {
+ return;
+ }
+ try {
+ String traceId = UUID.randomUUID().toString();
+ String now = Instant.now().toString();
+ String start = (startTs == null || "null".equals(startTs)) ? now : startTs;
+
+ Map usage = new LinkedHashMap<>();
+ usage.put("input", in);
+ usage.put("output", out);
+ usage.put("unit", "TOKENS");
+
+ Map genBody = new LinkedHashMap<>();
+ genBody.put("id", UUID.randomUUID().toString());
+ genBody.put("traceId", traceId);
+ genBody.put("name", "xlyAi-agent");
+ genBody.put("model", modelName);
+ genBody.put("startTime", start);
+ genBody.put("endTime", now);
+ genBody.put("usage", usage);
+ if (error != null) {
+ genBody.put("level", "ERROR");
+ genBody.put("statusMessage", error);
+ }
+
+ Map traceBody = new LinkedHashMap<>();
+ traceBody.put("id", traceId);
+ traceBody.put("name", "xlyAi-agent");
+ traceBody.put("timestamp", start);
+
+ List