103 lines
3.2 KiB
Python

"""Pi JSONL 到通用 FrameworkEvent 工件的确定性归一。"""
from __future__ import annotations
import json
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from framework.primitives.artifacts import payload_sha256, write_jsonl_atomic
_KIND_PHASE: dict[str, tuple[str, str]] = {
"session": ("session", "started"),
"turn_start": ("turn", "started"),
"turn_end": ("turn", "completed"),
"message_end": ("model", "completed"),
"tool_execution_start": ("tool", "started"),
"tool_execution_end": ("tool", "completed"),
"agent_start": ("agent", "started"),
"agent_end": ("agent", "completed"),
"agent_settled": ("agent", "completed"),
}
def _safe_details(event: dict[str, Any]) -> dict[str, Any]:
"""只保留执行观察字段,不把消息正文复制进通用事件。"""
keys = (
"type",
"provider",
"model",
"stopReason",
"toolName",
"toolCallId",
"isError",
"exitCode",
)
return {key: event[key] for key in keys if key in event}
def normalize_pi_transcript(
transcript_path: str | Path,
artifact_path: str | Path,
*,
run_id: str | None = None,
framework_version: str = "unknown",
) -> dict[str, Any]:
"""把 Pi 原始 JSONL 归一为可重放的通用事件工件。"""
source = Path(transcript_path)
try:
lines = source.read_text(encoding="utf-8").splitlines()
except (OSError, UnicodeError) as exc:
raise ValueError(f"Pi transcript 不可读: {source}") from exc
events: list[dict[str, Any]] = []
session_id: str | None = None
unknown: list[str] = []
observed_at = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
for raw_line in lines:
if not raw_line.strip():
continue
try:
raw = json.loads(raw_line)
except json.JSONDecodeError as exc:
raise ValueError("Pi transcript 含不可解析 JSON 行") from exc
if not isinstance(raw, dict):
raise ValueError("Pi transcript 事件必须是对象")
event_type = str(raw.get("type") or "unknown")
if event_type == "session":
session_id = str(raw.get("id") or "") or session_id
kind, phase = _KIND_PHASE.get(event_type, ("unknown", "progress"))
if kind == "unknown":
unknown.append(event_type)
events.append(
{
"framework": "pi",
"frameworkVersion": framework_version,
"sessionId": session_id,
"sourceSeq": len(events) + 1,
"sourceEventId": str(raw.get("id") or "") or None,
"runId": run_id,
"kind": kind,
"phase": phase,
"safeDetails": _safe_details(raw),
"payloadSha256": payload_sha256(raw),
"observedAt": observed_at,
}
)
if not events:
raise ValueError("Pi transcript 为空")
digest = write_jsonl_atomic(artifact_path, events)
return {
"path": str(artifact_path),
"sha256": digest,
"eventCount": len(events),
"unknownEventTypes": sorted(set(unknown)),
}
__all__ = ["normalize_pi_transcript"]