Skip to content

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)       # 默认 1s

run_once() 流程

加锁 → 加载 watcher state
  → 对每个 TrackedFile 调用 process_file()
  → 清理 idle timeout 的文件
  → 保存 state

3. 文件发现

命令行参数

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

6_streaming_mode.md

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.jsonWatcher 专用状态
~/.claude/state/langfuse_watcher_state.lock文件锁
~/.claude/state/langfuse_watcher.logWatcher 日志
~/.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

两者逻辑应保持同步;修改解析规则时需同时更新两个文件(或未来抽取共享包)。

下一章:6_streaming_mode.md