主题
5. Watcher 方案详解
入口:
langfuse_trace/watcher/claude_langfuse_watcher.py
1. 工作原理
独立守护进程启动
↓
每 1s 轮询 projects 目录(或指定 transcript)
↓
发现/跟踪主 session jsonl 文件
↓
read_new_jsonl() 增量读新增行
↓
build_turns(include_pending=True) 含进行中 Turn
↓
process_turn_streaming() 或 emit_main_turn() 上报
↓
Langfuse UI 实时更新生活类比:会计不等工作结束,而是每秒翻流水账,看到一行记一行。
2. JsonlWatcher 架构
python
class JsonlWatcher:
def __init__(self, langfuse, paths, poll_s=1.0, idle_timeout_s=300):
self._tracked: Dict[str, TrackedFile] = {} # 正在跟踪的文件
@dataclass
class TrackedFile:
path: Path
session_id: str
messages: List[Dict] # 累积的所有消息
last_activity: float # 最后有新行的时间
last_size: int
wait_for_create: bool # 文件尚未创建时等待主循环
python
def run_forever(self, rescan_roots):
while not self._stop:
self.run_once(rescan_roots) # 处理所有 tracked 文件
time.sleep(self.poll_s) # 默认 1srun_once() 流程
加锁 → 加载 watcher state
→ 对每个 TrackedFile 调用 process_file()
→ 清理 idle timeout 的文件
→ 保存 state3. 文件发现
命令行参数
bash
# 监听整个 projects 目录(推荐)
python3 claude_langfuse_watcher.py --projects-root ~/.claude/projects
# trpc-claudecode
python3 claude_langfuse_watcher.py --projects-root ~/.trpc-claudecode/projects
# 监听单个文件(不存在时会等待)
python3 claude_langfuse_watcher.py --transcript /path/to/session.jsonl
# 只跑一轮(测试)
python3 claude_langfuse_watcher.py --transcript demo.jsonl --once主 session 识别
python
# 只跟踪 UUID 命名的 jsonl,排除 agent-*.jsonl
_MAIN_JSONL_STEM_RE = r"^[0-9a-f]{8}-...-[0-9a-f]{12}$"新 session 自动发现
rescan() 定期扫描 projects_root,发现新的 <sessionId>.jsonl 自动加入跟踪列表。目录为空时不会退出,持续等待。
4. process_file() 双模式
python
def process_file(langfuse, global_state, tracked):
if streaming_mode_enabled(): # 默认 true
return process_file_streaming(...)
else:
return process_file_batch(...) # 内部逻辑,即下面的批量路径4.1 批量模式(CC_LANGFUSE_STREAMING=false)
python
all_turns = build_turns(tracked.messages)
done_turns = completed_turns(all_turns) # 只取已完成的
pending = done_turns[ss.turn_count:] # 跳过已上报的
for turn in pending:
emit_main_turn(...)
emitted += 1
if flush_ok and emitted == len(pending):
ss.turn_count += emitted与 Hook 行为类似:Turn 完整结束后才上报。
4.2 流式模式(默认,CC_LANGFUSE_STREAMING=true)
python
all_turns = build_turns(tracked.messages, include_pending=True)
active = all_turns[ss.turn_count] # 当前进行中的 Turn
finalize = is_turn_complete(active)
process_turn_streaming(..., finalize=finalize)
if finalize and flush_ok:
ss.turn_count += 1
ss.streaming = {} # 清空流式中间状态5. 状态管理
| 文件 | 用途 |
|---|---|
~/.claude/state/langfuse_watcher_state.json | Watcher 专用状态 |
~/.claude/state/langfuse_watcher_state.lock | 文件锁 |
~/.claude/state/langfuse_watcher.log | Watcher 日志 |
~/.claude/state/langfuse_hook.log | 解析/上报共享日志 |
流式模式下 state 额外保存 streaming 字段:
json
{
"streaming": {
"trace_id": "uuid-...",
"turn_num": 3,
"generations": {"msg_01ABC": {"id": "gen-uuid", "done": false}},
"tools": {"toolu_01XYZ": {"id": "span-uuid", "ended": false}},
"used_subagents": ["agent-id-1"]
}
}6. Idle Timeout
python
IDLE_TIMEOUT_S = float(os.environ.get("CC_LANGFUSE_WATCH_IDLE_S", "300"))文件超过 300 秒无新行 → 停止跟踪(释放内存)。下次 rescan 发现同一文件且有新内容 → 通过 _bootstrap_messages_if_needed() 从磁盘恢复完整消息列表。
设 CC_LANGFUSE_WATCH_IDLE_S=0 表示永不停止跟踪。
7. 启动方式
推荐:run_watcher.sh
bash
cd langfuse_trace/watcher
export TRACE_TO_LANGFUSE=true
export CC_LANGFUSE_PUBLIC_KEY=pk-lf-xxx
export CC_LANGFUSE_SECRET_KEY=sk-lf-xxx
./run_watcher.sh后台运行
bash
nohup ./run_watcher.sh > /tmp/langfuse_watcher.out 2>&1 &8. Watcher 方案的特点
| 优点 | 缺点 |
|---|---|
| 零侵入 Claude Code | 需常驻进程 |
| 流式上报,近实时 | 多占一点 CPU(1s 轮询) |
| 自动发现新 session | 与 Hook 不可同时开 |
支持 --once 做 CI 测试 | 自包含大文件(~2100 行) |
9. 与 langfuse_transcript.py 的关系
claude_langfuse_watcher.py 在文件顶部自包含了一份与 langfuse_transcript.py 同构的解析/上报代码(约 1-1600 行),然后在此基础上增加了:
JsonlWatcher轮询框架(1600 行以后)completed_turns()/TrackedFile/discover_transcripts()- Watcher 专用状态文件和 CLI 参数
设计原因:Watcher 可独立部署运行,不依赖 import common.langfuse_transcript。
两者逻辑应保持同步;修改解析规则时需同时更新两个文件(或未来抽取共享包)。