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
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -495,6 +495,7 @@ release without losing the old one while you verify):
| `POST /api/repos/toggle` | admin bearer token | `{name, ref, enabled}` | Flip searchability without touching indexed data |
| `POST /api/repos/reindex` | admin bearer token | `{name, ref}` | Re-run ingest for an existing entry (submits a job) |
| `POST /api/repos/delete` | admin bearer token | `{name, ref, purge}` | Remove the entry; `purge: true` also deletes its `kg_nodes`/`kg_edges` rows |
| `GET /api/ingest-log?limit=&entry=` | `corpus:view` | – | Persistent ingest history, newest first: every run (UI, CLI, kubectl) folded per `run_id` with status `ok` / `failed` (+ error) / `running` / `stale`, nodes, edges, sha, started/finished, duration |
| `GET /api/jobs` | admin bearer token | — | Last 50 ingest jobs (`queued`/`running`/`ok`/`failed`), newest first |
| `GET /api/jobs/{id}` | admin bearer token | — | Single job status |

Expand Down
2 changes: 1 addition & 1 deletion docs/FEATURES.md
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ Manage what the graph indexes at runtime, without editing files or restarting.
- **Enable / disable.** A one-row flip gates whether an entry is searchable — instant, reversible, with a short-TTL live cache so changes take effect without a restart.
- **Explicit delete with purge.** Removing an entry can optionally scrub its `kg_nodes`/`kg_edges` rows; ingest history is never purged.
- **Background jobs.** Add / reindex operations run on a background worker; job status is pollable.
- **Manage tab** in the web UI drives all of it (add, toggle, reindex, delete).
- **Manage tab** in the web UI drives all of it (add, toggle, reindex, delete), and shows the **ingest history** — every run, whether started from the UI, `tpk ingest` or `kubectl exec`: status (`running` / `ok` / `failed` with the error / `stale` for a run that never wrote its end row), nodes, edges, commit, duration. Backed by the append-only `kg_ingest_log` stream via `GET /api/ingest-log`.
- **Release-upgrade workflow:** add the new ref → ingest → verify → flip enabled; a legacy migration re-keys old bare-name rows on first ingest.
- **Export / import (issue #29).** `tpk export` dumps the ingested graph + corpus registry to a portable bundle (proton `FORMAT Parquet`, one file per stream + manifest); `tpk import` loads it into a new environment — no re-ingest, no LLM calls. Idempotent upsert by default (`--replace` to reset); auth streams excluded.

Expand Down
18 changes: 11 additions & 7 deletions src/tpk/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,14 @@
from datetime import datetime, timezone
from pathlib import Path

from fastapi import APIRouter, Depends, HTTPException
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel

import tpk.auth as auth_mod
from tpk import corpus, db
from tpk.auth import AuthLayer, User
from tpk.config import RepoConfig, Settings, config_path, entry_key, load_llm
from tpk import ingest as ingest_mod
from tpk.ingest import ingest_repo

REPOS_TOML = config_path()
Expand Down Expand Up @@ -326,12 +327,7 @@ def list_repos(actor: User = Depends(auth.require_cap(auth_mod.CAP_CORPUS_VIEW))
f"SELECT repo, count() FROM {db.latest(db.qualified('kg_nodes', prefix))} GROUP BY repo"
).result_rows
)
last: dict[str, tuple] = {}
for repo, sha, status, t in client.query(
f"SELECT repo, arg_max(git_sha, _tp_time), arg_max(status, _tp_time),"
f" max(_tp_time) FROM table({db.qualified('kg_ingest_log', prefix)}) GROUP BY repo"
).result_rows:
last[repo] = (sha, status, t)
last = ingest_mod.latest_status(client, prefix=prefix)
out = []
for e in entries:
key = entry_key(e)
Expand Down Expand Up @@ -374,6 +370,14 @@ def reindex_repo(body: EntryRef, actor: User = Depends(auth.require_cap(auth_mod
raise HTTPException(404, "no such corpus entry")
return {"job_id": jobs.submit(cfg)}

@router.get("/ingest-log")
def api_ingest_log(limit: int = Query(50, ge=1, le=200), entry: str | None = None,
actor: User = Depends(auth.require_cap(auth_mod.CAP_CORPUS_VIEW))):
"""Persistent ingest history (#15) -- every run, whether started from
the UI, the CLI or kubectl, folded per run_id; `running` / `stale`
for runs that have not written their end row."""
return {"runs": ingest_mod.list_runs(_client(), prefix=prefix, limit=limit, entry=entry)}

@router.get("/jobs")
def list_jobs(actor: User = Depends(auth.require_cap(auth_mod.CAP_CORPUS_VIEW))):
return sorted(jobs.snapshot(), key=lambda j: j["submitted_at"], reverse=True)[:50]
Expand Down
4 changes: 4 additions & 0 deletions src/tpk/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,10 @@ def ensure_schema(client, prefix: str = "") -> None:
status string
)
""")
# Ingest history (#15): a `started` row is written when a run begins and
# an ok/failed row when it ends (same run_id); `error` carries the failure
# reason. Backfilled as "" on streams created before this column existed.
_add_column_if_missing(client, qualified("kg_ingest_log", prefix), "error", "string")
client.command(_keyed_stream(prefix, "kg_repos", [
"name string", "ref string", "github string", "path string",
"enabled bool", "visibility string", "extraction string",
Expand Down
102 changes: 97 additions & 5 deletions src/tpk/ingest.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
import uuid
from collections.abc import Callable
from dataclasses import dataclass
from datetime import datetime, timezone
from datetime import datetime, timedelta, timezone
from pathlib import Path

from tpk import db
Expand Down Expand Up @@ -70,11 +70,11 @@ def _git_sha(repo_path: Path) -> str:
return proc.stdout.strip() if proc.returncode == 0 else "unknown"


def _log(client, prefix: str, result: IngestResult, git_sha: str) -> None:
def _log(client, prefix: str, result: IngestResult, git_sha: str, error: str = "") -> None:
client.insert(
db.qualified("kg_ingest_log", prefix),
[[result.repo, result.run_id, result.nodes, result.edges, git_sha, result.status]],
column_names=["repo", "run_id", "nodes", "edges", "git_sha", "status"],
[[result.repo, result.run_id, result.nodes, result.edges, git_sha, result.status, error[:2000]]],
column_names=["repo", "run_id", "nodes", "edges", "git_sha", "status", "error"],
)


Expand All @@ -98,6 +98,15 @@ def _report(phase: str, message: str = "", nodes: int = 0, edges: int = 0) -> No
if on_progress is not None:
on_progress(IngestProgress(phase, message, nodes, edges))

# A `started` row makes the run visible (UI ingest history, #15) while it
# runs -- including CLI / kubectl runs the API's JobManager never sees. Its
# ok/failed twin below shares the run_id. Best effort: a log failure must
# not stop the ingest.
try:
_log(client, prefix, IngestResult(key, run_id, 0, 0, "started"), "")
except Exception as exc: # noqa: BLE001
print(f"[tpk] could not log ingest start for {key}: {exc}")
error = ""
try:
if repo_cfg.github:
_report("fetch")
Expand Down Expand Up @@ -129,8 +138,91 @@ def _report(phase: str, message: str = "", nodes: int = 0, edges: int = 0) -> No
except Exception as exc: # per-repo isolation: never propagate, never touch prior rows
print(f"[tpk] ingest failed for {key}: {exc}")
result = IngestResult(key, run_id, 0, 0, "failed")
error = f"{type(exc).__name__}: {exc}"
sha = _git_sha(repo_path) if repo_path else "unknown"
if repo_cfg.ref:
sha = f"{sha} ({repo_cfg.ref})"
_log(client, prefix, result, sha)
_log(client, prefix, result, sha, error)
return result


# -- ingest history (#15) ---------------------------------------------------
# A `started` row with no ok/failed twin older than this is treated as
# abandoned: the process died mid-run (OOM-kill, pod restart) and will never
# write the end row.
STALE_AFTER = timedelta(hours=2)
_LOG_COLUMNS = ["repo", "run_id", "nodes", "edges", "git_sha", "status", "error", "_tp_time"]


def _iso(dt) -> str | None:
return dt.replace(tzinfo=timezone.utc).isoformat() if dt is not None else None


def list_runs(client, prefix: str = "", limit: int = 50, entry: str | None = None,
now: datetime | None = None) -> list[dict]:
"""Ingest runs, newest first, folding each run's `started` and ok/failed
rows (same run_id) into one record. Rows written before #15 have no
`started` twin: they come back with started_at / duration_s = None."""
now = now or datetime.now(timezone.utc)
clauses, params = [], {}
if entry:
clauses.append("repo = %(repo)s")
params["repo"] = entry
where = f" WHERE {' AND '.join(clauses)}" if clauses else ""
# Newest rows first; a run's two rows are seconds-to-minutes apart, so a
# generous window (4x the page) is enough to see both for every run on the
# page -- a truncated run at the very end just lacks its start time.
rows = client.query(
f"SELECT {', '.join(_LOG_COLUMNS)} FROM table({db.qualified('kg_ingest_log', prefix)})"
f"{where} ORDER BY _tp_time DESC LIMIT %(n)s",
parameters={**params, "n": max(limit * 4, 200)},
).result_rows
runs: dict[str, dict] = {}
order: list[str] = []
for repo, run_id, nodes, edges, sha, status, error, at in rows:
at = at.replace(tzinfo=timezone.utc)
rec = runs.get(run_id)
if rec is None:
rec = runs[run_id] = {"run_id": run_id, "entry_key": repo, "status": None, "nodes": 0,
"edges": 0, "git_sha": "", "error": "", "started_at": None,
"finished_at": None, "duration_s": None, "_sort": at}
order.append(run_id)
if status == "started":
rec["started_at"] = at
else:
rec.update(status=status, nodes=nodes, edges=edges, git_sha=sha, error=error or "",
finished_at=at)
out = []
for run_id in order[:limit]:
rec = runs[run_id]
if rec["status"] is None: # only a `started` row
age = now - rec["started_at"]
if age > STALE_AFTER:
rec["status"] = "stale"
mins = int(age.total_seconds() // 60)
ago = f"{mins // 60}h {mins % 60}m" if mins >= 120 else f"{mins} min"
rec["error"] = f"started {ago} ago and did not finish (no end row)"
else:
rec["status"] = "running"
if rec["started_at"] and rec["finished_at"]:
rec["duration_s"] = round((rec["finished_at"] - rec["started_at"]).total_seconds(), 1)
rec["started_at"], rec["finished_at"] = _iso(rec["started_at"]), _iso(rec["finished_at"])
rec.pop("_sort")
out.append(rec)
return out


def latest_status(client, prefix: str = "", now: datetime | None = None) -> dict[str, tuple]:
"""Per entry key: (git_sha, status, time) of the most recent row, with an
abandoned `started` reported as "stale" (same rule as list_runs)."""
now = now or datetime.now(timezone.utc)
out: dict[str, tuple] = {}
for repo, sha, status, t in client.query(
f"SELECT repo, arg_max(git_sha, _tp_time), arg_max(status, _tp_time),"
f" max(_tp_time) FROM table({db.qualified('kg_ingest_log', prefix)}) GROUP BY repo"
).result_rows:
if status == "started" and now - t.replace(tzinfo=timezone.utc) > STALE_AFTER:
status = "stale"
out[repo] = (sha, status, t)
return out

52 changes: 52 additions & 0 deletions tests/test_ingest.py
Original file line number Diff line number Diff line change
Expand Up @@ -277,3 +277,55 @@ def counts():
return dict(rows)

_eventually(lambda: counts(), lambda c: c.get("vrepo@v1.0.0") == 1 and "vrepo" not in c)


# -- ingest history (#15): a run is visible while it runs, and says why it failed

def _runs(client, prefix, repo):
return client.query(
f"SELECT run_id, status, error FROM table({prefix}kg_ingest_log)"
f" WHERE repo = %(r)s ORDER BY _tp_time", parameters={"r": repo}
).result_rows


def test_ingest_logs_a_started_row_then_the_outcome_with_the_same_run_id(tp, monkeypatch, tmp_path: Path):
from tpk import ingest as ingest_mod
from tpk.ingest import ingest_repo

client, prefix = tp
seen = {}

def fake_graphify(repo_path, out_dir, **kw):
# while extraction runs, the started row must already be there
seen["during"] = _eventually(lambda: _runs(client, prefix, "histrepo"),
lambda rows: len(rows) == 1)
gj = out_dir / "graphify-out" / "graph.json"
gj.parent.mkdir(parents=True, exist_ok=True)
gj.write_text('{"nodes": [], "links": []}')
return gj

monkeypatch.setattr(ingest_mod, "run_graphify", fake_graphify)
cfg = RepoConfig(name="histrepo", path=tmp_path, visibility="internal")
result = ingest_repo(client, cfg, prefix=prefix, out_root=tmp_path / "out")
assert result.status == "ok"
assert [r[1] for r in seen["during"]] == ["started"]
rows = _eventually(lambda: _runs(client, prefix, "histrepo"), lambda rows: len(rows) == 2)
assert [(r[1], r[2]) for r in rows] == [("started", ""), ("ok", "")]
assert rows[0][0] == rows[1][0] == result.run_id


def test_ingest_failure_row_carries_the_error(tp, monkeypatch, tmp_path: Path):
from tpk import ingest as ingest_mod
from tpk.ingest import ingest_repo

client, prefix = tp

def boom(repo_path, out_dir, **kw):
raise RuntimeError("graphify exploded: exit 137")

monkeypatch.setattr(ingest_mod, "run_graphify", boom)
cfg = RepoConfig(name="failrepo", path=tmp_path, visibility="internal")
ingest_repo(client, cfg, prefix=prefix, out_root=tmp_path / "out")
rows = _eventually(lambda: _runs(client, prefix, "failrepo"), lambda rows: len(rows) == 2)
assert [r[1] for r in rows] == ["started", "failed"]
assert "graphify exploded: exit 137" in rows[1][2]
125 changes: 125 additions & 0 deletions tests/test_ingest_log_api.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
"""GET /api/ingest-log (#15): persistent ingest history folded per run."""
import time
from datetime import datetime, timedelta, timezone

import pytest
from fastapi.testclient import TestClient

from conftest import requires_timeplus
from tpk import auth, db
from tpk.server import create_app

pytestmark = requires_timeplus


class _NoAgent:
async def astream_events(self, *a, **k):
yield # pragma: no cover


def _eventually(fn, predicate=bool, timeout=5.0, interval=0.1):
deadline = time.monotonic() + timeout
result = fn()
while not predicate(result) and time.monotonic() < deadline:
time.sleep(interval)
result = fn()
return result


def _log(client, prefix, repo, run_id, status, nodes=0, edges=0, sha="abc1234", error="", at=None):
cols = ["repo", "run_id", "nodes", "edges", "git_sha", "status", "error"]
row = [repo, run_id, nodes, edges, sha, status, error]
if at is not None: # backdate: the stream's _tp_time is writable on insert
cols.append("_tp_time"); row.append(at)
client.insert(db.qualified("kg_ingest_log", prefix), [row], column_names=cols)


@pytest.fixture()
def env(tp):
client, prefix = tp
c = TestClient(create_app(agent=_NoAgent(), stream_prefix=prefix))
auth.upsert_role(client, auth.Role("viewer", [], capabilities=[auth.CAP_CORPUS_VIEW]), prefix=prefix)
auth.upsert_role(client, auth.Role("chatter", [], capabilities=[auth.CAP_CHAT]), prefix=prefix)
for name, role in (("vi", "viewer"), ("ch", "chatter")):
auth.upsert_user(client, auth.User(name, auth.hash_password("password-1"), role), prefix=prefix)

def hdr(name):
login = lambda: c.post("/auth/login", json={"username": name, "password": "password-1"})
_eventually(lambda: login().status_code == 200)
return {"Authorization": f"Bearer {login().json()['token']}"}

return c, client, prefix, hdr


def _runs(c, h, **params):
r = c.get("/api/ingest-log", headers=h, params=params)
assert r.status_code == 200, r.text
return r.json()["runs"]


def test_requires_corpus_view(env):
c, _, _, hdr = env
assert c.get("/api/ingest-log").status_code == 401
assert c.get("/api/ingest-log", headers=hdr("ch")).status_code == 403
assert c.get("/api/ingest-log", headers=hdr("vi")).status_code == 200


def test_folds_start_and_end_rows_into_one_run_newest_first(env):
c, client, prefix, hdr = env
now = datetime.now(timezone.utc)
_log(client, prefix, "alpha@v1", "r1", "started", at=now - timedelta(minutes=10))
_log(client, prefix, "alpha@v1", "r1", "ok", nodes=120, edges=300, at=now - timedelta(minutes=6))
_log(client, prefix, "beta@v1", "r2", "started", at=now - timedelta(minutes=5))
_log(client, prefix, "beta@v1", "r2", "failed", error="RuntimeError: graphify exploded",
at=now - timedelta(minutes=4))
_log(client, prefix, "old@v0", "r0", "ok", nodes=5, edges=1, at=now - timedelta(days=3)) # pre-#15 row: no start
runs = _eventually(lambda: _runs(c, hdr("vi")), lambda v: len(v) == 3)
assert [r["run_id"] for r in runs] == ["r2", "r1", "r0"]
r1 = runs[1]
assert r1["entry_key"] == "alpha@v1" and r1["status"] == "ok"
assert (r1["nodes"], r1["edges"], r1["git_sha"], r1["error"]) == (120, 300, "abc1234", "")
assert r1["started_at"] < r1["finished_at"] and 235 <= r1["duration_s"] <= 245
r2 = runs[0]
assert r2["status"] == "failed" and "graphify exploded" in r2["error"]
r0 = runs[2]
assert r0["status"] == "ok" and r0["started_at"] is None and r0["duration_s"] is None


def test_running_and_stale(env):
c, client, prefix, hdr = env
now = datetime.now(timezone.utc)
_log(client, prefix, "live@v1", "live", "started", at=now - timedelta(minutes=2))
_log(client, prefix, "dead@v1", "dead", "started", at=now - timedelta(hours=5))
runs = _eventually(lambda: _runs(c, hdr("vi")), lambda v: len(v) == 2)
by = {r["run_id"]: r for r in runs}
assert by["live"]["status"] == "running" and by["live"]["finished_at"] is None
assert by["dead"]["status"] == "stale"
assert "no end" in by["dead"]["error"].lower() or "did not finish" in by["dead"]["error"].lower()


def test_filter_and_limit(env):
c, client, prefix, hdr = env
now = datetime.now(timezone.utc)
for i in range(5):
_log(client, prefix, "a@v1" if i % 2 else "b@v1", f"r{i}", "ok", at=now - timedelta(minutes=10 - i))
_eventually(lambda: _runs(c, hdr("vi")), lambda v: len(v) == 5)
assert [r["run_id"] for r in _runs(c, hdr("vi"), limit=2)] == ["r4", "r3"]
assert {r["entry_key"] for r in _runs(c, hdr("vi"), entry="a@v1")} == {"a@v1"}
assert len(_runs(c, hdr("vi"), entry="a@v1")) == 2
assert c.get("/api/ingest-log", headers=hdr("vi"), params={"limit": 0}).status_code == 422
assert c.get("/api/ingest-log", headers=hdr("vi"), params={"limit": 10_000}).status_code == 422


def test_last_ingest_column_treats_an_abandoned_start_as_stale(env):
from tpk import corpus
from tpk.config import RepoConfig

c, client, prefix, hdr = env
corpus.upsert_entry(client, RepoConfig(name="dead", github="org/dead", ref="v1", visibility="internal"),
prefix=prefix)
_log(client, prefix, "dead@v1", "d1", "ok", nodes=9, at=datetime.now(timezone.utc) - timedelta(days=1))
_log(client, prefix, "dead@v1", "d2", "started", at=datetime.now(timezone.utc) - timedelta(hours=5))
rows = _eventually(lambda: c.get("/api/repos", headers=hdr("vi")).json(),
lambda v: any(r["entry_key"] == "dead@v1" and r.get("last_status") for r in v))
row = next(r for r in rows if r["entry_key"] == "dead@v1")
assert row["last_status"] == "stale"
Loading
Loading