主题
2. 解析流水线
核心函数位于
langfuse_transcript.py(Watcher 中有同构副本)
1. 流水线总览
read_new_jsonl() # 增量读文件
↓
build_turns() # 切分为 Turn 列表
↓
completed_turns() # [Watcher] 判断哪些 Turn 可以上报
↓
discover_subagent_files() # 发现子代理 jsonl
↓
emit_main_turn() # 上报到 Langfuse2. 增量读取:read_new_jsonl()
2.1 SessionState 结构
python
@dataclass
class SessionState:
offset: int = 0 # 文件字节偏移
buffer: str = "" # 半行缓冲(未读完的 JSON 行)
turn_count: int = 0 # 已成功上报的 Turn 数
streaming: dict = {} # 流式模式下的中间状态状态持久化在 ~/.claude/state/langfuse_*_state.json,key 为 sha256(session_id::transcript_path)。
2.2 读取逻辑
python
def read_new_jsonl(transcript_path, ss):
# 1. 文件被截断(offset > 文件大小)→ 重置 offset
if size < ss.offset:
ss.offset = 0
ss.buffer = ""
# 2. 从 offset 读到文件末尾
f.seek(ss.offset)
chunk = f.read()
ss.offset = f.tell()
# 3. 拼接上次未完成的半行
combined = ss.buffer + chunk.decode()
lines = combined.split("\n")
ss.buffer = lines[-1] # 最后一行可能不完整,留到下次
# 4. 解析完整行
for line in lines[:-1]:
msgs.append(json.loads(line))生活类比:看书看到一半句子被分页切开——buffer 就是「上一页没读完的半句话」,下次接着读。
2.3 为什么需要半行缓冲
Claude Code 边写边刷 jsonl,一行 JSON 可能分多次 write()。没有缓冲会把半截 JSON 当坏行丢掉。
3. Turn 切分:build_turns()
3.1 Turn 数据结构
python
@dataclass
class Turn:
user_msg: Dict # 用户消息(Turn 起点)
assistant_msgs: List[Dict] # 本轮所有 assistant 行
tool_results_by_id: Dict # tool_use_id → 结果内容
tool_use_ts_by_id: Dict # tool_use_id → 发起时间
tool_use_inputs_by_id: Dict # tool_use_id → 输入参数
tool_use_names_by_id: Dict # tool_use_id → 工具名
tool_result_ts_by_id: Dict # tool_use_id → 结果时间3.2 切分规则
遍历 jsonl 消息:
tool_result 行
→ 提取 tool_use_id 和 content,存入 tool_results_by_id
→ 不新开 Turn
user 行(非 synthetic)
→ flush 上一个 Turn
→ 开始新 Turn
user 行(synthetic,如 Skill 注入)
→ 折叠进当前 Turn,不新开
assistant 行
→ 追加到当前 Turn 的 assistant_msgs3.3 状态机示意
[user: "写快排"] ← Turn 1 开始
[assistant: text+tool_use] ← 追加
[user: tool_result] ← 归档结果
[assistant: end_turn] ← 追加
[user: "再加注释"] ← Turn 1 结束,Turn 2 开始
[assistant: ...]3.4 Synthetic User 过滤
Skill 工具运行后,Claude Code 会把 Skill 指令以 user 角色写回 transcript:
"Base directory for this skill: /path/to/skill"is_synthetic_user() 识别这类消息,不拆 Turn——否则一次 Skill 调用会被误切成两轮对话。
4. Turn 完成判断:is_turn_complete()
python
def is_turn_complete(turn: Turn) -> bool:
last = turn.assistant_msgs[-1]
return last["message"]["stop_reason"] == "end_turn"| 场景 | 行为 |
|---|---|
| Hook 批量模式 | 每轮 Stop 触发时,新读到的 Turn 默认已完成 |
| Watcher 批量模式 | completed_turns() 只上报已完成的 Turn |
| Watcher 流式模式 | 进行中的 Turn 也处理,end_turn 时 finalize=True |
completed_turns() 逻辑:
python
# 除最后一个外,前面的 Turn 一定已完成(被下一个 user 截断)
out = all_turns[:-1]
# 最后一个仅在 stop_reason == end_turn 时包含
if is_turn_complete(all_turns[-1]):
out.append(all_turns[-1])5. 子代理发现与匹配
5.1 发现:discover_subagent_files()
扫描候选目录,收集所有 agent-*.jsonl:
python
@dataclass
class SubagentFile:
jsonl_path: Path
agent_id: str
agent_type: Optional[str] # 来自 meta.json 或首行
description: Optional[str]
first_ts: Optional[datetime] # 首条消息时间戳5.2 匹配:match_subagent()
当主 Turn 里出现 tool_use(name in {"Task", "Agent"}) 时,选最佳子代理文件:
评分公式(越低越好):
score = 0
+ 1000 (subagent_type 不匹配)
+ |Δt| (tool_use 时间与 subagent 首条消息的时间差)
+ 500 (超出时间窗口 SUBAGENT_WINDOW_S,默认 10s)python
requested_type = tool_use_input.get("subagent_type") # 如 "Explore"
sf.agent_type # 来自 meta.json,如 "Explore"生活类比:主厨喊「帮厨去做凉菜」(Task 工具),后厨有多个帮厨小本本,通过「工种 + 开始时间」找到对应那一本。
5.3 已用子代理去重
used_subagents: set 防止同一个 agent-*.jsonl 被多个 Task 工具重复匹配。
6. Assistant 消息分组:_group_assistant_calls()
解决 trpc 拆行问题,按 message.id 分组:
python
# 输入:一个 Turn 的所有 assistant 行
# 输出:每个实际模型调用一条记录
{
"mid": "msg_01ABC",
"text": "好的,我来...\n",
"start_dt": datetime,
"end_dt": datetime,
"model": "claude-sonnet-4-...",
"usage": {"input": 1200, "output": 85},
"tool_ids": ["toolu_01XYZ", "toolu_02ABC"]
}Token 处理:流式占位行的 usage 全零 → 视为未填充;多行重复非零 usage → 取最后一行(不叠加)。
7. 数据清洗
7.1 脱敏:redact()
正则匹配 API Key、Bearer Token、密码赋值、私钥等,替换为 [REDACTED]。
环境变量:CC_LANGFUSE_REDACT=true(默认开启)
7.2 截断:truncate_text()
超过 CC_LANGFUSE_MAX_CHARS(默认 20000)截断,metadata 记录原始长度和 sha256。
8. 解析流程图
下一章:3_langfuse_emit.md — 解析结果如何变成 Langfuse 对象。