Files
ai-agent/backend/indexer.py
Gabriel Vidal 6ae5290453 feat(ai-agent): effort select for claude runs, thinking toggle for pi
The composer's third toolbar chip is now harness-aware: claude models get an
`--effort` select (low|medium|high|xhigh|max) — the CLI's own scale — while pi,
which has no such flag, keeps the binary thinking toggle. The two are one knob
shown per harness, so only the active one is ever sent.

Effort is also a read-side tag: the claude CLI stamps an `effort` field on each
assistant record, so conversations.py mines it into `efforts` (busiest first,
like `models`) and EffortTags renders it on conversation rows and in the detail
header, behind a new "Effort tags" visibility switch.

- sidecar: `_effort_args` (claude-only, allow-list validated) on spawn + resume
- backend: effort on Spawn/ResumeBody, threaded through _send_message and
  stamped into the meta sidecar; resume reuses the level it last ran with
- parser: PARSER_VERSION 23 -> 24 to back-fill `efforts` on past transcripts
- mock: seed cycles every level (+ the untagged case); spawn/resume honour it

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-01 15:53:17 +02:00

425 lines
19 KiB
Python

"""
Background indexing for the ai-agent backend.
Two jobs, both run on a single daemon thread so they never block requests:
1. **Token + cost accounting** — scan the exposed context files, store their
metadata in SQLite, and compute the *real* Claude token count (via the
count_tokens endpoint) + dollar cost. To use the API sparingly we only count
a file when its content hash changed AND it has been stable for at least
``DEBOUNCE_SECS`` (so we don't burn a call on every keystroke of an edit).
2. **Skill-usage analytics** — mine the Claude Code transcripts (``*.jsonl``)
for ``Skill`` tool calls, counting invocations and last-used time per skill.
Parsing is incremental: a transcript is only re-read when its mtime/size
changed, and per-file contributions are summed into the ``skills`` table.
"""
import hashlib
import json
import math
import os
import pathlib
import shutil
import threading
import time
import conversations
import fsscan
from db import Store
# Bump when parse_conversation's stored summary shape changes (forces re-parse).
# v2: archive-backed indexing + byTool / projectsAuto fields in the summary.
# v3: lifecycleAuto (CI/CD stamps derived from transcript script/command usage).
# v4: doneMarker (last message ends with DONE → conversation finished).
# v5: memoriesRead (slugs of memories recalled into the transcript).
# v6: searchMsgs (ordered user/assistant text blocks for full-text search).
# v7: interruptedByUser (trailing "[Request interrupted by user]" → interrupted).
# v8: isSidechain/agentId (subagent transcripts flagged, hidden from lists,
# linked from the parent's Task/Agent card + usage rolled into the parent).
# v9: worktreeAuto (name of the git worktree mined from a new-worktree.sh call).
# v10: worktreesAuto (list of created/removed worktrees w/ create+remove dates).
# v11: worktree parse anchored to command position + name validation (drop noise).
# v12: projectsAuto/servicesAuto tag only projects a conversation *edited* (target
# path of Write/Edit/MultiEdit/NotebookEdit), not every cwd/read/bash mention.
# v13: byToolSub (second-level per-bucket usage breakdown).
# v14: pausedByUser (trailing "No response requested." → paused); interrupt +
# no-response markers hidden from the thread (rendered as compact tags).
# v15: models (every model that produced a turn, busiest first) → model tags on
# past conversations, not just the first model seen.
# v16: `<synthetic>` (Claude Code's placeholder model on messages it fabricates)
# never counts as a model the conversation ran on.
# v20: cost accumulates per assistant record — a pi record's runner-mirrored
# `costUSD` (real provider cost) beats the price-table estimate, and bare
# (local EVOX2) qwen ids price at 0.
# v21: first-turn context split — the first API call's input side divided into
# injected context (system prompt, CLAUDE.md, memories) vs the user's
# actual message: new top-level "context" activity bucket +
# firstContextTokens/firstMessageTokens + per-item turnContext/MessageTokens.
PARSER_VERSION = "24"
MODEL = os.environ.get("TOKEN_MODEL", "claude-opus-4-8")
# Input-token price in USD per million tokens (Opus 4.8 = $5.00 / 1M).
INPUT_PRICE_PER_MTOK = float(os.environ.get("INPUT_PRICE_PER_MTOK", "5.0"))
DEBOUNCE_SECS = float(os.environ.get("TOKEN_DEBOUNCE_SECS", "30"))
SCAN_INTERVAL = float(os.environ.get("INDEX_INTERVAL_SECS", "60"))
# How often the workspace is re-walked to *discover* exposed files. The walk
# visits every project dir to find ~280 files, so it is far too expensive to
# redo each SCAN_INTERVAL; between walks the cached list is just re-stat'd
# (which is what catches edits). A save forces a walk, so the UI never waits.
WALK_INTERVAL_SECS = float(os.environ.get("WALK_INTERVAL_SECS", "120"))
RETRY_ESTIMATE_SECS = 300.0 # retry API-failed (estimated) files at most this often
MD_EXTS = {".md", ".markdown"}
def _price(tokens: int) -> float:
return tokens / 1_000_000 * INPUT_PRICE_PER_MTOK
def _estimate_tokens(text: str) -> int:
"""Cheap fallback used before/instead of an API count (~chars/4)."""
return math.ceil(len(text) / 4) if text else 0
class Indexer:
def __init__(self, store: Store, workspace: pathlib.Path,
source_dir,
archive_dir: pathlib.Path | None, iter_files, read_text):
self.store = store
self.workspace = workspace
# Live transcript sources (read-only mounts): the Claude Code projects
# dir, plus the pi runner's mirror dir. A single path or a list — both
# are imported into the same archive (session UUIDs keep them distinct).
if source_dir is None:
self.source_dirs = []
elif isinstance(source_dir, (list, tuple)):
self.source_dirs = [d for d in source_dir if d]
else:
self.source_dirs = [source_dir]
# Our persistent copy: transcripts are imported here so deleted ones
# survive. Indexing + the detail view read from the archive.
self.archive_dir = archive_dir
# Back-compat alias used by the conversation detail endpoint.
self.transcripts_dir = archive_dir or (
self.source_dirs[0] if self.source_dirs else None)
self._iter_files = iter_files # () -> yields (abs_path, rel_posix)
self._read_text = read_text # (abs_path) -> str
self._client = None
self._client_tried = False
self._stop = threading.Event()
# Resumable parser states for *growing* transcripts: abs path →
# (byte offset of the next unread line, ParserState). A live session's
# transcript changes every couple of seconds; feeding only the appended
# lines keeps the per-tick parse cost flat instead of O(file size).
# In-memory only — a restart simply re-parses once. LRU-capped: only a
# handful of transcripts grow concurrently.
self._live: dict[str, tuple[int, "conversations.ParserState"]] = {}
self._live_cap = 8
# Cached result of the workspace walk (see _exposed_files).
self._files_cache: list | None = None
self._walked_at = 0.0
# ── Anthropic client (lazy) ──────────────────────────────────────────
def _anthropic(self):
if self._client_tried:
return self._client
self._client_tried = True
if not os.environ.get("ANTHROPIC_API_KEY"):
return None
try:
import anthropic
self._client = anthropic.Anthropic()
except Exception:
self._client = None
return self._client
def _count_tokens_api(self, text: str) -> int | None:
if not text.strip():
return 0
client = self._anthropic()
if client is None:
return None
try:
r = client.messages.count_tokens(
model=MODEL,
messages=[{"role": "user", "content": text}],
)
return int(r.input_tokens)
except Exception:
return None
# ── file scan + token accounting ─────────────────────────────────────
def _exposed_files(self, force_walk: bool) -> list:
"""The exposed context files, from a cached directory walk.
Discovering them means walking the whole workspace (every project dir),
which is far more expensive than the ~280 files it finds — and it was
being redone on every scan interval. The file *list* changes only when
a file is added or removed, so the walk is cached and refreshed at most
every ``WALK_INTERVAL_SECS`` (or on demand, e.g. right after a save)."""
now = time.time()
if (self._files_cache is None or force_walk
or (now - self._walked_at) >= WALK_INTERVAL_SECS):
self._files_cache = list(self._iter_files())
self._walked_at = now
return self._files_cache
def scan_files(self, force_walk: bool = False) -> None:
"""Record on-disk metadata for every exposed file.
Stat-gated: a file whose (mtime, size) matches its stored row is
skipped without being read, so a scan costs one stat per exposed file
— not a read+sha256 of the whole context tree. Word and char counts are
stored alongside (recomputed only on change) so ``/api/bundle`` can be
served from the DB without touching disk."""
seen: set[str] = set()
rows = self.store.all_files()
for ap, rel in self._exposed_files(force_walk):
relstr = str(rel)
try:
st = ap.stat()
except OSError:
continue # vanished since the walk — let the next walk prune it
seen.add(relstr)
row = rows.get(relstr)
if (row is not None and row["mtime"] == st.st_mtime
and row["size"] == st.st_size
and row["words"] is not None):
continue # unchanged on disk — nothing to record
try:
data = ap.read_bytes()
except OSError:
continue
digest = hashlib.sha256(data).hexdigest()
text = data.decode("utf-8", "ignore")
self.store.upsert_file_meta(
relstr, digest, st.st_size, st.st_mtime,
rel.suffix.lower() in MD_EXTS,
words=len(text.split()), chars=len(text),
)
self.store.prune_files(seen)
def count_pending(self) -> int:
"""Compute real token counts for changed/stale files past the debounce
window. Returns how many files were (re)counted."""
now = time.time()
done = 0
for path, row in self.store.all_files().items():
stale = bool(row["estimated"]) or row["tokens"] is None
if not stale:
continue
# debounce: only spend a call once the file has settled
if row["mtime"] and (now - row["mtime"]) < DEBOUNCE_SECS:
continue
# don't hammer the API on files that keep failing
if (row["tokens"] is not None and row["counted_at"]
and (now - row["counted_at"]) < RETRY_ESTIMATE_SECS):
continue
ap = self.workspace / path
text = self._read_text(ap)
tokens = self._count_tokens_api(text)
if tokens is None:
tokens = _estimate_tokens(text)
self.store.set_tokens(path, tokens, _price(tokens), estimated=True)
else:
self.store.set_tokens(path, tokens, _price(tokens), estimated=False)
done += 1
return done
# ── transcript import (persist live transcripts before Claude prunes) ──
def _copy_in(self, src: pathlib.Path, rel: pathlib.PurePath,
mtime: float, size: int, changed: list[str]) -> None:
"""Copy one source transcript into the archive if it's new or grew."""
dst = self.archive_dir / rel
try:
d = dst.stat()
if d.st_size >= size and d.st_mtime >= mtime:
return
except OSError:
pass # not in the archive yet
try:
dst.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(src, dst)
changed.append(rel.as_posix())
except OSError:
pass
def import_transcripts(self, only: list[str] | None = None) -> list[str]:
"""Copy new/grown source transcripts into the archive. Append-only, so
we copy when the source is larger or newer; we never delete archived
copies whose source has vanished. Returns the rel paths that changed.
Also copies each subagent's tiny `agent-<id>.meta.json` sidecar (it
carries the `toolUseId` that links a subagent back to the parent's
Task/Agent card, so the detail endpoint can cross-link them).
``only`` restricts the import to the given source-relative paths — the
watcher already knows exactly which transcripts changed, so the live
path doesn't re-walk thousands of files every couple of seconds. A
sibling `.meta.json` is pulled in alongside each named `.jsonl`; the
periodic full pass (``only=None``) catches anything else."""
if not (self.source_dirs and self.archive_dir):
return []
changed: list[str] = []
for source_dir in self.source_dirs:
if not source_dir.is_dir():
continue
if only is None:
for suffix in (".jsonl", ".meta.json"):
for path, mtime, size in fsscan.walk_files(source_dir, suffix):
src = pathlib.Path(path)
self._copy_in(src, src.relative_to(source_dir),
mtime, size, changed)
continue
for rel in only:
relp = pathlib.PurePosixPath(rel)
cands = [relp]
if relp.suffix == ".jsonl": # its subagent metadata sidecar
cands.append(relp.with_suffix(".meta.json"))
for cand in cands:
src = source_dir / cand
try:
st = src.stat()
except OSError:
continue
self._copy_in(src, cand, st.st_mtime, st.st_size, changed)
return changed
# ── transcript analytics (skills + per-conversation summaries) ────────
@staticmethod
def _feed_lines(state, ap, offset: int) -> int | None:
"""Feed the complete JSONL lines from ``offset`` to EOF into ``state``.
Returns the offset of the next unread byte — a trailing partial line
(a record still being written) is left for the next pass. ``None`` on
read errors. Mirrors ``conversations._iter_records`` semantics."""
try:
with ap.open("rb") as fh:
fh.seek(offset)
buf = fh.read()
except OSError:
return None
end = buf.rfind(b"\n")
if end == -1:
return offset # no complete new line yet
for raw in buf[:end].split(b"\n"):
line = raw.decode("utf-8", "ignore").strip()
if not line or line[0] != "{":
continue
try:
state.feed(json.loads(line))
except ValueError:
continue
return offset + end + 1
def _parse_summary(self, ap, size: int) -> dict:
"""Summary of a transcript — incremental when we saw it grow before."""
relstr = str(ap)
live = self._live.pop(relstr, None) # pop = LRU re-insert on success
if live is not None:
offset, state = live
if size >= offset:
new_off = self._feed_lines(state, ap, offset)
if new_off is not None:
self._live[relstr] = (new_off, state)
return state.summary()
# shrank/rewritten or unreadable → fall through to a full parse
state = conversations.ParserState(ap, full=False)
off = self._feed_lines(state, ap, 0)
if off is None:
return conversations.parse_conversation(ap, full=False)
self._live[relstr] = (off, state)
while len(self._live) > self._live_cap: # evict least-recently-grown
self._live.pop(next(iter(self._live)))
return state.summary()
def scan_transcripts(self, only: list[str] | None = None) -> None:
"""Import + (re)parse transcripts into the store.
``only`` is the source-relative paths the watcher saw change: the live
path then imports and re-parses just those, instead of re-walking every
source dir and the whole archive on each 2-second tick. The indexer's
own periodic pass (``only=None``) still does the full sweep, which is
what discovers files the watcher missed and prunes deleted ones."""
self.import_transcripts(only)
if not self.transcripts_dir or not self.transcripts_dir.is_dir():
return
# Force a full re-parse when the parser changes shape (new fields).
if self.store.get_meta("parser_version") != PARSER_VERSION:
self.store.prune_transcripts(set())
self.store.set_meta("parser_version", PARSER_VERSION)
self._live.clear()
only = None # everything must be re-parsed
rows = self.store.all_transcripts()
if only is not None:
targets = []
for rel in only:
ap = self.transcripts_dir / rel
try:
st = ap.stat()
except OSError:
continue
targets.append((str(ap), st.st_mtime, st.st_size))
else:
targets = list(fsscan.walk_files(self.transcripts_dir, ".jsonl"))
seen: set[str] = set()
changed = False
for relstr, mtime, size in targets:
seen.add(relstr)
row = rows.get(relstr)
if row and row["mtime"] == mtime and row["size"] == size:
continue
summary = self._parse_summary(pathlib.Path(relstr), size)
skills = summary.pop("skills", {})
self.store.upsert_transcript(relstr, mtime, size, skills, summary)
changed = True
# Only a full sweep knows what's gone; a targeted pass must not prune.
if only is None and rows.keys() - seen:
self.store.prune_transcripts(seen)
changed = True
if changed:
self._aggregate_skills()
def _aggregate_skills(self) -> None:
agg: dict[str, dict] = {}
for skills_json in self.store.all_transcript_skills():
try:
contrib = json.loads(skills_json or "{}")
except ValueError:
continue
for name, d in contrib.items():
e = agg.setdefault(name, {"count": 0, "last_used": None})
e["count"] += d.get("count", 0)
last = d.get("last")
if last and (e["last_used"] is None or last > e["last_used"]):
e["last_used"] = last
self.store.replace_skills(agg)
# ── lifecycle ────────────────────────────────────────────────────────
def run_once(self) -> None:
"""One full sweep: re-stat the context files, count what's pending, and
do the complete transcript discovery + prune pass."""
self.scan_files()
self.count_pending()
self.scan_transcripts()
def _loop(self) -> None:
# The watcher already handles live transcript changes (targeted, every
# couple of seconds). This loop is the slower full sweep: it re-stats
# the context files for token counting and re-walks the transcript dirs
# to discover/prune — work that is wasted at a 10-second cadence.
while not self._stop.is_set():
try:
self.run_once()
except Exception:
pass
self._stop.wait(SCAN_INTERVAL)
def start(self) -> None:
threading.Thread(target=self._loop, daemon=True).start()
def stop(self) -> None:
self._stop.set()