Files

1378 lines
71 KiB
Python

"""
Parse Claude Code transcripts (~/.claude/projects/<proj>/<session>.jsonl) into a
viewable conversation: a flat thread of text / thinking / tool-call / tool-result
items (tools meant to be collapsed in the UI) plus real per-turn token usage and
dollar cost taken from each assistant message's `usage` block.
Two depths:
* ``parse_summary`` — cheap rollup for the conversation list (title, model,
counts, aggregated usage/cost). Also returns per-skill usage so the indexer can
derive skill analytics from the same single read.
* ``parse_thread`` — the full message thread for the detail view.
"""
import datetime
import json
import pathlib
import re
from .pricing import (
add_usage,
norm_usage,
usage_cache_cost,
usage_cost,
usage_tokens,
zero_usage,
)
# Claude Code session ids (= main transcript file stems) are UUIDs.
_UUID_RE = re.compile(r"[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}"
r"-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}")
def _ms(iso: str | None) -> float | None:
"""Epoch milliseconds for a transcript ISO timestamp, or None."""
if not iso:
return None
try:
return datetime.datetime.fromisoformat(
iso.replace("Z", "+00:00")).timestamp() * 1000
except ValueError:
return None
#: Width of one activity bin, in minutes (the activity dashboard's histogram
#: resolution). Bins are absolute — indexed from the Unix epoch in UTC — so a
#: bin index is timezone-free and only the *reader* maps it to a local slot.
BIN_MINUTES = 5
_BIN_MS = BIN_MINUTES * 60 * 1000
def _bin5(iso: str | None) -> int | None:
"""The 5-minute bin index (epoch-based, UTC) a transcript timestamp lands in."""
ms = _ms(iso)
return None if ms is None else int(ms // _BIN_MS)
def _elapsed_ms(start: str | None, end: str | None) -> int | None:
"""Wall-clock milliseconds between two transcript timestamps."""
a, b = _ms(start), _ms(end)
if a is None or b is None or b < a:
return None
return int(b - a)
# Map a Claude Code tool name → the activity bucket we attribute its turn's
# tokens/cost to (so a conversation's spend can be broken down by what it was
# spent doing: plain messages, sub-agents, bash, file reads/edits, …).
TOOL_BUCKETS = {
"Read": "read", "NotebookRead": "read",
"Edit": "edit", "Write": "edit", "MultiEdit": "edit", "NotebookEdit": "edit",
"Bash": "bash", "BashOutput": "bash", "KillShell": "bash",
"Task": "agent", "Agent": "agent",
"Grep": "search", "Glob": "search", "LS": "search",
"Skill": "skill",
"WebFetch": "web", "WebSearch": "web",
"TodoWrite": "task", "TaskCreate": "task", "TaskUpdate": "task",
}
# All buckets the UI may show ("msg" = text-only turns, "context" = the first
# API call's injected context share, "tool" = catch-all).
TOOL_BUCKET_ORDER = ["msg", "context", "agent", "bash", "read", "edit", "search",
"skill", "web", "task", "tool"]
# The session's first API call pays for far more than the user's typed message:
# the harness front-loads the system prompt, tool definitions, CLAUDE.md files,
# memories and other <system-reminder> context into it. The API reports only the
# combined total, so we estimate the actual message's share from its length and
# attribute the rest of that first input side to context loading.
EST_CHARS_PER_TOKEN = 4
def _est_tokens(text: str | None) -> int:
"""Rough token estimate for a text (~4 chars/token), 0 for empty."""
return max(1, len(text) // EST_CHARS_PER_TOKEN) if text else 0
# A homelab project is `…/projects/<slug>`. Claude's own transcript store
# (`.claude/projects/<encoded-path>`) is excluded by the leading-dash filter
# below, since those encoded names always start with a dash.
_PROJECT_RE = re.compile(r"(?:^|[/\s])projects/([A-Za-z0-9._][\w.-]*)")
# A homelab service is `…/services/<slug>` — the infra side of the repo.
_SERVICE_RE = re.compile(r"(?:^|[/\s])services/([A-Za-z0-9._][\w.-]*)")
# A conversation is attributed to a project/service only when it actually
# **worked on** it — i.e. created or edited a file under `projects/<slug>` or
# `services/<slug>`. We therefore mine the auto-tags exclusively from the
# file-mutating tools' target paths (Write/Edit/MultiEdit → `file_path`,
# NotebookEdit → `notebook_path`), *not* from the cwd, plain reads, searches, or
# bash commands that merely mention a path. This keeps the tag list to projects
# the conversation genuinely changed, instead of everything it happened to look
# at or `cd` into.
_EDIT_TOOLS = {"Write", "Edit", "MultiEdit", "NotebookEdit"}
# A conversation that ran `scripts/new-worktree.sh <name>` spun up an isolated git
# worktree to do its work in (per the repo's CLAUDE.md convention), and one that
# ran `new-worktree.sh -r <path>` tore one down. We mine both out of the command
# so the viewer can flag which worktrees a conversation created (and when it
# removed them). To avoid matching mentions inside quoted strings / echo / grep
# (rife in a conversation *about* the script, like this one), we anchor to a
# command position: the script must sit at the start of the command or right
# after a shell separator (`;`, `&&`, `|`, `(`, newline) or a `bash`/`sudo`/…
# prefix — an optional path in front, then `new-worktree.sh` and its arg tail (up
# to the next separator). Events are teased out in `_worktree_events`.
_WORKTREE_RE = re.compile(
r"(?:^|[\n;&|(]|\b(?:bash|sh|sudo|exec|command)[ \t])[ \t]*"
r"(?:[\w./-]*/)?new-worktree\.sh\b([^\n|&;<>]*)")
# Management-mode flags that only list/help — they create and remove nothing.
_WORKTREE_NOOP_FLAGS = {"-l", "--list", "-h", "--help"}
# A plausible worktree name / directory basename (git-ref-ish). Anything with
# shell quoting or punctuation is parse noise, not a real invocation, and is
# dropped — a second line of defense behind the command-position anchor above.
_WORKTREE_NAME_RE = re.compile(r"^[A-Za-z0-9][\w./-]*$")
# A conversation signals it is finished by ending its last message with the
# literal marker ``DONE`` (see the "When work is done" rule in CLAUDE.md). We
# match it as an upper-case word so ordinary prose ("done", "well done") doesn't
# trip it. Combined with a recorded notification, this flips a stale "running"
# marker to finished in the viewer.
_DONE_RE = re.compile(r"\bDONE\b")
# The `complete` skill's end-of-task pass (summary published, memories, docs,
# artefacts, cleanup) continues the same session, so its own reply becomes the
# transcript's last message; it signs off with the distinct marker ``COMPLETED``
# rather than ``DONE``, which lets the viewer tell "the task finished" apart
# from "the task finished AND was tidied up after" (`meta.completedAt`).
_COMPLETED_RE = re.compile(r"\bCOMPLETED\b")
# When a user interrupts a session (Esc in an interactive session, or a SIGINT
# sent to a headless ``claude -p`` run — including the interrupt button in this
# viewer), Claude Code records a synthetic user message whose only text block is
# ``[Request interrupted by user]`` (or the ``… for tool use]`` variant). If such
# a marker is the *last* content-bearing message in the transcript, the run ended
# on an interruption, so we surface the conversation's state as ``interrupted``.
_INTERRUPT_RE = re.compile(r"^\s*\[Request interrupted by user")
# A follow-up the ai-agent delivered into a *live* run through the CLI's
# cross-session inbox socket (sidecar ``/message``) instead of killing the run.
# Claude Code records it as a peer message — "Another Claude session sent a
# message:\n<text>\n\nThis came from another Claude session — not typed by your
# user …" — and the sidecar prefixes the text with a bracketed note saying it
# *is* the user's (sidecar.INBOX_PREFIX). Shown as the plain user turn it is.
_INBOX_RE = re.compile(
r"\AAnother Claude session sent a message:\s*\n"
r"\[Follow-up from your user[^\]]*\]\s*\n+"
r"(.*?)"
r"(?:\n\n+This came from another Claude session\b.*)?\Z", re.S)
def _unwrap_inbox(c: str) -> str:
m = _INBOX_RE.match(c)
return m.group(1).strip() if m else c
# The sidecar's Stop guard (sidecar/stop_guard.py) held the turn while a plain
# background shell finished, then re-entered the model with this feedback. It
# is a runner event, not the user talking — the viewer shows a compact "held"
# tag (like "interrupted"), with the hold spelled out in its tooltip.
# (the CLI writes "Stop hook feedback:" then a newline before the reason)
_HELD_RE = re.compile(r"\A(?:Stop hook feedback:\s*)+Stop guard \(headless run\)")
_HELD_SECS_RE = re.compile(r"turn was held (\d+) s")
_HELD_MIN_RE = re.compile(r"turn was held (\d+) min")
_RELAY_RE = re.compile(r"<<<USER MESSAGE>>>\n(.*?)\n<<<END USER MESSAGE>>>", re.S)
def _relayed_message(c: str) -> str | None:
"""The user's text the Stop guard relayed (it arrived while the turn was
held), or None for a plain hold feedback."""
m = _RELAY_RE.search(c)
return m.group(1).strip() if m else None
def _held_label(c: str) -> str:
m = _HELD_SECS_RE.search(c)
if m:
return f"turn held {int(m.group(1))} s for a background shell"
m = _HELD_MIN_RE.search(c)
if m:
return f"background shell still running after {m.group(1)} min"
return "turn held for a background shell"
# A Task/Agent call spawned with ``run_in_background`` gets its tool_result the
# instant the agent launches ("Async agent launched successfully… agentId: …"),
# so — unlike a blocking agent, whose result *is* its report — that result says
# nothing about whether the agent is still working. Its real end signal is the
# ``<task-notification>`` Claude Code injects when the agent stops, which names
# the Task call it belongs to.
_AGENT_ASYNC_RE = re.compile(r"Async agent launched|^\s*agentId:", re.M)
_AGENT_DONE_RE = re.compile(
r"<task-notification>.*?<tool-use-id>\s*([^\s<]+)\s*</tool-use-id>", re.S)
# The canned assistant reply Claude Code records when a turn is interrupted (or a
# queued/resume prompt yields nothing to answer): a single text block reading
# "No response requested.". It carries no real content, so we hide the bubble and
# surface a small "paused" marker inline (and the conversation's state) instead.
_NO_RESPONSE_RE = re.compile(r"^\s*No response requested\.?\s*$")
# When Claude invokes a skill (the `Skill` tool), Claude Code injects the skill's
# SKILL.md into the transcript as a *user* message that opens with this marker
# line ("Base directory for this skill: <abs path>"). Left alone it renders as a
# giant "You" bubble and miscounts as a user turn — so we detect it, fold its body
# into the preceding Skill tool card, and drop it from the turn/search counts.
_SKILL_BODY_RE = re.compile(
r"^Base directory for this skill:[ \t]*(\S[^\n]*)\n+", re.S)
def _skill_body(text: str) -> tuple[str, str] | None:
"""If `text` is an injected skill body, return (base_dir, body) with the
marker line stripped; else None."""
m = _SKILL_BODY_RE.match(text or "")
if not m:
return None
return m.group(1).strip(), text[m.end():]
def _rel_skill_path(base_dir: str, cwd: str | None) -> str | None:
"""Repo-relative path to the skill's SKILL.md (for a link into the editor),
or None when the skill lives outside the repo (built-in/plugin skills)."""
if not base_dir:
return None
marker = ".claude/skills/"
idx = base_dir.find(marker)
if idx != -1:
return base_dir[idx:].rstrip("/") + "/SKILL.md"
if cwd:
root = cwd.rstrip("/") + "/"
if base_dir.startswith(root):
return base_dir[len(root):].rstrip("/") + "/SKILL.md"
return None
# A *recalled* memory is injected into the transcript (as a tool_result) with its
# full file content — frontmatter + body. We detect each one by its frontmatter
# (``name: <slug>`` shortly above ``node_type: memory``) and record the slug, so
# the viewer can show which memories a conversation read. The lazy span keeps the
# slug bound to the nearest ``node_type: memory`` below it; the session-start
# MEMORY.md index (plain ``- [..](..md)`` lines) carries no such frontmatter and
# so is correctly ignored.
_MEMORY_RE = re.compile(
r"name:\s*([A-Za-z0-9][\w.-]*)\s*\n[\s\S]{0,400}?node_type:\s*memory")
# ── lifecycle backfill ───────────────────────────────────────────────────────
# The CI/CD-style stamps (committed/pushed/merged/deployed/notified) are normally
# written live by the skill scripts (conv-meta.sh, called from the commit / deploy
# / notify skills) into the metadata sidecar. For conversations that ran before
# those hooks existed — or that committed/deployed via plain git/deploy commands —
# the sidecar is empty, so the status icons never light up. We *derive* the stamps
# from the transcript instead: scan each Skill/Bash tool call for evidence the
# action happened, and stamp it with that turn's timestamp. These derived values
# are merged non-destructively under `lifecycleAuto` (the live sidecar wins).
#
# The actual deploy commands (as opposed to the `conv-meta ... deployed` stamp
# below, which only *records* one after the fact) — reused to detect a deploy
# in flight, see `_is_deploy_call`.
_DEPLOY_CMD_RE = re.compile(
r"\bdeploy-html\b"
r"|\b(?:npm|pnpm|yarn)\s+(?:run\s+)?deploy\b"
r"|\bzipgo\s+deploy\b"
r"|\bmake\s+deploy[\w.]*\b"
r"|\bdeploy[\w.-]*\.sh\b")
# Each pattern matches a command/skill-name string; a hit lights its stage(s).
_LIFECYCLE_PATTERNS: list[tuple[str, re.Pattern]] = [
# commit-homelab / commit-project both commit *and* push to Gitea.
("committed", re.compile(r"\b(?:commit-homelab|commit-project)\b")),
("pushed", re.compile(r"\b(?:commit-homelab|commit-project)\b")),
("committed", re.compile(r"\bgit\b[^\n|&;]*\bcommit\b")),
("pushed", re.compile(r"\bgit\b[^\n|&;]*\bpush\b")),
("merged", re.compile(r"\bgit\s+merge\b")),
# deploy paths: skill/alias, npm/pnpm/yarn deploy, zipgo, make deploy, *.sh.
("deployed", _DEPLOY_CMD_RE),
("notified", re.compile(r"\b(?:notify-done|notify-ask|notify\.sh|notify\s+(?:send|ask|form))\b")),
# explicit conv-meta stamps (e.g. `conv-meta.sh committed`).
("committed", re.compile(r"\bconv-meta(?:\.sh)?\b[^\n|&;]*\bcommitted\b")),
("pushed", re.compile(r"\bconv-meta(?:\.sh)?\b[^\n|&;]*\bpushed\b")),
("merged", re.compile(r"\bconv-meta(?:\.sh)?\b[^\n|&;]*\bmerged\b")),
("deployed", re.compile(r"\bconv-meta(?:\.sh)?\b[^\n|&;]*\bdeployed\b")),
("notified", re.compile(r"\bconv-meta(?:\.sh)?\b[^\n|&;]*\bnotified\b")),
]
# Skill-tool invocations (matched on the `skill` input) → stage(s).
_SKILL_STAGES: dict[str, tuple[str, ...]] = {
"commit-homelab": ("committed", "pushed"),
"commit-project": ("committed", "pushed"),
"deploy-html": ("deployed",),
"notify-done": ("notified",),
}
# ── deployed URL ─────────────────────────────────────────────────────────────
# Every deploy path prints where the thing landed, but always *last* and often
# after a very long build log — and the result we ship to the UI is clipped to
# MAX_BLOCK chars. So the URL is pulled out of the raw output here, before the
# clip, and hung on the tool card as `deployedUrl` (the UI turns it into a live
# preview widget + the conversation's "open the deploy" gizmo).
_DEPLOY_URL_PATTERNS: list[re.Pattern] = [
re.compile(r"Deployed URL\s*[:=]\s*(https?://[^\s'\"<>]+)"),
re.compile(r"Live:\s*HTTP\s*\d+\s*at\s+(https?://[^\s'\"<>]+)"),
re.compile(r"(?:deployed|live|serving|available)\s+(?:at|on)\s+(https?://[^\s'\"<>]+)",
re.IGNORECASE),
re.compile(r"https://[\w.-]+\.(?:gabvdl\.xyz|gabriel\.vidal--ayrinhac\.xyz)[^\s'\"<>]*"),
]
def _deployed_url(rtext: str) -> str | None:
"""The URL a deploy command reported, best-effort (raw, unclipped output)."""
for pat in _DEPLOY_URL_PATTERNS:
m = pat.search(rtext)
if m:
url = m.group(1) if m.groups() else m.group(0)
return url.rstrip(".,;:)\"'")
return None
def _lifecycle_stages(name: str | None, inp: dict) -> set[str]:
"""Stages a tool call provides evidence for (best-effort, may over-report)."""
stages: set[str] = set()
if name == "Skill":
stages.update(_SKILL_STAGES.get((inp.get("skill") or ""), ()))
cmd = inp.get("command")
if isinstance(cmd, str) and cmd:
for stage, pat in _LIFECYCLE_PATTERNS:
if stage not in stages and pat.search(cmd):
stages.add(stage)
return stages
def _is_deploy_call(name: str | None, inp: dict) -> bool:
"""True if this tool call itself performs (or wraps) a deploy — used to
detect one in flight, unlike `_lifecycle_stages`'s "deployed" hit, which
also fires on the post-hoc `conv-meta ... deployed` stamp."""
if name == "Skill" and inp.get("skill") == "deploy-html":
return True
cmd = inp.get("command")
return isinstance(cmd, str) and bool(_DEPLOY_CMD_RE.search(cmd))
def _bucket(tool: str | None) -> str:
return TOOL_BUCKETS.get(tool or "", "tool")
# The leading program of a shell command — the sub-label a "bash" turn is broken
# down by (git / npm / docker / …). We skip common env/prefix wrappers so the
# real command shows through (e.g. `sudo docker …` → docker, `FOO=1 npm …` → npm).
_BASH_PREFIX = {"sudo", "env", "time", "nice", "nohup", "exec", "command",
"xargs", "then", "do", "if", "!"}
def _bash_program(command) -> str:
"""First real program name in a shell command (for the bash sub-breakdown)."""
if not isinstance(command, str):
return "?"
cmd = command.strip()
# `cd <dir> && <real cmd>` / `(cd <dir>; …)` is pervasive here — skip the cd
# prefix so the actual command shows through instead of "cd".
cmd = re.sub(r"^\(?\s*cd\s+[^\n&;|]+(?:&&|;|\|\|)\s*", "", cmd).lstrip("(").strip()
# Only look at the first pipeline/segment; strip a leading subshell/paren.
seg = re.split(r"[|;&\n]", cmd, 1)[0].strip().lstrip("(").strip()
for tok in seg.split():
# Skip VAR=value assignments and known wrapper prefixes.
if "=" in tok and not tok.startswith("-") and re.match(r"^\w+=", tok):
continue
if tok in _BASH_PREFIX:
continue
# Basename, drop any path so `/usr/bin/git` and `git` collapse together.
prog = tok.rsplit("/", 1)[-1]
return prog or "?"
return "?"
def _file_ext(inp: dict) -> str:
"""File extension sub-label for a read/edit turn (e.g. `.py`, `.tsx`)."""
p = inp.get("file_path") or inp.get("notebook_path") or ""
if not isinstance(p, str) or not p:
return "?"
base = p.rsplit("/", 1)[-1]
dot = base.rfind(".")
# Leading-dot names (`.env`) are extension-less; everything after the dot else.
if dot > 0:
return base[dot:].lower()
return "(no ext)"
def _subkey(bucket: str, tool: str | None, inp: dict) -> str | None:
"""The second-level breakdown label for a turn, given its activity bucket:
bash → program, read/edit → file extension, skill → constant "skills read",
any other tool bucket → the tool name. ``None`` for text-only (msg) turns,
which are split by prompt/reply at accumulation time instead."""
if bucket == "bash":
return _bash_program(inp.get("command"))
if bucket in ("read", "edit"):
return _file_ext(inp)
if bucket == "skill":
return "skills read"
if bucket == "msg":
return None
return tool or "tool"
# Valid task states we normalise to; anything else falls back to "pending".
_TASK_STATES = {"pending", "in_progress", "completed", "cancelled"}
_TASK_SUBJECT_MAX = 200
def _task_subject(s) -> str:
"""A short, single-line label for a task (bounded for the compact panel)."""
s = " ".join(str(s or "").split()).strip() or "task"
return s[:_TASK_SUBJECT_MAX]
def _projects_in(text) -> set[str]:
"""Pull `projects/<name>` slugs out of a path/command/cwd string."""
if not isinstance(text, str):
return set()
return {m.group(1) for m in _PROJECT_RE.finditer(text)
if not m.group(1).startswith("-")}
def _services_in(text) -> set[str]:
"""Pull `services/<name>` slugs out of a path/command/cwd string."""
if not isinstance(text, str):
return set()
return {m.group(1) for m in _SERVICE_RE.finditer(text)
if not m.group(1).startswith("-")}
def _worktree_events(text) -> list[dict]:
"""Worktree create/remove events in a `new-worktree.sh` command string.
Each event is ``{"kind": "create"|"remove", "name": str, "dir": str}`` where
``dir`` is the worktree *directory* basename — the stable key both a create
and its later removal share:
- create (``new-worktree.sh <name> [branch]`` / ``-p <project> <name>``):
``name`` is the ``<name>`` arg; ``dir`` is ``homelab-<name>`` (super mode)
or ``<project>-<name>`` (project mode) — i.e. the created directory.
- remove (``new-worktree.sh -r <path>``): ``dir`` is ``basename(<path>)``
(``name`` empty — the short name is recovered by pairing on ``dir``).
List/help modes create nothing and yield no events."""
events: list[dict] = []
if not isinstance(text, str):
return events
for m in _WORKTREE_RE.finditer(text):
args = m.group(1).split()
i = 0
name = project = remove_path = None
noop = False
while i < len(args):
a = args[i]
if a in _WORKTREE_NOOP_FLAGS:
noop = True
break
if a in ("-r", "--remove"):
remove_path = args[i + 1] if i + 1 < len(args) else None
i += 2
continue
if a in ("-p", "--project"): # takes a value (the project) — skip both
project = args[i + 1] if i + 1 < len(args) else None
i += 2
continue
if a.startswith("-"): # any other flag (e.g. --no-link)
i += 1
continue
if name is None:
name = a # first bare positional = worktree name
i += 1
if noop:
continue
if remove_path:
d = pathlib.PurePosixPath(remove_path.rstrip("/")).name
if _WORKTREE_NAME_RE.match(d):
events.append({"kind": "remove", "name": "", "dir": d})
elif name and _WORKTREE_NAME_RE.match(name) and (
project is None or _WORKTREE_NAME_RE.match(project)):
d = f"{project}-{name}" if project else f"homelab-{name}"
events.append({"kind": "create", "name": name, "dir": d})
return events
MAX_BLOCK = 8000 # cap a single text/result block so payloads stay bounded
def _real_model(model: str | None) -> bool:
"""False for Claude Code's own placeholder models (`<synthetic>`, …).
Messages the CLI fabricates itself (API-error notices, "No response
requested." fillers) carry a bracketed pseudo-model. It is not a model the
conversation ran on, so it must never become one of its model tags."""
return bool(model) and not model.startswith("<")
# The claude CLI's `error` field on an API-error turn → the viewer's kind (the
# same vocabulary as the sidecar's claude_cli.classify, which types launch
# failures). Text patterns cover older CLIs and `unknown`/unset fields.
_API_ERROR_FIELD = {"authentication_failed": "auth", "rate_limit": "rate_limit",
"server_error": "overloaded", "overloaded": "overloaded"}
_API_ERROR_TEXT = [
("auth", re.compile(r"not logged in|please run /login|failed to authenticate"
r"|oauth.*(expired|revoked|refresh)|invalid api key", re.I)),
("rate_limit", re.compile(r"session limit|usage limit|rate.?limit"
r"|weekly limit", re.I)),
("cli_outdated", re.compile(r"does not support this model|newer is required"
r"|claude update", re.I)),
("network", re.compile(r"connection (lost|error|refused)|ECONNRESET"
r"|ETIMEDOUT|fetch failed", re.I)),
("overloaded", re.compile(r"overloaded|API Error: 5\d\d", re.I)),
]
def api_error_kind(field: str | None, text: str | None) -> str:
"""The kind of a transcript API-error turn (``isApiErrorMessage``)."""
if field in _API_ERROR_FIELD:
return _API_ERROR_FIELD[field]
for kind, rx in _API_ERROR_TEXT:
if rx.search(text or ""):
return kind
return "unknown"
def _block_text(content) -> str:
"""Flatten a tool_result / message content into text."""
if content is None:
return ""
if isinstance(content, str):
return content
parts = []
if isinstance(content, list):
for b in content:
if isinstance(b, dict):
if b.get("type") == "text":
parts.append(b.get("text", ""))
elif b.get("type") == "image":
parts.append("[image]")
elif b.get("type") == "tool_reference":
# ToolSearch results are tool_reference blocks; keep one
# parseable line per match for the ToolSearch widget.
parts.append(f"tool_reference: {b.get('tool_name', '')}")
elif "text" in b:
parts.append(str(b["text"]))
else:
parts.append(str(b))
return "\n".join(parts)
def _clip(s: str) -> str:
if s and len(s) > MAX_BLOCK:
return s[:MAX_BLOCK] + f"\n… [{len(s) - MAX_BLOCK} more chars truncated]"
return s
# Per-message cap for the full-text search index. Search + snippet only need the
# meat of a message, and this text ships to the client (Fuse.js indexes it), so
# we bound each message to keep the search payload reasonable.
SEARCH_CLIP = 1000
def _search_clip(s: str) -> str:
s = " ".join((s or "").split()) # collapse whitespace/newlines for snippets
return s[:SEARCH_CLIP] if len(s) > SEARCH_CLIP else s
def _iter_records(path: pathlib.Path):
try:
with path.open(encoding="utf-8", errors="ignore") as fh:
for line in fh:
line = line.strip()
if not line or line[0] != "{":
continue
try:
yield json.loads(line)
except ValueError:
continue
except OSError:
return
class ParserState:
"""Resumable transcript parser: feed records one by one, snapshot anytime.
Holds exactly the cross-record state ``parse_conversation`` used to keep in
locals, so the indexer can parse a *growing* live transcript incrementally —
restore the state, feed only the appended JSONL lines, snapshot the summary
again — instead of re-reading the whole file on every watcher tick.
``summary()`` is a pure snapshot (it never mutates parser state), so it can
be called repeatedly between feeds."""
def __init__(self, path: pathlib.Path, full: bool = False):
self.path = path
self.full = full
# A main transcript is named `<sessionId>.jsonl`, and the file name is
# the only trustworthy identity: a resumed / forked / compacted session
# starts with history lines copied from its *source* session, whose
# `sessionId` fields still name that source. Seeding from the first
# record attributed those conversations to their origin — meta sidecar
# reads (title, state, projects) landed on the wrong conversation, so
# renaming one card appeared to rename another. Subagent transcripts
# (`agent-<id>.jsonl`, no UUID stem) keep the record fallback: their
# `sessionId` intentionally names the parent.
self.session_id = path.stem if _UUID_RE.fullmatch(path.stem) else None
self.cwd = None
self.git_branch = None
self.model = None
# Every model that produced a turn here, model id -> assistant records
# seen. A session can switch models mid-thread, so the viewer tags a
# conversation with all of them, busiest first.
self.model_turns: dict[str, int] = {}
# Same idea for the claude CLI's `--effort` level (low…max), which it
# stamps on each assistant record: a session can be resumed at another
# level, so we keep every one seen, busiest first.
self.effort_turns: dict[str, int] = {}
# When this run was actually working, as model id -> the set of
# 5-minute wall-clock bins (see `_bin5`) it produced a turn in. The
# activity dashboard counts *distinct agents per bin*, so presence is
# all that's needed — a bin is either worked in or it isn't. Collected
# here because the parser is already walking every record: deriving it
# later would mean a second pass over the whole archive.
self.active_bins: dict[str, set[int]] = {}
self.ai_title = None
self.first_user = None
self.started = None
self.ended = None
self.user_turns = 0
self.assistant_turns = 0
# A subagent (sidechain) transcript lives at
# `<convId>/subagents/agent-*.jsonl`; every record carries
# `isSidechain: true` and an `agentId` (its `sessionId` is the parent's).
self.is_sidechain = False
self.agent_id = None
self.agg = zero_usage()
# Total cost, accumulated per assistant record: the record's own
# `costUSD` (pi harness — the runner mirrors pi's real provider cost,
# 0 for the free local EVOX2 provider) when present, else the price
# table estimate for that turn's model. Summing per-turn also prices
# mixed-model sessions correctly (the old aggregate used one model).
self.cost_total = 0.0
# The conversation's current context size: the *last* assistant API
# call's full prompt+reply — how full the context window is now.
self.context_tokens = 0
self.context_model = None
# First-turn split: the first API call's input side divided into the
# injected context (system prompt, CLAUDE.md, memories, …) and the
# user's actual first message (estimated from its length). Filled once,
# on the first usage-bearing assistant record.
self.first_context_tokens = None
self.first_message_tokens = None
self._first_split_done = False
# The thread item of the user's first message (full mode): the split is
# stamped onto it once the first assistant usage arrives, so the viewer
# shows "context + message" right on that bubble. Patching an appended
# item later is the same pattern tool results use (`tool_items`).
self._first_user_item = None
self.by_tool: dict[str, dict] = {} # bucket -> accumulated usage
# …and its cost, accumulated per turn at *that turn's* model. Pricing the
# summed usage with one model at the end (what this used to do) overcharges
# a session that switched models mid-thread, so the breakdown no longer
# added up to `cost_total`. Same reason `cost_total` accumulates per turn.
self.by_tool_cost: dict[str, float] = {}
# Second-level breakdown: bucket -> sub-label -> usage (bash by program,
# read/edit by file extension, msg by prompt/replies).
self.by_tool_sub: dict[str, dict[str, dict]] = {}
self.by_tool_sub_cost: dict[str, dict[str, float]] = {}
self.projects: set[str] = set()
self.services: set[str] = set()
# Git worktrees created/removed, keyed by directory basename.
self.worktrees: dict[str, dict] = {}
self.skills: dict[str, dict] = {}
# Task list (TaskCreate/TaskUpdate or the older TodoWrite snapshots),
# keyed by sequential string id; insertion order ≈ creation order.
self.tasks: dict[str, dict] = {}
self.task_seq = 0 # running counter → next task id
self.lifecycle: dict[str, str] = {} # stage -> latest evidence ts
self.memories_read: set[str] = set()
self.last_assistant_text = "" # final assistant text (DONE marker)
# Trailing interrupt / "No response requested." markers — cleared by any
# later real activity; interrupt outranks paused.
self.interrupted_trailing = False
self.paused_trailing = False
# Claude parked itself: the last thing it did was schedule a wake-up
# (ScheduleWakeup) or block on a Monitor — the run sits idle until the
# trigger fires, so the viewer shows it "paused" rather than "running".
# Cleared by any later tool call, real user turn, or interrupt (NOT by
# the same turn's closing text — "I'll check back in 20 min" still ends
# with the run parked).
self.waiting_trailing = False
self.tool_name: dict[str, str] = {}
self.tool_items: dict[str, dict] = {} # tool_use_id -> thread item
# The thread item carrying the current assistant turn's cost (reset per
# assistant record) — its activity bucket is stamped on once known.
self._turn_item: dict | None = None
# Open Task/Agent calls: an entry here *is* a running subagent.
self.agent_calls: dict[str, dict] = {}
# Set while a deploy command (Bash or the deploy-html skill) has been
# issued and hasn't returned yet — cleared on its tool_result.
self.deploying_since: str | None = None
self._deploying_tool_use_id: str | None = None
# Skill card awaiting its injected SKILL.md body.
self.pending_skill: dict | None = None
# Timestamp of the previous record (assistant latency baseline).
self.prev_ts: str | None = None
self.thread: list[dict] = []
self.turns: list[dict] = []
# Full-text search index: one entry per user/assistant text block, in
# transcript order (`i` matches the thread items' `mi`).
self.search_msgs: list[dict] = []
def _add_search(self, role: str, text: str) -> int:
idx = len(self.search_msgs)
self.search_msgs.append(
{"i": idx, "role": role, "text": _search_clip(text)})
return idx
def feed(self, o: dict) -> None:
"""Fold one transcript record into the state."""
t = o.get("type")
self.session_id = self.session_id or o.get("sessionId")
if o.get("isSidechain") is True:
self.is_sidechain = True
self.agent_id = self.agent_id or o.get("agentId")
if o.get("cwd"):
# cwd only seeds the human-readable project *label* (_project_name);
# it does NOT auto-tag a project — a conversation is tagged only for
# projects/services it actually edited (see _EDIT_TOOLS below).
self.cwd = self.cwd or o.get("cwd")
self.git_branch = self.git_branch or o.get("gitBranch")
ts = o.get("timestamp")
rec_prev_ts = self.prev_ts
if ts:
self.started = self.started or ts
self.ended = ts
self.prev_ts = ts
if t == "ai-title" and o.get("aiTitle"):
self.ai_title = o["aiTitle"]
return
# A background agent's completion is announced by a `<task-notification>`
# injected as a queued-command attachment (its tool_result came back at
# launch, so that can't be the end signal — see `agent_calls`).
if t == "attachment":
m = _AGENT_DONE_RE.search((o.get("attachment") or {}).get("prompt") or "")
if m:
self.agent_calls.pop(m.group(1), None)
return
msg = o.get("message")
if not isinstance(msg, dict):
return
if t == "assistant":
# `<synthetic>` and friends are Claude Code's placeholders on messages
# it fabricates itself — not a model the conversation ran on, so they
# taint neither the primary model nor the model tags.
if _real_model(msg.get("model")):
self.model = self.model or msg["model"]
self.model_turns[msg["model"]] = self.model_turns.get(msg["model"], 0) + 1
# `effort` rides on the *record*, not the message. Only real
# model turns count, so synthetic records can't invent a level.
eff = o.get("effort")
if isinstance(eff, str) and eff.strip():
eff = eff.strip().lower()
self.effort_turns[eff] = self.effort_turns.get(eff, 0) + 1
# Mark this 5-minute bin as worked in, under the model that did
# the work — a mixed-model session lands in each model's own
# bins, so the histogram never attributes a turn to the wrong one.
b = _bin5(ts)
if b is not None:
self.active_bins.setdefault(msg["model"], set()).add(b)
u = msg.get("usage")
nu = None
out_tokens = 0
if isinstance(u, dict):
nu = norm_usage(u)
self.agg = add_usage(self.agg, nu)
self.assistant_turns += 1
out_tokens = nu["output"]
if self.full:
self.turns.append({"model": msg.get("model"), **nu})
blocks = msg.get("content")
if isinstance(blocks, str):
blocks = [{"type": "text", "text": blocks}]
first = True
turn_bucket = None # first tool this turn drives the cost attribution
turn_subkey = None # …and its sub-label (program / ext / tool name)
tmodel = msg.get("model") or self.model
# This assistant API call's own tokens/cost — surfaced as a per-tool
# "delta" on the first tool card of the turn (the call that made it).
turn_tokens = usage_tokens(nu) if nu is not None else 0
turn_cost = usage_cost(nu, tmodel) if nu is not None else 0.0
# A runner-mirrored pi record carries the real provider cost —
# authoritative, so it beats the price-table estimate.
real_cost = o.get("costUSD")
if isinstance(real_cost, (int, float)) and not isinstance(real_cost, bool):
turn_cost = float(real_cost)
self.cost_total += turn_cost
# …and how much of it is just the cached prompt being replayed, so the
# cost-growth chart can stack cache under fresh tokens.
turn_cache_tokens = int(nu["cacheRead"]) if nu is not None else 0
turn_cache_cost = usage_cache_cost(nu, tmodel) if nu is not None else 0.0
if nu is not None and turn_tokens:
self.context_tokens = turn_tokens
self.context_model = tmodel
# First usage-bearing record = the session's first API call: split
# its input side into injected context vs the user's actual message
# (see EST_CHARS_PER_TOKEN). `first_ctx_u` is the context share as a
# usage dict, peeled off this turn's bucket attribution below.
first_ctx_u = None
if nu is not None and not self._first_split_done:
self._first_split_done = True
in_side = nu["input"] + nu["cacheRead"] + nu["cacheWriteTokens"]
msg_tok = min(_est_tokens(self.first_user), in_side)
self.first_message_tokens = msg_tok
self.first_context_tokens = in_side - msg_tok
if in_side > 0 and self.first_context_tokens > 0:
f = self.first_context_tokens / in_side
first_ctx_u = {
"input": round(nu["input"] * f),
"output": 0,
"cacheRead": round(nu["cacheRead"] * f),
"cacheWriteTokens": round(nu["cacheWriteTokens"] * f),
"cacheWriteUnits": nu["cacheWriteUnits"] * f,
}
if self._first_user_item is not None and in_side > 0:
self._first_user_item["firstContextTokens"] = \
self.first_context_tokens
self._first_user_item["firstMessageTokens"] = \
self.first_message_tokens
delta_attached = False
# The one thread item this turn's cost is stamped on (the first tool
# card / text bubble of the turn, or the carrier item below when the
# turn rendered nothing). The bucket is stamped onto it once the
# attribution is known, at the end of the record.
self._turn_item = None
# How long this assistant API call took: from the previous transcript
# record (the tool result / user turn that unblocked it) to this one.
turn_ms = _elapsed_ms(rec_prev_ts, ts)
for b in blocks or []:
if not isinstance(b, dict):
continue
bt = b.get("type")
if bt == "tool_use":
self.interrupted_trailing = self.paused_trailing = False # real activity
self.tool_name[b.get("id")] = b.get("name")
inp = b.get("input") or {}
# A ScheduleWakeup (unless it's the stop:true loop-ender)
# parks the run until the wake-up fires; a Monitor blocks
# until its condition. Anything else is real work.
if b.get("name") == "ScheduleWakeup":
self.waiting_trailing = not inp.get("stop")
else:
self.waiting_trailing = b.get("name") == "Monitor"
if turn_bucket is None:
turn_bucket = _bucket(b.get("name"))
turn_subkey = _subkey(turn_bucket, b.get("name"), inp)
# Auto-tag projects/services only from files this turn
# *created or edited* — the target path of a file-mutating
# tool — not from reads, searches, or bash path mentions.
if b.get("name") in _EDIT_TOOLS:
for v in (inp.get("file_path"), inp.get("notebook_path")):
self.projects |= _projects_in(v)
self.services |= _services_in(v)
for ev in _worktree_events(inp.get("command")):
wt = self.worktrees.get(ev["dir"])
if wt is None:
wt = {"name": ev["name"] or ev["dir"], "dir": ev["dir"],
"createdAt": None, "removedAt": None}
self.worktrees[ev["dir"]] = wt
if ev["kind"] == "create":
if ev["name"]:
wt["name"] = ev["name"]
if ts and (wt["createdAt"] is None or ts < wt["createdAt"]):
wt["createdAt"] = ts
elif ts and (wt["removedAt"] is None or ts > wt["removedAt"]):
wt["removedAt"] = ts
for stage in _lifecycle_stages(b.get("name"), inp):
if ts and (stage not in self.lifecycle or ts > self.lifecycle[stage]):
self.lifecycle[stage] = ts
if b.get("id") and _is_deploy_call(b.get("name"), inp):
self.deploying_since = ts
self._deploying_tool_use_id = b["id"]
if b.get("name") == "Skill":
sk = (b.get("input") or {}).get("skill")
if sk:
e = self.skills.setdefault(sk, {"count": 0, "last": None})
e["count"] += 1
if ts and (e["last"] is None or ts > e["last"]):
e["last"] = ts
# Task-list tools. TaskCreate appends a task (id = its 1-based
# ordinal, the same id its tool_result reports and TaskUpdate
# cites); TaskUpdate mutates one by id. TodoWrite is the older
# tool that re-sends the *entire* list each call, so it resets
# our task map to that snapshot — the last snapshot wins.
tname = b.get("name")
if tname == "TaskCreate":
self.task_seq += 1
tid = str(self.task_seq)
self.tasks[tid] = {
"id": tid,
"subject": _task_subject(
inp.get("subject") or inp.get("content")),
"status": "pending",
}
elif tname == "TaskUpdate":
tid = str(inp.get("taskId") or inp.get("task_id") or "")
cur = self.tasks.get(tid)
if cur is not None:
st = inp.get("status")
if st in _TASK_STATES:
cur["status"] = st
if inp.get("subject"):
cur["subject"] = _task_subject(inp["subject"])
elif tname == "TodoWrite" and isinstance(
inp.get("todos"), list):
self.tasks.clear()
self.task_seq = 0
for td in inp["todos"]:
if not isinstance(td, dict):
continue
self.task_seq += 1
tid = str(self.task_seq)
st = td.get("status")
self.tasks[tid] = {
"id": tid,
"subject": _task_subject(
td.get("content") or td.get("activeForm")),
"status": st if st in _TASK_STATES else "pending",
}
if tname in ("Task", "Agent") and b.get("id"):
self.agent_calls[b["id"]] = {
"toolUseId": b["id"],
"agentType": inp.get("subagent_type"),
"description": _task_subject(inp["description"])
if inp.get("description") else None,
"startedAt": ts,
}
if self.full:
it = _item("assistant", "tool_use", first, out_tokens,
name=b.get("name"), input=b.get("input"), ts=ts,
model=tmodel)
# An Agent/Task call's block id links it to the spawned
# subagent transcript (its meta.json `toolUseId`), so the
# detail endpoint can attach the child conversation. Every
# other call carries its id too — it's the handle the tool
# result is matched on, and useful metadata in the UI.
it["toolUseId"] = b.get("id")
# Remember this Skill card so the SKILL.md body that Claude
# Code injects next (a "Base directory…" user message) folds
# into it instead of rendering as its own "You" bubble.
if b.get("name") == "Skill":
self.pending_skill = it
if not delta_attached and nu is not None:
it["turnTokens"] = turn_tokens
it["turnCost"] = turn_cost
it["turnCacheTokens"] = turn_cache_tokens
it["turnCacheCost"] = turn_cache_cost
delta_attached = True
self._turn_item = it
self.thread.append(it)
self.tool_items[b.get("id")] = it
first = False
elif bt == "text" and _NO_RESPONSE_RE.match(b.get("text", "")):
# Canned "No response requested." reply → hide the bubble and
# mark the run paused. Doesn't clear a trailing interrupt (an
# interrupt's own canned reply must stay "interrupted").
self.paused_trailing = True
if self.full:
self.thread.append(_item("assistant", "paused", first, 0, ts=ts))
first = False
elif bt == "text" and b.get("text", "").strip():
self.interrupted_trailing = self.paused_trailing = False # real activity
self.last_assistant_text = b["text"] # keep only the latest one
mi = self._add_search("assistant", b["text"])
if self.full:
it = _item("assistant", "text", first, out_tokens,
text=_clip(b["text"]), mi=mi, ts=ts,
model=tmodel, durationMs=turn_ms)
if o.get("isApiErrorMessage"):
# The CLI's own synthetic turn for a failed API
# call ("Not logged in · Please run /login", a
# usage limit…): the viewer draws an error card
# with the fix instead of a normal bubble.
it["apiError"] = api_error_kind(
o.get("error"), b["text"])
# Surface this turn's tokens/cost on its first text bubble
# (if no tool card already claimed the delta).
if not delta_attached and nu is not None:
it["turnTokens"] = turn_tokens
it["turnCost"] = turn_cost
it["turnCacheTokens"] = turn_cache_tokens
it["turnCacheCost"] = turn_cache_cost
delta_attached = True
self._turn_item = it
self.thread.append(it)
first = False
elif bt == "thinking" and b.get("thinking", "").strip():
self.interrupted_trailing = self.paused_trailing = False # real activity
if self.full:
self.thread.append(_item("assistant", "thinking", first, out_tokens,
text=_clip(b["thinking"]), ts=ts,
model=tmodel, durationMs=turn_ms))
first = False
if nu is not None:
bucket = turn_bucket or "msg"
turn_u = nu
# This turn's cost, split the same way its usage is. Accumulating
# the cost here — at this turn's model, and from `turn_cost` (which
# a pi record's real `costUSD` may have overridden) — is what keeps
# the breakdown summing to `cost_total`.
ctx_cost = 0.0
if first_ctx_u is not None:
# First turn: the injected-context share goes to its own
# top-level "context" bucket; only the remainder (the actual
# message + this turn's output) is attributed as usual.
turn_u = {k: nu[k] - first_ctx_u[k] for k in nu}
self.by_tool["context"] = add_usage(
self.by_tool.get("context", zero_usage()), first_ctx_u)
est = usage_cost(nu, tmodel)
share = (usage_cost(first_ctx_u, tmodel) / est) if est > 0 else 0.0
ctx_cost = turn_cost * share
self.by_tool_cost["context"] = \
self.by_tool_cost.get("context", 0.0) + ctx_cost
self.by_tool[bucket] = add_usage(self.by_tool.get(bucket, zero_usage()), turn_u)
own_cost = turn_cost - ctx_cost
self.by_tool_cost[bucket] = \
self.by_tool_cost.get(bucket, 0.0) + own_cost
if self.full:
# Every priced turn must carry its cost on exactly one thread
# item, or the cost-growth chart (which sums the items) lands
# below the conversation total. A turn slips through when it
# rendered nothing at all: a *redacted* thinking block is
# signature-only, so its `thinking` text is empty and no item is
# appended — yet the call was made and billed (mostly cache
# reads). Give those turns an invisible carrier item so the
# money still reaches the chart.
if not delta_attached:
it = _item("assistant", "turn", first, out_tokens, ts=ts,
model=tmodel, durationMs=turn_ms)
it["turnTokens"] = turn_tokens
it["turnCost"] = turn_cost
it["turnCacheTokens"] = turn_cache_tokens
it["turnCacheCost"] = turn_cache_cost
self.thread.append(it)
self._turn_item = it
delta_attached = True
first = False
# Colour the turn by the activity the breakdown actually charged
# it to, so the chart's legend can't drift from "cost by
# activity" (the frontend used to re-derive the bucket from the
# tool name, which has no way to know about the `context` split).
if self._turn_item is not None:
self._turn_item["turnBucket"] = bucket
if first_ctx_u is not None:
# The first call's injected-context share is its own
# bucket in the breakdown; hand the chart the same split
# (the very same numbers) so it can stack it under this
# turn rather than mis-attributing the whole call to
# `msg`.
self._turn_item["turnContextCost"] = ctx_cost
self._turn_item["turnContextTokens"] = usage_tokens(
first_ctx_u)
subs = self.by_tool_sub.setdefault(bucket, {})
sub_costs = self.by_tool_sub_cost.setdefault(bucket, {})
if bucket == "msg":
# A plain messaging turn: split its usage into the prompt side
# (everything fed in) and the reply (what the model produced) —
# and its cost the same way, so the subs sum to their bucket.
prompt_u = {**turn_u, "output": 0}
reply_u = {**zero_usage(), "output": turn_u["output"]}
subs["prompt"] = add_usage(subs.get("prompt", zero_usage()),
prompt_u)
subs["replies"] = add_usage(subs.get("replies", zero_usage()),
reply_u)
est = usage_cost(turn_u, tmodel)
reply_share = (usage_cost(reply_u, tmodel) / est) if est > 0 else 0.0
reply_cost = own_cost * reply_share
sub_costs["replies"] = sub_costs.get("replies", 0.0) + reply_cost
sub_costs["prompt"] = \
sub_costs.get("prompt", 0.0) + (own_cost - reply_cost)
else:
sk = turn_subkey or "?"
subs[sk] = add_usage(subs.get(sk, zero_usage()), turn_u)
sub_costs[sk] = sub_costs.get(sk, 0.0) + own_cost
elif t == "user":
c = msg.get("content")
if isinstance(c, str):
c = _unwrap_inbox(c)
relayed = _relayed_message(c) if _HELD_RE.match(c) else None
if relayed is not None:
c = relayed # the user's own words, handed over by the guard
done = _AGENT_DONE_RE.search(c) # replayed task-notification
if done:
self.agent_calls.pop(done.group(1), None)
sk = _skill_body(c)
if sk is not None:
# Injected SKILL.md — fold into the Skill card, not a turn.
if self.full and self.pending_skill is not None:
base_dir, body = sk
self.pending_skill["skillBody"] = _clip(body)
rel = _rel_skill_path(base_dir, self.cwd)
if rel:
self.pending_skill["skillPath"] = rel
self.pending_skill = None
elif _INTERRUPT_RE.match(c):
self.interrupted_trailing = True
self.paused_trailing = self.waiting_trailing = False
if self.full:
# Hide the raw "[Request interrupted by user]" bubble; the
# viewer renders a compact right-aligned "interrupted" tag.
self.thread.append(_item("user", "interrupted", False, 0, ts=ts))
elif _HELD_RE.match(c) and _relayed_message(c) is None:
# The Stop guard's feedback — a runner event, shown as a tag.
if self.full:
self.thread.append(_item("user", "held", False, 0, ts=ts,
text=_held_label(c)))
elif c.strip() and not c.startswith("<"):
self.interrupted_trailing = self.paused_trailing = False # real activity
self.waiting_trailing = False
self.user_turns += 1
if self.first_user is None:
self.first_user = c.strip()
mi = self._add_search("user", c)
if self.full:
it = _item("user", "text", False, 0,
text=_clip(c), mi=mi, ts=ts)
if self._first_user_item is None:
self._first_user_item = it
self.thread.append(it)
elif isinstance(c, list):
for b in c:
if not isinstance(b, dict):
continue
if b.get("type") == "tool_result":
rtext = _block_text(b.get("content"))
# A Monitor answering means its wait is over (a
# ScheduleWakeup's result only confirms the schedule —
# that one stays parked until the wake-up fires).
if self.tool_name.get(b.get("tool_use_id")) == "Monitor":
self.waiting_trailing = False
# A blocking agent's result is its report → it's done. A
# background one only reports its launch here, and ends on
# the `<task-notification>` handled above.
call = self.agent_calls.get(b.get("tool_use_id"))
if call is not None and not _AGENT_ASYNC_RE.search(rtext):
self.agent_calls.pop(b["tool_use_id"], None)
elif call is not None:
call["background"] = True
if b.get("tool_use_id") == self._deploying_tool_use_id:
self.deploying_since = None
self._deploying_tool_use_id = None
if "node_type: memory" in rtext:
self.memories_read.update(_MEMORY_RE.findall(rtext))
if self.full:
# Nest the result inside its originating tool_use card
# (the UI collapses it there) rather than as its own row.
tu = self.tool_items.get(b.get("tool_use_id"))
if tu is not None:
tu["result"] = _clip(rtext)
# How long the tool ran: the gap between the call
# and the record carrying its result.
tu["durationMs"] = _elapsed_ms(tu.get("ts"), ts)
tu["resultChars"] = len(rtext)
tu["resultLines"] = rtext.count("\n") + 1 if rtext else 0
if b.get("is_error"):
tu["isError"] = True
elif "deployed" in _lifecycle_stages(
tu.get("name"), tu.get("input") or {}):
url = _deployed_url(rtext)
if url:
tu["deployedUrl"] = url
else:
# Orphan result (call not captured) — keep a row.
self.thread.append(_item(
"tool", "tool_result", False, 0,
name=self.tool_name.get(b.get("tool_use_id"), "tool"),
text=_clip(rtext),
isError=bool(b.get("is_error")),
ts=ts,
))
elif b.get("type") == "text" and b.get("text", "").strip():
sk = _skill_body(b["text"])
if sk is not None:
if self.full and self.pending_skill is not None:
base_dir, body = sk
self.pending_skill["skillBody"] = _clip(body)
rel = _rel_skill_path(base_dir, self.cwd)
if rel:
self.pending_skill["skillPath"] = rel
self.pending_skill = None
elif _INTERRUPT_RE.match(b["text"]):
# Synthetic interrupt marker: mark the run as trailing-
# interrupted and show a compact tag inline, but don't
# treat it as a real user turn / title / searchable msg.
self.interrupted_trailing = True
self.paused_trailing = self.waiting_trailing = False
if self.full:
self.thread.append(_item("user", "interrupted", False, 0,
ts=ts))
elif not b["text"].startswith("<"):
self.interrupted_trailing = self.paused_trailing = False # real
self.waiting_trailing = False
self.user_turns += 1
if self.first_user is None:
self.first_user = b["text"].strip()
mi = self._add_search("user", b["text"])
if self.full:
it = _item("user", "text", False, 0,
text=_clip(b["text"]), mi=mi, ts=ts)
if self._first_user_item is None:
self._first_user_item = it
self.thread.append(it)
def summary(self) -> dict:
"""Snapshot the summary (and, when ``full``, the thread) so far."""
title = self.ai_title or (self.first_user[:90] if self.first_user else None) or "(untitled)"
summary = {
"sessionId": self.session_id,
"cwd": self.cwd,
"project": _project_name(self.cwd, self.path),
"gitBranch": self.git_branch,
"model": self.model or "claude-opus-4-8",
# All models that produced a turn, busiest first — the conversation's
# model tags. Usually one; more when the session switched mid-thread.
"models": [m for m, _ in sorted(self.model_turns.items(),
key=lambda kv: -kv[1])],
# The `--effort` level(s) the turns actually ran at, busiest first.
# Empty for pi runs and for transcripts predating the CLI flag.
"efforts": [e for e, _ in sorted(self.effort_turns.items(),
key=lambda kv: -kv[1])],
# model id -> the 5-minute bins this run produced a turn in, sorted.
# Small by construction (bins are 5 min wide and runs are short), so
# it rides along in the stored summary rather than costing a re-walk.
"activeBins": {m: sorted(b) for m, b in self.active_bins.items()},
"title": title,
"userTurns": self.user_turns,
"assistantTurns": self.assistant_turns,
"messages": self.user_turns + self.assistant_turns,
"startedAt": self.started,
"endedAt": self.ended,
"tokens": usage_tokens(self.agg),
"cost": self.cost_total,
"usage": self.agg,
"contextTokens": self.context_tokens or None,
"contextModel": self.context_model,
# First API call's input side, split into injected context
# (system prompt, CLAUDE.md, memories, …) vs the user's actual
# first message (length-estimated — the API only reports the total).
"firstContextTokens": self.first_context_tokens,
"firstMessageTokens": self.first_message_tokens,
# Costs come from `by_tool_cost` — accumulated per turn at that turn's
# own model — not from re-pricing the summed usage with `self.model`,
# which overcharged every mixed-model session and left the breakdown
# disagreeing with the conversation total.
"byTool": {
b: {"tokens": usage_tokens(u),
"cost": self.by_tool_cost.get(b, 0.0),
"output": u["output"]}
for b, u in sorted(
self.by_tool.items(),
key=lambda kv: TOOL_BUCKET_ORDER.index(kv[0])
if kv[0] in TOOL_BUCKET_ORDER else 99)
},
# Same buckets, one level deeper (bash→program, read/edit→ext, msg→
# prompt/replies, tool buckets→tool name), each sub sorted by tokens desc.
"byToolSub": {
b: {
sk: {"tokens": usage_tokens(u),
"cost": self.by_tool_sub_cost.get(b, {}).get(sk, 0.0),
"output": u["output"]}
for sk, u in sorted(subs.items(),
key=lambda kv: usage_tokens(kv[1]),
reverse=True)
}
for b, subs in sorted(
self.by_tool_sub.items(),
key=lambda kv: TOOL_BUCKET_ORDER.index(kv[0])
if kv[0] in TOOL_BUCKET_ORDER else 99)
},
"projectsAuto": sorted(self.projects),
"servicesAuto": sorted(self.services),
"worktreesAuto": list(self.worktrees.values()),
"tasks": list(self.tasks.values()),
"runningAgents": list(self.agent_calls.values()),
"deploying": bool(self.deploying_since),
"lifecycleAuto": self.lifecycle,
"memoriesRead": sorted(self.memories_read),
"doneMarker": bool(_DONE_RE.search(self.last_assistant_text)),
"completedMarker": bool(_COMPLETED_RE.search(self.last_assistant_text)),
"interruptedByUser": self.interrupted_trailing,
"pausedByUser": self.paused_trailing,
"waitingForTrigger": self.waiting_trailing,
"isSidechain": self.is_sidechain,
"agentId": self.agent_id,
"skills": self.skills,
"searchMsgs": self.search_msgs,
}
if self.full:
summary["thread"] = self.thread
summary["turns"] = self.turns
return summary
def parse_conversation(path: pathlib.Path, full: bool = False) -> dict:
"""Walk one transcript. Returns a summary; with full=True also the thread."""
st = ParserState(path, full=full)
for o in _iter_records(path):
st.feed(o)
return st.summary()
def find_cut_line(path: pathlib.Path, mi: int) -> dict | None:
"""Locate the raw transcript line holding user message ordinal ``mi``.
Forking a conversation at a message needs the JSONL line where that
message lives. ``mi`` is assigned by ``ParserState.feed`` (one per
user/assistant text block, in transcript order), so instead of
reimplementing that rule — and silently drifting from it — this walks the
file through a throwaway ``ParserState`` and watches ``search_msgs`` grow.
Returns ``{"line", "text", "userOrd"}`` for the record that produced
message ``mi`` — ``line`` is 0-based over *raw* file lines (so
``lines[:line]`` is the history strictly before it), ``text`` the
search-clipped message text (a canary for the sidecar's own slice), and
``userOrd`` this message's 0-based ordinal among *user* messages only
(how the pi runner finds the same point in pi's native session file).
Returns ``None`` when ``mi`` doesn't exist or belongs to an assistant
message (only user messages are fork points).
"""
st = ParserState(path, full=False)
try:
with path.open(encoding="utf-8", errors="ignore") as fh:
for i, line in enumerate(fh):
s = line.strip()
if not s or s[0] != "{":
continue
try:
o = json.loads(s)
except ValueError:
continue
before = len(st.search_msgs)
st.feed(o)
if before <= mi < len(st.search_msgs):
m = st.search_msgs[mi]
if m.get("role") != "user":
return None
return {"line": i, "text": m.get("text") or "",
"userOrd": sum(1 for e in st.search_msgs[:mi]
if e.get("role") == "user")}
except OSError:
return None
return None
def _item(role, kind, first, out_tokens, **rest) -> dict:
it = {"role": role, "kind": kind,
**{k: v for k, v in rest.items() if v is not None}}
if first and out_tokens:
it["out"] = out_tokens # output tokens for this assistant API call
return it
def _project_name(cwd: str | None, path: pathlib.Path) -> str:
"""A short human label for the conversation's project."""
if cwd:
return "/".join(p for p in cwd.split("/")[-2:] if p)
# fall back to the encoded directory name (…-homelab-projects-foo)
return path.parent.name