fix(langfuse): v4 가 옛 ingestion 을 거부 — OTLP/HTTP JSON(/api/public/otel/v1/traces)으로 전환, 로컬 Langfuse 실측 통과. 사용자·세션 속성은 자식 span 에도 복사

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
lee-hyeon-cheol
2026-09-22 13:53:20 +09:00
co-authored by Claude Fable 5.1
parent 0321050053
commit 55f2e8fb60
7 changed files with 175 additions and 68 deletions
+2
View File
@@ -8,3 +8,5 @@ staticfiles/
# 고객사 소스 위키(ABAP_INDEXING 결과물) — 고객 데이터라 git 밖
opencode/wiki/
deploy/langfuse/.env
deploy/langfuse/docker-compose.override.yml
+6 -5
View File
@@ -256,13 +256,14 @@ async def _finalize(session: ChatSession, state: TurnState, started: float, *, f
tid = f"{session.id}-{int(started * 1000)}"
user_email = await sync_to_async(lambda: session.user.email)()
meta = {"failed": failed} if failed else {}
t0 = datetime.fromtimestamp(started, timezone.utc).isoformat()
root = langfuse.trace(tid, "chat", userId=user_email, sessionId=session.id, input=state.user_text,
output=content, metadata=meta, tags=["codeassist"], start=t0, end=langfuse.now_iso())
langfuse.send_later([
langfuse.trace(tid, "chat", userId=user_email, sessionId=session.id, input=state.user_text,
output=content, metadata=meta, tags=["codeassist"]),
root,
langfuse.generation(
tid, "opencode-turn",
startTime=datetime.fromtimestamp(started, timezone.utc).isoformat(),
endTime=langfuse.now_iso(),
tid, "opencode-turn", parent=root,
startTime=t0, endTime=langfuse.now_iso(),
input=state.user_text, output=content,
usage=langfuse.usage_of(usage["input"], usage["output"]) if usage else None,
level="ERROR" if failed else "DEFAULT", statusMessage=failed or "",
+103 -20
View File
@@ -1,5 +1,6 @@
"""Langfuse 로 trace/generation 쏘기 — SDK 없이 HTTP 한 방 (docs-lib/langfuse.md).
"""Langfuse 로 trace/generation 쏘기 — SDK 없이 OTLP/HTTP JSON 한 방 (docs-lib/langfuse.md).
v4 는 옛 /api/public/ingestion 이 막혀서(score 만) OTel 엔드포인트로 감. 트레이스 하나 = 루트 span(trace 속성) + 자식 span(generation).
LANGFUSE_HOST 가 비어 있으면 전부 no-op. 보내는 건 fire-and-forget: 실패해도 채팅엔 영향 0, 로그만.
두 군데서 부름:
- apps/chat/stream.py _finalize → 사용자 단위 trace(누가·어느 세션·질문·최종 답·토큰·시간)
@@ -10,7 +11,10 @@ LANGFUSE_HOST 가 비어 있으면 전부 no-op. 보내는 건 fire-and-forget:
from __future__ import annotations
import asyncio
import hashlib
import json
import logging
import time
import uuid
from datetime import datetime, timezone
@@ -31,16 +35,86 @@ def now_iso() -> str:
return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
def _event(kind: str, body: dict) -> dict:
return {"id": str(uuid.uuid4()), "timestamp": now_iso(), "type": kind, "body": body}
def _trace_hex(trace_id: str) -> str:
"""아무 문자열 → OTel traceId(16바이트 hex). 같은 문자열이면 같은 trace 로 묶임."""
return hashlib.sha256(trace_id.encode()).hexdigest()[:32]
def trace(trace_id: str, name: str, **body) -> dict:
return _event("trace-create", {"id": trace_id, "name": name, **body})
def _nanos(iso_or_none: str | None) -> str:
if not iso_or_none:
return str(time.time_ns())
dt = datetime.fromisoformat(iso_or_none.replace("Z", "+00:00"))
return str(int(dt.timestamp() * 1_000_000_000))
def generation(trace_id: str, name: str, **body) -> dict:
return _event("generation-create", {"id": str(uuid.uuid4()), "traceId": trace_id, "name": name, **body})
def _attr(k: str, v) -> dict | None:
if v is None or v == "" or v == {} or v == []:
return None
if isinstance(v, bool):
return {"key": k, "value": {"boolValue": v}}
if isinstance(v, int):
return {"key": k, "value": {"intValue": str(v)}}
if isinstance(v, float):
return {"key": k, "value": {"doubleValue": v}}
if isinstance(v, str):
return {"key": k, "value": {"stringValue": v}}
return {"key": k, "value": {"stringValue": json.dumps(v, ensure_ascii=False)}}
def _span(trace_id: str, name: str, attrs: dict, *, start: str | None, end: str | None, parent: str | None, error: str | None) -> dict:
span = {
"traceId": _trace_hex(trace_id),
"spanId": uuid.uuid4().hex[:16],
"name": name,
"kind": 1, # INTERNAL
"startTimeUnixNano": _nanos(start),
"endTimeUnixNano": _nanos(end),
"attributes": [a for a in (_attr(k, v) for k, v in attrs.items()) if a],
"status": {"code": 2, "message": error} if error else {"code": 1},
}
if parent:
span["parentSpanId"] = parent
return span
def trace(trace_id: str, name: str, *, userId: str = "", sessionId: str = "", input=None, output=None,
metadata: dict | None = None, tags: list[str] | None = None, start: str | None = None, end: str | None = None) -> dict:
"""루트 span. trace 속성(이름·사용자·세션·태그)은 여기 실림."""
attrs = {
"langfuse.trace.name": name,
"langfuse.observation.type": "span",
"langfuse.user.id": userId,
"langfuse.session.id": sessionId,
"langfuse.trace.input": input,
"langfuse.trace.output": output,
"langfuse.trace.tags": tags,
**{f"langfuse.trace.metadata.{k}": v for k, v in (metadata or {}).items()},
}
return _span(trace_id, name, attrs, start=start, end=end, parent=None, error=None)
def generation(trace_id: str, name: str, *, model: str = "", input=None, output=None, usage: dict | None = None,
level: str = "DEFAULT", statusMessage: str = "", metadata: dict | None = None,
modelParameters: dict | None = None, startTime: str | None = None, endTime: str | None = None,
parent: dict | None = None) -> dict:
"""generation span. parent 로 trace() 결과를 주면 그 밑에 붙고, 없으면 같은 traceId 의 루트.
사용자·세션·이름·태그는 자식에도 복사 — Langfuse 가 필터·집계할 때 span 마다 보기 때문(로컬 실측: 안 하면 빈 값)."""
inherited = {a["key"]: a["value"].get("stringValue") for a in (parent or {}).get("attributes", [])
if a["key"] in ("langfuse.user.id", "langfuse.session.id", "langfuse.trace.name", "langfuse.trace.tags")}
attrs = {
**inherited,
"langfuse.observation.type": "generation",
"langfuse.observation.model.name": model,
"langfuse.observation.input": input,
"langfuse.observation.output": output,
"langfuse.observation.usage_details": usage,
"langfuse.observation.model_parameters": modelParameters,
"langfuse.observation.level": level,
"langfuse.observation.status_message": statusMessage,
**{f"langfuse.observation.metadata.{k}": v for k, v in (metadata or {}).items()},
}
return _span(trace_id, name, attrs, start=startTime, end=endTime,
parent=parent["spanId"] if parent else None, error=statusMessage if level == "ERROR" else None)
def usage_of(inp: int | None, out: int | None) -> dict | None:
@@ -49,36 +123,45 @@ def usage_of(inp: int | None, out: int | None) -> dict | None:
return {"input": inp or 0, "output": out or 0, "total": (inp or 0) + (out or 0)}
async def send(events: list[dict]) -> bool:
"""배치 하나 전송. 207 이면 성공(부분 실패는 로그)."""
def _otlp_body(spans: list[dict]) -> dict:
return {"resourceSpans": [{
"resource": {"attributes": [_attr("service.name", "codeassist-backend")]},
"scopeSpans": [{"scope": {"name": "codeassist"}, "spans": spans}],
}]}
async def send(spans: list[dict]) -> bool:
"""span 묶음 하나 전송. OTLP 는 200 + 빈 JSON 이면 성공, partialSuccess 있으면 일부 실패."""
cfg = settings.LANGFUSE
if not cfg.get("host") or not events:
if not cfg.get("host") or not spans:
return False
try:
async with httpx.AsyncClient(timeout=5.0, transport=TRANSPORT) as client:
resp = await client.post(
cfg["host"].rstrip("/") + "/api/public/ingestion",
json={"batch": events},
cfg["host"].rstrip("/") + "/api/public/otel/v1/traces",
json=_otlp_body(spans),
headers={"x-langfuse-ingestion-version": "4"},
auth=(cfg.get("public_key", ""), cfg.get("secret_key", "")),
)
if resp.status_code not in (200, 207):
if resp.status_code != 200:
log.warning("langfuse ← %s %s", resp.status_code, resp.text[:200])
return False
errors = (resp.json() or {}).get("errors") or []
if errors:
log.warning("langfuse 일부 실패: %s", errors[:3])
return not errors
partial = (resp.json() or {}).get("partialSuccess") or {}
if partial.get("rejectedSpans"):
log.warning("langfuse 일부 거부: %s", partial)
return False
return True
except Exception as e: # noqa: BLE001 — 관측용이라 절대 본 흐름 안 깨뜨림
log.warning("langfuse 전송 실패: %s: %s", type(e).__name__, e)
return False
def send_later(events: list[dict]) -> None:
def send_later(spans: list[dict]) -> None:
"""지금 흐름 안 막고 백그라운드로. 이벤트 루프 없으면(동기 테스트) 조용히 버림."""
if not enabled() or not events:
if not enabled() or not spans:
return
try:
task = asyncio.get_running_loop().create_task(send(events))
task = asyncio.get_running_loop().create_task(send(spans))
except RuntimeError:
return
_pending.add(task)
+5 -3
View File
@@ -101,11 +101,13 @@ def _observe(model_id: str, body: dict, col: _Collect, started: float, status: i
return
u = col.usage or {}
tid = str(uuid.uuid4())
t0 = datetime.fromtimestamp(started, timezone.utc).isoformat()
root = langfuse.trace(tid, "fabrix", tags=["gateway"], metadata={"parts": kinds}, start=t0, end=langfuse.now_iso())
langfuse.send_later([
langfuse.trace(tid, "fabrix", tags=["gateway"], metadata={"parts": kinds}),
root,
langfuse.generation(
tid, "fabrix.chat", model=model_id,
startTime=datetime.fromtimestamp(started, timezone.utc).isoformat(), endTime=langfuse.now_iso(),
tid, "fabrix.chat", model=model_id, parent=root,
startTime=t0, endTime=langfuse.now_iso(),
input=body.get("messages"), output="".join(col.text),
usage=langfuse.usage_of(u.get("prompt_tokens"), u.get("completion_tokens")),
level="ERROR" if status >= 400 else "DEFAULT", statusMessage="" if status < 400 else f"upstream {status}",
+21 -19
View File
@@ -2,30 +2,32 @@
출처: `https://cloud.langfuse.com/generated/api/openapi.yml`, `langfuse.com/self-hosting`. SDK 안 쓰고 HTTP 로 직접 쏨(오프라인 wheels 반입 줄이려고).
## POST /api/public/ingestion
## v4 는 OTLP 로 (2026-09-22 로컬 실측)
- 인증: Basic auth — user = `pk-lf-…`(public key), password = `sk-lf-…`(secret key)
- 응답: **207** (배치 부분 성공). 각 이벤트 결과가 `successes[]`/`errors[]`
- Langfuse Cloud 에선 2026-11-16 폐기 예정이지만 **self-host 는 계속 지원**. OTel 엔드포인트(`/api/public/otel/v1/traces`)는 v3+ 만.
`POST /api/public/ingestion` 은 v4 기본(`LANGFUSE_MIGRATION_V4_WRITE_MODE=events_only`)에서 **trace/generation 을 400 으로 거부**함(score 만 받음).
`dual` 로 열 수 있지만 다음 메이저에서 없어짐 → 우리는 OTLP감.
### POST /api/public/otel/v1/traces
- OTLP/HTTP **JSON** 도 받음(protobuf 안 써도 됨). gRPC 는 없음.
- 헤더: `Authorization: Basic base64(pk-lf-…:sk-lf-…)`, `x-langfuse-ingestion-version: 4`(직접 쓰기, 15분 지연 없음)
- 응답 200 `{}`. 일부 거부면 `partialSuccess.rejectedSpans`.
- traceId = 16바이트 hex(32자), spanId = 8바이트 hex(16자). 시간은 `startTimeUnixNano`/`endTimeUnixNano` 문자열.
```json
{
"batch": [
{"id": "evt-1", "timestamp": "2026-09-22T01:00:00Z", "type": "trace-create",
"body": {"id": "trace-1", "name": "chat", "userId": "u@x.com", "sessionId": "ses_…", "input": "…", "output": "", "metadata": {}, "tags": []}},
{"id": "evt-2", "timestamp": "2026-09-22T01:00:05Z", "type": "generation-create",
"body": {"id": "gen-1", "traceId": "trace-1", "name": "fabrix", "model": "581",
"startTime": "…", "endTime": "", "input": [...messages], "output": "…",
"usage": {"input": 120, "output": 30, "total": 150}, "level": "DEFAULT", "statusMessage": "", "metadata": {}}}
]
}
{"resourceSpans":[{"resource":{"attributes":[{"key":"service.name","value":{"stringValue":"codeassist-backend"}}]},
"scopeSpans":[{"scope":{"name":"codeassist"},"spans":[
{"traceId":"<32hex>","spanId":"<16hex>","name":"chat","kind":1,"startTimeUnixNano":"…","endTimeUnixNano":"…",
"attributes":[{"key":"langfuse.trace.name","value":{"stringValue":"chat"}},{"key":"langfuse.user.id","value":{"stringValue":"u@x.com"}}],
"status":{"code":1}},
{"traceId":"<같은>","spanId":"<16hex>","parentSpanId":"<루트 spanId>","name":"fabrix.chat","kind":1,,
"attributes":[{"key":"langfuse.observation.type","value":{"stringValue":"generation"}},]}]}]}]}
```
- envelope: `id`(중복 제거용, 유일), `timestamp`(ISO 8601), `type`, `body`
- 이벤트 타입: `trace-create`, `generation-create`, `span-create`, `event-create`, `score-create`, `*-update`
- TraceBody: id·name·userId·sessionId·input·output·metadata·tags·version·release·environment·public
- CreateGenerationBody: id·traceId·name·startTime·endTime·completionStartTime·model·modelParameters·input·output·usage{input,output,total,unit,inputCost,outputCost,totalCost}·usageDetails·costDetails·level(DEBUG/DEFAULT/WARNING/ERROR)·statusMessage·metadata·parentObservationId
- 같은 traceId 로 trace-create 를 여러 번 보내도 됨(upsert). generation 먼저 와도 trace 가 나중에 붙음.
### 속성 → Langfuse 필드
- 관측(span) 단위: `langfuse.observation.type`(generation/span), `.model.name`, `.input`, `.output`, `.usage_details`(JSON `{"input","output","total"}`), `.model_parameters`, `.level`(DEBUG/DEFAULT/WARNING/ERROR), `.status_message`, `.metadata.<키>`, `.cost_details`
(대안: `gen_ai.request.model`, `gen_ai.prompt`/`gen_ai.completion`, `gen_ai.usage.*`)
- 트레이스 단위(아무 span 에나 실으면 트레이스 전체에 적용, 필터 쓰려면 루트에): `langfuse.trace.name`(없으면 루트 span 이름), `langfuse.user.id`, `langfuse.session.id`, `langfuse.trace.tags`(JSON 배열), `langfuse.trace.input`/`.output`, `langfuse.trace.metadata.<키>`, `langfuse.environment`, `langfuse.version`
- 값 인코딩: 문자열은 stringValue, 객체/배열은 JSON 문자열로 넣으면 Langfuse 가 파싱해서 보여줌.
## 자체 호스팅 (v4 docker-compose, 2026-09 기준)
+6 -3
View File
@@ -133,13 +133,16 @@ async def test_stream_reports_generation_to_langfuse(upstream, monkeypatch):
from apps.gateway import langfuse
got = []
monkeypatch.setattr(langfuse, "TRANSPORT", httpx.MockTransport(lambda r: (got.append(json.loads(r.content)), httpx.Response(207, json={}))[1]))
monkeypatch.setattr(langfuse, "TRANSPORT", httpx.MockTransport(lambda r: (got.append(json.loads(r.content)), httpx.Response(200, json={}))[1]))
upstream(lambda r: _sse(b'data: {"choices":[{"delta":{"content":"hi"}}],"usage":{"prompt_tokens":2,"completion_tokens":1}}\n\ndata: [DONE]\n\n'))
resp = await _post({"model": "581", "stream": True, "messages": [{"role": "user", "content": "q"}]})
b"".join([c async for c in resp.streaming_content])
await asyncio.gather(*list(langfuse._pending))
gen = got[0]["batch"][1]["body"]
assert gen["model"] == "581" and gen["output"] == "hi" and gen["usage"]["total"] == 3 and gen["input"][0]["content"] == "q"
spans = got[0]["resourceSpans"][0]["scopeSpans"][0]["spans"]
ga = {a["key"]: a["value"].get("stringValue") for a in spans[1]["attributes"]}
assert ga["langfuse.observation.model.name"] == "581" and ga["langfuse.observation.output"] == "hi"
assert json.loads(ga["langfuse.observation.usage_details"])["total"] == 3
assert json.loads(ga["langfuse.observation.input"])[0]["content"] == "q"
@override_settings(FABRIX_ENV=FULL_ENV)
+32 -18
View File
@@ -1,4 +1,4 @@
"""Langfuse 전송 — 이벤트 모양·인증·게이트웨이/턴 마무리 훅. 서버는 MockTransport."""
"""Langfuse 전송(OTLP/HTTP JSON) — span 모양·인증·게이트웨이/턴 마무리 훅. 서버는 MockTransport."""
import asyncio
import base64
@@ -19,27 +19,38 @@ def lf_server(monkeypatch):
def handler(req: httpx.Request) -> httpx.Response:
got.append({"url": str(req.url), "auth": req.headers.get("authorization"), "body": json.loads(req.content)})
return httpx.Response(207, json={"successes": [], "errors": []})
return httpx.Response(200, json={})
monkeypatch.setattr(langfuse, "TRANSPORT", httpx.MockTransport(handler))
return got
def _attrs(span: dict) -> dict:
out = {}
for a in span["attributes"]:
v = a["value"]
out[a["key"]] = v.get("stringValue", v.get("intValue", v.get("boolValue")))
return out
@override_settings(LANGFUSE=LF)
async def test_send_batch_shape_and_basic_auth(lf_server):
ok = await langfuse.send([
langfuse.trace("t1", "chat", userId="u@x.com", sessionId="s1"),
langfuse.generation("t1", "gen", model="581", usage=langfuse.usage_of(10, 5)),
])
assert ok and len(lf_server) == 1
async def test_send_otlp_shape_and_basic_auth(lf_server):
root = langfuse.trace("t1", "chat", userId="u@x.com", sessionId="s1", tags=["x"])
gen = langfuse.generation("t1", "gen", model="581", usage=langfuse.usage_of(10, 5), parent=root)
assert await langfuse.send([root, gen]) and len(lf_server) == 1
req = lf_server[0]
assert req["url"] == "http://lf.test/api/public/ingestion"
assert req["url"] == "http://lf.test/api/public/otel/v1/traces"
assert req["auth"] == "Basic " + base64.b64encode(b"pk-lf-a:sk-lf-b").decode()
batch = req["body"]["batch"]
assert [e["type"] for e in batch] == ["trace-create", "generation-create"]
assert batch[0]["body"] == {"id": "t1", "name": "chat", "userId": "u@x.com", "sessionId": "s1"}
assert batch[1]["body"]["traceId"] == "t1" and batch[1]["body"]["usage"] == {"input": 10, "output": 5, "total": 15}
assert batch[0]["id"] != batch[1]["id"] and batch[0]["timestamp"].endswith("Z")
spans = req["body"]["resourceSpans"][0]["scopeSpans"][0]["spans"]
assert len(spans) == 2 and spans[0]["traceId"] == spans[1]["traceId"] and len(spans[0]["traceId"]) == 32
assert spans[1]["parentSpanId"] == spans[0]["spanId"] and "parentSpanId" not in spans[0]
ra, ga = _attrs(spans[0]), _attrs(spans[1])
assert ra["langfuse.trace.name"] == "chat" and ra["langfuse.user.id"] == "u@x.com" and ra["langfuse.session.id"] == "s1"
assert json.loads(ra["langfuse.trace.tags"]) == ["x"]
assert ga["langfuse.observation.type"] == "generation" and ga["langfuse.observation.model.name"] == "581"
assert ga["langfuse.user.id"] == "u@x.com" and ga["langfuse.session.id"] == "s1" # 자식에도 복사
assert json.loads(ga["langfuse.observation.usage_details"]) == {"input": 10, "output": 5, "total": 15}
assert int(spans[1]["endTimeUnixNano"]) >= int(spans[1]["startTimeUnixNano"])
@override_settings(LANGFUSE={"host": "", "public_key": "", "secret_key": ""})
@@ -80,8 +91,11 @@ async def test_finalize_sends_user_trace(lf_server, django_user_model, monkeypat
await _finalize(session, TurnState("ses_1", "질문"), 1.0)
await asyncio.gather(*list(langfuse._pending))
assert len(lf_server) == 1
batch = lf_server[0]["body"]["batch"]
t, g = batch[0]["body"], batch[1]["body"]
assert t["userId"] == "u@x.com" and t["sessionId"] == "ses_1" and t["input"] == "질문" and t["output"] == ""
assert g["traceId"] == t["id"] and g["usage"] == {"input": 3, "output": 5, "total": 8} and g["level"] == "DEFAULT"
spans = lf_server[0]["body"]["resourceSpans"][0]["scopeSpans"][0]["spans"]
t, g = _attrs(spans[0]), _attrs(spans[1])
assert t["langfuse.user.id"] == "u@x.com" and t["langfuse.session.id"] == "ses_1"
assert t["langfuse.trace.input"] == "질문" and t["langfuse.trace.output"] == ""
assert spans[1]["parentSpanId"] == spans[0]["spanId"]
assert json.loads(g["langfuse.observation.usage_details"]) == {"input": 3, "output": 5, "total": 8}
assert g["langfuse.observation.level"] == "DEFAULT" and spans[1]["status"]["code"] == 1
assert stream_mod is not None