13 Commits

Author SHA1 Message Date
Claude 370adcde6f test: trigger after reopen 2026-08-31 18:36:07 +00:00
Claude 84325e61c1 test: trigger review after .pr-review.json merged to main 2026-08-31 18:34:18 +00:00
gitea_admin 1644c7f6b1 chore: enable pragent pilot on this repo (.pr-review.json on PR branch) 2026-08-31 18:33:06 +00:00
Claude 3543d15677 test: re-trigger after dedupe window 2026-08-31 18:31:32 +00:00
Claude 72ac0f76bc test: retrigger review after eval rule wiring 2026-08-31 18:25:57 +00:00
Claude 5d44121b28 feat(eval): LLM-as-judge evaluators for finding actionability and review self-consistency
Two llm_as_judge evaluators score the review generation directly: a
NUMERIC 0-1 on finding actionability, a BOOLEAN on whether the summary
agrees with the findings. Both run on every observation whose trace
name is pr-review or opencode-review.

The judge is kimi-k2.7-code through the headroom hub. Local Ollama
returns Anthropic-format responses but the thinking blocks lack the
signature field Langfuse Zod schema requires; the evaluator preflight
fails as Invalid JSON response. A small judge-proxy pod on 8802
forwards to the hub and patches every thinking block with a synthetic
signature before returning.

Trace + generation output now includes the findings themselves
(capped at 25) rather than just the count, so a judge has something
to grade. generation input/output mirrors the trace so an
observation-level evaluator can read them.

Idempotent: existing evaluators and rules are skipped on re-run,
not duplicated. The connection is upserted on provider.
2026-08-31 17:17:16 +00:00
Marcos 2e1ad817e7 feat(eval): filterable item metadata and dataset runs for the Experiments tab
The filter bar matches on metadata only — not on input and not on the item id —
so a dataset seeded with repo/pr in `input` alone could not be sliced by repo
at all. Every facet worth filtering on is now a flat primitive in `metadata`:
repo, owner, repo_name, pr, head_sha, finding_count, has_findings,
max_severity, reviews_run and the review timestamp both ways. `owner` is split
out because a filter on the joined repo matches one repo, never a whole org,
and `max_severity` is "none" rather than absent because an absent key matches
no filter.

`eval_experiment.py` links reviews that already ran into a dataset run, one run
per model, so the Experiments tab is populated without re-running the reviewer.
One trace per (run, item), the most recent: a PR re-reviewed on every push has
many traces and a run is one output per input.

It posts to the deprecated /api/public/dataset-run-items — the notice exempts
self-hosted v3 from the cutoff date and the pilot is stdlib-only by design.
Revisit at v4.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-31 15:53:35 +00:00
Marcos 1534a99630 fix(eval): dataset item ids that survive a URL path
Items were keyed `{repo}#{pr}`, e.g. `netcracker/interview#29`. Both
characters break the UI's item route: the `/` in `owner/repo` splits into
extra path segments, and everything after the `#` is a fragment the browser
never sends. Items were created successfully and then 404'd when opened.

Ids are now `{owner}__{repo}__pr{n}`, which needs no percent-encoding. The
real repo and pr stay intact in `input`, so nothing downstream reads the id
back apart. Session ids elsewhere keep the `{repo}#{pr}` form — those are
never path segments and feedback_scores depends on that shape.

The 28 existing items were unusable and are regenerable from feedback.db;
they were deleted and recreated under the new ids.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-31 15:38:44 +00:00
gitea_admin d9eb4d9822 Merge PR #12: langfuse evaluation layer 2026-08-31 14:48:06 +00:00
Claude 2f96e66aab feat(pilot): behavioural scorers, feedback ground truth, and an eval dataset
Adds the evaluation layer on top of the review traces: five deterministic
scores describing how the reviewer behaved, a bridge that turns human reactions
into ground truth, and a dataset seeded from the reviews already run.

The two are kept apart on purpose. feedback.db has recorded 113 reviews and
zero reactions, resolutions or replies — nobody has ever responded to a bot
comment — so an accuracy metric cannot be built yet. The scorers therefore
measure behaviour, which is computable from data in hand, and feedback_scores
turns verdicts into scores the moment any arrive.

eval_scores.py emits finding_rate, severity_info_ratio, severity_max,
dropped_findings and cost_per_finding into the same ingestion batch as the
trace. Undefined values are omitted rather than reported as zero: an info ratio
over a silent review is undefined, and charting it as 0 would read as perfect
calibration.

dropped_findings needed a parser change. Both parsers silently discard findings
with an unusable path/line, which made a model emitting garbage locations
indistinguishable from one that found nothing. last_parse_dropped() exposes the
delta, read at parse time — after apply_repo_config the drops are the config
working as intended, not the model misbehaving.

feedback_scores.py scores the session ("{repo}#{pr}"), because feedback arrives
days later against a PR and nothing records which re-run produced which
comment. review_acceptance is absent rather than 0 when nothing was engaged.

eval_bootstrap.py registers the score configs, seeds the pragent-reviews
dataset, and can backfill scores onto traces that predate the scorers.
expectedOutput is the reviewer's own prior output, flagged
labelled_by_human: false — a regression baseline, not verified truth.

Also fixes a silent telemetry failure: the ingestion endpoint answers 207 when
only some events succeed, so a batch with every event rejected still looked
like success. Score events were missing the required per-event timestamp and
ingested nothing while reporting 207. _warn_on_rejected_events now logs the
per-event errors under LANGFUSE_DEBUG.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-31 14:22:55 +00:00
gitea_admin a4c35a4472 Merge pull request 'feat(pilot): emit per-review Langfuse traces' (#11) from feat/langfuse-tracing into main 2026-08-31 13:03:14 +00:00
Claude 92283c44e8 feat(pilot): emit per-review Langfuse traces
Ship token spend, latency and equivalent cost for every review to the
self-hosted Langfuse so per-model behaviour is queryable as a trend rather
than one PR comment at a time.

langfuse_trace.py is stdlib-only and emits via the public ingestion API.
Traces split into `ollama` and `claude` environments keyed off the bare model
name, not the provider: both paths go through the same headroom proxy, so the
provider prefix says nothing about which spend story a review belongs to. The
pilot's own path bills $0, so the reported cost is the equivalent price from
cost_model.PRICES.

ai_review.py calls _emit_langfuse on both token-spending exit paths (the
normal post and the salvage path). Import and emission are wrapped in a
blanket except: with no LANGFUSE_HOST or key pair the whole thing is a silent
no-op, and a telemetry failure must never fail a review.

These files were previously deployed only by way of the image build's
`COPY . /app`, so a clean checkout would have silently dropped tracing.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-31 12:45:35 +00:00
Claude 51b81def98 feat(opencode): configure vllm-qwen38 provider for local Qwen3.8-27B 2026-08-24 22:49:50 +00:00
24 changed files with 3641 additions and 22 deletions
+2
View File
@@ -8,3 +8,5 @@ dist/
__pycache__/
*.pyc
.worktrees/
.claude/
.opencode/package-lock.json
+5
View File
@@ -0,0 +1,5 @@
{
"enabled": true,
"model": "headroom/MiniMax-M2.7",
"static_message": "PR-Agent pilot on this repo. Comments are LLM-generated; treat as suggestions, not mandates."
}
+5
View File
@@ -94,6 +94,11 @@ path is in [`pilot/README.md`](pilot/README.md).
The model endpoint is supplied at runtime via `PRAGENT_MODEL_BASE_URL`; the
committed `opencode.json` carries a placeholder.
Per-review token spend, latency and equivalent cost are shipped to a
self-hosted Langfuse, split into `ollama` and `claude` environments so the two
spend stories stay separate: [`pilot/README-langfuse.md`](pilot/README-langfuse.md).
Emission is a silent no-op unless `LANGFUSE_HOST` and the key pair are set.
## Extending it
The review "factory" is [`.opencode/`](.opencode/README.md) — agent definitions
+82
View File
@@ -0,0 +1,82 @@
# Judge-side think-block patcher. Stands between Langfuse evaluators and the
# headroom-ollama hub (port 8790). Local Ollama does not emit the `signature`
# field that Langfuse's Anthropic adapter's Zod schema requires on every
# `thinking` content block — without it, the evaluator preflight fails as
# "Invalid JSON response". The proxy forwards /v1/* verbatim and adds a dummy
# signature to each thinking block before returning.
apiVersion: v1
kind: ConfigMap
metadata:
name: judge-proxy
namespace: pragent
data:
proxy.py: |
#!/usr/bin/env python3
"""Judge proxy: forward to headroom-ollama, fix thinking blocks."""
import json, sys, urllib.request, urllib.error
from http.server import BaseHTTPRequestHandler, HTTPServer
from socketserver import ThreadingMixIn
UPSTREAM = "http://100.74.17.70:8790"
DUMMY_SIG = "kimi-local-judge-no-signature"
class H(BaseHTTPRequestHandler):
def _proxy(self):
n = int(self.headers.get("Content-Length", 0))
body = self.rfile.read(n) if n else b""
h = {k: v for k, v in self.headers.items() if k.lower() not in ("host", "content-length")}
req = urllib.request.Request(UPSTREAM + self.path, data=body, headers=h, method=self.command)
try:
with urllib.request.urlopen(req, timeout=120) as r:
resp_body = r.read(); status = r.status; rh = dict(r.headers)
except urllib.error.HTTPError as e:
resp_body = e.read(); status = e.code; rh = dict(e.headers)
ct = rh.get("content-type", "")
if status == 200 and "application/json" in ct and self.path.startswith("/v1/messages"):
try:
obj = json.loads(resp_body)
patched = 0
for blk in obj.get("content") or []:
if isinstance(blk, dict) and blk.get("type") == "thinking" and "signature" not in blk:
blk["signature"] = DUMMY_SIG; patched += 1
if patched:
resp_body = json.dumps(obj).encode("utf-8")
rh["content-length"] = str(len(resp_body))
print(f"judge-proxy: patched {patched} thinking block(s)", file=sys.stderr, flush=True)
except Exception as e:
print(f"judge-proxy: patch failed: {e}", file=sys.stderr, flush=True)
self.send_response(status)
for k, v in rh.items():
if k.lower() not in ("transfer-encoding", "content-length", "connection"):
self.send_header(k, v)
self.send_header("Content-Length", str(len(resp_body)))
self.end_headers(); self.wfile.write(resp_body)
def do_POST(self): self._proxy()
def do_GET(self): self._proxy()
def log_message(self, *a, **k): pass
class S(ThreadingMixIn, HTTPServer): daemon_threads = True
S(("0.0.0.0", 8802), H).serve_forever()
---
apiVersion: v1
kind: Pod
metadata:
name: judge-proxy
namespace: pragent
labels:
app: judge-proxy
spec:
nodeSelector:
kubernetes.io/hostname: kubernets
hostNetwork: true
dnsPolicy: ClusterFirstWithHostNet
restartPolicy: Always
containers:
- name: p
image: python:3.12-alpine
command: ["sh","-c","apk add --no-cache ca-certificates >/dev/null && python3 -u /etc/cfg/proxy.py"]
volumeMounts:
- {name: cfg, mountPath: /etc/cfg}
ports:
- {containerPort: 8802, hostPort: 8802}
volumes:
- name: cfg
configMap:
name: judge-proxy
+18 -5
View File
@@ -21,19 +21,32 @@
}
}
},
"local": {
"vllm-qwen38": {
"npm": "@ai-sdk/openai-compatible",
"name": "Local AI workstation (qwen3.8-27b)",
"name": "Qwen3.8-27B vLLM (RTX 3090, MTP spec-decode)",
"options": {
"baseURL": "http://192.168.1.79:18020/v1",
"apiKey": "PLACEHOLDER_REPLACED_AT_RUNTIME"
"apiKey": "PLACEHOLDER_REPLACED_AT_RUNTIME",
"timeout": 300000,
"chunkTimeout": 30000
},
"models": {
"qwen3.8-27b": {
"name": "Qwen 3.8 27B (local)",
"name": "Qwen3.8-27B (vLLM, MTP, 150k ctx)",
"tools": true,
"thinking": true,
"attachments": false,
"limit": {
"context": 57344,
"context": 150000,
"output": 8192
},
"options": {
"temperature": 0.3,
"topP": 0.8,
"topK": 20,
"repetitionPenalty": 1.05,
"frequencyPenalty": 0,
"presencePenalty": 0
}
}
}
+224
View File
@@ -0,0 +1,224 @@
# Evaluation — scorers, ground truth, and the dataset
Langfuse already receives one trace per review (`README-langfuse.md`). This is
the layer on top: numbers attached to those traces that say how the reviewer
*behaved*, and the beginnings of a ground-truth signal that says whether it was
*right*.
Those two things are deliberately kept apart, because only one of them exists
yet.
## What could and could not be built
`feedback.db` has recorded 113 reviews across 4 repos. It has recorded **zero**
reactions, zero thread resolutions and zero replies. The harvester, the schema
and the daily analyzer are all working; nobody has ever reacted to a bot
comment.
That rules out an accuracy metric today. Correctness needs labels, and a
judge scored against no labels is theatre. So the scorers here measure
behaviour, which is computable from data already in hand, and a separate
bridge exists to turn human reactions into scores the moment any arrive.
## The five behavioural scores
Emitted with every review by `eval_scores.py`, folded into the same ingestion
batch as the trace so they cost no extra request.
| score | type | what a change in it means |
|---|---|---|
| `finding_rate` | NUMERIC | Findings posted. 0 is the restraint case — good on clean code, a failure when the run degraded. Only the rate over time separates those. |
| `severity_info_ratio` | NUMERIC 01 | Share of findings the model rated `info`/`trivial`. Rising = the model is hedging rather than committing. `None` when the review was silent: a ratio over an empty set is undefined, and charting it as 0 would read as perfect calibration. |
| `severity_max` | CATEGORICAL | Highest severity surfaced, `none` when silent. Categorical because "did this ever surface something serious" is the real question, and a mean of severity ranks answers nothing. |
| `dropped_findings` | NUMERIC | Findings the model emitted that the parser rejected for an unusable `path`/`line`. This is the only score here that measures the model's raw output. |
| `cost_per_finding` | NUMERIC | Equivalent USD per finding. A cheaper model that finds nothing is not cheaper. |
### Why `dropped_findings` needed a change to the parser
`parse_findings` and `parse_review_output` discard any finding with a missing or
unusable location. That happens silently, so a model emitting ten findings at
invalid locations was indistinguishable from a model that found nothing — both
produce an empty list. `ai_review.last_parse_dropped()` exposes the delta,
recorded at parse time.
It must be read at parse time specifically: by the time findings reach
`_emit_langfuse`, `apply_repo_config` has already filtered them by
`severity_threshold` and `max_findings`, and those drops are the config working
as intended, not the model misbehaving.
## Ground truth: `feedback_scores.py`
Turns `feedback.db` into two session-level scores, keyed on `"{repo}#{pr}"`
(which is what `langfuse_trace` already sets as `sessionId`).
| score | meaning |
|---|---|
| `review_engagement` | Share of a PR's findings that drew any human reaction, resolution or reply. **Watch this first** — every quality number is vapour until it moves off 0. |
| `review_acceptance` | Net verdict over engaged findings, 1 to +1. Absent, not 0, when nothing was engaged: zero would claim humans judged the review neutral, when the truth is nobody looked. |
Session-level rather than trace-level because feedback arrives days later
against a PR, and nothing in `feedback.db` records which re-run of the reviewer
produced which comment. The session is both the available join and the honest
granularity.
Score ids are `uuid5(namespace, repo#pr#name)`, so the daily backfill updates
rather than duplicates.
## The dataset
`pragent-reviews`, one item per PR the reviewer has run on, seeded by
`eval_bootstrap.py` from `feedback.db`.
`expectedOutput` is **the reviewer's own prior output**, not human-verified
truth — every item carries `metadata.labelled_by_human: false`. Read it as a
regression baseline: re-run a candidate model over these PRs and the diff
against this column is the behaviour change. Promoting an item to real ground
truth means a human editing it in the dataset view after re-reading the PR.
### Item ids
`{owner}__{repo}__pr{n}`. The obvious `{repo}#{pr}` cannot be used: items are
routed as `/datasets/{id}/items/{item_id}`, so the `/` in `owner/repo` splits
into extra path segments and everything after `#` is a fragment the browser
never sends — the item is created fine by the API and then 404s when opened.
Session ids elsewhere keep `{repo}#{pr}`; those are never path segments.
### Filterable metadata
The filter bar matches on `metadata` only — not on `input`, and not on the item
id — so every facet worth slicing on is a flat, primitive key in `metadata`
even where it duplicates `input`:
| key | why it is there |
| --- | --- |
| `repo`, `owner`, `repo_name` | `owner` exists because a filter on the joined `repo` matches one repo, never a whole org |
| `pr`, `head_sha` | jump from a filtered row back to the actual PR |
| `finding_count`, `has_findings` | isolate the silent reviews, which are the interesting negatives |
| `max_severity` | `"none"` rather than absent — an absent key matches no filter |
| `reviews_run` | how churny the PR was; high values skew per-item averages |
| `last_reviewed_at` / `_iso` | epoch sorts, ISO reads |
| `labelled_by_human` | `false` everywhere today; the flag to filter on before trusting any of it |
Nested objects and lists are deliberately absent: the filter bar cannot reach
into them.
`max_severity` is derived from `feedback.db`, whose `severity` column is
re-parsed out of the rendered comment by `feedback_harvest._parse_severity` and
defaults to `INFO` when its regex misses the badge. Trust the `severity_max`
**score** (read from the model's structured output) over this facet.
## Experiments
`eval_experiment.py` links reviews that already ran into a dataset run, so the
Experiments tab is populated without re-running anything. Runs are grouped by
model — the comparison the pilot actually needs is the same PRs under a
candidate model with `finding_rate` and `cost_per_finding` side by side. A new
model produces a new run automatically on the next invocation.
One trace per (run, item), the most recent: a PR re-reviewed on every push has
many traces, and a run is one output per input.
It uses `POST /api/public/dataset-run-items`, which is deprecated in favour of
the SDK experiment runner and disappears in Langfuse v4. The deprecation notice
exempts self-hosted v3 from the cutoff date, and this pilot is stdlib-only by
design. Revisit when this deployment moves to v4.
Coverage is bounded by the dataset, not by the traces: items only exist for PRs
with a row in `feedback.db`, and a review that posted no comment leaves a trace
but no row. That is why a run links fewer items than there are traces.
## Evaluators: `eval_judges.py`
Behaviour scores answer "how many, how severe, how much" — computable from data
already in hand. Two things they cannot answer:
- **Was the finding any good?** Specificity vs. hedge, generic advice vs.
fix-it-now advice — the difference between a useful review and one a
developer scrolls past.
- **Did the summary match the findings?** Claiming "no issues" above two
criticals, or describing a problem in prose that never became a finding.
These need a judge. `eval_judges.py` registers two `llm_as_judge` evaluators
against the trace names this project emits (`pr-review`, `opencode-review`)
and wires a sampling=1 rule per evaluator. Both run on every observation in a
matching trace; the only observations in those traces are the review itself.
| evaluator | output | what it answers |
|---|---|---|
| `finding_actionability` | NUMERIC 01 | How specific and fixable is each finding? |
| `review_self_consistency` | BOOLEAN | Does the summary agree with the findings? |
The judge is a different model from the reviewer (`kimi-k2.7-code` through the
headroom hub). A model grading its own output agrees with itself for reasons
that have nothing to do with quality. The judges are also asked only what they
can answer from the review itself — never whether a finding is correct, since
that needs the diff the trace does not carry.
### Why the judge goes through `judge-proxy` (port 8802)
The headroom hub in front of local Ollama returns Anthropic-format responses,
but every `thinking` content block is missing the `signature` field real
Claude emits. Langfuse's Zod schema requires it; the omission fails the
evaluator preflight as `Invalid JSON response`. The `judge-proxy` pod sits in
front of the hub on `100.74.17.70:8802` and patches every thinking block with
a synthetic signature before forwarding the response. The model is unchanged;
only the wire shape is fixed.
```bash
python3 pilot/eval_judges.py --dry-run # show what would be created
python3 pilot/eval_judges.py # create the LLM connection, evaluators, rules
```
Idempotent: existing evaluators and rules are skipped, not duplicated. The
connection is upserted on `provider` so re-runs return the same record.
## Running it
```bash
# once per project: score configs + dataset (+ score historical traces)
python3 pilot/eval_bootstrap.py --db /data/feedback.db --backfill-traces
# ship feedback verdicts (runs daily from the feedback CronJob)
python3 pilot/feedback_scores.py --db /data/feedback.db
# link already-traced reviews into a dataset run per model
python3 pilot/eval_experiment.py --dry-run
python3 pilot/eval_experiment.py
```
Both need `LANGFUSE_HOST`, `LANGFUSE_PUBLIC_KEY`, `LANGFUSE_SECRET_KEY`. In
cluster they come from the `pragent-langfuse` Secret and point at the ClusterIP
— never the NodePort, whose oauth2-proxy 302s ingestion to Logto and drops it.
## Gotcha: HTTP 207 is not success
The ingestion endpoint answers `207 Multi-Status` when *some* events failed, so
a batch where **every** event was rejected still returns 207. An early version
of these scorers omitted the required per-event `timestamp` and silently
ingested nothing while reporting success. `langfuse_trace._warn_on_rejected_events`
now logs the per-event errors under `LANGFUSE_DEBUG=1`. If scores are missing,
check that before anything else.
## What the first run showed
Backfilled over 42 existing traces and 13 PRs:
```
cost_per_finding n=42 mean=0.3133 min=0.0880 max=0.9042
finding_rate n=42 mean=0.4762 min=0.0000 max=4.0000
severity_info_ratio n=14 mean=0.0000
review_engagement n=14 mean=0.0000
severity_max {none: 28, medium: 11, high: 1, critical: 2}
```
Two things worth keeping:
- **The reviewer is not info-heavy.** `feedback.db` shows 61 of 62 findings at
`INFO`, which looked like a badly calibrated model. It is not: `severity_max`
reads `medium`/`high`/`critical` on every trace that found anything, and
`severity_info_ratio` is flat 0. The `INFO` in the DB comes from
`feedback_harvest._parse_severity`, which defaults to `INFO` when its regex
misses the severity badge in the rendered comment. The DB severity is a
re-parse artifact; the score reads the model's structured output directly.
- **28 of 42 reviews found nothing** (67%), and **engagement is flat zero**. The
first is not yet interpretable without the second.
+122
View File
@@ -0,0 +1,122 @@
# pragent → Langfuse
Every review the pilot runs ships one **trace** to a self-hosted Langfuse. The
review body already prints a usage table, but that table lives and dies inside
one Gitea PR. Langfuse is where the same numbers become a trend: tokens per
review, latency per model, equivalent cost per repo, and how those move when
the model or the tiering changes.
## The ollama / claude split
Both paths route through the same headroom proxy, so the provider prefix does
not distinguish them — `headroom/claude-sonnet-5` is Claude spend,
`headroom/glm-5.2:cloud` is not. The split is keyed off the **bare model name**
and lands on the trace's `environment`:
| resolved model | environment |
| --------------------------- | ----------- |
| `headroom/claude-sonnet-5` | `claude` |
| `claude-opus-5` | `claude` |
| `headroom/glm-5.2:cloud` | `ollama` |
| `headroom/MiniMax-M2.7` | `ollama` |
| `vllm-qwen38/qwen3.8-27b` | `ollama` |
Langfuse takes an environment selector on every dashboard, filter and cost
breakdown, so the two spend stories stay separate inside one project — one key
pair to rotate instead of two. Tags carry the finer cut:
`provider:headroom`, `model:<bare>`, `engine:opencode`, `repo:<owner/name>`,
`lens:<id>` per fan-out lens.
To split into two *projects* later, point `LANGFUSE_PUBLIC_KEY` /
`LANGFUSE_SECRET_KEY` at the second project on whichever deployment runs the
Claude path. Nothing in the code needs to change.
## What a trace carries
- **trace** `pr-review``sessionId` = `owner/repo#index`, so every push to one
PR groups together. Input is the PR identity; output is the summary + finding
count; metadata carries steps, duration, severity counts and the provider's
own reported cost.
- **generation** `opencode-review``model`, `usageDetails`, `costDetails`.
`usageDetails.input` is the **uncached** input. opencode reports `cache_read`
*inside* `input`, and Langfuse sums the keys it is given, so passing both
verbatim would bill the resent prefix twice.
### How cost is priced
Langfuse has no price table of its own here — we compute the number and ship it
as `costDetails.total`, so what Langfuse charts is exactly what
`cost_model.PRICES` says.
A model that genuinely bills (`claude-*`, `gpt-*`, `gemini-*`, `grok-*`) is
priced **as itself**: basis `actual`.
A model that costs nothing through the headroom proxy is priced against a
**comparison target** instead: basis `equivalent:<target>`. That covers the
models absent from `PRICES` (`MiniMax-M2.7` — which is what the webhook
actually runs — and `glm-5.2:cloud`) as well as entries priced at all zeros
(the self-hosted vLLM `qwen3.8-27b`). Without this the dashboard would be a
flat $0.00 line, since the pilot's own path is free.
The target follows the same precedence as the review body, so the PR and the
dashboard never disagree:
.pr-review.json:cost_target > PRAGENT_PRICE_TARGET > claude-sonnet-5
An equivalent cost is a hypothetical, not money spent, so every trace is tagged
`cost:actual` or `cost:equivalent:<target>` and the generation metadata carries
`cost_basis`. Filter on it before reading any cost chart as spend.
If the comparison target itself is unknown, the trace ships usage with **no**
cost block — better no number than a wrong one.
Anthropic prices in `cost_model.PRICES` were fetched 2026-08-18; re-check them
before quoting anything externally.
## Configuration
| env | meaning |
| --------------------- | --------------------------------------------------------- |
| `LANGFUSE_HOST` | `http://langfuse-web.langfuse.svc.cluster.local:3000` |
| `LANGFUSE_PUBLIC_KEY` | `pk-lf-…` |
| `LANGFUSE_SECRET_KEY` | `sk-lf-…` |
| `LANGFUSE_TIMEOUT` | seconds, default `5` |
| `LANGFUSE_DEBUG` | `1` to log ingestion failures to stderr |
Unset host or either key ⇒ emission is a silent no-op. That is the default, so
a checkout without Langfuse behaves exactly as before.
## Fail-open
`langfuse_trace` is stdlib-only (`urllib`) and every entry point swallows its
own exceptions; `_emit_langfuse` in `ai_review.py` wraps even the import. A
Langfuse outage cannot fail, delay past `LANGFUSE_TIMEOUT`, or alter a review.
Both token-spending exit paths emit — the normal post **and** the salvage path
where the agent produced unparseable output. That run cost the same as a clean
one, and is precisely the failure worth trending.
## Deployment
Cluster side lives outside this repo: `~/k8s/langfuse.yaml` (ClickHouse +
web + worker, reusing the gitea postgres, gitea valkey and minio),
`~/k8s/oauth2-proxy-langfuse.yaml` (the Logto gate), and
`~/k8s/langfuse-setup.sh`, which provisions the database, the bucket, the
secrets, and wires `pragent-webhook` with the three env vars above.
The UI is at **https://langfuse.marcospaulo.dev.br**:
browser -> Caddy (VPS, TLS, DNS-01) -> tailscale
-> 100.74.17.70:30361 -> oauth2-proxy (Logto, email allowlist)
-> langfuse-web (ClusterIP)
Logto sits at *both* layers off one app (`langfuse`, two redirect URIs): the
proxy gates the domain, and Langfuse's own NextAuth uses the same Logto as a
custom OIDC provider, so the inner login is a silent redirect rather than a
second password.
pragent does **not** go through any of that. It posts to
`langfuse-web.langfuse.svc.cluster.local:3000` from inside the cluster, on
API-key auth — putting ingestion behind an interactive SSO gate would break it
on the first review.
+83 -3
View File
@@ -489,7 +489,7 @@ def _resolve_display_model(base_model: str, config: dict | None) -> str:
like `claude-sonnet-5` or `qwen3.8-27b` is safe. Re-prefixed with
the model's `provider` field from `cost_model.Price` (default
`headroom`) so the opencode subprocess routes correctly — e.g.
`qwen3.8-27b` → `local/qwen3.8-27b` (local AI workstation on
`qwen3.8-27b` → `vllm-qwen38/qwen3.8-27b` (vLLM on RTX 3090 at
192.168.1.79:18020), `claude-sonnet-5` → `headroom/claude-sonnet-5`
(Anthropic pricing proxy).
3. Default — `f"headroom/{base_model}"` where `base_model` is the bare
@@ -505,8 +505,8 @@ def _resolve_display_model(base_model: str, config: dict | None) -> str:
cfg_model = (config or {}).get("model")
if isinstance(cfg_model, str) and cfg_model.strip():
# Look up the provider from PRICES so the opencode subprocess routes
# through the right provider block (local vs headroom). Lazy import —
# the ollama path doesn't touch cost_model.
# through the right provider block (vllm-qwen38 vs headroom). Lazy
# import — the ollama path doesn't touch cost_model.
from cost_model import PRICES
provider = PRICES.get(cfg_model.strip())
if provider is not None:
@@ -678,6 +678,19 @@ def _strip_path_prefix(p: str) -> str:
# ---------------------------------------------------------------------------
# How many raw findings the last `parse_review_output` / `parse_findings` call
# rejected for an unusable path/line. A side channel rather than a return value
# because both parsers already return fixed-width tuples that several callers
# and their tests unpack positionally; widening them to carry a telemetry
# number would be a breaking change for a fail-open signal.
_LAST_PARSE_DROPPED: dict[str, int] = {"n": 0}
def last_parse_dropped() -> int:
"""Findings the last parse discarded. Read it immediately after parsing."""
return int(_LAST_PARSE_DROPPED.get("n") or 0)
def _normalize_finding(f: dict) -> dict | None:
"""Validate + normalize one raw finding dict. Returns None if it's unusable
(missing path/line). Normalises severity, keeps `reference` (default "")."""
@@ -754,6 +767,7 @@ def parse_findings(text: str) -> list[dict]:
Also accepts a bare JSON array as the outer value: ``[{...}, {...}]`` —
some agents skip the ``{"summary":..., "findings":[...]}`` wrapper.
"""
_LAST_PARSE_DROPPED["n"] = 0
data = _parse_json_tolerant(text)
if isinstance(data, dict):
findings = data.get("findings")
@@ -768,6 +782,7 @@ def parse_findings(text: str) -> list[dict]:
n = _normalize_finding(f)
if n is not None:
out.append(n)
_LAST_PARSE_DROPPED["n"] = len(findings) - len(out)
return out
@@ -823,6 +838,7 @@ def parse_review_output(
block), with a tolerant fallback that scans for the last balanced
object/array in the prose tail. Never raises.
"""
_LAST_PARSE_DROPPED["n"] = 0
blob = _last_json_block(text)
if blob is None:
return "", [], [], [], [], "", ""
@@ -856,6 +872,12 @@ def parse_review_output(
n = _normalize_finding(f)
if n is not None:
out.append(n)
# A model that emits findings at unusable locations is indistinguishable
# from one that found nothing, because both end up with an empty `out`.
# Stash the delta so the caller can score it (see `eval_scores`).
_LAST_PARSE_DROPPED["n"] = len(findings_raw) - len(out)
else:
_LAST_PARSE_DROPPED["n"] = 0
return summary, out, summary_changes, risks, walkthrough, risk_verdict, test_coverage
@@ -2068,6 +2090,50 @@ def _need(name: str) -> str:
return v
def _emit_langfuse(
*,
repo: str,
index: str,
sha: str,
title: str,
model: str,
usage: dict | None,
findings: list[dict],
summary: str,
engine: str,
config: dict | None = None,
dropped_count: float | None = None,
) -> None:
"""Ship this review's usage to Langfuse, if one is configured.
Called on both exit paths that spent tokens — the normal post and the
salvage path — because an unparseable run costs the same as a clean one and
is exactly the kind of thing worth trending.
Local import + blanket except: `langfuse_trace` is stdlib-only but optional,
and telemetry is never allowed to fail a review (see the fail-open contract
in `review_pr`). The trace's `environment` is `claude` or `ollama`, so the
two spend stories stay separated in every Langfuse view.
"""
try:
import langfuse_trace
# Same comparison model the review body prices against, so the number
# in Langfuse and the number in the PR agree. Free/unknown models
# (MiniMax, glm, self-hosted qwen) are priced against it; a paid model
# is priced as itself.
price_target, _err = _resolve_price_target(config)
langfuse_trace.emit_review_trace(
repo=repo, index=index, sha=sha, title=title, model=model,
usage=usage, findings=findings, summary=summary or "",
engine=engine, lenses=(usage or {}).get("lenses"),
price_target=price_target, dropped_count=dropped_count,
)
except Exception as e:
print(f"pragent: langfuse emit skipped: {e}", file=sys.stderr)
def review_pr(
api: str,
repo: str,
@@ -2192,6 +2258,7 @@ def review_pr(
additional_context=additional_context,
)
review_summary, findings, summary_changes, risks, _walkthrough, _risk_verdict, _test_coverage = parse_review_output(stdout)
parse_dropped = last_parse_dropped()
if not findings and not review_summary:
# The findings JSON was missing or malformed. Don't discard the
# run: salvage the prose, keep the usage report (the tokens were
@@ -2208,11 +2275,18 @@ def review_pr(
salvaged or "AI review produced no parseable output.",
display_model, sha, usage_section=usage_section,
static_message=(config or {}).get("static_message", "")))
_emit_langfuse(
repo=repo, index=index, sha=sha, title=title,
model=display_model, usage=usage, findings=[],
summary=salvaged, engine=engine, config=config,
dropped_count=parse_dropped,
)
return True
else:
user_prompt = build_user_prompt(title, body + compression_note, diff, config, prior, additional_context)
raw_findings = call_model(ollama_url, model, SYSTEM_PROMPT, user_prompt, max_tokens)
findings = parse_findings(raw_findings)
parse_dropped = last_parse_dropped()
usage = None
# Filter / cap findings per `.pr-review.json` (style, threshold, max,
@@ -2291,6 +2365,12 @@ def review_pr(
)
post_inline_review(api, repo, index, token, summary_body, anchored)
_emit_langfuse(
repo=repo, index=index, sha=sha, title=title,
model=display_model, usage=usage, findings=findings,
summary=review_summary, engine=engine, config=config,
dropped_count=parse_dropped,
)
print(
f"pragent: reviewed {repo}#{index} sha={sha[:8]} "
f"engine={engine} findings={len(findings)} inline={len(anchored)}",
+8 -7
View File
@@ -57,7 +57,7 @@ class Price:
"""Per-MTok prices. `cache_write` and `cache_read` are absolute rates, not
multipliers, so providers with different cache economics stay comparable.
`provider` is the opencode provider name (`headroom`, `local`, ...). It
`provider` is the opencode provider name (`headroom`, `vllm-qwen38`, ...). It
doubles as the dispatch key for `.pr-review.json:model` overrides — when
a per-repo override is set, `_resolve_display_model` returns
`f"{provider}/{key}"` so the opencode subprocess routes correctly.
@@ -98,12 +98,13 @@ PRICES: dict[str, Price] = {
# xAI Grok — cache_write = input
"grok-4.5": Price("Grok 4.5", 2.00, 6.00, 2.00, 0.30),
"grok-4.3": Price("Grok 4.3", 1.25, 2.50, 1.25, 0.20),
# Self-hosted — local AI workstation, no per-token charge. provider="local"
# so the opencode subprocess routes via the `local` provider block in
# opencode.json (baseURL=http://192.168.1.79:18020/v1). Equivalent-cost
# column will read $0 — the cost-comparison signal is that the same work
# would bill $X on a paid model.
"qwen3.8-27b": Price("Qwen 3.8 27B (local)", 0.0, 0.0, 0.0, 0.0, provider="local"),
# Self-hosted — AI workstation RTX 3090, vLLM + DFlash2 spec-decode, no
# per-token charge. provider="vllm-qwen38" so the opencode subprocess
# routes via the matching provider block in opencode.json
# (baseURL=http://192.168.1.79:18020/v1). Equivalent-cost column reads $0
# — the cost-comparison signal is that the same work would bill $X on a
# paid model.
"qwen3.8-27b": Price("Qwen3.8-27B (vLLM, MTP, 150k ctx)", 0.0, 0.0, 0.0, 0.0, provider="vllm-qwen38"),
}
+350
View File
@@ -0,0 +1,350 @@
#!/usr/bin/env python3
"""pragent pilot — one-time Langfuse project setup for evaluation.
Three jobs, each idempotent so it can be re-run after any change:
1. **Score configs.** Registers the schema for every score pragent emits
(`eval_scores.SCORE_CONFIGS` + `feedback_scores.SCORE_CONFIGS`). Without
these the scores still ingest, but nothing stops a later scorer writing
`severity_max="HIGH"` beside today's `"high"` and quietly splitting one
series into two. Configs are immutable in Langfuse — a name that already
exists is left alone rather than updated.
2. **Dataset.** Seeds `pragent-reviews` from `feedback.db`: one item per PR
the reviewer has actually run on, carrying the repo/PR/sha as input and
the findings it posted as `expectedOutput`.
Read `expectedOutput` here as "what the reviewer said last time", not "what
is correct" — no human has labelled any of it. It is a regression baseline:
re-run a candidate model over these PRs and the diff against this column is
the behaviour change. Promoting an item to real ground truth means a human
editing it after reviewing the PR, which is what the dataset view is for.
3. **Trace backfill** (`--backfill-traces`). Scores only ride along with new
reviews, so without this the charts stay empty until the next PR lands.
Every trace `langfuse_trace` has ever written already carries the finding
count, the severity histogram and the cost in its metadata, which is
everything four of the five scorers need. `dropped_findings` is absent from
historical traces and is left unscored rather than backfilled as zero.
4. **Reports** what it found, so the gap between "reviews recorded" and
"reviews with human feedback" is visible rather than assumed.
Usage:
LANGFUSE_HOST=... LANGFUSE_PUBLIC_KEY=... LANGFUSE_SECRET_KEY=... \\
python3 eval_bootstrap.py --db /data/feedback.db
"""
from __future__ import annotations
import argparse
import base64
import json
import os
import sqlite3
import sys
import urllib.error
import urllib.request
from datetime import datetime, timezone
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import eval_scores # noqa: E402
import feedback_scores # noqa: E402
DATASET_NAME = "pragent-reviews"
def _conf() -> tuple[str, str, str]:
host = (os.environ.get("LANGFUSE_HOST") or "").strip().rstrip("/")
pk = (os.environ.get("LANGFUSE_PUBLIC_KEY") or "").strip()
sk = (os.environ.get("LANGFUSE_SECRET_KEY") or "").strip()
if not host or not pk or not sk:
raise SystemExit("LANGFUSE_HOST / LANGFUSE_PUBLIC_KEY / LANGFUSE_SECRET_KEY must be set")
return host, pk, sk
def _call(method: str, path: str, body: dict | None = None, timeout: float = 20.0):
host, pk, sk = _conf()
auth = base64.b64encode(f"{pk}:{sk}".encode()).decode("ascii")
data = json.dumps(body).encode() if body is not None else None
req = urllib.request.Request(
host + path,
data=data,
headers={
"Content-Type": "application/json",
"Authorization": f"Basic {auth}",
"User-Agent": "pragent-pilot/1.0",
},
method=method,
)
try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
raw = resp.read()
return resp.status, (json.loads(raw) if raw else None)
except urllib.error.HTTPError as e:
return e.code, e.read()[:400].decode("utf-8", "replace")
# ---------------------------------------------------------------------------
# 1. Score configs
# ---------------------------------------------------------------------------
def ensure_score_configs() -> dict:
status, existing = _call("GET", "/api/public/score-configs?limit=100")
have = set()
if status == 200 and isinstance(existing, dict):
have = {c.get("name") for c in existing.get("data", [])}
created, skipped, failed = [], [], []
for cfg in list(eval_scores.SCORE_CONFIGS) + list(feedback_scores.SCORE_CONFIGS):
if cfg["name"] in have:
skipped.append(cfg["name"])
continue
st, resp = _call("POST", "/api/public/score-configs", cfg)
if st in (200, 201):
created.append(cfg["name"])
else:
failed.append({"name": cfg["name"], "status": st, "error": resp})
return {"created": created, "already_present": skipped, "failed": failed}
# ---------------------------------------------------------------------------
# 2. Dataset from recorded reviews
# ---------------------------------------------------------------------------
def item_id(repo: str, pr) -> str:
"""A dataset-item id that survives being put in a URL path.
The obvious `{repo}#{pr}` is unusable: the UI routes items as
`/datasets/{id}/items/{item_id}`, so the `/` in `owner/repo` splits into
extra path segments and everything after the `#` is a fragment the browser
never sends. The item is created fine and then 404s when opened.
Session ids elsewhere keep the `{repo}#{pr}` form — those are never path
segments, and `feedback_scores` depends on that shape.
"""
return f"{repo.replace('/', '__')}__pr{pr}"
def _item_metadata(*, repo, pr, head_sha, reviews_run, last_seen, findings) -> dict:
"""Filterable facets for one dataset item.
Kept flat and primitive: the filter bar matches a metadata key against a
literal, so a nested object or a list is not reachable from the UI.
"""
owner, _, repo_name = str(repo).partition("/")
sevs = [str(f["severity"] or "").lower() for f in findings]
ranked = [s for s in sevs if s in eval_scores.SEVERITY_RANK]
return {
"repo": repo,
"owner": owner or repo,
"repo_name": repo_name or repo,
"pr": int(pr),
"head_sha": head_sha,
"reviews_run": reviews_run,
"last_reviewed_at": last_seen,
"last_reviewed_iso": datetime.fromtimestamp(last_seen, timezone.utc).isoformat(),
"finding_count": len(findings),
"has_findings": bool(findings),
# "none" rather than omitting the key: a filter for silent reviews needs
# something to match, and an absent key matches nothing.
"max_severity": (
max(ranked, key=lambda s: eval_scores.SEVERITY_RANK[s]) if ranked else "none"
),
# Flags that this row is the reviewer's own past output, not a human
# judgement. Filter on it before anyone treats the dataset as truth.
"labelled_by_human": False,
}
def read_review_items(db_path: str) -> list[dict]:
"""One dataset item per (repo, pr) the reviewer has run on.
Keyed on the PR rather than on each individual review row: the same PR is
re-reviewed on every push, and 113 rows over 26 PRs would make a benchmark
that is 4x redundant and weighted towards whichever PR churned most.
"""
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
try:
prs = conn.execute(
"""
SELECT repo, pr, MAX(posted_at) AS last_seen, COUNT(*) AS reviews,
MAX(head_sha) AS head_sha
FROM review GROUP BY repo, pr ORDER BY repo, pr
"""
).fetchall()
items = []
for row in prs:
findings = conn.execute(
"""
SELECT path, line, severity, problem, fix
FROM inline_finding WHERE repo = ? AND pr = ?
ORDER BY path, line
""",
(row["repo"], row["pr"]),
).fetchall()
items.append(
{
"id": item_id(row["repo"], row["pr"]),
"input": {
"repo": row["repo"],
"pr": int(row["pr"]),
"head_sha": row["head_sha"],
},
"expectedOutput": {
"findings": [dict(f) for f in findings],
"finding_count": len(findings),
},
# The UI's filter bar reads metadata and nothing else, so
# anything worth slicing on is a top-level key here even
# where it duplicates `input`. `owner` and `repo_name` are
# split out because a filter on the joined `repo` can only
# match one repo at a time, never a whole org.
"metadata": _item_metadata(
repo=row["repo"],
pr=row["pr"],
head_sha=row["head_sha"],
reviews_run=int(row["reviews"]),
last_seen=int(row["last_seen"]),
findings=findings,
),
}
)
return items
finally:
conn.close()
def ensure_dataset(items: list[dict], name: str = DATASET_NAME) -> dict:
st, _ = _call(
"POST",
"/api/public/datasets",
{
"name": name,
"description": (
"PRs the pragent pilot has reviewed, seeded from feedback.db. "
"expectedOutput is the reviewer's own prior output — a regression "
"baseline, not human-verified ground truth."
),
"metadata": {"source": "feedback.db", "seeded_by": "eval_bootstrap.py"},
},
)
# A duplicate name is fine: the dataset already exists from an earlier run.
dataset_ok = st in (200, 201, 409)
created, failed = 0, []
for item in items:
body = {
"datasetName": name,
"id": item["id"], # idempotent: same PR updates rather than duplicates
"input": item["input"],
"expectedOutput": item["expectedOutput"],
"metadata": item["metadata"],
}
ist, resp = _call("POST", "/api/public/dataset-items", body)
if ist in (200, 201):
created += 1
else:
failed.append({"item": item["id"], "status": ist, "error": resp})
return {"dataset": name, "dataset_created": dataset_ok, "items_upserted": created, "failed": failed}
# ---------------------------------------------------------------------------
# 3. Backfill scores onto traces that predate the scorers
# ---------------------------------------------------------------------------
def _synth_findings(severities: dict) -> list[dict]:
"""Rebuild a findings list from a trace's severity histogram.
Only severity matters to the scorers, and that is all the histogram kept.
Reconstructing placeholders is honest here because every scorer being
backfilled reads nothing else off a finding.
"""
out = []
for sev, count in (severities or {}).items():
out.extend({"severity": sev} for _ in range(int(count)))
return out
def backfill_traces(limit_pages: int = 20) -> dict:
import eval_scores as es
scored, skipped, events = 0, 0, []
page = 1
while page <= limit_pages:
st, resp = _call("GET", f"/api/public/traces?limit=50&page={page}&name=pr-review")
if st != 200 or not isinstance(resp, dict):
break
rows = resp.get("data") or []
if not rows:
break
for tr in rows:
meta = tr.get("metadata") or {}
severities = meta.get("severities") or {}
count = meta.get("findings")
if count is None:
skipped += 1
continue
findings = _synth_findings(severities)
# The histogram is authoritative when present; a trace that recorded
# a count but no histogram still scores its rate.
if not findings and count:
findings = [{"severity": "medium"} for _ in range(int(count))]
batch = es.build_scores(
trace_id=tr["id"],
findings=findings,
environment=tr.get("environment") or "default",
cost_usd=(tr.get("totalCost") or meta.get("provider_cost_usd")),
timestamp=tr.get("timestamp"),
comment="backfilled from trace metadata",
)
events.extend(batch)
scored += 1
page += 1
posted = False
status = None
if events:
import langfuse_trace
host, pk, sk = _conf()
# Chunked: one 2000-event POST is refused, and a partial backfill that
# reports success is worse than a slow one.
for i in range(0, len(events), 200):
status = langfuse_trace._post(host, pk, sk, events[i:i + 200], 30.0)
posted = status in (200, 201, 207)
if not posted:
break
return {"traces_scored": scored, "traces_skipped": skipped, "scores": len(events),
"posted": posted, "http_status": status}
def main() -> int:
ap = argparse.ArgumentParser(description="Bootstrap Langfuse evaluation for the pragent pilot")
ap.add_argument("--db", default=os.environ.get("PRAGENT_FEEDBACK_DB", "/data/feedback.db"))
ap.add_argument("--skip-dataset", action="store_true")
ap.add_argument("--skip-configs", action="store_true")
ap.add_argument("--backfill-traces", action="store_true",
help="score traces written before the scorers existed")
args = ap.parse_args()
out: dict = {}
if not args.skip_configs:
out["score_configs"] = ensure_score_configs()
if not args.skip_dataset:
items = read_review_items(args.db)
out["dataset"] = ensure_dataset(items)
out["dataset"]["items_read"] = len(items)
if args.backfill_traces:
out["trace_backfill"] = backfill_traces()
print(json.dumps(out, indent=2))
failed = (out.get("score_configs", {}).get("failed") or []) + (
out.get("dataset", {}).get("failed") or []
)
return 1 if failed else 0
if __name__ == "__main__":
raise SystemExit(main())
+212
View File
@@ -0,0 +1,212 @@
#!/usr/bin/env python3
"""pragent pilot — populate the Experiments tab from reviews already traced.
An "experiment" in Langfuse is a dataset run: a set of (dataset item, trace)
links under one run name. The Experiments tab then shows one row per item with
its scores, and lets two runs be diffed side by side.
Nothing here re-runs the reviewer. Every PR in `pragent-reviews` has already
been reviewed, and each of those reviews left a trace carrying its findings,
cost and scores. This links what exists, which is what makes the tab useful on
day one instead of after the next N pushes.
Runs are grouped by **model** by default, because that is the comparison the
pilot actually needs to make: the same PRs reviewed by MiniMax vs whatever
replaces it, with `finding_rate` and `cost_per_finding` side by side. Group by
`none` for a single "all traces" run.
One trace per (run, item) — the most recent. A PR re-reviewed on every push has
many traces, and a dataset run is defined as one output per input; feeding it
the other five would make the per-run averages meaningless.
Note on the endpoint: `POST /api/public/dataset-run-items` is deprecated in
favour of the SDK experiment runner / OTel ingestion, and disappears in
Langfuse v4. This instance is self-hosted v3, which the deprecation notice
explicitly exempts from the cutoff date, and the pilot is stdlib-only by
design. Revisit when this deployment moves to v4.
Usage:
LANGFUSE_HOST=... LANGFUSE_PUBLIC_KEY=... LANGFUSE_SECRET_KEY=... \\
python3 eval_experiment.py --dry-run
"""
from __future__ import annotations
import argparse
import json
import os
import sys
import urllib.parse
from collections import defaultdict
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import eval_bootstrap as eb # noqa: E402
TRACE_NAME = "pr-review"
# ---------------------------------------------------------------------------
# Reading what already exists
# ---------------------------------------------------------------------------
def fetch_traces(name: str = TRACE_NAME, limit: int = 100, max_pages: int = 50) -> list[dict]:
"""Every review trace, newest first."""
out: list[dict] = []
for page in range(1, max_pages + 1):
q = urllib.parse.urlencode({"name": name, "limit": limit, "page": page})
st, body = eb._call("GET", f"/api/public/traces?{q}")
if st != 200 or not isinstance(body, dict):
raise SystemExit(f"listing traces failed: {st} {body}")
data = body.get("data") or []
out.extend(data)
meta = body.get("meta") or {}
if page * meta.get("limit", limit) >= meta.get("totalItems", 0):
break
return out
def fetch_item_ids(dataset: str) -> set[str]:
"""Ids present in the dataset, so runs never reference a missing item."""
ids: set[str] = set()
for page in range(1, 51):
q = urllib.parse.urlencode({"datasetName": dataset, "limit": 100, "page": page})
st, body = eb._call("GET", f"/api/public/dataset-items?{q}")
if st != 200 or not isinstance(body, dict):
raise SystemExit(f"listing dataset items failed: {st} {body}")
ids.update(i["id"] for i in body.get("data") or [])
meta = body.get("meta") or {}
if page * meta.get("limit", 100) >= meta.get("totalItems", 0):
break
return ids
# ---------------------------------------------------------------------------
# Grouping traces into runs
# ---------------------------------------------------------------------------
def trace_model(trace: dict) -> str:
"""The model that produced a review, from its `model:` tag."""
for tag in trace.get("tags") or []:
if tag.startswith("model:"):
return tag[len("model:"):] or "unknown"
return "unknown"
def trace_item_id(trace: dict) -> str | None:
"""The dataset item a trace belongs to, or None if it is not a PR review."""
md = trace.get("metadata") or {}
repo, pr = md.get("repo"), md.get("pr")
if not repo or pr in (None, ""):
return None
return eb.item_id(str(repo), pr)
def _sort_key(trace: dict):
return (trace.get("timestamp") or "", trace.get("id") or "")
def plan_runs(traces: list[dict], known_items: set[str], group_by: str = "model") -> dict:
"""Map run name -> {item id: trace}, keeping only the newest trace per item.
Traces whose PR is not in the dataset are dropped: `feedback.db` is the
source for both, but a review can be traced without its row landing (the
posting step can fail after the model ran), and a run item pointing at a
non-existent dataset item is rejected.
"""
runs: dict[str, dict[str, dict]] = defaultdict(dict)
skipped_no_item, skipped_unknown = 0, 0
for tr in traces:
iid = trace_item_id(tr)
if iid is None:
skipped_unknown += 1
continue
if iid not in known_items:
skipped_no_item += 1
continue
run = "all-traces" if group_by == "none" else trace_model(tr)
prev = runs[run].get(iid)
if prev is None or _sort_key(tr) > _sort_key(prev):
runs[run][iid] = tr
return {
"runs": dict(runs),
"skipped_not_in_dataset": skipped_no_item,
"skipped_not_a_review": skipped_unknown,
}
def run_name(prefix: str, key: str) -> str:
return f"{prefix}-{key}" if prefix else key
# ---------------------------------------------------------------------------
# Writing the runs
# ---------------------------------------------------------------------------
def create_run(name: str, items: dict[str, dict], description: str = "") -> dict:
"""Link each (item, trace) pair into the named run. Idempotent per pair."""
created, failed = 0, []
for iid, tr in sorted(items.items()):
md = tr.get("metadata") or {}
body = {
"runName": name,
"runDescription": description,
"datasetItemId": iid,
"traceId": tr["id"],
"metadata": {
"model": trace_model(tr),
"engine": md.get("engine"),
"findings": md.get("findings"),
"duration_s": md.get("duration_s"),
"cost_basis": md.get("cost_basis"),
"linked_by": "eval_experiment.py",
},
}
st, resp = eb._call("POST", "/api/public/dataset-run-items", body)
if st in (200, 201):
created += 1
else:
failed.append({"item": iid, "status": st, "error": resp})
return {"run": name, "items_linked": created, "failed": failed}
def main(argv: list[str] | None = None) -> int:
ap = argparse.ArgumentParser(description=__doc__)
ap.add_argument("--dataset", default=eb.DATASET_NAME)
ap.add_argument("--group-by", choices=("model", "none"), default="model")
ap.add_argument("--prefix", default="baseline",
help="run name prefix; '' for the bare group key")
ap.add_argument("--dry-run", action="store_true")
args = ap.parse_args(argv)
traces = fetch_traces()
items = fetch_item_ids(args.dataset)
plan = plan_runs(traces, items, group_by=args.group_by)
report = {
"traces_read": len(traces),
"dataset_items": len(items),
"skipped_not_in_dataset": plan["skipped_not_in_dataset"],
"skipped_not_a_review": plan["skipped_not_a_review"],
"runs": {},
}
for key, mapping in sorted(plan["runs"].items()):
name = run_name(args.prefix, key)
if args.dry_run:
report["runs"][name] = {"items_would_link": len(mapping)}
continue
report["runs"][name] = create_run(
name,
mapping,
description=(
"Reviews already run by the pilot, linked after the fact. "
"Scores come from the traces; expectedOutput is the reviewer's "
"own prior output, not human-verified ground truth."
),
)
report["dry_run"] = args.dry_run
print(json.dumps(report, indent=2))
return 0
if __name__ == "__main__":
raise SystemExit(main())
+308
View File
@@ -0,0 +1,308 @@
#!/usr/bin/env python3
"""pragent pilot — LLM-as-a-judge evaluators for the reviewer.
The deterministic scorers in `eval_scores.py` measure *behaviour*: how many
findings, how severe, how much they cost. None of them can say whether a
finding was any good. With no human labels in `feedback.db`, a judge is the
only thing that can — so these two ask the questions that need no ground truth,
only the review itself:
`finding_actionability` — is each finding concrete enough to act on? A
reviewer that says "consider improving error handling" at file level is
indistinguishable from a useful one by finding count alone. This is the
failure mode a cheap model degrades into first.
`review_self_consistency` — does the summary agree with the findings it
posted? Claiming "no issues found" above a list of two criticals, or
describing a problem in prose that never became a finding, is a defect the
reviewer can commit entirely on its own.
Neither judge is asked whether a finding is *correct*. That needs the diff,
which these traces do not carry, and a judge asked to rule on correctness from
a summary alone will confabulate. Accuracy stays an open question until humans
start labelling — which is what `feedback_scores.py` is there to capture.
**The judge is a different model from the reviewer.** The reviewer runs
MiniMax-M2.7; the judge runs kimi-k2.7-code through the same headroom hub. A
model grading its own output agrees with itself for reasons that have nothing
to do with quality.
Evaluators score *observations*, and their variable mapping reads the
observation's own input/output — which is why `langfuse_trace` now writes the
review onto the generation and not just onto the trace.
Usage:
LANGFUSE_HOST=... LANGFUSE_PUBLIC_KEY=... LANGFUSE_SECRET_KEY=... \\
python3 eval_judges.py --dry-run
"""
from __future__ import annotations
import argparse
import json
import os
import sys
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import eval_bootstrap as eb # noqa: E402
# The headroom hub in front of the local Ollama, plus a small pass-through
# proxy (`judge-proxy` on 8802) that patches every `thinking` content block
# to carry the `signature` field Langfuse's Anthropic adapter requires. The
# underlying model is kimi-k2.7-code through the hub on 8790; the proxy fixes
# the shape so Mastra's Zod parse stops failing.
JUDGE_PROVIDER = "headroom-ollama"
JUDGE_BASE_URL = os.environ.get("PRAGENT_JUDGE_BASE_URL", "http://100.74.17.70:8802")
JUDGE_API_KEY = os.environ.get("PRAGENT_JUDGE_API_KEY", "ollama")
JUDGE_MODEL = os.environ.get("PRAGENT_JUDGE_MODEL", "kimi-k2.7-code:cloud")
# The trace names this project emits (`pr-review` on the trace, `opencode-review`
# on the generation). Filter on `traceName` rather than observation `name` — the
# observation-rule schema only exposes `traceName` as a stringOptions column, and
# every observation inside these traces is the review itself, so the narrowness
# is the same.
REVIEW_TRACE_NAMES = ["pr-review", "opencode-review"]
def _model_config() -> dict:
return {"provider": JUDGE_PROVIDER, "model": JUDGE_MODEL}
JUDGES = [
{
"name": "finding_actionability",
"prompt": (
"You are auditing the output of an automated code reviewer.\n\n"
"PR under review:\n{{input}}\n\n"
"What the reviewer produced:\n{{output}}\n\n"
"Rate how ACTIONABLE the findings are, from 0 to 1. A finding is "
"actionable when a developer could act on it without asking a "
"follow-up question: it points at a specific location, names a "
"concrete problem, and proposes a fix that could be applied.\n\n"
"Score 1.0 when every finding is specific and fixable. Score around "
"0.5 when findings identify a real area but leave the developer to "
"work out what to change. Score near 0.0 when findings are generic "
"advice that would apply to almost any pull request.\n\n"
"Judge only specificity and actionability. You cannot see the diff, "
"so do NOT attempt to judge whether a finding is factually correct, "
"and do not penalise a finding for being one you cannot verify.\n\n"
"If the reviewer reported no findings at all, return 1.0 and say in "
"your reasoning that there was nothing to judge — a silent review is "
"measured by finding_rate, not here."
),
"outputDefinition": {
"dataType": "NUMERIC",
"minValue": 0,
"maxValue": 1,
"reasoning": {
"description": (
"Name the least actionable finding and say what it would "
"need in order to be acted on."
)
},
"score": {"description": "0 = generic advice, 1 = every finding is specific and fixable."},
},
},
{
"name": "review_self_consistency",
"prompt": (
"You are auditing the output of an automated code reviewer.\n\n"
"PR under review:\n{{input}}\n\n"
"What the reviewer produced:\n{{output}}\n\n"
"The output contains a prose `summary` and a list of `findings`. "
"Decide whether the summary is CONSISTENT with the findings.\n\n"
"Inconsistent means, for example: the summary says no issues were "
"found while findings are listed; the summary describes a problem "
"that never became a finding; the summary characterises the severity "
"of the findings in a way the findings themselves contradict; or the "
"summary refers to files that appear in no finding and in no part of "
"the PR description.\n\n"
"A summary that adds context beyond the findings is NOT inconsistent "
"as long as nothing in it contradicts them. A review that found "
"nothing and says so is consistent.\n\n"
"You cannot see the diff. Judge the summary against the findings and "
"the PR title only — never against what you imagine the code does."
),
"outputDefinition": {
"dataType": "BOOLEAN",
"reasoning": {
"description": "Quote the part of the summary that conflicts with the findings, if any."
},
"score": {"description": "true = summary agrees with the findings, false = it contradicts them."},
},
},
]
# Both judges read the observation's own input/output.
MAPPING = [
{"variable": "input", "source": "input"},
{"variable": "output", "source": "output"},
]
# ---------------------------------------------------------------------------
# LLM connection
# ---------------------------------------------------------------------------
def ensure_llm_connection() -> dict:
"""Point the project at the judge model. Upserted on `provider`."""
body = {
"provider": JUDGE_PROVIDER,
"adapter": "anthropic",
"baseURL": JUDGE_BASE_URL,
"secretKey": JUDGE_API_KEY,
"customModels": [JUDGE_MODEL],
# The hub serves two local models and none of Anthropic's, so the
# default catalogue would be a list of models that all fail on use.
"withDefaultModels": False,
}
st, resp = eb._call("PUT", "/api/public/llm-connections", body)
return {"status": st, "ok": st in (200, 201), "provider": JUDGE_PROVIDER,
"error": None if st in (200, 201) else resp}
# ---------------------------------------------------------------------------
# Evaluators
# ---------------------------------------------------------------------------
def existing_evaluators() -> dict[str, str]:
"""name -> id for evaluators already in the project."""
out: dict[str, str] = {}
st, body = eb._call("GET", "/api/public/unstable/evaluators?limit=100")
if st == 200 and isinstance(body, dict):
for ev in body.get("data") or []:
out[ev.get("name")] = ev.get("id")
return out
def ensure_evaluators() -> dict:
"""Create each judge if no version exists for the name yet.
POST /evaluators with a name that already exists creates a new version, not
a no-op — re-running this script would pile up versions until the page
listing them is unreadable. Skip when an evaluator of that name is present.
"""
created, skipped, failed = {}, [], []
existing = set(existing_evaluators())
for judge in JUDGES:
if judge["name"] in existing:
skipped.append(judge["name"])
continue
body = {
"type": "llm_as_judge",
"name": judge["name"],
"prompt": judge["prompt"],
"outputDefinition": judge["outputDefinition"],
"modelConfig": _model_config(),
}
st, resp = eb._call("POST", "/api/public/unstable/evaluators", body, timeout=60.0)
if st in (200, 201) and isinstance(resp, dict):
created[judge["name"]] = resp.get("id")
else:
failed.append({"name": judge["name"], "status": st, "error": resp})
return {"created": created, "skipped": skipped, "failed": failed}
# ---------------------------------------------------------------------------
# Rules — what gets judged, and how often
# ---------------------------------------------------------------------------
def rule_body(name: str, judge_name: str, sampling: float) -> dict:
"""POST /evaluation-rules shape for an LLM-as-judge observation rule.
The judge is referenced by `name`+`scope`, not by id — ids name specific
versions, names name the evaluator across versions. Mapping is required at
both the rule root (the server validates it there) and inside `evaluator`
(the API echoes it back). Filter is on `traceName` because that is the only
stringOptions column the observation-rule schema exposes.
"""
return {
"name": name,
"enabled": True,
"target": "observation",
"sampling": sampling,
"filter": [
{"column": "traceName", "operator": "any of",
"value": REVIEW_TRACE_NAMES, "type": "stringOptions"},
],
"evaluator": {
"name": judge_name,
"scope": "project",
"variableMapping": MAPPING,
},
"mapping": MAPPING,
}
def ensure_rules(evaluator_ids: dict[str, str], sampling: float) -> dict:
"""Idempotent: existing rules with the same name are skipped, not duplicated.
The API has no `name`-keyed upsert; the convention is to POST once and
re-run the script to verify the response. A duplicate POST raises 409.
"""
created, failed, skipped = [], [], []
existing = existing_rule_names()
for name, eid in evaluator_ids.items():
if not eid:
continue
rule_name = f"{name}-on-reviews"
if rule_name in existing:
skipped.append(name)
continue
st, resp = eb._call(
"POST", "/api/public/unstable/evaluation-rules",
rule_body(rule_name, name, sampling), timeout=60.0,
)
if st in (200, 201):
created.append(name)
else:
failed.append({"rule": name, "status": st, "error": resp})
return {"created": created, "failed": failed, "skipped": skipped}
def existing_rule_names() -> set[str]:
"""Names of observation-target rules already in the project."""
out: set[str] = set()
st, body = eb._call("GET", "/api/public/unstable/evaluation-rules?limit=100")
if st == 200 and isinstance(body, dict):
for r in body.get("data") or []:
if r.get("target") == "observation":
out.add(r.get("name"))
return out
def main(argv: list[str] | None = None) -> int:
ap = argparse.ArgumentParser(description=__doc__)
ap.add_argument("--sampling", type=float, default=1.0,
help="fraction of matching observations to judge (default: all)")
ap.add_argument("--skip-connection", action="store_true")
ap.add_argument("--dry-run", action="store_true")
args = ap.parse_args(argv)
if args.dry_run:
print(json.dumps({
"would_connect": {"provider": JUDGE_PROVIDER, "baseURL": JUDGE_BASE_URL,
"model": JUDGE_MODEL},
"would_create": [j["name"] for j in JUDGES],
"existing_evaluators": sorted(existing_evaluators()),
"sampling": args.sampling,
}, indent=2))
return 0
report = {}
if not args.skip_connection:
report["llm_connection"] = ensure_llm_connection()
report["evaluators"] = ensure_evaluators()
ids = dict(report["evaluators"]["created"])
# Fall back to whatever is already registered, so a re-run still wires rules.
for name, eid in existing_evaluators().items():
ids.setdefault(name, eid)
report["rules"] = ensure_rules(
{j["name"]: ids.get(j["name"]) for j in JUDGES}, args.sampling
)
print(json.dumps(report, indent=2))
return 0
if __name__ == "__main__":
raise SystemExit(main())
+233
View File
@@ -0,0 +1,233 @@
#!/usr/bin/env python3
"""pragent pilot — deterministic review scorers.
Four numbers computed from a review that already happened, shipped to Langfuse
as scores on the review's trace. All are derived from data the reviewer already
has in hand: no LLM judge, no ground truth, no extra token spend.
Why these four and not `helpfulness`/`quality`
----------------------------------------------
They come from what the recorded reviews actually did, not from a generic eval
checklist:
* `severity_info_ratio` — of the findings ever posted to a PR, effectively all
landed at `info`. Either the model will not commit to a severity or the
per-repo `severity_threshold` is filtering the rest out. Trending the ratio
per model says which.
* `finding_rate` — most reviews post nothing at all. Silence on clean code is
the goal; silence because the run degraded is a failure. Same output, two
causes, and only the rate over time separates them.
* `dropped_findings` — `ai_review.parse_findings` discards any finding whose
`path`/`line` is unusable. That happens silently, so a model that emits ten
findings at invalid locations is indistinguishable from one that found
nothing. This is the only signal here that measures the *model's* output
rather than the review's.
* `cost_per_finding` — the equivalent-cost number is already trended per
review; per finding is what actually compares two models, since a cheaper
model that finds nothing is not cheaper.
None of these say whether a finding was *correct*. That needs labels, and the
labels come from `feedback_scores.py` once maintainers start reacting to review
comments. Read these as behavioural drift detectors, not as accuracy.
Fail-open, like every other telemetry path here: a scorer that raises returns no
score rather than failing the review.
"""
from __future__ import annotations
import uuid
from datetime import datetime, timezone
# Mirrors ai_review.SEVERITY_RANK. Duplicated rather than imported because this
# module is also run standalone (backfill) where ai_review's import side effects
# are unwanted.
SEVERITY_RANK = {"info": -1, "trivial": 0, "low": 1, "medium": 2, "high": 3, "critical": 4}
# Findings at or below this rank are "the model declined to commit". `trivial`
# and `info` are advisory by the reviewer's own prompt contract.
_ADVISORY_MAX_RANK = 0
# Score names. Named for what is measured, not for the mechanism producing it —
# these land on every trace and become the axis of every chart.
FINDING_RATE = "finding_rate"
SEVERITY_INFO_RATIO = "severity_info_ratio"
SEVERITY_MAX = "severity_max"
DROPPED_FINDINGS = "dropped_findings"
COST_PER_FINDING = "cost_per_finding"
def _sev(f: dict) -> str:
return str(f.get("severity") or "medium").strip().lower()
def finding_rate(findings: list[dict] | None) -> float:
"""How many findings this review posted. 0.0 is the restraint case."""
return float(len(findings or []))
def severity_info_ratio(findings: list[dict] | None) -> float | None:
"""Share of findings the model rated advisory (`info`/`trivial`).
`None` for a review with no findings — a ratio over an empty set is not 0,
it is undefined, and charting it as 0 would read as "perfectly calibrated".
"""
fs = findings or []
if not fs:
return None
advisory = sum(1 for f in fs if SEVERITY_RANK.get(_sev(f), 2) <= _ADVISORY_MAX_RANK)
return round(advisory / len(fs), 4)
def severity_max(findings: list[dict] | None) -> str:
"""Highest severity present, or `none` when the review was silent.
Categorical on purpose: the useful question is "did this review ever surface
something serious", and an average of severity ranks answers nothing.
"""
fs = findings or []
if not fs:
return "none"
top = max(fs, key=lambda f: SEVERITY_RANK.get(_sev(f), 2))
sev = _sev(top)
return sev if sev in SEVERITY_RANK else "medium"
def dropped_findings(raw_count: int | None, kept_count: int | None) -> float | None:
"""Findings the model emitted that the parser could not use.
`raw_count` is what came back in the JSON; `kept_count` is what survived
`_normalize_finding`. `None` when the caller could not determine the raw
count — better no score than a fabricated zero.
"""
if raw_count is None or kept_count is None:
return None
return float(max(0, int(raw_count) - int(kept_count)))
def cost_per_finding(cost_usd: float | None, findings: list[dict] | None) -> float | None:
"""Equivalent USD spent per finding posted.
`None` when nothing could be priced. A silent review divides by one, not by
zero: the run still cost money, and attributing that whole cost to "found
nothing" is the honest reading.
"""
if cost_usd is None:
return None
try:
c = float(cost_usd)
except (TypeError, ValueError):
return None
return round(c / max(1, len(findings or [])), 6)
def build_scores(
*,
trace_id: str,
findings: list[dict] | None,
environment: str,
cost_usd: float | None = None,
dropped_count: float | None = None,
timestamp: str | None = None,
comment: str = "",
) -> list[dict]:
"""The `score-create` ingestion events for one review.
`dropped_count` must be measured at parse time, not here: by the time
`findings` reaches this function the per-repo config has already filtered it
by severity threshold and `max_findings`, and those drops are the config
working as intended, not the model emitting garbage.
Returns [] rather than raising if something is unscoreable — scores are
telemetry and must never cost a review.
"""
# The ingestion envelope requires a timestamp on every event; omitting it
# gets the whole batch rejected with an HTTP 207 whose per-event 400s are
# easy to mistake for success.
ts = timestamp or datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
out: list[dict] = []
def add(name: str, value, data_type: str) -> None:
if value is None:
return
body = {
"id": str(uuid.uuid4()),
"traceId": trace_id,
"name": name,
"dataType": data_type,
"environment": environment,
}
if data_type == "CATEGORICAL":
body["value"] = str(value)
else:
body["value"] = float(value)
if comment:
body["comment"] = comment
out.append(
{
"id": str(uuid.uuid4()),
"type": "score-create",
"timestamp": ts,
"body": body,
}
)
try:
add(FINDING_RATE, finding_rate(findings), "NUMERIC")
add(SEVERITY_INFO_RATIO, severity_info_ratio(findings), "NUMERIC")
add(SEVERITY_MAX, severity_max(findings), "CATEGORICAL")
add(DROPPED_FINDINGS, dropped_count, "NUMERIC")
add(COST_PER_FINDING, cost_per_finding(cost_usd, findings), "NUMERIC")
except Exception: # pragma: no cover - defensive
return out
return out
# ---------------------------------------------------------------------------
# Score configs — the schema these scores must comply with
# ---------------------------------------------------------------------------
# Registered once per project via `eval_bootstrap.py`. Without configs the
# scores still ingest, but nothing constrains a future scorer from writing
# `severity_max="HIGH"` next to today's `"high"` and silently splitting the
# series in two.
SCORE_CONFIGS = [
{
"name": FINDING_RATE,
"dataType": "NUMERIC",
"minValue": 0,
"description": "Findings posted by one review. 0 = the reviewer stayed silent.",
},
{
"name": SEVERITY_INFO_RATIO,
"dataType": "NUMERIC",
"minValue": 0,
"maxValue": 1,
"description": "Share of a review's findings rated info/trivial. High = the model is not committing to a severity.",
},
{
"name": SEVERITY_MAX,
"dataType": "CATEGORICAL",
"categories": [
{"label": "none", "value": 0},
{"label": "info", "value": 1},
{"label": "trivial", "value": 2},
{"label": "low", "value": 3},
{"label": "medium", "value": 4},
{"label": "high", "value": 5},
{"label": "critical", "value": 6},
],
"description": "Highest severity surfaced by one review; 'none' when it posted nothing.",
},
{
"name": DROPPED_FINDINGS,
"dataType": "NUMERIC",
"minValue": 0,
"description": "Findings the model emitted that the parser rejected for an unusable path/line.",
},
{
"name": COST_PER_FINDING,
"dataType": "NUMERIC",
"minValue": 0,
"description": "Equivalent USD per finding posted. Silent reviews divide by 1, not 0.",
},
]
+247
View File
@@ -0,0 +1,247 @@
#!/usr/bin/env python3
"""pragent pilot — feedback DB to Langfuse scores.
`feedback.db` already records every reaction, thread resolution and reply a
maintainer leaves on a bot comment. That is the only ground truth pragent has
about whether a finding was any good, and until now it went to a markdown report
nobody reads and nowhere else. This ships it to Langfuse as session-level
scores, so "was the reviewer right" sits on the same axis as "what did it cost".
Session, not trace
------------------
`langfuse_trace` sets `sessionId` to `"{repo}#{pr}"` and lets the trace id be a
fresh uuid per review. Feedback arrives days later against a PR, not against one
particular re-run of the reviewer, and nothing in `feedback.db` records which
trace produced which comment. Scoring the session is therefore both the
available join and the honest granularity: this is feedback on the review of
this PR, not on one invocation.
Two scores, deliberately separated
----------------------------------
* `review_engagement` — the share of a PR's findings that got any human
response at all. This is a signal about the *feedback loop*, not the
reviewer: at the time of writing it is 0.0 across all 113 recorded reviews,
which is exactly the fact that makes an accuracy metric impossible today.
It must be watched first, because every other quality number is vapour
until it moves.
* `review_acceptance` — net verdict over the findings that *did* get a
response: (upvotes + resolved) - (downvotes + negation replies), normalised
to -1..1. Computed only over engaged findings, so an ignored review scores
`None` rather than 0. Zero would read as "humans judged this exactly
neutral"; the truth is nobody looked.
Fail-open and idempotent. Score ids are derived from (repo, pr, name) so a
re-run overwrites rather than duplicates.
"""
from __future__ import annotations
import argparse
import json
import os
import sqlite3
import sys
import uuid
from datetime import datetime, timezone
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from feedback_harvest import classify_reaction, _is_negation_reply # noqa: E402
REVIEW_ENGAGEMENT = "review_engagement"
REVIEW_ACCEPTANCE = "review_acceptance"
# Stable namespace so the same (repo, pr, score) always produces the same score
# id — Langfuse treats a repeated id as an update, which is what a backfill of a
# still-accumulating PR should do.
_NS = uuid.UUID("6f1d9c2e-4a77-4f2a-9c1a-0d3b5e8a7c41")
def _score_id(repo: str, pr: int, name: str) -> str:
return str(uuid.uuid5(_NS, f"{repo}#{pr}#{name}"))
def collect_pr_feedback(conn: sqlite3.Connection, repo: str, pr: int) -> dict:
"""Tally one PR's findings and the human responses attached to them.
Returns counts only — the scoring maths lives in `score_pr` so it can be
tested without a database.
"""
rows = conn.execute(
"SELECT id, comment_id FROM inline_finding WHERE repo = ? AND pr = ?",
(repo, pr),
).fetchall()
total = len(rows)
engaged = 0
positive = 0
negative = 0
for row in rows:
fid = row["id"] if isinstance(row, sqlite3.Row) else row[0]
cid = row["comment_id"] if isinstance(row, sqlite3.Row) else row[1]
pos = neg = 0
if cid is not None:
for r in conn.execute(
"SELECT content FROM reaction WHERE comment_id = ?", (cid,)
):
kind = classify_reaction(r[0])
if kind == "positive":
pos += 1
elif kind == "negative":
neg += 1
for r in conn.execute(
"SELECT resolved FROM thread_state WHERE finding_id = ?", (fid,)
):
# A resolved thread means the maintainer acted on the finding.
if r[0]:
pos += 1
# A reply counts as engagement either way; only a negation phrase makes
# it a vote against. A neutral reply ("done", "good catch, but…") is
# deliberately not a positive vote — it says someone looked, not that
# they agreed.
replied = 0
for r in conn.execute(
"SELECT body FROM reply WHERE finding_id = ?", (fid,)
):
replied += 1
if _is_negation_reply(r[0]):
neg += 1
if pos or neg or replied:
engaged += 1
positive += pos
negative += neg
return {"total": total, "engaged": engaged, "positive": positive, "negative": negative}
def score_pr(tally: dict) -> dict:
"""Turn one PR's tally into score values.
`review_acceptance` is `None` when nothing was engaged — see the module
docstring on why that is not 0.
"""
total = int(tally.get("total") or 0)
engaged = int(tally.get("engaged") or 0)
pos = int(tally.get("positive") or 0)
neg = int(tally.get("negative") or 0)
engagement = round(engaged / total, 4) if total else None
acceptance = None
if pos or neg:
acceptance = round((pos - neg) / (pos + neg), 4)
return {REVIEW_ENGAGEMENT: engagement, REVIEW_ACCEPTANCE: acceptance}
def build_score_events(
repo: str, pr: int, values: dict, environment: str = "default",
timestamp: str | None = None,
) -> list[dict]:
"""`score-create` events for one PR's feedback.
Every event carries a timestamp: the ingestion endpoint rejects those that
do not, and it reports the rejection as a per-event 400 inside an HTTP 207,
which reads as success to a caller that only checks the status code.
"""
ts = timestamp or datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
events = []
for name, value in values.items():
if value is None:
continue
events.append(
{
"id": str(uuid.uuid4()),
"type": "score-create",
"timestamp": ts,
"body": {
"id": _score_id(repo, pr, name),
"sessionId": f"{repo}#{pr}",
"name": name,
"value": float(value),
"dataType": "NUMERIC",
"environment": environment,
"comment": f"from feedback.db · {repo}#{pr}",
},
}
)
return events
SCORE_CONFIGS = [
{
"name": REVIEW_ENGAGEMENT,
"dataType": "NUMERIC",
"minValue": 0,
"maxValue": 1,
"description": "Share of a PR's findings that drew any human reaction, resolution or reply. 0 = nobody engaged with the review.",
},
{
"name": REVIEW_ACCEPTANCE,
"dataType": "NUMERIC",
"minValue": -1,
"maxValue": 1,
"description": "Net human verdict over engaged findings: +1 all accepted, -1 all rejected. Absent when nothing was engaged.",
},
]
def iter_prs(conn: sqlite3.Connection):
for row in conn.execute(
"SELECT DISTINCT repo, pr FROM inline_finding ORDER BY repo, pr"
):
yield row[0], int(row[1])
def backfill(db_path: str, *, environment: str = "default", dry_run: bool = False) -> dict:
"""Score every PR in the feedback DB. Returns a summary dict."""
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
events: list[dict] = []
scanned = 0
engaged_prs = 0
try:
for repo, pr in iter_prs(conn):
scanned += 1
tally = collect_pr_feedback(conn, repo, pr)
values = score_pr(tally)
if (values.get(REVIEW_ENGAGEMENT) or 0) > 0:
engaged_prs += 1
events.extend(build_score_events(repo, pr, values, environment))
finally:
conn.close()
summary = {"prs_scanned": scanned, "prs_with_engagement": engaged_prs, "scores": len(events)}
if dry_run or not events:
summary["posted"] = False
return summary
import langfuse_trace
conf = langfuse_trace._enabled()
if conf is None:
summary["posted"] = False
summary["error"] = "Langfuse not configured (LANGFUSE_HOST / keys unset)"
return summary
host, pk, sk = conf
status = langfuse_trace._post(host, pk, sk, events, 15.0)
summary["posted"] = status in (200, 201, 207)
summary["http_status"] = status
return summary
def main() -> int:
ap = argparse.ArgumentParser(description="Ship feedback.db verdicts to Langfuse as scores")
ap.add_argument("--db", default=os.environ.get("PRAGENT_FEEDBACK_DB", "/data/feedback.db"))
ap.add_argument("--environment", default="default")
ap.add_argument("--dry-run", action="store_true")
args = ap.parse_args()
summary = backfill(args.db, environment=args.environment, dry_run=args.dry_run)
print(json.dumps(summary, indent=2))
return 0 if summary.get("posted") or args.dry_run else 1
if __name__ == "__main__":
raise SystemExit(main())
+468
View File
@@ -0,0 +1,468 @@
#!/usr/bin/env python3
"""pragent pilot — Langfuse trace emission.
Ships one trace per PR review to a self-hosted Langfuse (v3) so the reviewer's
token spend, latency and per-model behaviour are queryable outside the review
body. The review body already renders a usage table; that table is per-PR and
disappears into Gitea. This is the same numbers, aggregated.
Why hand-rolled instead of the `langfuse` SDK: the pilot image is stdlib-only
(see pilot/Dockerfile — no requirements.txt anywhere in the repo), and the
ingestion API is a single authenticated POST of a JSON batch. Pulling an SDK
plus its otel dependency tree into a fail-open telemetry side-path is a bad
trade.
Provider split
--------------
`environment` on every trace is either `ollama` or `claude`, derived from the
resolved display model (`resolve_environment`). That is what keeps the two
spend stories separate in Langfuse: every dashboard, filter and cost breakdown
takes an environment selector, so "what did the local/self-hosted path cost"
and "what did the Claude path cost" are two views of one project rather than
two projects with two key pairs to rotate. Tags carry the finer split
(`provider:headroom`, `model:...`, `engine:opencode`).
Cost
----
The pilot's own path bills $0 (headroom proxy, no per-token charge), so the
`cost` reported to Langfuse is the *equivalent* cost from `cost_model` — what
the same tokens would bill on the comparison model. That is the number worth
trending; a chart of $0.00 is not.
A model is "free" when `cost_model.PRICES` has no entry for it (MiniMax-M2.7,
glm-5.2:cloud) or when its entry is all zeros (the self-hosted vLLM qwen). In
both cases the reported cost is priced against the comparison target instead —
same precedence the review body uses: `.pr-review.json:cost_target` >
`PRAGENT_PRICE_TARGET` > `claude-sonnet-5`. A paid model is priced as itself.
Because a hypothetical and a real charge must never be read as the same
number, every trace is tagged `cost:actual` or `cost:equivalent:<target>`, and
the generation's metadata carries `cost_basis`.
Fail-open: every entry point swallows its own exceptions. Telemetry must never
cost a review.
Env:
LANGFUSE_HOST e.g. http://langfuse-web.langfuse.svc.cluster.local:3000
LANGFUSE_PUBLIC_KEY pk-lf-...
LANGFUSE_SECRET_KEY sk-lf-...
LANGFUSE_TIMEOUT seconds, default 5
LANGFUSE_DEBUG 1 to log ingestion failures to stderr
Disabled (silently) when host or either key is unset.
"""
from __future__ import annotations
import base64
import json
import os
import sys
import time
import urllib.error
import urllib.request
import uuid
from datetime import datetime, timezone
INGESTION_PATH = "/api/public/ingestion"
# Model-key prefixes that mean "this review ran against Anthropic-shaped
# billing". Everything else (glm, MiniMax, qwen, local vLLM) is the ollama /
# self-hosted side of the split.
_CLAUDE_PREFIXES = ("claude-", "anthropic/")
def _now_iso() -> str:
return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
def _enabled() -> tuple[str, str, str] | None:
host = (os.environ.get("LANGFUSE_HOST") or "").strip().rstrip("/")
pk = (os.environ.get("LANGFUSE_PUBLIC_KEY") or "").strip()
sk = (os.environ.get("LANGFUSE_SECRET_KEY") or "").strip()
if not host or not pk or not sk:
return None
return host, pk, sk
def _debug(msg: str) -> None:
if os.environ.get("LANGFUSE_DEBUG"):
print(f"pragent/langfuse: {msg}", file=sys.stderr, flush=True)
def strip_provider(model: str) -> str:
"""`headroom/claude-sonnet-5` -> `claude-sonnet-5`. Bare names pass through."""
return model.split("/", 1)[1] if "/" in model else model
def provider_of(model: str) -> str:
"""The opencode provider block a display model routes through."""
return model.split("/", 1)[0] if "/" in model else "headroom"
def resolve_environment(model: str) -> str:
"""Which spend story this review belongs to: `claude` or `ollama`.
Keyed off the bare model name, not the provider, because both paths route
through the same `headroom` proxy — `headroom/claude-sonnet-5` is Claude
spend, `headroom/glm-5.2:cloud` is not.
"""
bare = strip_provider(model).lower()
return "claude" if bare.startswith(_CLAUDE_PREFIXES) else "ollama"
def _usage_details(usage: dict) -> dict:
"""opencode's usage dict -> Langfuse `usageDetails`.
Langfuse sums every key except the ones it knows are derived, so `input`
here is the *uncached* portion: reporting both `input` (which opencode
reports as the full input, cache included) and `cache_read_input_tokens`
would double-count.
"""
inp = int(usage.get("input") or 0)
cache_read = int(usage.get("cache_read") or 0)
cache_write = int(usage.get("cache_write") or 0)
details = {
"input": max(0, inp - cache_read),
"output": int(usage.get("output") or 0),
}
if cache_read:
details["cache_read_input_tokens"] = cache_read
if cache_write:
details["cache_write_input_tokens"] = cache_write
reasoning = int(usage.get("reasoning") or 0)
if reasoning:
details["reasoning"] = reasoning
return details
DEFAULT_PRICE_TARGET = "claude-sonnet-5"
def resolve_price_target(price_target: str | None = None) -> str:
"""The model to price free/unknown runs against.
Mirrors `ai_review._resolve_price_target`: an explicit target (which the
caller reads from `.pr-review.json:cost_target`) wins, then
`PRAGENT_PRICE_TARGET`, then Sonnet.
"""
if price_target and price_target.strip():
return price_target.strip()
env = os.environ.get("PRAGENT_PRICE_TARGET", "").strip()
return env or DEFAULT_PRICE_TARGET
def _is_free(price) -> bool:
"""A price entry that charges nothing — self-hosted or proxied at no cost."""
return price.input == 0 and price.output == 0
def _cost_details(usage: dict, model: str, price_target: str | None = None) -> tuple[dict, str]:
"""USD for this usage plus the basis it was computed on.
Returns `({"total": …}, basis)` where basis is `actual` for a model that
genuinely bills, or `equivalent:<target>` for one that does not. `({}, "")`
when nothing can be priced at all — better no number than a wrong one.
Local import + broad except: `cost_model` is only present on the opencode
path, and an unknown model key must not break telemetry.
"""
try:
from cost_model import PRICES, Usage, cost
bare = strip_provider(model)
price = PRICES.get(bare)
basis = "actual"
if price is None or _is_free(price):
# MiniMax / glm / self-hosted qwen: $0 through the proxy, so the
# useful number is what these tokens would have billed elsewhere.
target = resolve_price_target(price_target)
price = PRICES.get(target)
if price is None:
_debug(f"comparison target {target!r} not in PRICES")
return {}, ""
basis = f"equivalent:{target}"
u = Usage(
uncached_input=max(0, int(usage.get("input") or 0) - int(usage.get("cache_read") or 0)),
cached_input=int(usage.get("cache_read") or 0),
cache_writes=int(usage.get("cache_write") or 0),
output=int(usage.get("output") or 0),
)
return {"total": round(cost(u, price), 6)}, basis
except Exception as e: # pragma: no cover - defensive
_debug(f"cost lookup failed for {model!r}: {e}")
return {}, ""
def _severity_counts(findings: list[dict] | None) -> dict:
counts: dict[str, int] = {}
for f in findings or []:
sev = str(f.get("severity") or "unknown").lower()
counts[sev] = counts.get(sev, 0) + 1
return counts
def build_batch(
*,
repo: str,
index: str,
sha: str,
title: str,
model: str,
usage: dict | None,
findings: list[dict] | None = None,
summary: str = "",
engine: str = "opencode",
tier: str = "",
lenses: list[str] | None = None,
trace_id: str | None = None,
release: str = "",
price_target: str | None = None,
dropped_count: float | None = None,
) -> list[dict]:
"""The ingestion batch for one review: a trace, a generation, and scores.
Split out from `emit_review_trace` so the shape is testable without a
Langfuse to POST to.
`dropped_count` is how many findings the parser rejected for an unusable
`path`/`line`, measured where the model output was parsed. Passing it turns
on the `dropped_findings` score; leaving it `None` omits that score rather
than reporting a zero the caller never measured.
"""
usage = usage or {}
tid = trace_id or str(uuid.uuid4())
ts = _now_iso()
env = resolve_environment(model)
duration = float(usage.get("duration_s") or 0.0)
started = datetime.fromtimestamp(
time.time() - duration, tz=timezone.utc
).isoformat().replace("+00:00", "Z")
tags = [
f"provider:{provider_of(model)}",
f"model:{strip_provider(model)}",
f"engine:{engine}",
f"repo:{repo}",
]
if tier:
tags.append(f"tier:{tier}")
for lens in lenses or []:
tags.append(f"lens:{lens}")
costs, cost_basis = _cost_details(usage, model, price_target) if usage else ({}, "")
if cost_basis:
# Filterable in Langfuse, so an equivalent-cost chart can never be
# mistaken for money actually spent.
tags.append(f"cost:{cost_basis}")
metadata = {
"repo": repo,
"pr": index,
"sha": sha,
"engine": engine,
"steps": usage.get("steps"),
"duration_s": duration or None,
"findings": len(findings or []),
"severities": _severity_counts(findings),
"provider_cost_usd": usage.get("cost"),
"cost_basis": cost_basis or None,
}
if lenses:
metadata["lenses"] = lenses
if tier:
metadata["tier"] = tier
metadata = {k: v for k, v in metadata.items() if v not in (None, {}, [])}
trace_body = {
"id": tid,
"name": "pr-review",
"timestamp": ts,
"environment": env,
"sessionId": f"{repo}#{index}",
"input": _review_input(repo, index, sha, title),
"output": _review_output(summary, findings),
"metadata": metadata,
"tags": tags,
}
if release:
trace_body["release"] = release
events = [
{
"id": str(uuid.uuid4()),
"type": "trace-create",
"timestamp": ts,
"body": trace_body,
}
]
if usage:
gen_body = {
"id": str(uuid.uuid4()),
"traceId": tid,
"type": "GENERATION",
"name": f"{engine}-review",
"environment": env,
"startTime": started,
"endTime": ts,
"model": strip_provider(model),
"usageDetails": _usage_details(usage),
"metadata": metadata,
"level": "DEFAULT",
# Repeated from the trace on purpose: an evaluator's variable
# mapping reads the *observation's* input/output, so a generation
# left blank cannot be judged at all.
"input": _review_input(repo, index, sha, title),
"output": _review_output(summary, findings),
}
if costs:
gen_body["costDetails"] = costs
events.append(
{
"id": str(uuid.uuid4()),
"type": "generation-create",
"timestamp": ts,
"body": gen_body,
}
)
events.extend(
_score_events(
trace_id=tid,
findings=findings,
environment=env,
cost_usd=costs.get("total"),
dropped_count=dropped_count,
timestamp=ts,
cost_basis=cost_basis,
)
)
return events
MAX_JUDGED_FINDINGS = 25
_FIELD_CAP = 600
def _review_input(repo: str, index, sha: str, title: str) -> dict:
return {"repo": repo, "pr": index, "sha": sha, "title": title}
def _review_output(summary: str, findings) -> dict:
"""What the reviewer actually said, in a shape an evaluator can read.
The findings themselves are included, not just their count. A judge given
only `{"summary": ..., "findings": 3}` can say nothing about whether those
three findings are specific, actionable, or consistent with the summary —
which is the whole question worth asking of a reviewer that has no ground
truth to check against.
Capped rather than complete: this rides in every ingestion batch, and a
review with 80 findings would push the payload past what is reasonable to
store per trace. `finding_count` stays exact so nothing reading the count
is misled by the cap.
"""
items = list(findings or [])
return {
"summary": summary[:2000],
"finding_count": len(items),
"findings_truncated": len(items) > MAX_JUDGED_FINDINGS,
"findings": [
{
"path": f.get("path"),
"line": f.get("line"),
"severity": f.get("severity"),
"problem": str(f.get("problem") or "")[:_FIELD_CAP],
"fix": str(f.get("fix") or "")[:_FIELD_CAP],
}
for f in items[:MAX_JUDGED_FINDINGS]
],
}
def _score_events(*, cost_basis: str, **kwargs) -> list[dict]:
"""Deterministic scores for this review, or [] if the scorer is missing.
Local import + blanket except for the same reason the rest of this module
swallows: `eval_scores` is optional, and a scoring bug must not cost the
trace it was supposed to annotate.
"""
try:
import eval_scores
# The cost score is only meaningful next to its basis — a $/finding
# figure computed from an equivalent price is not money that was spent.
comment = f"cost basis: {cost_basis}" if cost_basis else ""
return eval_scores.build_scores(comment=comment, **kwargs)
except Exception as e: # pragma: no cover - defensive
_debug(f"scoring failed: {e}")
return []
def _post(host: str, pk: str, sk: str, batch: list[dict], timeout: float) -> int:
payload = json.dumps({"batch": batch}).encode("utf-8")
auth = base64.b64encode(f"{pk}:{sk}".encode("utf-8")).decode("ascii")
req = urllib.request.Request(
host + INGESTION_PATH,
data=payload,
headers={
"Content-Type": "application/json",
"Authorization": f"Basic {auth}",
"User-Agent": "pragent-pilot/1.0",
},
method="POST",
)
with urllib.request.urlopen(req, timeout=timeout) as resp:
_warn_on_rejected_events(resp.read())
return resp.status
def _warn_on_rejected_events(raw: bytes) -> None:
"""Surface per-event rejections hiding inside a 207.
The ingestion endpoint answers 207 Multi-Status when *some* events failed,
so a caller that only checks the status code reads a batch where every
single event was rejected as a success. That failure mode is invisible
exactly when it matters — the traces simply never appear.
"""
try:
body = json.loads(raw or b"{}")
errors = body.get("errors") or []
if errors:
first = errors[0]
_debug(
f"{len(errors)} event(s) rejected by ingestion; "
f"first: status={first.get('status')} {first.get('error')}"
)
except Exception: # pragma: no cover - never let logging break emission
pass
def emit_review_trace(**kwargs) -> bool:
"""Ship one review's trace. Returns True if Langfuse accepted it.
No-op (False) when Langfuse is unconfigured. Never raises — a telemetry
outage must not turn into a failed review.
"""
conf = _enabled()
if conf is None:
return False
host, pk, sk = conf
try:
timeout = float(os.environ.get("LANGFUSE_TIMEOUT", "5"))
except ValueError:
timeout = 5.0
try:
batch = build_batch(**kwargs)
status = _post(host, pk, sk, batch, timeout)
if status not in (200, 201, 207):
_debug(f"ingestion returned HTTP {status}")
return False
return True
except urllib.error.HTTPError as e:
_debug(f"ingestion HTTP {e.code}: {e.read()[:300]!r}")
except Exception as e:
_debug(f"ingestion failed: {e}")
return False
+4 -4
View File
@@ -397,7 +397,7 @@ def install_config(src: str, dst: str) -> bool:
private-network addresses. Real values are supplied at runtime and patched
in here.
Env var convention (case-sensitive provider name `headroom`, `local`):
Env var convention (case-sensitive provider name `headroom`, `vllm-qwen38`):
PRAGENT_<NAME>_BASE_URL per-provider endpoint override
PRAGENT_<NAME>_API_KEY per-provider API key override
@@ -406,9 +406,9 @@ def install_config(src: str, dst: str) -> bool:
PRAGENT_MODEL_API_KEY legacy catchall (same)
Per-provider wins over the catchall. The first 2 win when the operator
needs a different endpoint per upstream (e.g. headroom MiniMax, local
ai-workstation). The catchall keeps the single-provider deploys from
needing any env config.
needs a different endpoint per upstream (e.g. headroom MiniMax,
vllm-qwen38 ai-workstation). The catchall keeps the single-provider
deploys from needing any env config.
This is done in Python rather than with opencode's own `{env:VAR}` config
templating because the reviewer subprocess runs with an allow-listed
+4 -3
View File
@@ -462,11 +462,12 @@ def test_resolve_display_model_precedence(monkeypatch):
ai_review._resolve_display_model("MiniMax-M2.7", {"model": "claude-sonnet-5"})
== "headroom/claude-sonnet-5"
)
# Self-hosted models carry provider="local" → routes to the `local`
# provider block in opencode.json (AI workstation on 192.168.1.79:18020).
# Self-hosted models carry provider="vllm-qwen38" → routes to the
# matching provider block in opencode.json (AI workstation on
# 192.168.1.79:18020).
assert (
ai_review._resolve_display_model("MiniMax-M2.7", {"model": "qwen3.8-27b"})
== "local/qwen3.8-27b"
== "vllm-qwen38/qwen3.8-27b"
)
# 3. Env wins over config
monkeypatch.setenv("OPENCODE_MODEL", "headroom/MiniMax-M2.7")
+158
View File
@@ -0,0 +1,158 @@
"""Tests for the eval bootstrap's dataset-item construction."""
import os
import sqlite3
import sys
import urllib.parse
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", "pilot"))
import eval_bootstrap as eb # noqa: E402
# --- item_id --------------------------------------------------------------
def test_item_id_has_no_path_separator():
"""A `/` would split the UI's item route into extra path segments."""
assert "/" not in eb.item_id("netcracker/interview", 29)
def test_item_id_has_no_fragment_marker():
"""Everything after a `#` is a fragment the browser never sends."""
assert "#" not in eb.item_id("netcracker/interview", 29)
def test_item_id_survives_a_url_round_trip():
"""The id must appear verbatim in a path, needing no percent-encoding."""
ident = eb.item_id("netcracker/interview", 29)
assert urllib.parse.quote(ident, safe="") == ident
def test_item_id_keeps_repo_and_pr_readable():
assert eb.item_id("netcracker/interview", 29) == "netcracker__interview__pr29"
def test_item_id_is_unique_per_pr():
assert eb.item_id("o/r", 1) != eb.item_id("o/r", 2)
def test_item_id_is_unique_per_repo():
assert eb.item_id("o/one", 1) != eb.item_id("o/two", 1)
def test_item_id_accepts_a_string_pr():
assert eb.item_id("o/r", "29") == eb.item_id("o/r", 29)
# --- read_review_items ----------------------------------------------------
def _db(tmp_path, rows, findings=()):
path = str(tmp_path / "feedback.db")
conn = sqlite3.connect(path)
conn.execute(
"CREATE TABLE review (repo TEXT, pr INTEGER, posted_at INTEGER, head_sha TEXT)"
)
conn.execute(
"CREATE TABLE inline_finding (repo TEXT, pr INTEGER, path TEXT, line INTEGER,"
" severity TEXT, problem TEXT, fix TEXT)"
)
conn.executemany("INSERT INTO review VALUES (?,?,?,?)", rows)
conn.executemany("INSERT INTO inline_finding VALUES (?,?,?,?,?,?,?)", findings)
conn.commit()
conn.close()
return path
def test_items_use_url_safe_ids(tmp_path):
path = _db(tmp_path, [("netcracker/interview", 29, 100, "abc")])
items = eb.read_review_items(path)
assert [i["id"] for i in items] == ["netcracker__interview__pr29"]
def test_item_input_keeps_the_real_repo_name(tmp_path):
"""The id is mangled for the URL; the payload must stay faithful."""
path = _db(tmp_path, [("netcracker/interview", 29, 100, "abc")])
item = eb.read_review_items(path)[0]
assert item["input"]["repo"] == "netcracker/interview"
assert item["input"]["pr"] == 29
def test_one_item_per_pr_not_per_review(tmp_path):
path = _db(
tmp_path,
[
("o/r", 1, 100, "a"),
("o/r", 1, 200, "b"),
("o/r", 2, 300, "c"),
],
)
items = eb.read_review_items(path)
assert [i["id"] for i in items] == ["o__r__pr1", "o__r__pr2"]
assert items[0]["metadata"]["reviews_run"] == 2
def test_items_are_not_flagged_as_human_labelled(tmp_path):
path = _db(tmp_path, [("o/r", 1, 100, "a")])
assert eb.read_review_items(path)[0]["metadata"]["labelled_by_human"] is False
# --- metadata facets ------------------------------------------------------
def _md(findings=(), repo="netcracker/interview", pr=29):
return eb._item_metadata(
repo=repo, pr=pr, head_sha="abc", reviews_run=2, last_seen=1788189422,
findings=[{"severity": s} for s in findings],
)
def test_metadata_carries_the_repo_for_filtering():
assert _md()["repo"] == "netcracker/interview"
def test_metadata_splits_owner_from_repo_name():
"""A filter on the joined repo can match one repo; owner matches an org."""
md = _md()
assert md["owner"] == "netcracker"
assert md["repo_name"] == "interview"
def test_owner_falls_back_when_the_repo_is_unqualified():
md = _md(repo="standalone")
assert md["owner"] == "standalone"
assert md["repo_name"] == "standalone"
def test_metadata_values_are_filterable_primitives():
"""Nested objects and lists are not reachable from the filter bar."""
for key, value in _md(["high"]).items():
assert isinstance(value, (str, int, float, bool)), key
def test_max_severity_is_the_worst_finding():
assert _md(["low", "critical", "medium"])["max_severity"] == "critical"
def test_max_severity_is_none_not_absent_for_a_silent_review():
md = _md([])
assert md["max_severity"] == "none"
assert md["has_findings"] is False
def test_unknown_severity_does_not_win_the_max():
assert _md(["banana", "low"])["max_severity"] == "low"
def test_severity_comparison_ignores_case():
assert _md(["HIGH"])["max_severity"] == "high"
def test_finding_count_matches_the_findings():
md = _md(["low", "low"])
assert md["finding_count"] == 2
assert md["has_findings"] is True
def test_last_reviewed_is_exposed_both_ways():
"""The epoch sorts; the ISO string is what a human reads in a filter."""
md = _md()
assert md["last_reviewed_at"] == 1788189422
assert md["last_reviewed_iso"].startswith("2026-08-31T")
+159
View File
@@ -0,0 +1,159 @@
"""Tests for linking existing review traces into dataset runs."""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", "pilot"))
import eval_experiment as ex # noqa: E402
def trace(tid, repo="o/r", pr=1, model="M2", ts="2026-08-01T00:00:00Z", **md):
meta = {"repo": repo, "pr": pr}
meta.update(md)
return {
"id": tid,
"timestamp": ts,
"tags": [f"model:{model}", "engine:opencode"],
"metadata": meta,
}
# --- trace_model ----------------------------------------------------------
def test_model_read_from_tag():
assert ex.trace_model(trace("t1", model="MiniMax-M2.7")) == "MiniMax-M2.7"
def test_model_falls_back_when_untagged():
assert ex.trace_model({"tags": ["engine:opencode"]}) == "unknown"
def test_model_falls_back_when_tags_absent():
assert ex.trace_model({}) == "unknown"
# --- trace_item_id --------------------------------------------------------
def test_item_id_matches_the_bootstrap_scheme():
assert ex.trace_item_id(trace("t1", repo="netcracker/interview", pr=29)) == \
"netcracker__interview__pr29"
def test_trace_without_repo_is_not_an_item():
assert ex.trace_item_id({"metadata": {"pr": 1}}) is None
def test_trace_without_pr_is_not_an_item():
assert ex.trace_item_id({"metadata": {"repo": "o/r"}}) is None
def test_trace_without_metadata_is_not_an_item():
assert ex.trace_item_id({}) is None
# --- plan_runs ------------------------------------------------------------
ITEMS = {"o__r__pr1", "o__r__pr2"}
def test_traces_group_by_model():
plan = ex.plan_runs(
[trace("a", pr=1, model="x"), trace("b", pr=2, model="y")], ITEMS
)
assert set(plan["runs"]) == {"x", "y"}
def test_group_by_none_collapses_to_one_run():
plan = ex.plan_runs(
[trace("a", pr=1, model="x"), trace("b", pr=2, model="y")],
ITEMS,
group_by="none",
)
assert list(plan["runs"]) == ["all-traces"]
def test_only_the_newest_trace_per_item_is_kept():
"""A re-reviewed PR has many traces; a run takes one output per input."""
plan = ex.plan_runs(
[
trace("old", pr=1, ts="2026-08-01T00:00:00Z"),
trace("new", pr=1, ts="2026-08-09T00:00:00Z"),
],
ITEMS,
)
assert plan["runs"]["M2"]["o__r__pr1"]["id"] == "new"
def test_newest_wins_regardless_of_input_order():
older = trace("old", pr=1, ts="2026-08-01T00:00:00Z")
newer = trace("new", pr=1, ts="2026-08-09T00:00:00Z")
for order in ([older, newer], [newer, older]):
plan = ex.plan_runs(order, ITEMS)
assert plan["runs"]["M2"]["o__r__pr1"]["id"] == "new"
def test_trace_for_a_pr_outside_the_dataset_is_skipped():
plan = ex.plan_runs([trace("a", pr=99)], ITEMS)
assert plan["runs"] == {}
assert plan["skipped_not_in_dataset"] == 1
def test_non_review_trace_is_counted_separately():
plan = ex.plan_runs([{"id": "x", "metadata": {}}], ITEMS)
assert plan["skipped_not_a_review"] == 1
assert plan["skipped_not_in_dataset"] == 0
def test_same_pr_different_models_lands_in_both_runs():
plan = ex.plan_runs([trace("a", pr=1, model="x"), trace("b", pr=1, model="y")], ITEMS)
assert plan["runs"]["x"]["o__r__pr1"]["id"] == "a"
assert plan["runs"]["y"]["o__r__pr1"]["id"] == "b"
# --- run_name -------------------------------------------------------------
def test_run_name_prefixed():
assert ex.run_name("baseline", "MiniMax-M2.7") == "baseline-MiniMax-M2.7"
def test_empty_prefix_leaves_the_key_bare():
assert ex.run_name("", "MiniMax-M2.7") == "MiniMax-M2.7"
# --- create_run -----------------------------------------------------------
def test_create_run_posts_one_item_per_pair(monkeypatch):
calls = []
def fake_call(method, path, body=None, timeout=20.0):
calls.append((method, path, body))
return 201, {}
monkeypatch.setattr(ex.eb, "_call", fake_call)
res = ex.create_run("run-1", {"o__r__pr1": trace("t1"), "o__r__pr2": trace("t2", pr=2)})
assert res["items_linked"] == 2
assert res["failed"] == []
assert {c[1] for c in calls} == {"/api/public/dataset-run-items"}
assert {c[2]["runName"] for c in calls} == {"run-1"}
def test_create_run_links_the_trace_to_the_item(monkeypatch):
seen = {}
def fake_call(method, path, body=None, timeout=20.0):
seen.update(body)
return 201, {}
monkeypatch.setattr(ex.eb, "_call", fake_call)
ex.create_run("run-1", {"o__r__pr1": trace("t1")})
assert seen["datasetItemId"] == "o__r__pr1"
assert seen["traceId"] == "t1"
assert seen["metadata"]["model"] == "M2"
def test_create_run_reports_rejected_items(monkeypatch):
monkeypatch.setattr(ex.eb, "_call", lambda *a, **k: (400, "nope"))
res = ex.create_run("run-1", {"o__r__pr1": trace("t1")})
assert res["items_linked"] == 0
assert res["failed"][0]["item"] == "o__r__pr1"
assert res["failed"][0]["status"] == 400
+126
View File
@@ -0,0 +1,126 @@
"""Tests for the LLM-as-judge evaluator bootstrap."""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", "pilot"))
import eval_judges as ej # noqa: E402
# --- rule_body ------------------------------------------------------------
def test_rule_body_targets_observations():
"""Trace-level rules wouldn't see observation input/output."""
body = ej.rule_body("rule-x", "finding_actionability", 1.0)
assert body["target"] == "observation"
assert body["enabled"] is True
def test_rule_body_filters_on_trace_name():
"""`name` isn't a stringOptions column; only `traceName` is."""
body = ej.rule_body("rule-x", "finding_actionability", 1.0)
f = body["filter"][0]
assert f["column"] == "traceName"
assert f["operator"] == "any of"
assert f["type"] == "stringOptions"
assert "pr-review" in f["value"]
def test_rule_body_references_evaluator_by_name():
"""Ids are version-specific; rules must name the evaluator across versions."""
body = ej.rule_body("rule-x", "finding_actionability", 1.0)
assert body["evaluator"]["name"] == "finding_actionability"
assert body["evaluator"]["scope"] == "project"
def test_rule_body_maps_input_and_output():
"""Both judges read the observation's own input/output."""
body = ej.rule_body("rule-x", "any", 1.0)
sources = {m["source"] for m in body["mapping"]}
assert sources == {"input", "output"}
def test_rule_body_carries_mapping_at_both_levels():
"""The server validates `mapping` at the rule root and echoes it on the evaluator."""
body = ej.rule_body("rule-x", "any", 1.0)
assert body["mapping"]
assert body["evaluator"]["variableMapping"] == body["mapping"]
def test_rule_body_passes_sampling_through():
assert ej.rule_body("r", "any", 0.25)["sampling"] == 0.25
# --- ensure_evaluators idempotency ---------------------------------------
def test_ensure_evaluators_skips_existing(monkeypatch):
seen = []
def fake_call(method, path, body=None, timeout=20.0):
seen.append(path)
return 200, {}
monkeypatch.setattr(ej.eb, "_call", fake_call)
monkeypatch.setattr(ej, "existing_evaluators",
lambda: {"finding_actionability": "id-1", "review_self_consistency": "id-2"})
res = ej.ensure_evaluators()
assert res["created"] == {}
assert sorted(res["skipped"]) == ["finding_actionability", "review_self_consistency"]
assert res["failed"] == []
assert seen == []
def test_ensure_evaluators_records_failures(monkeypatch):
def fake_call(method, path, body=None, timeout=20.0):
return 422, "boom"
monkeypatch.setattr(ej.eb, "_call", fake_call)
monkeypatch.setattr(ej, "existing_evaluators", lambda: {})
res = ej.ensure_evaluators()
assert res["created"] == {}
assert res["failed"][0]["status"] == 422
# --- ensure_rules idempotency --------------------------------------------
def test_ensure_rules_skips_existing(monkeypatch):
calls = []
monkeypatch.setattr(ej.eb, "_call",
lambda *a, **k: calls.append(a) or (200, {}))
monkeypatch.setattr(ej, "existing_evaluators",
lambda: {"finding_actionability": "id-1",
"review_self_consistency": "id-2"})
monkeypatch.setattr(ej, "existing_rule_names",
lambda: {"finding_actionability-on-reviews",
"review_self_consistency-on-reviews"})
res = ej.ensure_rules({"finding_actionability": "id-1",
"review_self_consistency": "id-2"}, 1.0)
assert res["created"] == []
assert sorted(res["skipped"]) == ["finding_actionability", "review_self_consistency"]
assert calls == []
def test_ensure_rules_creates_when_missing(monkeypatch):
calls = []
monkeypatch.setattr(ej.eb, "_call",
lambda *a, **k: calls.append(a) or (201, {}))
monkeypatch.setattr(ej, "existing_rule_names", lambda: set())
res = ej.ensure_rules({"finding_actionability": "id-1"}, 1.0)
assert res["created"] == ["finding_actionability"]
assert calls[0][0] == "POST"
assert calls[0][1] == "/api/public/unstable/evaluation-rules"
# --- judge shape ----------------------------------------------------------
def test_judges_have_required_keys():
for j in ej.JUDGES:
assert j["prompt"]
assert j["outputDefinition"]["dataType"] in ("NUMERIC", "BOOLEAN", "CATEGORICAL")
def test_default_base_url_points_at_the_thinking_patch_proxy():
"""`8802` is the judge-proxy that adds a `signature` to thinking blocks."""
assert "8802" in ej.JUDGE_BASE_URL
+200
View File
@@ -0,0 +1,200 @@
"""Tests for the deterministic review scorers."""
import os
import sys
import pytest
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", "pilot"))
import eval_scores as es # noqa: E402
def f(sev, path="a.py", line=1):
return {"severity": sev, "path": path, "line": line, "problem": "p", "fix": ""}
# --- finding_rate ---------------------------------------------------------
def test_finding_rate_counts_findings():
assert es.finding_rate([f("high"), f("low")]) == 2.0
def test_finding_rate_zero_for_silent_review():
assert es.finding_rate([]) == 0.0
assert es.finding_rate(None) == 0.0
# --- severity_info_ratio --------------------------------------------------
def test_info_ratio_all_advisory():
assert es.severity_info_ratio([f("info"), f("trivial")]) == 1.0
def test_info_ratio_mixed():
assert es.severity_info_ratio([f("info"), f("high")]) == 0.5
def test_info_ratio_none_when_no_findings():
# Undefined, not zero — zero would read as perfectly calibrated.
assert es.severity_info_ratio([]) is None
def test_info_ratio_unknown_severity_treated_as_medium():
# Matches _normalize_finding's fallback, so an odd severity is not
# silently counted as advisory.
assert es.severity_info_ratio([f("bogus")]) == 0.0
# --- severity_max ---------------------------------------------------------
def test_severity_max_picks_highest():
assert es.severity_max([f("info"), f("critical"), f("low")]) == "critical"
def test_severity_max_none_when_silent():
assert es.severity_max([]) == "none"
def test_severity_max_case_insensitive():
assert es.severity_max([f("HIGH")]) == "high"
# --- dropped_findings -----------------------------------------------------
def test_dropped_findings_delta():
assert es.dropped_findings(5, 2) == 3.0
def test_dropped_findings_never_negative():
assert es.dropped_findings(1, 3) == 0.0
def test_dropped_findings_none_when_unknown():
assert es.dropped_findings(None, 2) is None
# --- cost_per_finding -----------------------------------------------------
def test_cost_per_finding_divides():
assert es.cost_per_finding(1.0, [f("high"), f("low")]) == 0.5
def test_cost_per_finding_silent_review_divides_by_one():
# The run still cost money; attributing all of it to "found nothing" is
# the honest reading, and it avoids a division by zero.
assert es.cost_per_finding(0.25, []) == 0.25
def test_cost_per_finding_none_when_unpriced():
assert es.cost_per_finding(None, [f("high")]) is None
def test_cost_per_finding_none_on_garbage():
assert es.cost_per_finding("abc", [f("high")]) is None
# --- build_scores ---------------------------------------------------------
def _by_name(events):
return {e["body"]["name"]: e["body"] for e in events}
def test_build_scores_emits_expected_set():
events = es.build_scores(
trace_id="t1", findings=[f("high"), f("info")], environment="claude",
cost_usd=0.5, dropped_count=2, timestamp="2026-01-01T00:00:00Z",
)
names = _by_name(events)
assert set(names) == {
es.FINDING_RATE, es.SEVERITY_INFO_RATIO, es.SEVERITY_MAX,
es.DROPPED_FINDINGS, es.COST_PER_FINDING,
}
assert names[es.FINDING_RATE]["value"] == 2.0
assert names[es.SEVERITY_MAX]["value"] == "high"
assert names[es.DROPPED_FINDINGS]["value"] == 2.0
assert names[es.COST_PER_FINDING]["value"] == 0.25
def test_build_scores_all_events_are_score_create_on_the_trace():
events = es.build_scores(
trace_id="t9", findings=[f("low")], environment="ollama", cost_usd=1.0,
)
assert all(e["type"] == "score-create" for e in events)
assert all(e["body"]["traceId"] == "t9" for e in events)
assert all(e["body"]["environment"] == "ollama" for e in events)
def test_build_scores_omits_undefined_scores():
# No cost and no drop count measured -> those scores are absent, not zero.
events = es.build_scores(trace_id="t2", findings=[], environment="ollama")
names = set(_by_name(events))
assert es.COST_PER_FINDING not in names
assert es.DROPPED_FINDINGS not in names
assert es.SEVERITY_INFO_RATIO not in names
assert names == {es.FINDING_RATE, es.SEVERITY_MAX}
def test_build_scores_categorical_value_is_string():
events = es.build_scores(trace_id="t3", findings=[f("high")], environment="claude")
sev = _by_name(events)[es.SEVERITY_MAX]
assert sev["dataType"] == "CATEGORICAL"
assert isinstance(sev["value"], str)
def test_build_scores_numeric_values_are_floats():
events = es.build_scores(
trace_id="t4", findings=[f("high")], environment="claude", cost_usd=1,
)
for name, body in _by_name(events).items():
if body["dataType"] == "NUMERIC":
assert isinstance(body["value"], float), name
def test_build_scores_comment_propagates():
events = es.build_scores(
trace_id="t5", findings=[f("high")], environment="claude",
cost_usd=1.0, comment="cost basis: equivalent:claude-sonnet-5",
)
assert all("equivalent" in e["body"]["comment"] for e in events)
# --- score configs --------------------------------------------------------
def test_every_emitted_score_has_a_config():
configured = {c["name"] for c in es.SCORE_CONFIGS}
events = es.build_scores(
trace_id="t6", findings=[f("high")], environment="claude",
cost_usd=1.0, dropped_count=0,
)
assert set(_by_name(events)) <= configured
def test_severity_max_config_covers_every_severity_it_can_emit():
labels = {c["label"] for c in
next(c for c in es.SCORE_CONFIGS if c["name"] == es.SEVERITY_MAX)["categories"]}
assert set(es.SEVERITY_RANK) | {"none"} == labels
# --- ingestion envelope ---------------------------------------------------
def test_every_event_carries_a_timestamp():
# Ingestion rejects events without one, and reports the rejection as a
# per-event 400 inside an HTTP 207 that reads as success.
events = es.build_scores(
trace_id="t7", findings=[f("high")], environment="claude", cost_usd=1.0,
)
assert events
assert all(e.get("timestamp") for e in events)
def test_timestamp_defaults_when_caller_omits_it():
events = es.build_scores(trace_id="t8", findings=[f("low")], environment="claude")
assert all(isinstance(e["timestamp"], str) and e["timestamp"].endswith("Z") for e in events)
def test_explicit_timestamp_is_used():
events = es.build_scores(
trace_id="t9", findings=[f("low")], environment="claude",
timestamp="2026-01-02T03:04:05Z",
)
assert all(e["timestamp"] == "2026-01-02T03:04:05Z" for e in events)
+207
View File
@@ -0,0 +1,207 @@
"""Tests for the feedback.db -> Langfuse score bridge."""
import os
import sys
import pytest
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "..", "pilot"))
import feedback # noqa: E402
import feedback_scores as fs # noqa: E402
@pytest.fixture
def db(tmp_path):
conn = feedback.init(str(tmp_path / "fb.db"))
yield conn
conn.close()
def _seed_finding(conn, repo="o/r", pr=1, comment_id=100, path="a.py", line=1):
cur = conn.execute(
"INSERT INTO review (repo, pr, head_sha, posted_at) VALUES (?,?,?,?)",
(repo, pr, "deadbeef", 1000),
)
review_id = cur.lastrowid
cur = conn.execute(
"""INSERT INTO inline_finding
(review_id, repo, pr, path, line, severity, problem, comment_id, posthash, posted_at)
VALUES (?,?,?,?,?,?,?,?,?,?)""",
(review_id, repo, pr, path, line, "HIGH", "problem", comment_id, f"h{comment_id}", 1000),
)
conn.commit()
return cur.lastrowid
# --- score_pr maths -------------------------------------------------------
def test_engagement_zero_when_nobody_responded():
v = fs.score_pr({"total": 4, "engaged": 0, "positive": 0, "negative": 0})
assert v[fs.REVIEW_ENGAGEMENT] == 0.0
def test_acceptance_absent_when_nobody_engaged():
# Not 0.0 — zero would claim humans judged it neutral.
v = fs.score_pr({"total": 4, "engaged": 0, "positive": 0, "negative": 0})
assert v[fs.REVIEW_ACCEPTANCE] is None
def test_engagement_is_a_share_of_findings():
v = fs.score_pr({"total": 4, "engaged": 1, "positive": 1, "negative": 0})
assert v[fs.REVIEW_ENGAGEMENT] == 0.25
def test_acceptance_all_positive():
v = fs.score_pr({"total": 2, "engaged": 2, "positive": 3, "negative": 0})
assert v[fs.REVIEW_ACCEPTANCE] == 1.0
def test_acceptance_all_negative():
v = fs.score_pr({"total": 2, "engaged": 2, "positive": 0, "negative": 2})
assert v[fs.REVIEW_ACCEPTANCE] == -1.0
def test_acceptance_mixed_is_normalised():
v = fs.score_pr({"total": 4, "engaged": 4, "positive": 3, "negative": 1})
assert v[fs.REVIEW_ACCEPTANCE] == 0.5
def test_engagement_absent_when_no_findings_at_all():
v = fs.score_pr({"total": 0, "engaged": 0, "positive": 0, "negative": 0})
assert v[fs.REVIEW_ENGAGEMENT] is None
# --- collect_pr_feedback over a real sqlite ------------------------------
def test_collect_counts_nothing_on_untouched_findings(db):
_seed_finding(db)
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally == {"total": 1, "engaged": 0, "positive": 0, "negative": 0}
def test_collect_counts_positive_reaction(db):
_seed_finding(db, comment_id=101)
db.execute(
"INSERT INTO reaction (comment_id, user, content, created_at) VALUES (?,?,?,?)",
(101, "alice", "+1", 1),
)
db.commit()
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally["positive"] == 1 and tally["engaged"] == 1
def test_collect_counts_negative_reaction(db):
_seed_finding(db, comment_id=102)
db.execute(
"INSERT INTO reaction (comment_id, user, content, created_at) VALUES (?,?,?,?)",
(102, "bob", "-1", 1),
)
db.commit()
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally["negative"] == 1 and tally["engaged"] == 1
def test_resolved_thread_counts_positive(db):
fid = _seed_finding(db, comment_id=103)
db.execute(
"INSERT INTO thread_state (finding_id, resolved, checked_at) VALUES (?,?,?)",
(fid, 1, 1),
)
db.commit()
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally["positive"] == 1 and tally["engaged"] == 1
def test_unresolved_thread_is_not_a_vote(db):
fid = _seed_finding(db, comment_id=104)
db.execute(
"INSERT INTO thread_state (finding_id, resolved, checked_at) VALUES (?,?,?)",
(fid, 0, 1),
)
db.commit()
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally == {"total": 1, "engaged": 0, "positive": 0, "negative": 0}
def test_negation_reply_counts_negative(db):
fid = _seed_finding(db, comment_id=105)
db.execute(
"INSERT INTO reply (finding_id, author, body, created_at) VALUES (?,?,?,?)",
(fid, "carol", "this is a false positive", 1),
)
db.commit()
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally["negative"] == 1 and tally["engaged"] == 1
def test_neutral_reply_is_engagement_but_not_a_vote(db):
fid = _seed_finding(db, comment_id=106)
db.execute(
"INSERT INTO reply (finding_id, author, body, created_at) VALUES (?,?,?,?)",
(fid, "dave", "done", 1),
)
db.commit()
tally = fs.collect_pr_feedback(db, "o/r", 1)
assert tally["engaged"] == 1
assert tally["positive"] == 0 and tally["negative"] == 0
# --- event shape ----------------------------------------------------------
def test_build_score_events_shape():
events = fs.build_score_events("o/r", 7, {fs.REVIEW_ENGAGEMENT: 0.5}, "claude")
assert len(events) == 1
body = events[0]["body"]
assert events[0]["type"] == "score-create"
assert body["sessionId"] == "o/r#7"
assert body["value"] == 0.5
assert body["environment"] == "claude"
def test_build_score_events_skips_none():
events = fs.build_score_events("o/r", 7, {fs.REVIEW_ACCEPTANCE: None})
assert events == []
def test_score_ids_are_stable_across_runs():
# A backfill re-run must update, not duplicate.
a = fs.build_score_events("o/r", 7, {fs.REVIEW_ENGAGEMENT: 0.5})[0]["body"]["id"]
b = fs.build_score_events("o/r", 7, {fs.REVIEW_ENGAGEMENT: 0.9})[0]["body"]["id"]
assert a == b
def test_score_ids_differ_per_pr_and_name():
e1 = fs.build_score_events("o/r", 7, {fs.REVIEW_ENGAGEMENT: 1})[0]["body"]["id"]
e2 = fs.build_score_events("o/r", 8, {fs.REVIEW_ENGAGEMENT: 1})[0]["body"]["id"]
e3 = fs.build_score_events("o/r", 7, {fs.REVIEW_ACCEPTANCE: 1})[0]["body"]["id"]
assert len({e1, e2, e3}) == 3
def test_backfill_dry_run_reports_without_posting(db, tmp_path):
_seed_finding(db, comment_id=107)
db.commit()
path = db.execute("PRAGMA database_list").fetchone()[2]
summary = fs.backfill(path, dry_run=True)
assert summary["prs_scanned"] == 1
assert summary["prs_with_engagement"] == 0
assert summary["posted"] is False
def test_every_emitted_score_has_a_config():
configured = {c["name"] for c in fs.SCORE_CONFIGS}
assert {fs.REVIEW_ENGAGEMENT, fs.REVIEW_ACCEPTANCE} == configured
def test_every_event_carries_a_timestamp():
# Without one the ingestion endpoint 400s the event inside a 207 that the
# caller reads as success.
events = fs.build_score_events("o/r", 1, {fs.REVIEW_ENGAGEMENT: 0.0})
assert events
assert all(e.get("timestamp") for e in events)
def test_explicit_timestamp_is_used():
events = fs.build_score_events(
"o/r", 1, {fs.REVIEW_ENGAGEMENT: 0.0}, timestamp="2026-01-02T03:04:05Z"
)
assert events[0]["timestamp"] == "2026-01-02T03:04:05Z"
+326
View File
@@ -0,0 +1,326 @@
"""Unit tests for Langfuse trace emission. No network.
`_post` is monkeypatched everywhere a POST would happen; a test that reaches
the real network is a bug in the test, not a slow test.
"""
import json
import os
import sys
HERE = os.path.dirname(os.path.abspath(__file__))
ROOT = os.path.abspath(os.path.join(HERE, "..", ".."))
sys.path.insert(0, os.path.join(ROOT, "pilot"))
import langfuse_trace as lt # noqa: E402
USAGE = {
"input": 2_000_000,
"output": 17_000,
"reasoning": 500,
"cache_read": 400_000,
"cache_write": 50_000,
"total": 2_017_000,
"cost": 0.0,
"steps": 28,
"duration_s": 348.3,
}
BASE = dict(
repo="techspark/pragent",
index="42",
sha="2613b3e1122334455",
title="Harden the review path",
usage=USAGE,
findings=[
{"severity": "critical", "path": "a.py"},
{"severity": "minor", "path": "b.py"},
{"severity": "minor", "path": "c.py"},
],
summary="Three findings.",
)
# ---------------------------------------------------------------------------
# model -> environment split (the whole point of the integration)
# ---------------------------------------------------------------------------
def test_claude_models_land_in_the_claude_environment():
assert lt.resolve_environment("headroom/claude-sonnet-5") == "claude"
assert lt.resolve_environment("claude-opus-5") == "claude"
def test_everything_else_lands_in_the_ollama_environment():
for m in (
"headroom/glm-5.2:cloud",
"headroom/MiniMax-M2.7",
"vllm-qwen38/qwen3.8-27b",
"gpt-5",
):
assert lt.resolve_environment(m) == "ollama", m
def test_provider_and_bare_model_are_split_on_the_first_slash_only():
assert lt.provider_of("vllm-qwen38/qwen3.8-27b") == "vllm-qwen38"
assert lt.strip_provider("headroom/glm-5.2:cloud") == "glm-5.2:cloud"
# A bare name has no provider prefix; default to the pilot's proxy.
assert lt.provider_of("glm-5.2:cloud") == "headroom"
assert lt.strip_provider("glm-5.2:cloud") == "glm-5.2:cloud"
# ---------------------------------------------------------------------------
# usage accounting
# ---------------------------------------------------------------------------
def test_cache_reads_are_subtracted_from_input_not_added():
# Langfuse sums usageDetails keys; opencode reports cache_read *inside*
# input, so reporting both raw would bill the prefix twice.
d = lt._usage_details(USAGE)
assert d["input"] == 2_000_000 - 400_000
assert d["cache_read_input_tokens"] == 400_000
assert d["cache_write_input_tokens"] == 50_000
assert d["output"] == 17_000
assert d["reasoning"] == 500
def test_zero_cache_fields_are_omitted_rather_than_sent_as_zero():
d = lt._usage_details({"input": 100, "output": 10})
assert d == {"input": 100, "output": 10}
def test_a_paid_model_is_priced_as_itself():
costs, basis = lt._cost_details(USAGE, "headroom/claude-sonnet-5")
assert costs["total"] > 0
assert basis == "actual"
def test_minimax_is_priced_against_the_comparison_target_not_zero():
# MiniMax-M2.7 is the model the webhook actually runs and it is absent from
# PRICES; charting it at $0 would make the whole dashboard a flat line.
costs, basis = lt._cost_details(USAGE, "headroom/MiniMax-M2.7")
assert costs["total"] > 0
assert basis == "equivalent:claude-sonnet-5"
def test_glm_is_priced_against_the_comparison_target():
costs, basis = lt._cost_details(USAGE, "headroom/glm-5.2:cloud")
assert costs["total"] > 0
assert basis.startswith("equivalent:")
def test_an_all_zero_price_entry_counts_as_free_not_as_priced():
# The self-hosted vLLM qwen IS in PRICES, at 0.00 across the board.
costs, basis = lt._cost_details(USAGE, "vllm-qwen38/qwen3.8-27b")
assert costs["total"] > 0
assert basis.startswith("equivalent:")
def test_explicit_price_target_wins_over_the_default():
costs, basis = lt._cost_details(USAGE, "headroom/MiniMax-M2.7", "claude-opus-5")
assert basis == "equivalent:claude-opus-5"
sonnet, _ = lt._cost_details(USAGE, "headroom/MiniMax-M2.7", "claude-sonnet-5")
assert costs["total"] > sonnet["total"]
def test_env_overrides_the_default_target(monkeypatch):
monkeypatch.setenv("PRAGENT_PRICE_TARGET", "claude-haiku-4-5")
assert lt.resolve_price_target() == "claude-haiku-4-5"
# An explicit argument still beats the env.
assert lt.resolve_price_target("gpt-5") == "gpt-5"
def test_unknown_comparison_target_yields_no_cost_block_rather_than_a_wrong_one():
costs, basis = lt._cost_details(USAGE, "headroom/MiniMax-M2.7", "not-a-real-model")
assert costs == {}
assert basis == ""
# ---------------------------------------------------------------------------
# batch shape
# ---------------------------------------------------------------------------
def test_batch_has_a_trace_and_a_generation_linked_by_trace_id():
batch = lt.build_batch(model="headroom/claude-sonnet-5", **BASE)
types = [e["type"] for e in batch]
# Scores ride in the same batch; the trace and generation lead it.
assert types[:2] == ["trace-create", "generation-create"]
trace, gen = batch[0], batch[1]
assert gen["body"]["traceId"] == trace["body"]["id"]
assert trace["body"]["environment"] == gen["body"]["environment"] == "claude"
def test_batch_without_usage_has_no_generation():
batch = lt.build_batch(model="headroom/glm-5.2:cloud", **{**BASE, "usage": None})
types = [e["type"] for e in batch]
assert "generation-create" not in types
assert types[0] == "trace-create"
def test_trace_carries_repo_pr_session_and_severity_counts():
batch = lt.build_batch(model="headroom/glm-5.2:cloud", **BASE)
body = batch[0]["body"]
assert body["sessionId"] == "techspark/pragent#42"
assert body["metadata"]["severities"] == {"critical": 1, "minor": 2}
assert body["metadata"]["findings"] == 3
assert "provider:headroom" in body["tags"]
assert "model:glm-5.2:cloud" in body["tags"]
def test_lens_names_become_tags():
batch = lt.build_batch(
model="headroom/glm-5.2:cloud", lenses=["security", "tests"], **BASE
)
assert "lens:security" in batch[0]["body"]["tags"]
assert "lens:tests" in batch[0]["body"]["tags"]
def test_cost_basis_is_tagged_so_equivalent_is_never_read_as_spend():
batch = lt.build_batch(model="headroom/MiniMax-M2.7", **BASE)
trace = batch[0]["body"]
assert "cost:equivalent:claude-sonnet-5" in trace["tags"]
assert trace["metadata"]["cost_basis"] == "equivalent:claude-sonnet-5"
paid = lt.build_batch(model="headroom/claude-sonnet-5", **BASE)
assert "cost:actual" in paid[0]["body"]["tags"]
def test_minimax_generation_carries_a_nonzero_cost():
batch = lt.build_batch(model="headroom/MiniMax-M2.7", **BASE)
assert batch[1]["body"]["costDetails"]["total"] > 0
def test_batch_is_json_serializable():
batch = lt.build_batch(model="headroom/claude-sonnet-5", **BASE)
json.dumps({"batch": batch})
# ---------------------------------------------------------------------------
# emit_review_trace — config gate and fail-open
# ---------------------------------------------------------------------------
def _configure(monkeypatch):
monkeypatch.setenv("LANGFUSE_HOST", "http://langfuse.test:3000/")
monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-lf-test")
monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-lf-test")
def test_no_config_means_no_post_and_no_error(monkeypatch):
for k in ("LANGFUSE_HOST", "LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY"):
monkeypatch.delenv(k, raising=False)
calls = []
monkeypatch.setattr(lt, "_post", lambda *a, **k: calls.append(a) or 200)
assert lt.emit_review_trace(model="headroom/glm-5.2:cloud", **BASE) is False
assert calls == []
def test_configured_emit_posts_to_the_ingestion_endpoint(monkeypatch):
_configure(monkeypatch)
seen = {}
def fake_post(host, pk, sk, batch, timeout):
seen.update(host=host, pk=pk, sk=sk, batch=batch, timeout=timeout)
return 207
monkeypatch.setattr(lt, "_post", fake_post)
assert lt.emit_review_trace(model="headroom/claude-sonnet-5", **BASE) is True
# Trailing slash stripped so the path is not doubled.
assert seen["host"] == "http://langfuse.test:3000"
kinds = [e["type"] for e in seen["batch"]]
assert kinds[:2] == ["trace-create", "generation-create"]
assert "score-create" in kinds
def test_transport_failure_is_swallowed(monkeypatch):
_configure(monkeypatch)
def boom(*a, **k):
raise OSError("connection refused")
monkeypatch.setattr(lt, "_post", boom)
assert lt.emit_review_trace(model="headroom/glm-5.2:cloud", **BASE) is False
def test_non_success_status_reports_failure_without_raising(monkeypatch):
_configure(monkeypatch)
monkeypatch.setattr(lt, "_post", lambda *a, **k: 401)
assert lt.emit_review_trace(model="headroom/glm-5.2:cloud", **BASE) is False
# ---------------------------------------------------------------------------
# Scores folded into the review batch (added with eval_scores)
# ---------------------------------------------------------------------------
def _scores(events):
return {e["body"]["name"]: e["body"] for e in events if e["type"] == "score-create"}
def test_build_batch_appends_scores():
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t",
model="headroom/claude-sonnet-5",
usage={"input": 100, "output": 10},
findings=[{"severity": "high", "path": "a.py", "line": 1}],
)
names = set(_scores(events))
assert "finding_rate" in names
assert "severity_max" in names
def test_scores_attach_to_the_same_trace():
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t", model="m",
usage={"input": 1, "output": 1}, findings=[], trace_id="fixed-id",
)
for body in _scores(events).values():
assert body["traceId"] == "fixed-id"
def test_scores_inherit_the_trace_environment():
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t",
model="headroom/glm-5.2:cloud",
usage={"input": 1, "output": 1}, findings=[],
)
for body in _scores(events).values():
assert body["environment"] == "ollama"
def test_dropped_findings_scored_when_provided():
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t", model="m",
usage={"input": 1, "output": 1}, findings=[], dropped_count=3,
)
assert _scores(events)["dropped_findings"]["value"] == 3.0
def test_dropped_findings_absent_when_not_measured():
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t", model="m",
usage={"input": 1, "output": 1}, findings=[],
)
assert "dropped_findings" not in _scores(events)
def test_cost_score_carries_its_basis_in_the_comment():
# An equivalent-cost $/finding must never be read as money spent.
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t",
model="headroom/glm-5.2:cloud",
usage={"input": 1000, "output": 100}, findings=[{"severity": "low", "path": "a", "line": 1}],
)
cpf = _scores(events).get("cost_per_finding")
if cpf is not None: # only when cost_model could price the comparison target
assert "equivalent" in cpf["comment"]
def test_batch_without_usage_still_scores_findings():
# A run with no usage report still produced findings worth scoring.
events = lt.build_batch(
repo="o/r", index="1", sha="abc", title="t", model="m",
usage=None, findings=[{"severity": "critical", "path": "a", "line": 2}],
)
assert _scores(events)["severity_max"]["value"] == "critical"
+90
View File
@@ -0,0 +1,90 @@
"""The parse-time drop counter feeding the `dropped_findings` score.
A model that emits findings at unusable locations produces an empty findings
list, exactly like a model that found nothing. These tests pin the signal that
tells the two apart.
"""
import json
import os
import sys
HERE = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, os.path.abspath(os.path.join(HERE, "..", "..", "pilot")))
import ai_review # noqa: E402
def _payload(findings):
return "```json\n" + json.dumps({"summary": "s", "findings": findings}) + "\n```"
GOOD = {"severity": "high", "path": "a.py", "line": 3, "problem": "p", "fix": "f"}
NO_PATH = {"severity": "high", "line": 3, "problem": "p"}
NO_LINE = {"severity": "high", "path": "a.py", "problem": "p"}
BAD_LINE = {"severity": "high", "path": "a.py", "line": 0, "problem": "p"}
def test_no_drops_on_clean_output():
_, findings, *_ = ai_review.parse_review_output(_payload([GOOD, GOOD]))
assert len(findings) == 2
assert ai_review.last_parse_dropped() == 0
def test_counts_findings_missing_path():
_, findings, *_ = ai_review.parse_review_output(_payload([GOOD, NO_PATH]))
assert len(findings) == 1
assert ai_review.last_parse_dropped() == 1
def test_counts_findings_missing_line():
_, findings, *_ = ai_review.parse_review_output(_payload([NO_LINE, NO_LINE]))
assert findings == []
assert ai_review.last_parse_dropped() == 2
def test_counts_findings_with_unusable_line():
_, findings, *_ = ai_review.parse_review_output(_payload([BAD_LINE]))
assert findings == []
assert ai_review.last_parse_dropped() == 1
def test_all_dropped_is_distinguishable_from_found_nothing():
ai_review.parse_review_output(_payload([NO_PATH, NO_PATH, NO_PATH]))
all_dropped = ai_review.last_parse_dropped()
ai_review.parse_review_output(_payload([]))
found_nothing = ai_review.last_parse_dropped()
assert all_dropped == 3 and found_nothing == 0
def test_counter_resets_on_unparseable_output():
# Otherwise a salvage-path review inherits the previous review's count.
ai_review.parse_review_output(_payload([NO_PATH, NO_PATH]))
assert ai_review.last_parse_dropped() == 2
ai_review.parse_review_output("no json here at all")
assert ai_review.last_parse_dropped() == 0
def test_counter_resets_on_malformed_json():
ai_review.parse_review_output(_payload([NO_PATH]))
ai_review.parse_review_output("```json\n{not valid json,,,}\n```")
assert ai_review.last_parse_dropped() == 0
def test_parse_findings_tracks_drops_too():
# The non-opencode path must be scored on the same basis.
findings = ai_review.parse_findings(json.dumps({"findings": [GOOD, NO_PATH]}))
assert len(findings) == 1
assert ai_review.last_parse_dropped() == 1
def test_parse_findings_resets_on_garbage():
ai_review.parse_findings(json.dumps({"findings": [NO_PATH]}))
assert ai_review.last_parse_dropped() == 1
ai_review.parse_findings("not json")
assert ai_review.last_parse_dropped() == 0
def test_bare_array_output_is_counted():
_, findings, *_ = ai_review.parse_review_output("```json\n" + json.dumps([GOOD, NO_PATH]) + "\n```")
assert len(findings) == 1
assert ai_review.last_parse_dropped() == 1