Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 14 additions & 3 deletions backend/app/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -116,9 +116,11 @@ def _approval_review_origin_from_text(text: str | None) -> str | None:
-- One row per event that used an MCP server or loaded a skill, derived from
-- events by extension_usage.sync() so the Extensions view never scans the
-- payload-heavy events table. Names are what the event said, and they resolve to
-- installed extensions at read time, since installs change.
-- installed extensions at read time, since installs change. One shell command
-- can read several SKILL.md files, so ord numbers the uses within an event.
CREATE TABLE IF NOT EXISTS extension_uses (
event_id INTEGER PRIMARY KEY REFERENCES events(id) ON DELETE CASCADE,
event_id INTEGER NOT NULL REFERENCES events(id) ON DELETE CASCADE,
ord INTEGER NOT NULL DEFAULT 0,
session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
source TEXT NOT NULL,
kind TEXT NOT NULL,
Expand All @@ -129,7 +131,8 @@ def _approval_review_origin_from_text(text: str | None) -> str | None:
phase TEXT,
status TEXT,
duration_ms INTEGER,
ts TEXT NOT NULL
ts TEXT NOT NULL,
PRIMARY KEY (event_id, ord)
);
CREATE INDEX IF NOT EXISTS idx_extension_uses_session ON extension_uses(session_id);

Expand Down Expand Up @@ -361,6 +364,14 @@ def _migrate(conn: sqlite3.Connection) -> None:
"CREATE INDEX IF NOT EXISTS idx_sessions_parent ON sessions(parent_session_id)"
)

if "ord" not in _table_columns(conn, "extension_uses"):
# Keyed by event_id alone before ord existed. The table is derived from
# events, so rebuild it; extension_usage.sync() refills it.
conn.execute("DROP TABLE extension_uses")
for statement in SCHEMA.split(";"):
if "extension_uses" in statement:
conn.execute(statement)

if "source" in session_cols:
conn.execute(
"UPDATE sessions SET source = ("
Expand Down
86 changes: 69 additions & 17 deletions backend/app/extension_usage.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,16 @@
``extension_uses`` table, past a watermark, so reads never scan the
payload-heavy events table. ``overview()`` and ``detail()`` join that usage to
the on-disk inventory from :mod:`extensions`, resolving what each event said
(``mcp__github__create_issue``, ``Skill pptx``, a ``SKILL.md`` read, a
``/command`` prompt) to an installed extension.
(``mcp__github__create_issue``, ``Skill pptx``, a ``SKILL.md`` read or
``cat``, a ``/command`` prompt) to an installed extension.
"""

from __future__ import annotations

import json
import os
import re
import shlex
import sqlite3
import statistics
import threading
Expand All @@ -25,14 +26,18 @@
from . import extensions, marketplace_scan, store, timeutil

# Bump to rebuild extension_uses from scratch (e.g. when classify() changes).
USAGE_VERSION = "1"
USAGE_VERSION = "2"
_WATERMARK_KEY = "extension_uses_watermark"
_VERSION_KEY = "extension_uses_version"

_SLASH = re.compile(r"^\s*/([A-Za-z0-9][\w.:-]*)(?=\s|$)")
_COMMAND_TAG = re.compile(r"<command-name>/?([\w.:-]+)</command-name>")
# A skill read right after the Skill tool loaded it is the same load.
_FOLD_WINDOW_S = 180
# Codex has no Skill tool: it loads a skill by printing its SKILL.md.
_SHELL_READERS = {"cat", "head", "tail", "sed", "awk", "less", "more", "bat", "nl", "read"}
_SHELL_WRAPPERS = {"rtk", "proxy", "sudo", "command", "env", "exec"}
_SHELLS = {"sh", "bash", "zsh"}

_sync_lock = threading.Lock()

Expand All @@ -47,8 +52,52 @@ def _parse_mcp_tool(tool: str) -> tuple[str | None, str | None]:
return None, None


def classify(row: dict[str, Any]) -> dict[str, Any] | None:
"""What one event says about extension use, or None."""
def _shell_skill_reads(command: str) -> list[str]:
"""SKILL.md paths a shell command prints, e.g. ``cat a/SKILL.md b/SKILL.md``."""
lexer = shlex.shlex(command, posix=True, punctuation_chars=True)
lexer.whitespace_split = True
try:
tokens = list(lexer)
except ValueError: # unbalanced quotes
tokens = command.split()
segments: list[list[str]] = [[]]
for tok in tokens:
if tok and set(tok) <= set(";&|()"):
segments.append([])
else:
segments[-1].append(tok)
paths: list[str] = []
for words in segments:
while words and (words[0] in _SHELL_WRAPPERS or re.match(r"^\w+=", words[0])):
words = words[1:]
if not words:
continue
prog = os.path.basename(words[0])
if prog in _SHELLS:
for i, w in enumerate(words[1:-1], 1):
if w.startswith("-") and "c" in w:
paths.extend(_shell_skill_reads(words[i + 1]))
break
elif prog in _SHELL_READERS:
paths.extend(os.path.expanduser(w) for w in words[1:] if w.lower().endswith("/skill.md"))
return list(dict.fromkeys(paths))


def classify(row: dict[str, Any]) -> list[dict[str, Any]]:
"""What one event says about extension use: usually one use, or none."""
found = _classify_one(row)
if found is not None:
return [found]
if row.get("category") == "shell":
command = row.get("command") or row.get("target") or ""
return [
{"kind": "skill", "name": Path(p).parent.name, "path": p, "via": "read"}
for p in _shell_skill_reads(command)
]
return []


def _classify_one(row: dict[str, Any]) -> dict[str, Any] | None:
tool = row.get("tool") or ""
target = row.get("target") or ""
category = row.get("category") or ""
Expand Down Expand Up @@ -88,9 +137,14 @@ def classify(row: dict[str, Any]) -> dict[str, Any] | None:
_SYNC_SELECT = (
"SELECT id, session_id, source, tool, category, target, phase, status, duration_ms, ts,"
" CASE WHEN tool = 'CallMcpTool' THEN payload END AS payload,"
" CASE WHEN category = 'prompt' THEN substr(detail, 1, 400) END AS detail"
" CASE WHEN category = 'prompt' THEN substr(detail, 1, 400) END AS detail,"
# target is clipped at 120 chars; the full shell command is in detail.
" CASE WHEN category = 'shell' AND target LIKE '%…' AND json_valid(detail)"
" THEN json_extract(detail, '$.command') END AS command"
" FROM events"
" WHERE category IN ('mcp', 'memory', 'web', 'context_read', 'file_read', 'prompt') AND id > ?"
" WHERE category IN ('mcp', 'memory', 'web', 'context_read', 'file_read', 'prompt', 'shell') AND id > ?"
" AND (category != 'shell' OR target LIKE '%skill.md%'"
" OR (target LIKE '%…' AND detail LIKE '%skill.md%'))"
)


Expand Down Expand Up @@ -131,22 +185,20 @@ def sync(batch: int = 2000) -> int:
chunk = rows[start : start + batch]
values = []
for r in chunk:
c = classify(r)
if c is None:
continue
values.append((
r["id"], r["session_id"], r["source"], c["kind"], c.get("name"), c.get("tool"),
c.get("path"), c["via"], r["phase"], r["status"], r["duration_ms"], r["ts"],
))
for ord_, c in enumerate(classify(r)):
values.append((
r["id"], ord_, r["session_id"], r["source"], c["kind"], c.get("name"), c.get("tool"),
c.get("path"), c["via"], r["phase"], r["status"], r["duration_ms"], r["ts"],
))
with store.write() as conn:
if values:
conn.executemany(
"INSERT OR REPLACE INTO extension_uses (event_id, session_id, source, kind, name,"
"INSERT OR REPLACE INTO extension_uses (event_id, ord, session_id, source, kind, name,"
" tool, path, via, phase, status, duration_ms, ts)"
" SELECT ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?"
" SELECT ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?"
" WHERE EXISTS (SELECT 1 FROM events WHERE id = ?)"
" AND EXISTS (SELECT 1 FROM sessions WHERE id = ?)",
[(*v, v[0], v[1]) for v in values],
[(*v, v[0], v[2]) for v in values],
)
_settings_set(conn, _WATERMARK_KEY, str(chunk[-1]["id"]))
added += len(values)
Expand Down
28 changes: 28 additions & 0 deletions backend/tests/test_extensions.py
Original file line number Diff line number Diff line change
Expand Up @@ -190,6 +190,34 @@ def test_usage_attribution(homes):
assert used == {"mcp:github", "skill:tidy", "mcp:plugin_myplug_srv", "skill:myplug:deploy"}


def test_codex_shell_skill_reads_count_as_loads(homes):
agents = homes["home"] / ".agents" / "skills"
_skill(agents, "code-review")
_skill(agents, "unslop")
extensions.clear_cache()
_session("x1", "codex", "/p")
# Codex loads skills by printing SKILL.md; the stored target is clipped.
command = (f"cat {homes['home']}/.codex/RTK.md {agents}/code-review/SKILL.md"
f" {agents}/unslop/SKILL.md {agents}/tidy/SKILL.md")
for phase in ("start", "end"):
_event("x1", "codex", tool="exec_command", category="shell", phase=phase,
target=command[:119] + "…", detail=json.dumps({"command": command}))
_event("x1", "codex", tool="exec_command", category="shell",
target=f"rtk proxy sh -c 'sed -n 1,80p {agents}/unslop/SKILL.md'", ts="2026-10-01T12:00:00+00:00")
# Not loads: writing a skill, listing skills, searching one.
_event("x1", "codex", tool="exec_command", category="shell", target=f"vim {agents}/tidy/SKILL.md")
_event("x1", "codex", tool="exec_command", category="shell", target=f"ls {agents}/tidy/SKILL.md")
_event("x1", "codex", tool="exec_command", category="shell", target=f"grep name {agents}/tidy/SKILL.md")

items = _by_key(extension_usage.overview(refresh=True))
assert items["skill:code-review"]["usage"]["calls"] == 1
assert items["skill:code-review"]["usage"]["agents"] == {"codex": 1}
assert items["skill:unslop"]["usage"]["calls"] == 2
assert items["skill:tidy"]["usage"]["calls"] == 1
used = {e["key"] for e in extension_usage.session_extensions("x1")}
assert used == {"skill:code-review", "skill:unslop", "skill:tidy"}


def test_sync_is_incremental_and_cascades(homes):
_session("c1", "claude", "/p")
_event("c1", "claude", tool="mcp__github__x", category="mcp", target="github/x")
Expand Down
Loading