diff --git a/.agents/skills/icm-cctv-ops/SKILL.md b/.agents/skills/icm-cctv-ops/SKILL.md new file mode 100644 index 000000000..3c6e7b957 --- /dev/null +++ b/.agents/skills/icm-cctv-ops/SKILL.md @@ -0,0 +1,16 @@ +--- +name: icm-cctv-ops +description: ICM-CCTV visualization projection for ops dashboards. Load on Stepie 2087 step 10023. +--- + +# Skill: icm-cctv-ops + +Emit lane counts + PR rows as JSON. No secrets. Do not reorder primary Stepie goal 2149. + +```text +python3 -m ml.pipelines.cli cctv +``` + +Contract: `docs/ops/ICM-CCTV.md` +Session 2026-09-24 14:02 PDT. Live master `03ffb33b`. +Agent-Identity: Grok (Administrator) diff --git a/.agents/skills/ml-pipeline-ops/SKILL.md b/.agents/skills/ml-pipeline-ops/SKILL.md new file mode 100644 index 000000000..2808036d7 --- /dev/null +++ b/.agents/skills/ml-pipeline-ops/SKILL.md @@ -0,0 +1,37 @@ +--- +name: ml-pipeline-ops +description: Keep-alive ML pipeline DAG for termux-monorepo. Load on Issue #175 ML lanes. Never wholesale-merge mega ML PRs. +--- + +# Skill: ml-pipeline-ops + +Session 2026-09-24 14:02 PDT. Live master `03ffb33b`. Dual-gate GREEN on tip (repo-gate 36058704829, termux-smoke 36058704728). Keep-alive child rebases here; #787 / #746 / #682 SUPERSEDE-or-EXTRACT. + +## Hard rules + +1. Extract-only for mega ML PRs. Slim tree is `ml/pipelines/`. +2. Dual gates before promote: hygiene+portability + agentic termux smoke. Bind the named jobs, not combined status. +3. No secrets, no Class 3/4 artifacts, no GPU runners. +4. GitLab / Vercel hobby rate-limit are **non-gate** (#772). +5. Dirty + files>40 + minesweeper overlap → HOLD or EXTRACT. +6. `master_staging_base` (#48, #788) and stacked feature bases never retarget to master. +7. Do **not** restamp `docs/ops/LANE-MATRIX.md` (policy SSOT). Live board is generated. +8. Do **not** pulse-comment Issue #175. Edit the issue body only when intent changes. + +## Commands + +```text +python3 -m ml.pipelines.cli status +python3 -m ml.pipelines.cli lanes +python3 -m ml.pipelines.cli run +python3 -m ml.pipelines.cli cctv +python3 -m ml.pipelines.cli explain 707 +python3 -m ml.pipelines.cli gate 48 +python3 -m unittest discover -s ml/pipelines -p 'test_*.py' +``` + +## Cycle + +RECON → INGEST → FEATURES → TRAIN → EVALUATE → (WAIT|HOLD|EXTRACT|SUPERSEDE|PROMOTE) → MONITOR + +Operator ACTIVE. Agent-Identity: Grok (Administrator) diff --git a/.agents/skills/ml-pipeline-ops/references/dual-gate.md b/.agents/skills/ml-pipeline-ops/references/dual-gate.md new file mode 100644 index 000000000..666e6cfeb --- /dev/null +++ b/.agents/skills/ml-pipeline-ops/references/dual-gate.md @@ -0,0 +1,8 @@ +# Dual-gate (promote authority) + +1. `hygiene + portability gate` / `repo gate` SUCCESS on **this** SHA +2. `agentic termux smoke` / `termux smoke` SUCCESS on **this** SHA +3. Vercel rate-limits are **non-gate** (#772) +4. Copilot / CodeRabbit / Qodo / Devin = advisory only +5. Dual-gate SUCCESS on an older head does not authorize a newer SHA +6. Combined commit status is not dual-gate; bind the named jobs diff --git a/.agents/skills/ml-pipeline-ops/references/minesweeper.md b/.agents/skills/ml-pipeline-ops/references/minesweeper.md new file mode 100644 index 000000000..84b6dd3e9 --- /dev/null +++ b/.agents/skills/ml-pipeline-ops/references/minesweeper.md @@ -0,0 +1,7 @@ +# Minesweeper + +Concurrent agents (Jules, Sentinel, Bolt, Devin, Copilot) open overlapping PRs. + +- Do not overwrite a WAIT peer's branch. +- Bot author + files>40 → EXTRACT, never wholesale merge. +- Classify with `ml.pipelines.lanes.minesweeper_rules`. diff --git a/.github/workflows/ml-pipelines-keepalive.yml b/.github/workflows/ml-pipelines-keepalive.yml new file mode 100644 index 000000000..2f8798379 --- /dev/null +++ b/.github/workflows/ml-pipelines-keepalive.yml @@ -0,0 +1,32 @@ +name: ml pipelines keep-alive + +on: + pull_request: + paths: + - "ml/pipelines/**" + - ".github/workflows/ml-pipelines-keepalive.yml" + push: + branches: [master] + paths: + - "ml/pipelines/**" + workflow_dispatch: + +permissions: + contents: read + +concurrency: + group: ml-keepalive-${{ github.ref }} + cancel-in-progress: true + +jobs: + unittest: + name: keep-alive unit tests + runs-on: ubuntu-latest + timeout-minutes: 8 + steps: + - uses: actions/checkout@v4 + with: + fetch-depth: 1 + submodules: false + - name: stdlib unittest + run: python3 -m unittest discover -s ml/pipelines -p "test_*.py" diff --git a/SKILLS.md b/SKILLS.md index d30fdec0e..0b398fdff 100644 --- a/SKILLS.md +++ b/SKILLS.md @@ -8,6 +8,8 @@ | **Primary agent entry (BIUDL)** | [`CLAUDE.md`](CLAUDE.md) | | Adaptive wait / feedback | `.agents/skills/adaptive-feedback-cycle/SKILL.md` | | Admin ops | `.agents/skills/evidence-led-monorepo-ops/SKILL.md` | +| ML keep-alive | `.agents/skills/ml-pipeline-ops/SKILL.md` | +| ICM-CCTV projection | `.agents/skills/icm-cctv-ops/SKILL.md` | | External contribute (help-wanted) | `.agents/skills/help-wanted-lane/SKILL.md` | | Production WAIT → VALIDATE | `.github/skills/production-reconciliation/SKILL.md` | @@ -29,7 +31,7 @@ Every skill directory must contain a `SKILL.md`. Inventory lists **all** of them | Role | Load first | |------|------------| | Collaborator | `adaptive-feedback-cycle` → dual-gate | -| Admin / Grok | `evidence-led-monorepo-ops` + `adaptive-wait` | +| Admin / Grok | `evidence-led-monorepo-ops` + `adaptive-wait` + `ml-pipeline-ops` | | Oversight / external PR | `help-wanted-lane` | | Evaluation / DOE | `multivariate-doe` + `blind-agent-evaluation` | diff --git a/docs/icm/cards/icm-cctv-10023.md b/docs/icm/cards/icm-cctv-10023.md new file mode 100644 index 000000000..2e47dd7ae --- /dev/null +++ b/docs/icm/cards/icm-cctv-10023.md @@ -0,0 +1,5 @@ +# Card: icm-cctv-10023 + +Stepie 2087 step 10023 visualization surface. + +`ml/pipelines/viz` emits JSON for dashboards. Does not steal focus from primary goal 2149. diff --git a/docs/icm/cards/issue-175-operator.md b/docs/icm/cards/issue-175-operator.md new file mode 100644 index 000000000..8d649d430 --- /dev/null +++ b/docs/icm/cards/issue-175-operator.md @@ -0,0 +1,6 @@ +# Card: issue-175-operator + +Priority hub. Dual-gate SSOT. Vercel #772 is non-gate. + +Do not YOLO merge. Do not pulse-comment. Edit the issue body when intent changes. +Live board: `docs/ops/generated/lane-matrix-status.md`. diff --git a/docs/icm/cards/mlp-keep-001.md b/docs/icm/cards/mlp-keep-001.md new file mode 100644 index 000000000..e799d0ba4 --- /dev/null +++ b/docs/icm/cards/mlp-keep-001.md @@ -0,0 +1,9 @@ +# Card: MLP-KEEP-001 + +Keep-alive ML DAG under `ml/pipelines/`. + +- Extract-only vs #682 / #432 / #549 / #601 / #746 / #787 +- Dual-gate before promote +- MoneyBall scorer is transparent weights, not a black box +- Lane short-circuit: EXTRACT / HOLD / WAIT / OBSERVE / SUPERSEDE before PROMOTE +- Session 2026-09-24 14:02 PDT tip `03ffb33b` diff --git a/docs/icm/cards/session-ssot-pulse.md b/docs/icm/cards/session-ssot-pulse.md new file mode 100644 index 000000000..6f4a3ae4a --- /dev/null +++ b/docs/icm/cards/session-ssot-pulse.md @@ -0,0 +1,6 @@ +# Card: session-ssot-pulse + +`docs/ops/LANE-MATRIX.md` is **policy**, not a living pulse. + +Generated board: `docs/ops/generated/lane-matrix-status.md`. +Session recon stays in chat. Product PRs carry code. diff --git a/docs/ops/ICM-CCTV.md b/docs/ops/ICM-CCTV.md new file mode 100644 index 000000000..6faef05e6 --- /dev/null +++ b/docs/ops/ICM-CCTV.md @@ -0,0 +1,14 @@ +# ICM-CCTV ops surface + +Stepie goal 2087 step 10023: wire icm-cctv && visualization surfaces into Ops Dashboards. + +This is the **projection contract**, not a runtime GitHub Pages app. + +- Producer: `python3 -m ml.pipelines.cli cctv` +- Schema: `docs/schemas/icm-cctv.json` +- Code: `ml/pipelines/viz/cctv.py` +- Consumed by operator dashboards (Command Center + future GitHub Pages) + +No secrets. Counts + lane labels + master SHA only. +Primary Stepie goal 2149 is not reordered by this surface. +Live master at extract: `03ffb33b`. diff --git a/docs/ops/ML-PIPELINES.md b/docs/ops/ML-PIPELINES.md new file mode 100644 index 000000000..f46a495dd --- /dev/null +++ b/docs/ops/ML-PIPELINES.md @@ -0,0 +1,27 @@ +# ML Pipelines keep-alive + +Issue #175 forbids wholesale merge of #432 / #549 / #601 / #682. +This document maps the slim extract in `ml/pipelines/` re-based onto live master `03ffb33b`. + +- DAG: `ml/pipelines/cli.py` + `ml/pipelines/stages/` +- Ranking: `ml/pipelines/moneyball/scorer.py` +- Lane short-circuit: `ml/pipelines/lanes/` +- Gate: `ml/pipelines/contracts/gate.py` +- Replay adapters: `ml/pipelines/replay/` +- ICM-CCTV: `ml/pipelines/viz/cctv.py` (Stepie 2087 step 10023) +- Latest fixture: `ml/pipelines/fixtures/session_20260924.json` + +## Commands + +```text +python3 -m ml.pipelines.cli status +python3 -m ml.pipelines.cli lanes +python3 -m ml.pipelines.cli run +python3 -m ml.pipelines.cli cctv +python3 -m ml.pipelines.cli explain 48 +python3 -m ml.pipelines.cli gate 707 +python3 -m unittest discover -s ml/pipelines -p 'test_*.py' +``` + +Promote path remains **dual-gate only**. Vercel / GitLab / Copilot / CodeRabbit / Qodo / Devin are non-gate. +Do not restamp `docs/ops/LANE-MATRIX.md`. diff --git a/docs/ops/SKILLS-INVENTORY.md b/docs/ops/SKILLS-INVENTORY.md index dc1236d52..45193baf4 100644 --- a/docs/ops/SKILLS-INVENTORY.md +++ b/docs/ops/SKILLS-INVENTORY.md @@ -1,6 +1,6 @@ # Skills Inventory (termux-monorepo) -**Version:** 2026-09-19 17:26 PDT · **Master tip:** `826dc1e4` +**Version:** 2026-09-24 14:02 PDT · **Master tip:** `03ffb33b` **Primary agent entry:** [`CLAUDE.md`](../../CLAUDE.md) **Collaborator short entry:** [`SKILLS.md`](../../SKILLS.md) @@ -40,6 +40,8 @@ Local agent mirrors (e.g. `.grok/skills/`) are convenience only — **not** a se | gemini-performance-psychology | `.agents/skills/gemini-performance-psychology/SKILL.md` | | **github-pages-operator** | `.agents/skills/github-pages-operator/SKILL.md` | | help-wanted-lane | `.agents/skills/help-wanted-lane/SKILL.md` | +| **icm-cctv-ops** | `.agents/skills/icm-cctv-ops/SKILL.md` | +| **ml-pipeline-ops** | `.agents/skills/ml-pipeline-ops/SKILL.md` | | multivariate-doe | `.agents/skills/multivariate-doe/SKILL.md` | | review-loop | `.agents/skills/review-loop/SKILL.md` | | termux-monorepo | `.agents/skills/termux-monorepo/SKILL.md` | diff --git a/docs/proposals/active/actions-refinements/ITEMS.md b/docs/proposals/active/actions-refinements/ITEMS.md index 7ea52613a..292a34d44 100644 --- a/docs/proposals/active/actions-refinements/ITEMS.md +++ b/docs/proposals/active/actions-refinements/ITEMS.md @@ -22,6 +22,9 @@ | AR-18 | Establish an observe-mode Capability–Scope–Specialist decision spine that normalizes declared model/provider capabilities, applies hard provenance/availability/quota/scope gates before scoring, emits a secret-free decision envelope, and records repository-local performance evidence ahead of public leaderboard priors. | P0 | Manus AI | executing | Approved Evolutionary Capability Spine plan 2026-08-22. Extends RL-05, RL-11–RL-14, and RL-17 without replacing event classification, peer evidence, command dispatch, or existing branch-write confirmation. Initial delivery is read-only/observe mode with an immediate feature-gate rollback; B3, B4/AR-04, B5/A-14, and PR #276 remain held. | | AR-19 | Add a read-only Context Relationship Evidence Matrix decision-support record that reuses the canonical graph and lead/lag audit surfaces to classify authority, provenance, freshness, risk, coverage gaps, and outcome bindings. | P1 | Manus AI | executing | Approved 2026-08-25 as the Issue #192 half of a tightly coupled CRG-14 + AR-19 PR. No provider invocation, routing promotion, branch write, command composition, raw-body collection, or generated-index hand edit is authorized. The live audit parsing warning is a collector-correctness remediation; B3, B4/AR-04, B5/A-14, and PR #276 remain held. | | AR-20 | Repair the Continuous Team Evaluation live-invocation admission gate so dynamic candidate eligibility is emitted as an explicit step output, stale/unavailable candidates are recorded as skips, and eligible lanes reach the provider invocation while preserving the existing immutable MVT ledger and no-source-mutation boundary. | P1 | ChatGPT | executing | Live run `34069136232` at `4460cf4ff8db8860dfb1b477b64de3a39b7faf1c` discovered 22 eligible candidates but every experiment invocation step was skipped; the candidate script printed `eligible=...` to stdout without writing `$GITHUB_OUTPUT`. Fix is limited to workflow admission/observability; provider secrets remain Actions secrets and output remains uncapped at the existing 131072 request value. | +| AR-21 | Provider-neutral evaluation lane registry + GitHub App capability probe (Checks-publication failure mode). Names-only for Issue #184. | P1 | Grok | executing | PR #806. Rebase onto live master before promote. | +| AR-22 | GAMUT remote evaluation + all-repo wiki knowledge fabric. Observational. No secret values. No private Wiki bodies. | P1 | Grok | executing | PR #809. Rebase onto live master before promote. | +| AR-23 | ML keep-alive DAG re-extract onto live master (`ml/pipelines/`). Extract-only vs wholesale #432/#549/#601/#682/#746/#787. ICM-CCTV projection. Do not restamp LANE-MATRIX policy. | P0 | Grok | executing | Implements MLP-KEEP-001. Dual-gate on this SHA. Vercel non-gate (#772). | ## Batch Rules diff --git a/docs/proposals/active/ml-keep-alive/ITEMS.md b/docs/proposals/active/ml-keep-alive/ITEMS.md new file mode 100644 index 000000000..5a8bce58e --- /dev/null +++ b/docs/proposals/active/ml-keep-alive/ITEMS.md @@ -0,0 +1,8 @@ +# ITEMS — ML keep-alive extract (Issue #175) + +| ID | Work | Priority | Owner | Status | Evidence / Boundary | +|---|---|---|---|---|---| +| MLP-KEEP-001 | Slim `ml/pipelines/` DAG + MoneyBall + lane short-circuit onto live master. Never wholesale-merge #432/#549/#601/#682. | P0 | Grok | executing | Dual-gate on this SHA. Vercel non-gate (#772). | +| MLP-KEEP-002 | ICM-CCTV projection (`cli cctv`) for ops dashboards. No secrets. | P1 | Grok | executing | Schema `docs/schemas/icm-cctv.json`. | +| MLP-KEEP-003 | Minesweeper + supersede classifiers so Jules/Sentinel/Bolt WAIT peers are not overwritten. | P1 | Grok | executing | `lanes/minesweeper_rules.py`, `lanes/supersede_rules.py`. | +| MLP-KEEP-004 | Latest-session fixture resolver so keep-alive rebases do not restamp LANE-MATRIX policy. | P1 | Grok | executing | `lib/latest.py` + `fixtures/session_20260924.json`. | diff --git a/docs/proposals/active/ml-keep-alive/MANIFEST.md b/docs/proposals/active/ml-keep-alive/MANIFEST.md new file mode 100644 index 000000000..f49fb4594 --- /dev/null +++ b/docs/proposals/active/ml-keep-alive/MANIFEST.md @@ -0,0 +1,7 @@ +# Manifest + +Package: `ml/pipelines/` +Issue: #175 +Item: MLP-KEEP-001 +Master at cut: `03ffb33b54b86b52cf8a11cc2f7db6d2d8061772` +Agent-Identity: Grok (Administrator) diff --git a/docs/proposals/active/ml-keep-alive/source.md b/docs/proposals/active/ml-keep-alive/source.md new file mode 100644 index 000000000..0986b725e --- /dev/null +++ b/docs/proposals/active/ml-keep-alive/source.md @@ -0,0 +1,6 @@ +# ML keep-alive proposal + +Implements: MLP-KEEP-001 + +Extract-only vs wholesale ML PRs. Stdlib Python. No GPU. No secrets. +Promote only when dual-gate is green on **this** SHA. diff --git a/docs/proposals/registry.yaml b/docs/proposals/registry.yaml index 7b15d1333..947c94f50 100644 --- a/docs/proposals/registry.yaml +++ b/docs/proposals/registry.yaml @@ -1,9 +1,27 @@ # ArchW1z proposal registry — agents read this first version: 1 -updated_at: "2026-09-23T18:15:00Z" +updated_at: "2026-09-24T21:10:00Z" updated_by: grok-administrator proposals: + - id: ml-keep-alive + title: "ML keep-alive DAG extract onto live master (Issue #175)" + author: grok + status: executing + priority: P0 + path: active/ml-keep-alive/ + reviewers: + - id: grok + role: author+executor + status: executing + - id: timerloggedout-spec + role: operator-authorizer + status: requested + related_issues: [175, 772] + related_prs: [787, 746, 682, 432] + related_branches: [feat/ml-pipelines-keepalive-20260924] + gates_required: [repo-gate, termux-smoke] + - id: approxination-integration title: "Approxination skill search/create/contribute + A/B/C/D evaluation cohort" author: grok diff --git a/docs/schemas/icm-cctv.json b/docs/schemas/icm-cctv.json new file mode 100644 index 000000000..b08e6380d --- /dev/null +++ b/docs/schemas/icm-cctv.json @@ -0,0 +1,25 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "title": "ICM-CCTV projection", + "type": "object", + "required": [ + "master_sha", + "issue", + "lanes" + ], + "properties": { + "master_sha": { + "type": "string", + "minLength": 7 + }, + "issue": { + "const": 175 + }, + "lanes": { + "type": "object" + }, + "operator": { + "const": "ACTIVE" + } + } +} diff --git a/ml/__init__.py b/ml/__init__.py new file mode 100644 index 000000000..5cece41ac --- /dev/null +++ b/ml/__init__.py @@ -0,0 +1 @@ +"""ML keep-alive package. Extract-only vs wholesale #682/#432/#549/#601.""" diff --git a/ml/pipelines/README.md b/ml/pipelines/README.md new file mode 100644 index 000000000..36c62792f --- /dev/null +++ b/ml/pipelines/README.md @@ -0,0 +1,16 @@ +# ml/pipelines keep-alive + +Slim extract of the ML keep-alive DAG. Issue #175 forbids wholesale merge of #432 / #549 / #601 / #682. + +Re-extracted onto live master `03ffb33b` after #787 waited on `ba5f6b6d`. + +```text +python3 -m ml.pipelines.cli status +python3 -m ml.pipelines.cli lanes +python3 -m ml.pipelines.cli run +python3 -m ml.pipelines.cli cctv +python3 -m unittest discover -s ml/pipelines -p 'test_*.py' +``` + +Operator ACTIVE. Dual-gate remains promote authority. Vercel non-gate. +Do not restamp LANE-MATRIX policy. diff --git a/ml/pipelines/VERSION b/ml/pipelines/VERSION new file mode 100644 index 000000000..1d0ba9ea1 --- /dev/null +++ b/ml/pipelines/VERSION @@ -0,0 +1 @@ +0.4.0 diff --git a/ml/pipelines/__init__.py b/ml/pipelines/__init__.py new file mode 100644 index 000000000..a782dfb35 --- /dev/null +++ b/ml/pipelines/__init__.py @@ -0,0 +1,7 @@ +"""ml.pipelines: Keep-alive ML DAG for operator ranking. + +Implements: MLP-KEEP-001 +""" +from __future__ import annotations + +__version__ = "0.4.0" diff --git a/ml/pipelines/__main__.py b/ml/pipelines/__main__.py new file mode 100644 index 000000000..2dcaa29ee --- /dev/null +++ b/ml/pipelines/__main__.py @@ -0,0 +1,5 @@ +from ml.pipelines.cli import main +import sys + +if __name__ == '__main__': + sys.exit(main()) diff --git a/ml/pipelines/cli.py b/ml/pipelines/cli.py new file mode 100644 index 000000000..38d081d7a --- /dev/null +++ b/ml/pipelines/cli.py @@ -0,0 +1,110 @@ +"""cli: python3 -m ml.pipelines.cli """ +from __future__ import annotations + +import argparse +import json +import sys +from typing import Any + +from ml.pipelines.contracts.gate import assert_promotable +from ml.pipelines.lanes.classify import classify_pr +from ml.pipelines.lib.engine import run_dag, summarize +from ml.pipelines.lib.errors import GateBlocked +from ml.pipelines.lib.io import load_json +from ml.pipelines.lib.latest import latest_session_path +from ml.pipelines.moneyball.explain import explain +from ml.pipelines.moneyball.scorer import score +from ml.pipelines.stages import STAGES +from ml.pipelines.viz.projection import project + + +def _fixture() -> dict[str, Any]: + return load_json(latest_session_path()) + + +def cmd_status(_: argparse.Namespace) -> int: + payload = _fixture() + print(json.dumps({ + "master_sha": payload["master_sha"], + "issue": 175, + "open_prs": len(payload["prs"]), + "operator": "ACTIVE", + "session": payload.get("session"), + "fixture": latest_session_path().name, + }, indent=2)) + return 0 + + +def cmd_run(_: argparse.Namespace) -> int: + context: dict[str, Any] = {"snapshot": _fixture()} + results = run_dag(STAGES, context) + print(json.dumps({"summary": summarize(results), "lanes": context.get("lanes", [])}, indent=2)) + return 0 + + +def cmd_lanes(_: argparse.Namespace) -> int: + payload = _fixture() + rows = [] + for pr in payload["prs"]: + points = score(pr) + lane = classify_pr(pr) + rows.append({"number": pr["number"], "lane": lane.value, "score": points, "why": pr.get("why")}) + print(json.dumps(rows, indent=2)) + return 0 + + +def cmd_cctv(_: argparse.Namespace) -> int: + payload = _fixture() + lanes = [ + {"number": pr["number"], "lane": classify_pr(pr).value, "score": score(pr), "why": pr.get("why")} + for pr in payload["prs"] + ] + print(json.dumps(project(payload, lanes), indent=2)) + return 0 + + +def cmd_explain(ns: argparse.Namespace) -> int: + payload = _fixture() + wanted = int(ns.number) + for pr in payload["prs"]: + if int(pr["number"]) == wanted: + print(json.dumps(explain(pr), indent=2)) + return 0 + return 1 + + +def cmd_gate(ns: argparse.Namespace) -> int: + payload = _fixture() + wanted = int(ns.number) + for pr in payload["prs"]: + if int(pr["number"]) == wanted: + lane = classify_pr(pr) + try: + assert_promotable(pr, lane) + print(json.dumps({"number": wanted, "ok": True, "lane": lane.value})) + return 0 + except GateBlocked as exc: + print(json.dumps({"number": wanted, "ok": False, "lane": lane.value, "reason": str(exc)})) + return 2 + return 1 + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(prog="ml.pipelines.cli") + sub = parser.add_subparsers(dest="cmd", required=True) + sub.add_parser("status").set_defaults(func=cmd_status) + sub.add_parser("run").set_defaults(func=cmd_run) + sub.add_parser("lanes").set_defaults(func=cmd_lanes) + sub.add_parser("cctv").set_defaults(func=cmd_cctv) + p_ex = sub.add_parser("explain") + p_ex.add_argument("number", type=int) + p_ex.set_defaults(func=cmd_explain) + p_g = sub.add_parser("gate") + p_g.add_argument("number", type=int) + p_g.set_defaults(func=cmd_gate) + args = parser.parse_args(argv) + return int(args.func(args)) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/ml/pipelines/contracts/__init__.py b/ml/pipelines/contracts/__init__.py new file mode 100644 index 000000000..68526a6e0 --- /dev/null +++ b/ml/pipelines/contracts/__init__.py @@ -0,0 +1,2 @@ +from .gate import assert_promotable +from .non_gate import is_non_gate, NON_GATE diff --git a/ml/pipelines/contracts/dual_gate.py b/ml/pipelines/contracts/dual_gate.py new file mode 100644 index 000000000..bfd12cfe0 --- /dev/null +++ b/ml/pipelines/contracts/dual_gate.py @@ -0,0 +1,11 @@ +"""Named dual-gate checks. Both must be SUCCESS.""" +from __future__ import annotations + +from typing import Mapping + +HYGIENE = "hygiene + portability gate" +SMOKE = "agentic termux smoke" + + +def dual_gate_green(checks: Mapping[str, str]) -> bool: + return checks.get(HYGIENE) == "success" and checks.get(SMOKE) == "success" diff --git a/ml/pipelines/contracts/gate.py b/ml/pipelines/contracts/gate.py new file mode 100644 index 000000000..a147af637 --- /dev/null +++ b/ml/pipelines/contracts/gate.py @@ -0,0 +1,24 @@ +"""gate: Promote only when dual gates are green and lane is promote.""" +from __future__ import annotations + +from typing import Any, Mapping + +from ml.pipelines.lib.errors import GateBlocked +from ml.pipelines.lib.types import Lane + + +def assert_promotable(pr: Mapping[str, Any], lane: Lane) -> None: + if lane is not Lane.PROMOTE: + raise GateBlocked(f"lane {lane.value} is not promote") + if pr.get("dual_gate") != "green": + raise GateBlocked("dual_gate is not green") + if pr.get("mergeable_state") != "clean": + raise GateBlocked("mergeable_state is not clean") + if pr.get("gitlab_required"): + raise GateBlocked("gitlab must remain non-blocking") + if pr.get("hitl_risk"): + raise GateBlocked("HITL YOLO blocked") + if pr.get("master_staging_base"): + raise GateBlocked("master-staging is not a master promote path") + if pr.get("stacked_feature_base"): + raise GateBlocked("stacked feature base is not a master promote path") diff --git a/ml/pipelines/contracts/non_gate.py b/ml/pipelines/contracts/non_gate.py new file mode 100644 index 000000000..41a676bfa --- /dev/null +++ b/ml/pipelines/contracts/non_gate.py @@ -0,0 +1,19 @@ +"""Advisory / non-gate signals. Never block promote on these alone.""" +from __future__ import annotations + +NON_GATE = frozenset( + { + "vercel_rate_limit", + "gitlab", + "mintlify", + "copilot", + "coderabbit", + "qodo", + "devin", + "age_days", + } +) + + +def is_non_gate(signal: str) -> bool: + return signal in NON_GATE diff --git a/ml/pipelines/contracts/promote.py b/ml/pipelines/contracts/promote.py new file mode 100644 index 000000000..c4871f7de --- /dev/null +++ b/ml/pipelines/contracts/promote.py @@ -0,0 +1,22 @@ +"""Promote packet: dual-gate + clean mergeable + extract-sized.""" +from __future__ import annotations + +from typing import Any, Mapping + +from ml.pipelines.contracts.gate import assert_promotable +from ml.pipelines.lib.errors import GateBlocked +from ml.pipelines.lib.types import Lane +from ml.pipelines.lib.weights import FILES_OVER_40 + + +def promote_packet(pr: Mapping[str, Any], lane: Lane) -> dict[str, Any]: + assert_promotable(pr, lane) + files = int(pr.get("changed_files") or 0) + if files > FILES_OVER_40: + raise GateBlocked("extract required: changed_files over 40") + return { + "number": pr.get("number"), + "lane": lane.value, + "decision": "promote", + "files": files, + } diff --git a/ml/pipelines/contracts/test_dual_gate.py b/ml/pipelines/contracts/test_dual_gate.py new file mode 100644 index 000000000..6c3a79716 --- /dev/null +++ b/ml/pipelines/contracts/test_dual_gate.py @@ -0,0 +1,11 @@ +import unittest + +from ml.pipelines.contracts.dual_gate import HYGIENE, SMOKE, dual_gate_green + + +class DualGateTests(unittest.TestCase): + def test_both(self) -> None: + self.assertTrue(dual_gate_green({HYGIENE: "success", SMOKE: "success"})) + + def test_one(self) -> None: + self.assertFalse(dual_gate_green({HYGIENE: "success", SMOKE: "failure"})) diff --git a/ml/pipelines/contracts/test_gate.py b/ml/pipelines/contracts/test_gate.py new file mode 100644 index 000000000..4d855a767 --- /dev/null +++ b/ml/pipelines/contracts/test_gate.py @@ -0,0 +1,28 @@ +import unittest + +from ml.pipelines.contracts.gate import assert_promotable +from ml.pipelines.lib.errors import GateBlocked +from ml.pipelines.lib.types import Lane + + +class GateTests(unittest.TestCase): + def test_blocks_extract(self) -> None: + with self.assertRaises(GateBlocked): + assert_promotable({"dual_gate": "green", "mergeable_state": "clean"}, Lane.EXTRACT) + + def test_allows_promote(self) -> None: + assert_promotable({"dual_gate": "green", "mergeable_state": "clean"}, Lane.PROMOTE) + + def test_blocks_staging(self) -> None: + with self.assertRaises(GateBlocked): + assert_promotable( + {"dual_gate": "green", "mergeable_state": "clean", "master_staging_base": True}, + Lane.PROMOTE, + ) + + def test_blocks_hitl(self) -> None: + with self.assertRaises(GateBlocked): + assert_promotable( + {"dual_gate": "green", "mergeable_state": "clean", "hitl_risk": True}, + Lane.PROMOTE, + ) diff --git a/ml/pipelines/contracts/test_non_gate.py b/ml/pipelines/contracts/test_non_gate.py new file mode 100644 index 000000000..c65af6ebd --- /dev/null +++ b/ml/pipelines/contracts/test_non_gate.py @@ -0,0 +1,11 @@ +import unittest + +from ml.pipelines.contracts.non_gate import is_non_gate + + +class NonGateTests(unittest.TestCase): + def test_vercel(self) -> None: + self.assertTrue(is_non_gate("vercel_rate_limit")) + + def test_dual_gate_is_gate(self) -> None: + self.assertFalse(is_non_gate("hygiene + portability gate")) diff --git a/ml/pipelines/contracts/test_promote.py b/ml/pipelines/contracts/test_promote.py new file mode 100644 index 000000000..f215456f9 --- /dev/null +++ b/ml/pipelines/contracts/test_promote.py @@ -0,0 +1,21 @@ +import unittest + +from ml.pipelines.contracts.promote import promote_packet +from ml.pipelines.lib.errors import GateBlocked +from ml.pipelines.lib.types import Lane + + +class PromoteTests(unittest.TestCase): + def test_packet(self) -> None: + packet = promote_packet( + {"number": 707, "changed_files": 3, "dual_gate": "green", "mergeable_state": "clean"}, + Lane.PROMOTE, + ) + self.assertEqual(packet["decision"], "promote") + + def test_oversize(self) -> None: + with self.assertRaises(GateBlocked): + promote_packet( + {"number": 48, "changed_files": 73, "dual_gate": "green", "mergeable_state": "clean"}, + Lane.PROMOTE, + ) diff --git a/ml/pipelines/features/__init__.py b/ml/pipelines/features/__init__.py new file mode 100644 index 000000000..4c7902ec9 --- /dev/null +++ b/ml/pipelines/features/__init__.py @@ -0,0 +1,4 @@ +from .staleness import flag as stale +from .security import flag as security +from .docs_only import flag as docs_only +from .bot_authors import flag as bot_author diff --git a/ml/pipelines/features/bot_authors.py b/ml/pipelines/features/bot_authors.py new file mode 100644 index 000000000..1cb5d630a --- /dev/null +++ b/ml/pipelines/features/bot_authors.py @@ -0,0 +1,6 @@ +"""Jules/bot authorship is a mild penalty, not a HOLD by itself.""" +from __future__ import annotations +from typing import Any, Mapping + +def flag(pr: Mapping[str, Any]) -> bool: + return "bot" in str(pr.get("author") or "").lower() diff --git a/ml/pipelines/features/docs_only.py b/ml/pipelines/features/docs_only.py new file mode 100644 index 000000000..879e1f2e8 --- /dev/null +++ b/ml/pipelines/features/docs_only.py @@ -0,0 +1,6 @@ +"""docs_only pulses stay WAIT/SUPERSEDE; they do not replace keep-alive extracts.""" +from __future__ import annotations +from typing import Any, Mapping + +def flag(pr: Mapping[str, Any]) -> bool: + return bool(pr.get("docs_only")) diff --git a/ml/pipelines/features/dual_gate.py b/ml/pipelines/features/dual_gate.py new file mode 100644 index 000000000..8a029a128 --- /dev/null +++ b/ml/pipelines/features/dual_gate.py @@ -0,0 +1,10 @@ +"""Named dual-gate feature. Combined commit status is not dual-gate.""" +from __future__ import annotations +from typing import Any, Mapping + +def dual_gate_green(pr: Mapping[str, Any]) -> bool: + return pr.get("dual_gate") == "green" + +def vercel_is_nongate(pr: Mapping[str, Any]) -> bool: + """Issue #772: Vercel rate-limit / mergeable_state=unstable is non-gate.""" + return True diff --git a/ml/pipelines/features/nongate.py b/ml/pipelines/features/nongate.py new file mode 100644 index 000000000..7a06ce40b --- /dev/null +++ b/ml/pipelines/features/nongate.py @@ -0,0 +1,16 @@ +"""Surfaces that must never block promote: Vercel, GitLab, Copilot, CodeRabbit, Qodo, Devin.""" +from __future__ import annotations + +NON_GATE_CONTEXTS = ( + "Vercel", + "GitLab", + "Copilot", + "CodeRabbit", + "Qodo", + "Devin Review", + "Mintlify", +) + +def is_nongate_context(name: str) -> bool: + lowered = name.lower() + return any(token.lower() in lowered for token in NON_GATE_CONTEXTS) diff --git a/ml/pipelines/features/security.py b/ml/pipelines/features/security.py new file mode 100644 index 000000000..f96a87974 --- /dev/null +++ b/ml/pipelines/features/security.py @@ -0,0 +1,6 @@ +"""security_fix adds score; still requires dual-gate + clean mergeable.""" +from __future__ import annotations +from typing import Any, Mapping + +def flag(pr: Mapping[str, Any]) -> bool: + return bool(pr.get("security")) diff --git a/ml/pipelines/features/staleness.py b/ml/pipelines/features/staleness.py new file mode 100644 index 000000000..d8a82c5f9 --- /dev/null +++ b/ml/pipelines/features/staleness.py @@ -0,0 +1,6 @@ +"""stale_base is a WAIT signal when dual-gate is green, never a promote by itself.""" +from __future__ import annotations +from typing import Any, Mapping + +def flag(pr: Mapping[str, Any]) -> bool: + return bool(pr.get("stale_base")) diff --git a/ml/pipelines/features/test_bot_authors.py b/ml/pipelines/features/test_bot_authors.py new file mode 100644 index 000000000..4f89d81ee --- /dev/null +++ b/ml/pipelines/features/test_bot_authors.py @@ -0,0 +1,16 @@ +import unittest +from ml.pipelines.features.bot_authors import flag + +class FeatureFlagTests(unittest.TestCase): + def test_true(self) -> None: + key = { + "staleness": "stale_base", + "security": "security", + "docs_only": "docs_only", + "bot_authors": "author", + }["bot_authors"] + sample = {key: True} if key != "author" else {"author": "google-labs-jules[bot]"} + self.assertTrue(flag(sample)) + + def test_false(self) -> None: + self.assertFalse(flag({})) diff --git a/ml/pipelines/features/test_docs_only.py b/ml/pipelines/features/test_docs_only.py new file mode 100644 index 000000000..d69b2c0da --- /dev/null +++ b/ml/pipelines/features/test_docs_only.py @@ -0,0 +1,16 @@ +import unittest +from ml.pipelines.features.docs_only import flag + +class FeatureFlagTests(unittest.TestCase): + def test_true(self) -> None: + key = { + "staleness": "stale_base", + "security": "security", + "docs_only": "docs_only", + "bot_authors": "author", + }["docs_only"] + sample = {key: True} if key != "author" else {"author": "google-labs-jules[bot]"} + self.assertTrue(flag(sample)) + + def test_false(self) -> None: + self.assertFalse(flag({})) diff --git a/ml/pipelines/features/test_dual_gate.py b/ml/pipelines/features/test_dual_gate.py new file mode 100644 index 000000000..f0b7a660d --- /dev/null +++ b/ml/pipelines/features/test_dual_gate.py @@ -0,0 +1,12 @@ +import unittest +from ml.pipelines.features.dual_gate import dual_gate_green, vercel_is_nongate + +class DualGateFeature(unittest.TestCase): + def test_green(self): + self.assertTrue(dual_gate_green({"dual_gate": "green"})) + + def test_pending(self): + self.assertFalse(dual_gate_green({"dual_gate": "pending"})) + + def test_vercel_nongate(self): + self.assertTrue(vercel_is_nongate({"mergeable_state": "unstable"})) diff --git a/ml/pipelines/features/test_nongate.py b/ml/pipelines/features/test_nongate.py new file mode 100644 index 000000000..eb2cb365d --- /dev/null +++ b/ml/pipelines/features/test_nongate.py @@ -0,0 +1,12 @@ +import unittest +from ml.pipelines.features.nongate import is_nongate_context + +class Nongate(unittest.TestCase): + def test_vercel(self): + self.assertTrue(is_nongate_context("Vercel – termux-monorepo")) + + def test_coderabbit(self): + self.assertTrue(is_nongate_context("CodeRabbit")) + + def test_repo_gate_is_gate(self): + self.assertFalse(is_nongate_context("hygiene + portability gate")) diff --git a/ml/pipelines/features/test_security.py b/ml/pipelines/features/test_security.py new file mode 100644 index 000000000..d9b455455 --- /dev/null +++ b/ml/pipelines/features/test_security.py @@ -0,0 +1,16 @@ +import unittest +from ml.pipelines.features.security import flag + +class FeatureFlagTests(unittest.TestCase): + def test_true(self) -> None: + key = { + "staleness": "stale_base", + "security": "security", + "docs_only": "docs_only", + "bot_authors": "author", + }["security"] + sample = {key: True} if key != "author" else {"author": "google-labs-jules[bot]"} + self.assertTrue(flag(sample)) + + def test_false(self) -> None: + self.assertFalse(flag({})) diff --git a/ml/pipelines/features/test_staleness.py b/ml/pipelines/features/test_staleness.py new file mode 100644 index 000000000..14f2d764e --- /dev/null +++ b/ml/pipelines/features/test_staleness.py @@ -0,0 +1,16 @@ +import unittest +from ml.pipelines.features.staleness import flag + +class FeatureFlagTests(unittest.TestCase): + def test_true(self) -> None: + key = { + "staleness": "stale_base", + "security": "security", + "docs_only": "docs_only", + "bot_authors": "author", + }["staleness"] + sample = {key: True} if key != "author" else {"author": "google-labs-jules[bot]"} + self.assertTrue(flag(sample)) + + def test_false(self) -> None: + self.assertFalse(flag({})) diff --git a/ml/pipelines/fixtures/lane_cases.json b/ml/pipelines/fixtures/lane_cases.json new file mode 100644 index 000000000..2546c8419 --- /dev/null +++ b/ml/pipelines/fixtures/lane_cases.json @@ -0,0 +1,27 @@ +{ + "promote": [ + 707 + ], + "wait": [ + 724, + 745, + 717, + 714 + ], + "hold": [ + 48, + 69, + 739 + ], + "extract": [ + 682, + 630 + ], + "observe": [ + 740, + 735 + ], + "supersede": [ + 744 + ] +} diff --git a/ml/pipelines/fixtures/prs/432.json b/ml/pipelines/fixtures/prs/432.json new file mode 100644 index 000000000..c1fe0f06d --- /dev/null +++ b/ml/pipelines/fixtures/prs/432.json @@ -0,0 +1,8 @@ +{ + "number": 432, + "title": "observe-mode GitHub ML pipelines", + "changed_files": 140, + "mergeable_state": "dirty", + "ml_wholesale": true, + "why": "wholesale no-go" +} diff --git a/ml/pipelines/fixtures/prs/48.json b/ml/pipelines/fixtures/prs/48.json new file mode 100644 index 000000000..c68fd9747 --- /dev/null +++ b/ml/pipelines/fixtures/prs/48.json @@ -0,0 +1,8 @@ +{ + "number": 48, + "title": "llm-api-hub", + "changed_files": 73, + "mergeable_state": "dirty", + "master_staging_base": true, + "why": "do not retarget to master" +} diff --git a/ml/pipelines/fixtures/prs/543.json b/ml/pipelines/fixtures/prs/543.json new file mode 100644 index 000000000..17ebf79b5 --- /dev/null +++ b/ml/pipelines/fixtures/prs/543.json @@ -0,0 +1,9 @@ +{ + "number": 543, + "title": "skill quality lane", + "changed_files": 21, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "extract quality checks" +} diff --git a/ml/pipelines/fixtures/prs/587.json b/ml/pipelines/fixtures/prs/587.json new file mode 100644 index 000000000..a4645a5af --- /dev/null +++ b/ml/pipelines/fixtures/prs/587.json @@ -0,0 +1,8 @@ +{ + "number": 587, + "title": "timing quotas", + "changed_files": 9, + "mergeable_state": "unstable", + "author": "google-labs-jules[bot]", + "why": "jules peer" +} diff --git a/ml/pipelines/fixtures/prs/617.json b/ml/pipelines/fixtures/prs/617.json new file mode 100644 index 000000000..dca09f22c --- /dev/null +++ b/ml/pipelines/fixtures/prs/617.json @@ -0,0 +1,9 @@ +{ + "number": 617, + "title": "proposal registry gate", + "changed_files": 5, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase then dual-gate" +} diff --git a/ml/pipelines/fixtures/prs/630.json b/ml/pipelines/fixtures/prs/630.json new file mode 100644 index 000000000..e91242086 --- /dev/null +++ b/ml/pipelines/fixtures/prs/630.json @@ -0,0 +1,9 @@ +{ + "number": 630, + "title": "Jules dashboard rich UI", + "changed_files": 54, + "mergeable_state": "dirty", + "minesweeper": true, + "author": "google-labs-jules[bot]", + "why": "minesweeper" +} diff --git a/ml/pipelines/fixtures/prs/67.json b/ml/pipelines/fixtures/prs/67.json new file mode 100644 index 000000000..7d53554d2 --- /dev/null +++ b/ml/pipelines/fixtures/prs/67.json @@ -0,0 +1,8 @@ +{ + "number": 67, + "title": "PR scope discipline", + "changed_files": 2, + "docs_only": true, + "stale_base": true, + "why": "ancient docs" +} diff --git a/ml/pipelines/fixtures/prs/682.json b/ml/pipelines/fixtures/prs/682.json new file mode 100644 index 000000000..bab28c290 --- /dev/null +++ b/ml/pipelines/fixtures/prs/682.json @@ -0,0 +1,9 @@ +{ + "number": 682, + "title": "ML wholesale DAG", + "changed_files": 130, + "mergeable_state": "dirty", + "ml_wholesale": true, + "tests": true, + "why": "size != quality" +} diff --git a/ml/pipelines/fixtures/prs/69.json b/ml/pipelines/fixtures/prs/69.json new file mode 100644 index 000000000..47725efa7 --- /dev/null +++ b/ml/pipelines/fixtures/prs/69.json @@ -0,0 +1,12 @@ +{ + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "pr": { + "number": 69, + "title": "debate dock stacked", + "changed_files": 22, + "mergeable_state": "dirty", + "stacked_feature_base": true, + "why": "slice already in #784" + } +} diff --git a/ml/pipelines/fixtures/prs/707.json b/ml/pipelines/fixtures/prs/707.json new file mode 100644 index 000000000..3dd6c7395 --- /dev/null +++ b/ml/pipelines/fixtures/prs/707.json @@ -0,0 +1,9 @@ +{ + "number": 707, + "title": "historical promote specimen", + "changed_files": 3, + "mergeable_state": "clean", + "dual_gate": "green", + "tests": true, + "why": "historical promote specimen" +} diff --git a/ml/pipelines/fixtures/prs/725.json b/ml/pipelines/fixtures/prs/725.json new file mode 100644 index 000000000..c269a3c60 --- /dev/null +++ b/ml/pipelines/fixtures/prs/725.json @@ -0,0 +1,8 @@ +{ + "number": 725, + "title": "evidence JSONL SeekLog", + "changed_files": 11, + "mergeable_state": "unstable", + "tests": true, + "why": "independent extract" +} diff --git a/ml/pipelines/fixtures/prs/746.json b/ml/pipelines/fixtures/prs/746.json new file mode 100644 index 000000000..8f08305b9 --- /dev/null +++ b/ml/pipelines/fixtures/prs/746.json @@ -0,0 +1,10 @@ +{ + "number": 746, + "title": "ML slim keep-alive on stale tip", + "changed_files": 90, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "supersede": true, + "why": "superseded by this extract" +} diff --git a/ml/pipelines/fixtures/prs/753.json b/ml/pipelines/fixtures/prs/753.json new file mode 100644 index 000000000..a581f0698 --- /dev/null +++ b/ml/pipelines/fixtures/prs/753.json @@ -0,0 +1,12 @@ +{ + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "pr": { + "number": 753, + "title": "vercel sitemap", + "changed_files": 3, + "mergeable_state": "unstable", + "docs_only": true, + "why": "vercel non-gate" + } +} diff --git a/ml/pipelines/fixtures/prs/762.json b/ml/pipelines/fixtures/prs/762.json new file mode 100644 index 000000000..872065e88 --- /dev/null +++ b/ml/pipelines/fixtures/prs/762.json @@ -0,0 +1,7 @@ +{ + "number": 762, + "title": "plugin connector parity", + "changed_files": 12, + "draft": true, + "why": "keep draft" +} diff --git a/ml/pipelines/fixtures/prs/764.json b/ml/pipelines/fixtures/prs/764.json new file mode 100644 index 000000000..de967e1b1 --- /dev/null +++ b/ml/pipelines/fixtures/prs/764.json @@ -0,0 +1,8 @@ +{ + "number": 764, + "title": "Termux hub Tailscale", + "changed_files": 18, + "draft": true, + "hitl_risk": true, + "why": "human device edge" +} diff --git a/ml/pipelines/fixtures/prs/768.json b/ml/pipelines/fixtures/prs/768.json new file mode 100644 index 000000000..9b90ed545 --- /dev/null +++ b/ml/pipelines/fixtures/prs/768.json @@ -0,0 +1,13 @@ +{ + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "pr": { + "number": 768, + "title": "Sentinel harden", + "changed_files": 4, + "mergeable_state": "unstable", + "security": true, + "author": "google-labs-jules[bot]", + "why": "observe security slice" + } +} diff --git a/ml/pipelines/fixtures/prs/781.json b/ml/pipelines/fixtures/prs/781.json new file mode 100644 index 000000000..df25c2116 --- /dev/null +++ b/ml/pipelines/fixtures/prs/781.json @@ -0,0 +1,14 @@ +{ + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "pr": { + "number": 781, + "title": "session pulse 10:37", + "changed_files": 6, + "mergeable_state": "unstable", + "docs_only": true, + "supersede": true, + "stale_base": true, + "why": "older pulse" + } +} diff --git a/ml/pipelines/fixtures/prs/783.json b/ml/pipelines/fixtures/prs/783.json new file mode 100644 index 000000000..3f44fb7ea --- /dev/null +++ b/ml/pipelines/fixtures/prs/783.json @@ -0,0 +1,14 @@ +{ + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "pr": { + "number": 783, + "title": "session pulse 11:24", + "changed_files": 8, + "mergeable_state": "unstable", + "docs_only": true, + "supersede": true, + "stale_base": true, + "why": "older pulse" + } +} diff --git a/ml/pipelines/fixtures/prs/784.json b/ml/pipelines/fixtures/prs/784.json new file mode 100644 index 000000000..771327290 --- /dev/null +++ b/ml/pipelines/fixtures/prs/784.json @@ -0,0 +1,14 @@ +{ + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "pr": { + "number": 784, + "title": "debate dock extract of #69", + "changed_files": 16, + "mergeable_state": "unstable", + "dual_gate": "pending", + "stale_base": true, + "why": "corrective extract, base behind tip", + "tests": true + } +} diff --git a/ml/pipelines/fixtures/prs/785.json b/ml/pipelines/fixtures/prs/785.json new file mode 100644 index 000000000..8af3170b0 --- /dev/null +++ b/ml/pipelines/fixtures/prs/785.json @@ -0,0 +1,9 @@ +{ + "number": 785, + "title": "Wingman wait-loop skills", + "changed_files": 10, + "mergeable_state": "unstable", + "dual_gate": "pending", + "tests": true, + "why": "wait-loop extract, do not collide" +} diff --git a/ml/pipelines/fixtures/prs/787.json b/ml/pipelines/fixtures/prs/787.json new file mode 100644 index 000000000..8905c3d7f --- /dev/null +++ b/ml/pipelines/fixtures/prs/787.json @@ -0,0 +1,11 @@ +{ + "number": 787, + "title": "ML keep-alive DAG ICM-CCTV", + "changed_files": 103, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "keep_alive_parent": true, + "supersede": true, + "why": "re-extract onto 03ffb33b" +} diff --git a/ml/pipelines/fixtures/prs/788.json b/ml/pipelines/fixtures/prs/788.json new file mode 100644 index 000000000..bc7960ea2 --- /dev/null +++ b/ml/pipelines/fixtures/prs/788.json @@ -0,0 +1,9 @@ +{ + "number": 788, + "title": "Linear client/sync coverage TER-15", + "changed_files": 22, + "mergeable_state": "unknown", + "master_staging_base": true, + "tests": true, + "why": "wrong-base master-staging" +} diff --git a/ml/pipelines/fixtures/prs/803.json b/ml/pipelines/fixtures/prs/803.json new file mode 100644 index 000000000..919a933a7 --- /dev/null +++ b/ml/pipelines/fixtures/prs/803.json @@ -0,0 +1,8 @@ +{ + "number": 803, + "title": "Bolt log parse / registry regex", + "changed_files": 6, + "mergeable_state": "unstable", + "author": "google-labs-jules[bot]", + "why": "observe perf slice" +} diff --git a/ml/pipelines/fixtures/prs/804.json b/ml/pipelines/fixtures/prs/804.json new file mode 100644 index 000000000..3a2499756 --- /dev/null +++ b/ml/pipelines/fixtures/prs/804.json @@ -0,0 +1,9 @@ +{ + "number": 804, + "title": "Sentinel context collector path traversal", + "changed_files": 4, + "mergeable_state": "unstable", + "security": true, + "author": "google-labs-jules[bot]", + "why": "observe security slice" +} diff --git a/ml/pipelines/fixtures/prs/806.json b/ml/pipelines/fixtures/prs/806.json new file mode 100644 index 000000000..c01aea359 --- /dev/null +++ b/ml/pipelines/fixtures/prs/806.json @@ -0,0 +1,9 @@ +{ + "number": 806, + "title": "evaluation lanes + GitHub App capability evidence", + "changed_files": 5, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase onto live master" +} diff --git a/ml/pipelines/fixtures/prs/809.json b/ml/pipelines/fixtures/prs/809.json new file mode 100644 index 000000000..74330f5e4 --- /dev/null +++ b/ml/pipelines/fixtures/prs/809.json @@ -0,0 +1,9 @@ +{ + "number": 809, + "title": "GAMUT remote + wiki knowledge fabric", + "changed_files": 6, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase onto 03ffb33b then dual-gate" +} diff --git a/ml/pipelines/fixtures/replay_events.json b/ml/pipelines/fixtures/replay_events.json new file mode 100644 index 000000000..3a60481ae --- /dev/null +++ b/ml/pipelines/fixtures/replay_events.json @@ -0,0 +1,23 @@ +[ + { + "experiment_id": "742", + "sha": "cf8d12473e47b19655ea6eb1693225bcb9750926", + "gain": 0.4, + "promoted": true, + "kind": "action_effect" + }, + { + "experiment_id": "741", + "sha": "9a5beae60169df0157cea480bae5c51e8ce9eef3", + "gain": 0.3, + "promoted": true, + "kind": "corpus" + }, + { + "experiment_id": "695", + "sha": "187e2f9894282e931812bcab76f768105b96e49c", + "gain": 0.5, + "promoted": true, + "kind": "paper2agent" + } +] diff --git a/ml/pipelines/fixtures/session_20260922.json b/ml/pipelines/fixtures/session_20260922.json new file mode 100644 index 000000000..e60776839 --- /dev/null +++ b/ml/pipelines/fixtures/session_20260922.json @@ -0,0 +1,116 @@ +{ + "captured_at": "2026-09-22T19:18:00Z", + "master_sha": "d10a7a5472261a524bbdccf28db3702f82ea52e4", + "issue": 175, + "open_issues": 109, + "open_prs": 91, + "stepie": { + "primary": 2149, + "adjacent": 2087, + "step": 10023, + "progress": "1/12" + }, + "prs": [ + { + "number": 682, + "changed_files": 130, + "mergeable_state": "dirty", + "ml_wholesale": true, + "why": "wholesale ML" + }, + { + "number": 630, + "changed_files": 89, + "mergeable_state": "dirty", + "minesweeper": true, + "author": "google-labs-jules[bot]", + "why": "minesweeper dashboard" + }, + { + "number": 48, + "changed_files": 73, + "mergeable_state": "dirty", + "master_staging_base": true, + "why": "llm-api-hub on master-staging" + }, + { + "number": 69, + "changed_files": 24, + "mergeable_state": "dirty", + "stacked_feature_base": true, + "why": "debate dock stacked feature base" + }, + { + "number": 724, + "changed_files": 18, + "mergeable_state": "unstable", + "dual_gate": "green", + "tests": true, + "stale_base": true, + "why": "ML slim extract wait" + }, + { + "number": 745, + "changed_files": 4, + "mergeable_state": "unstable", + "docs_only": true, + "dual_gate": "green", + "why": "12:10 pulse" + }, + { + "number": 744, + "changed_files": 4, + "mergeable_state": "unstable", + "docs_only": true, + "supersede": true, + "why": "stale pulse" + }, + { + "number": 739, + "changed_files": 22, + "mergeable_state": "unstable", + "hitl_risk": true, + "why": "HITL control plane" + }, + { + "number": 740, + "changed_files": 16, + "mergeable_state": "unstable", + "draft": true, + "why": "System One draft" + }, + { + "number": 717, + "changed_files": 3, + "mergeable_state": "unstable", + "tests": true, + "dual_gate": "green", + "stale_base": true, + "why": "Gantt extract" + }, + { + "number": 714, + "changed_files": 6, + "mergeable_state": "unstable", + "docs_only": true, + "dual_gate": "green", + "stale_base": true, + "why": "Stepie skill" + }, + { + "number": 735, + "changed_files": 5, + "mergeable_state": "unstable", + "author": "google-labs-jules[bot]", + "why": "Sentinel symlink" + }, + { + "number": 707, + "changed_files": 3, + "mergeable_state": "clean", + "dual_gate": "green", + "tests": true, + "why": "historical promote specimen" + } + ] +} diff --git a/ml/pipelines/fixtures/session_20260923.json b/ml/pipelines/fixtures/session_20260923.json new file mode 100644 index 000000000..593951d55 --- /dev/null +++ b/ml/pipelines/fixtures/session_20260923.json @@ -0,0 +1,174 @@ +{ + "master_sha": "ba5f6b6de4589d27fb4f438ac19f94074038fc51", + "issue": 175, + "captured_at": "2026-09-23T20:14:00Z", + "operator": "ACTIVE", + "session": "2026-09-23-1314-PDT", + "prs": [ + { + "number": 785, + "title": "Wingman wait-loop skills", + "changed_files": 10, + "mergeable_state": "unstable", + "dual_gate": "pending", + "why": "just opened wait-loop extract", + "tests": true + }, + { + "number": 784, + "title": "debate dock extract of #69", + "changed_files": 16, + "mergeable_state": "unstable", + "dual_gate": "pending", + "stale_base": true, + "why": "corrective extract, base behind tip", + "tests": true + }, + { + "number": 783, + "title": "session pulse 11:24", + "changed_files": 8, + "mergeable_state": "unstable", + "docs_only": true, + "supersede": true, + "stale_base": true, + "why": "older pulse" + }, + { + "number": 781, + "title": "session pulse 10:37", + "changed_files": 6, + "mergeable_state": "unstable", + "docs_only": true, + "supersede": true, + "stale_base": true, + "why": "older pulse" + }, + { + "number": 746, + "title": "ML slim keep-alive on stale tip", + "changed_files": 90, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "re-extract onto live tip" + }, + { + "number": 48, + "title": "llm-api-hub", + "changed_files": 73, + "mergeable_state": "dirty", + "master_staging_base": true, + "why": "do not retarget to master" + }, + { + "number": 69, + "title": "debate dock stacked", + "changed_files": 22, + "mergeable_state": "dirty", + "stacked_feature_base": true, + "why": "slice already in #784" + }, + { + "number": 630, + "title": "Jules dashboard rich UI", + "changed_files": 54, + "mergeable_state": "dirty", + "minesweeper": true, + "author": "google-labs-jules[bot]", + "why": "minesweeper" + }, + { + "number": 768, + "title": "Sentinel harden", + "changed_files": 4, + "mergeable_state": "unstable", + "security": true, + "author": "google-labs-jules[bot]", + "why": "observe security slice" + }, + { + "number": 762, + "title": "plugin connector parity", + "changed_files": 12, + "draft": true, + "why": "keep draft" + }, + { + "number": 764, + "title": "Termux hub Tailscale", + "changed_files": 18, + "draft": true, + "hitl_risk": true, + "why": "human device edge" + }, + { + "number": 617, + "title": "proposal registry gate", + "changed_files": 5, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase then dual-gate" + }, + { + "number": 587, + "title": "timing quotas", + "changed_files": 9, + "mergeable_state": "unstable", + "author": "google-labs-jules[bot]", + "why": "jules peer" + }, + { + "number": 543, + "title": "skill quality lane", + "changed_files": 21, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "extract quality checks" + }, + { + "number": 725, + "title": "evidence JSONL SeekLog", + "changed_files": 11, + "mergeable_state": "unstable", + "tests": true, + "why": "independent extract" + }, + { + "number": 753, + "title": "vercel sitemap", + "changed_files": 3, + "mergeable_state": "unstable", + "docs_only": true, + "why": "vercel non-gate" + }, + { + "number": 67, + "title": "PR scope discipline", + "changed_files": 2, + "docs_only": true, + "stale_base": true, + "why": "ancient docs" + }, + { + "number": 682, + "title": "ML wholesale DAG", + "changed_files": 130, + "mergeable_state": "dirty", + "ml_wholesale": true, + "tests": true, + "why": "size != quality" + }, + { + "number": 707, + "title": "historical promote specimen", + "changed_files": 3, + "mergeable_state": "clean", + "dual_gate": "green", + "tests": true, + "why": "historical promote specimen" + } + ] +} diff --git a/ml/pipelines/fixtures/session_20260924.json b/ml/pipelines/fixtures/session_20260924.json new file mode 100644 index 000000000..81da055f9 --- /dev/null +++ b/ml/pipelines/fixtures/session_20260924.json @@ -0,0 +1,189 @@ +{ + "master_sha": "03ffb33b54b86b52cf8a11cc2f7db6d2d8061772", + "issue": 175, + "captured_at": "2026-09-24T21:10:00Z", + "operator": "ACTIVE", + "session": "2026-09-24-1402-PDT", + "dual_gate_on_master": { + "repo_gate": "success", + "termux_smoke": "success", + "run_repo_gate": 36058704829, + "run_termux_smoke": 36058704728 + }, + "prs": [ + { + "number": 809, + "title": "GAMUT remote + wiki knowledge fabric", + "changed_files": 6, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase onto 03ffb33b then dual-gate" + }, + { + "number": 806, + "title": "evaluation lanes + GitHub App capability evidence", + "changed_files": 5, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase onto live master" + }, + { + "number": 787, + "title": "ML keep-alive DAG ICM-CCTV", + "changed_files": 103, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "keep_alive_parent": true, + "supersede": true, + "why": "re-extract onto 03ffb33b" + }, + { + "number": 746, + "title": "ML slim keep-alive on stale tip", + "changed_files": 90, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "supersede": true, + "why": "superseded by this extract" + }, + { + "number": 682, + "title": "ML wholesale DAG", + "changed_files": 130, + "mergeable_state": "dirty", + "ml_wholesale": true, + "tests": true, + "why": "size != quality" + }, + { + "number": 48, + "title": "llm-api-hub", + "changed_files": 73, + "mergeable_state": "dirty", + "master_staging_base": true, + "why": "do not retarget to master" + }, + { + "number": 788, + "title": "Linear client/sync coverage TER-15", + "changed_files": 22, + "mergeable_state": "unknown", + "master_staging_base": true, + "tests": true, + "why": "wrong-base master-staging" + }, + { + "number": 630, + "title": "Jules dashboard rich UI", + "changed_files": 54, + "mergeable_state": "dirty", + "minesweeper": true, + "author": "google-labs-jules[bot]", + "why": "minesweeper" + }, + { + "number": 804, + "title": "Sentinel context collector path traversal", + "changed_files": 4, + "mergeable_state": "unstable", + "security": true, + "author": "google-labs-jules[bot]", + "why": "observe security slice" + }, + { + "number": 803, + "title": "Bolt log parse / registry regex", + "changed_files": 6, + "mergeable_state": "unstable", + "author": "google-labs-jules[bot]", + "why": "observe perf slice" + }, + { + "number": 785, + "title": "Wingman wait-loop skills", + "changed_files": 10, + "mergeable_state": "unstable", + "dual_gate": "pending", + "tests": true, + "why": "wait-loop extract, do not collide" + }, + { + "number": 617, + "title": "proposal registry gate", + "changed_files": 5, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "rebase then dual-gate" + }, + { + "number": 543, + "title": "skill quality lane", + "changed_files": 21, + "mergeable_state": "unstable", + "stale_base": true, + "tests": true, + "why": "extract quality checks" + }, + { + "number": 587, + "title": "timing quotas", + "changed_files": 9, + "mergeable_state": "unstable", + "author": "google-labs-jules[bot]", + "why": "jules peer" + }, + { + "number": 762, + "title": "plugin connector parity", + "changed_files": 12, + "draft": true, + "why": "keep draft" + }, + { + "number": 764, + "title": "Termux hub Tailscale", + "changed_files": 18, + "draft": true, + "hitl_risk": true, + "why": "human device edge" + }, + { + "number": 725, + "title": "evidence JSONL SeekLog", + "changed_files": 11, + "mergeable_state": "unstable", + "tests": true, + "why": "independent extract" + }, + { + "number": 67, + "title": "PR scope discipline", + "changed_files": 2, + "docs_only": true, + "stale_base": true, + "why": "ancient docs" + }, + { + "number": 432, + "title": "observe-mode GitHub ML pipelines", + "changed_files": 140, + "mergeable_state": "dirty", + "ml_wholesale": true, + "why": "wholesale no-go" + }, + { + "number": 707, + "title": "historical promote specimen", + "changed_files": 3, + "mergeable_state": "clean", + "dual_gate": "green", + "tests": true, + "why": "historical promote specimen" + } + ] +} diff --git a/ml/pipelines/lanes/__init__.py b/ml/pipelines/lanes/__init__.py new file mode 100644 index 000000000..55c91a67c --- /dev/null +++ b/ml/pipelines/lanes/__init__.py @@ -0,0 +1 @@ +from .classify import classify_pr diff --git a/ml/pipelines/lanes/classify.py b/ml/pipelines/lanes/classify.py new file mode 100644 index 000000000..f5b411559 --- /dev/null +++ b/ml/pipelines/lanes/classify.py @@ -0,0 +1,27 @@ +"""Lane router: extract/hold/wait/observe/supersede short-circuit before scorer promote.""" +from __future__ import annotations +from typing import Any, Mapping +from ml.pipelines.lib.types import Lane +from ml.pipelines.moneyball.scorer import classify as moneyball_classify, score +from .extract_rules import should_extract +from .hold_rules import should_hold +from .observe_rules import should_observe +from .wait_rules import should_wait +from .supersede_rules import should_supersede +from .minesweeper_rules import is_minesweeper + +def classify_pr(pr: Mapping[str, Any]) -> Lane: + if should_supersede(pr): + return Lane.SUPERSEDE + if should_extract(pr) or (is_minesweeper(pr) and int(pr.get("changed_files") or 0) > 40): + if is_minesweeper(pr) and int(pr.get("changed_files") or 0) > 40: + return Lane.EXTRACT + if should_extract(pr): + return Lane.EXTRACT + if should_hold(pr): + return Lane.HOLD + if should_observe(pr): + return Lane.OBSERVE + if should_wait(pr): + return Lane.WAIT + return moneyball_classify(score(pr), pr) diff --git a/ml/pipelines/lanes/extract_rules.py b/ml/pipelines/lanes/extract_rules.py new file mode 100644 index 000000000..a04283880 --- /dev/null +++ b/ml/pipelines/lanes/extract_rules.py @@ -0,0 +1,12 @@ +from __future__ import annotations +from typing import Any, Mapping + +def should_extract(pr: Mapping[str, Any]) -> bool: + files = int(pr.get("changed_files") or 0) + if pr.get("ml_wholesale"): + return True + if files > 80 and str(pr.get("mergeable_state") or "") == "dirty": + return True + if pr.get("minesweeper") and files > 40: + return True + return False diff --git a/ml/pipelines/lanes/hold_rules.py b/ml/pipelines/lanes/hold_rules.py new file mode 100644 index 000000000..59fe9f2cf --- /dev/null +++ b/ml/pipelines/lanes/hold_rules.py @@ -0,0 +1,12 @@ +from __future__ import annotations +from typing import Any, Mapping + +def should_hold(pr: Mapping[str, Any]) -> bool: + if pr.get("master_staging_base") or pr.get("stacked_feature_base"): + return True + if pr.get("hitl_risk"): + return True + if str(pr.get("mergeable_state") or "") == "dirty" and int(pr.get("changed_files") or 0) <= 80: + if not pr.get("ml_wholesale"): + return True + return False diff --git a/ml/pipelines/lanes/minesweeper_rules.py b/ml/pipelines/lanes/minesweeper_rules.py new file mode 100644 index 000000000..aebeea8ee --- /dev/null +++ b/ml/pipelines/lanes/minesweeper_rules.py @@ -0,0 +1,19 @@ +"""Concurrent-agent overlap. Classify, do not overwrite WAIT peers.""" +from __future__ import annotations +from typing import Any, Mapping + +BOT_AUTHORS = ( + "google-labs-jules[bot]", + "devin-ai-integration[bot]", + "coderabbitai[bot]", + "github-actions[bot]", +) + +def is_minesweeper(pr: Mapping[str, Any]) -> bool: + if pr.get("minesweeper"): + return True + author = str(pr.get("author") or "") + files = int(pr.get("changed_files") or 0) + if author in BOT_AUTHORS and files > 40: + return True + return False diff --git a/ml/pipelines/lanes/observe_rules.py b/ml/pipelines/lanes/observe_rules.py new file mode 100644 index 000000000..65d8f30b5 --- /dev/null +++ b/ml/pipelines/lanes/observe_rules.py @@ -0,0 +1,10 @@ +from __future__ import annotations +from typing import Any, Mapping + +def should_observe(pr: Mapping[str, Any]) -> bool: + if pr.get("draft"): + return True + author = str(pr.get("author") or "") + if "jules" in author and not pr.get("minesweeper"): + return True + return False diff --git a/ml/pipelines/lanes/supersede_rules.py b/ml/pipelines/lanes/supersede_rules.py new file mode 100644 index 000000000..398c5fbe5 --- /dev/null +++ b/ml/pipelines/lanes/supersede_rules.py @@ -0,0 +1,13 @@ +"""Session pulses and stale keep-alive parents are SUPERSEDE, not WAIT.""" +from __future__ import annotations +from typing import Any, Mapping + +def should_supersede(pr: Mapping[str, Any]) -> bool: + if pr.get("supersede"): + return True + title = str(pr.get("title") or "").lower() + if "session pulse" in title or "lane-matrix" in title and "session" in title: + return True + if pr.get("keep_alive_parent") and pr.get("stale_base"): + return True + return False diff --git a/ml/pipelines/lanes/test_classify.py b/ml/pipelines/lanes/test_classify.py new file mode 100644 index 000000000..5cbf9c192 --- /dev/null +++ b/ml/pipelines/lanes/test_classify.py @@ -0,0 +1,13 @@ +import unittest +from ml.pipelines.lanes.classify import classify_pr +from ml.pipelines.lib.types import Lane + +class ClassifyPrTests(unittest.TestCase): + def test_promote(self) -> None: + self.assertEqual( + classify_pr({"changed_files": 3, "mergeable_state": "clean", "dual_gate": "green", "tests": True}), + Lane.PROMOTE, + ) + + def test_supersede(self) -> None: + self.assertEqual(classify_pr({"supersede": True}), Lane.SUPERSEDE) diff --git a/ml/pipelines/lanes/test_extract_rules.py b/ml/pipelines/lanes/test_extract_rules.py new file mode 100644 index 000000000..cf64a2cc0 --- /dev/null +++ b/ml/pipelines/lanes/test_extract_rules.py @@ -0,0 +1,9 @@ +import unittest +from ml.pipelines.lanes.extract_rules import should_extract + +class ExtractTests(unittest.TestCase): + def test_682(self) -> None: + self.assertTrue(should_extract({"ml_wholesale": True, "changed_files": 130})) + + def test_small(self) -> None: + self.assertFalse(should_extract({"changed_files": 4, "mergeable_state": "clean"})) diff --git a/ml/pipelines/lanes/test_hold_rules.py b/ml/pipelines/lanes/test_hold_rules.py new file mode 100644 index 000000000..aa1cc9f9a --- /dev/null +++ b/ml/pipelines/lanes/test_hold_rules.py @@ -0,0 +1,9 @@ +import unittest +from ml.pipelines.lanes.hold_rules import should_hold + +class HoldTests(unittest.TestCase): + def test_48(self) -> None: + self.assertTrue(should_hold({"number": 48, "master_staging_base": True, "changed_files": 73, "mergeable_state": "dirty"})) + + def test_69(self) -> None: + self.assertTrue(should_hold({"stacked_feature_base": True})) diff --git a/ml/pipelines/lanes/test_minesweeper_rules.py b/ml/pipelines/lanes/test_minesweeper_rules.py new file mode 100644 index 000000000..9cfef804c --- /dev/null +++ b/ml/pipelines/lanes/test_minesweeper_rules.py @@ -0,0 +1,12 @@ +import unittest +from ml.pipelines.lanes.minesweeper_rules import is_minesweeper + +class MinesweeperRules(unittest.TestCase): + def test_flag(self): + self.assertTrue(is_minesweeper({"minesweeper": True})) + + def test_jules_large(self): + self.assertTrue(is_minesweeper({"author": "google-labs-jules[bot]", "changed_files": 54})) + + def test_jules_small_not(self): + self.assertFalse(is_minesweeper({"author": "google-labs-jules[bot]", "changed_files": 4})) diff --git a/ml/pipelines/lanes/test_observe_rules.py b/ml/pipelines/lanes/test_observe_rules.py new file mode 100644 index 000000000..758e39050 --- /dev/null +++ b/ml/pipelines/lanes/test_observe_rules.py @@ -0,0 +1,6 @@ +import unittest +from ml.pipelines.lanes.observe_rules import should_observe + +class ObserveTests(unittest.TestCase): + def test_draft(self) -> None: + self.assertTrue(should_observe({"draft": True})) diff --git a/ml/pipelines/lanes/test_supersede_rules.py b/ml/pipelines/lanes/test_supersede_rules.py new file mode 100644 index 000000000..a9057f444 --- /dev/null +++ b/ml/pipelines/lanes/test_supersede_rules.py @@ -0,0 +1,15 @@ +import unittest +from ml.pipelines.lanes.supersede_rules import should_supersede + +class SupersedeRules(unittest.TestCase): + def test_flag(self): + self.assertTrue(should_supersede({"supersede": True})) + + def test_pulse_title(self): + self.assertTrue(should_supersede({"title": "ops(session): LANE-MATRIX pulse"})) + + def test_stale_keep_alive_parent(self): + self.assertTrue(should_supersede({"keep_alive_parent": True, "stale_base": True})) + + def test_negative(self): + self.assertFalse(should_supersede({"title": "feat(ml): keep-alive extract", "tests": True})) diff --git a/ml/pipelines/lanes/test_wait_rules.py b/ml/pipelines/lanes/test_wait_rules.py new file mode 100644 index 000000000..ac3e88a30 --- /dev/null +++ b/ml/pipelines/lanes/test_wait_rules.py @@ -0,0 +1,9 @@ +import unittest +from ml.pipelines.lanes.wait_rules import should_wait + +class WaitTests(unittest.TestCase): + def test_724(self) -> None: + self.assertTrue(should_wait({"mergeable_state": "unstable", "dual_gate": "green"})) + + def test_hitl_not_wait(self) -> None: + self.assertFalse(should_wait({"mergeable_state": "unstable", "hitl_risk": True})) diff --git a/ml/pipelines/lanes/wait_rules.py b/ml/pipelines/lanes/wait_rules.py new file mode 100644 index 000000000..a913d4ebc --- /dev/null +++ b/ml/pipelines/lanes/wait_rules.py @@ -0,0 +1,10 @@ +from __future__ import annotations +from typing import Any, Mapping + +def should_wait(pr: Mapping[str, Any]) -> bool: + if str(pr.get("mergeable_state") or "") == "unstable": + if not pr.get("hitl_risk") and not pr.get("draft"): + return True + if pr.get("stale_base") and pr.get("dual_gate") == "green": + return True + return False diff --git a/ml/pipelines/lib/__init__.py b/ml/pipelines/lib/__init__.py new file mode 100644 index 000000000..34848388f --- /dev/null +++ b/ml/pipelines/lib/__init__.py @@ -0,0 +1,2 @@ +from .engine import run_dag, summarize +from .types import Lane, StageResult, StageStatus diff --git a/ml/pipelines/lib/context.py b/ml/pipelines/lib/context.py new file mode 100644 index 000000000..010c45ec2 --- /dev/null +++ b/ml/pipelines/lib/context.py @@ -0,0 +1,20 @@ +"""Shared DAG context helpers.""" +from __future__ import annotations + +from typing import Any, Mapping, MutableMapping + + +def snapshot(ctx: Mapping[str, Any]) -> dict[str, Any]: + raw = ctx.get("snapshot") or {} + if not isinstance(raw, dict): + return {} + return raw + + +def prs(ctx: Mapping[str, Any]) -> list[dict[str, Any]]: + items = snapshot(ctx).get("prs") or [] + return [item for item in items if isinstance(item, dict)] + + +def put(ctx: MutableMapping[str, Any], key: str, value: Any) -> None: + ctx[key] = value diff --git a/ml/pipelines/lib/engine.py b/ml/pipelines/lib/engine.py new file mode 100644 index 000000000..4befbed13 --- /dev/null +++ b/ml/pipelines/lib/engine.py @@ -0,0 +1,25 @@ +"""engine: Sequential DAG runner with halt-on-fail.""" +from __future__ import annotations + +from typing import Any, Callable, Iterable, Mapping, MutableMapping + +from .types import StageResult, StageStatus + +StageFn = Callable[[MutableMapping[str, Any]], StageResult] + + +def run_dag(stages: Iterable[tuple[str, StageFn]], context: MutableMapping[str, Any]) -> list[StageResult]: + results: list[StageResult] = [] + for stage_id, fn in stages: + result = fn(context) + if result.stage_id != stage_id: + result = StageResult(stage_id=stage_id, status=result.status, artifacts=result.artifacts, notes=result.notes) + results.append(result) + context.setdefault("results", {})[stage_id] = result + if result.status is StageStatus.FAILED: + break + return results + + +def summarize(results: Iterable[StageResult]) -> Mapping[str, str]: + return {item.stage_id: item.status.value for item in results} diff --git a/ml/pipelines/lib/errors.py b/ml/pipelines/lib/errors.py new file mode 100644 index 000000000..d1002b431 --- /dev/null +++ b/ml/pipelines/lib/errors.py @@ -0,0 +1,10 @@ +class GateBlocked(RuntimeError): + """Promote packet refused.""" + + +class SnapshotError(RuntimeError): + """Fixture / snapshot unreadable.""" + + +class StageHalted(RuntimeError): + """DAG halted on FAILED stage.""" diff --git a/ml/pipelines/lib/io.py b/ml/pipelines/lib/io.py new file mode 100644 index 000000000..8a9c927a0 --- /dev/null +++ b/ml/pipelines/lib/io.py @@ -0,0 +1,19 @@ +from __future__ import annotations + +import json +from pathlib import Path +from typing import Any + +from .errors import SnapshotError + + +def load_json(path: Path) -> dict[str, Any]: + try: + return json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as exc: + raise SnapshotError(str(path)) from exc + + +def dump_json(path: Path, payload: Any) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n", encoding="utf-8") diff --git a/ml/pipelines/lib/latest.py b/ml/pipelines/lib/latest.py new file mode 100644 index 000000000..ed9c1a795 --- /dev/null +++ b/ml/pipelines/lib/latest.py @@ -0,0 +1,17 @@ +"""Resolve the newest keep-alive session fixture. No network.""" +from __future__ import annotations + +from pathlib import Path + +FIXTURES = Path(__file__).resolve().parents[1] / "fixtures" + + +def session_paths() -> list[Path]: + return sorted(FIXTURES.glob("session_*.json")) + + +def latest_session_path() -> Path: + paths = session_paths() + if not paths: + raise FileNotFoundError("no session_*.json fixtures") + return paths[-1] diff --git a/ml/pipelines/lib/test_context.py b/ml/pipelines/lib/test_context.py new file mode 100644 index 000000000..0e6c5050b --- /dev/null +++ b/ml/pipelines/lib/test_context.py @@ -0,0 +1,15 @@ +import unittest + +from ml.pipelines.lib.context import prs, put, snapshot + + +class ContextTests(unittest.TestCase): + def test_empty(self) -> None: + self.assertEqual(snapshot({}), {}) + self.assertEqual(prs({}), []) + + def test_put(self) -> None: + ctx: dict = {"snapshot": {"prs": [{"number": 1}]}} + put(ctx, "n", 3) + self.assertEqual(ctx["n"], 3) + self.assertEqual(len(prs(ctx)), 1) diff --git a/ml/pipelines/lib/test_engine.py b/ml/pipelines/lib/test_engine.py new file mode 100644 index 000000000..613ff2401 --- /dev/null +++ b/ml/pipelines/lib/test_engine.py @@ -0,0 +1,28 @@ +import unittest +from typing import Any, MutableMapping + +from ml.pipelines.lib.engine import run_dag, summarize +from ml.pipelines.lib.types import StageResult, StageStatus + + +def _ok(stage_id: str) -> StageResult: + return StageResult(stage_id=stage_id, status=StageStatus.OK) + + +def _fail(stage_id: str) -> StageResult: + return StageResult(stage_id=stage_id, status=StageStatus.FAILED) + + +class EngineTests(unittest.TestCase): + def test_runs_in_order(self) -> None: + ctx: MutableMapping[str, Any] = {} + results = run_dag([("a", lambda c: _ok("a")), ("b", lambda c: _ok("b"))], ctx) + self.assertEqual(summarize(results), {"a": "ok", "b": "ok"}) + + def test_halts_on_fail(self) -> None: + ctx: MutableMapping[str, Any] = {} + results = run_dag( + [("a", lambda c: _fail("a")), ("b", lambda c: _ok("b"))], + ctx, + ) + self.assertEqual([item.stage_id for item in results], ["a"]) diff --git a/ml/pipelines/lib/test_io.py b/ml/pipelines/lib/test_io.py new file mode 100644 index 000000000..48ddfea02 --- /dev/null +++ b/ml/pipelines/lib/test_io.py @@ -0,0 +1,19 @@ +import json +import tempfile +import unittest +from pathlib import Path + +from ml.pipelines.lib.errors import SnapshotError +from ml.pipelines.lib.io import dump_json, load_json + + +class IoTests(unittest.TestCase): + def test_roundtrip(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + path = Path(tmp) / "x.json" + dump_json(path, {"n": 1}) + self.assertEqual(load_json(path)["n"], 1) + + def test_missing(self) -> None: + with self.assertRaises(SnapshotError): + load_json(Path("/no/such/file.json")) diff --git a/ml/pipelines/lib/test_latest.py b/ml/pipelines/lib/test_latest.py new file mode 100644 index 000000000..930a83c5b --- /dev/null +++ b/ml/pipelines/lib/test_latest.py @@ -0,0 +1,13 @@ +import unittest +from ml.pipelines.lib.latest import latest_session_path, session_paths + +class LatestFixture(unittest.TestCase): + def test_at_least_one(self): + self.assertGreaterEqual(len(session_paths()), 1) + + def test_latest_is_newest_name(self): + paths = session_paths() + self.assertEqual(latest_session_path(), paths[-1]) + + def test_latest_is_20260924(self): + self.assertTrue(latest_session_path().name.startswith("session_20260924")) diff --git a/ml/pipelines/lib/test_weights.py b/ml/pipelines/lib/test_weights.py new file mode 100644 index 000000000..b08fb6fcf --- /dev/null +++ b/ml/pipelines/lib/test_weights.py @@ -0,0 +1,11 @@ +import unittest + +from ml.pipelines.lib.weights import THRESHOLDS, WEIGHTS + + +class WeightTests(unittest.TestCase): + def test_promote_threshold(self) -> None: + self.assertGreater(THRESHOLDS["promote"], THRESHOLDS["wait"]) + + def test_wholesale_penalty(self) -> None: + self.assertLess(WEIGHTS["ml_wholesale"], -1) diff --git a/ml/pipelines/lib/types.py b/ml/pipelines/lib/types.py new file mode 100644 index 000000000..17357d6be --- /dev/null +++ b/ml/pipelines/lib/types.py @@ -0,0 +1,31 @@ +"""types: Canonical records for the keep-alive DAG.""" +from __future__ import annotations + +from dataclasses import dataclass, field +from enum import Enum +from typing import Any, Mapping + + +class Lane(str, Enum): + PROMOTE = "promote" + WAIT = "wait" + HOLD = "hold" + EXTRACT = "extract" + OBSERVE = "observe" + SUPERSEDE = "supersede" + + +class StageStatus(str, Enum): + PENDING = "pending" + RUNNING = "running" + OK = "ok" + SKIPPED = "skipped" + FAILED = "failed" + + +@dataclass(frozen=True) +class StageResult: + stage_id: str + status: StageStatus + artifacts: Mapping[str, Any] = field(default_factory=dict) + notes: tuple[str, ...] = () diff --git a/ml/pipelines/lib/weights.py b/ml/pipelines/lib/weights.py new file mode 100644 index 000000000..9e7a009a1 --- /dev/null +++ b/ml/pipelines/lib/weights.py @@ -0,0 +1,28 @@ +"""Static MoneyBall weights. Keep transparent; never hide a term.""" +from __future__ import annotations + +WEIGHTS = { + "dual_gate_green": 5.0, + "mergeable_clean": 3.0, + "files_small": 1.5, + "tests_present": 1.0, + "security_fix": 2.0, + "docs_only": 0.5, + "stale_base": -2.0, + "mergeable_dirty": -8.0, + "mergeable_unstable": -1.0, + "changed_files_over_40": -3.0, + "minesweeper_overlap": -6.0, + "ml_wholesale": -10.0, + "jules_bot": -0.25, + "master_staging_base": -4.0, + "stacked_feature_base": -3.5, + "hitl_risk": -5.0, + "draft": -1.5, + "replay_landed": 0.75, +} + +THRESHOLDS = {"promote": 4.0, "wait": 0.0, "hold": -3.0} +FILES_SMALL = 8 +FILES_OVER_40 = 40 +FILES_WHOLESALE = 80 diff --git a/ml/pipelines/moneyball/__init__.py b/ml/pipelines/moneyball/__init__.py new file mode 100644 index 000000000..af492c40e --- /dev/null +++ b/ml/pipelines/moneyball/__init__.py @@ -0,0 +1,2 @@ +from .scorer import classify, score +from .explain import explain diff --git a/ml/pipelines/moneyball/explain.py b/ml/pipelines/moneyball/explain.py new file mode 100644 index 000000000..4d75060d9 --- /dev/null +++ b/ml/pipelines/moneyball/explain.py @@ -0,0 +1,20 @@ +"""Human-readable score breakdown.""" +from __future__ import annotations + +from typing import Any, Mapping + +from ml.pipelines.lib.weights import WEIGHTS +from ml.pipelines.moneyball.features import features +from ml.pipelines.moneyball.scorer import classify, score + + +def explain(pr: Mapping[str, Any]) -> dict[str, Any]: + feats = features(pr) + terms = {key: round(WEIGHTS[key] * value, 4) for key, value in feats.items() if value} + points = score(pr) + return { + "number": pr.get("number"), + "score": points, + "lane": classify(points, pr).value, + "terms": terms, + } diff --git a/ml/pipelines/moneyball/features.py b/ml/pipelines/moneyball/features.py new file mode 100644 index 000000000..74d5a6fa6 --- /dev/null +++ b/ml/pipelines/moneyball/features.py @@ -0,0 +1,31 @@ +"""Feature vector for a PR packet.""" +from __future__ import annotations + +from typing import Any, Mapping + +from ml.pipelines.lib.weights import FILES_OVER_40, FILES_SMALL + + +def features(pr: Mapping[str, Any]) -> dict[str, float]: + files = int(pr.get("changed_files") or 0) + state = str(pr.get("mergeable_state") or "unknown") + return { + "dual_gate_green": 1.0 if pr.get("dual_gate") == "green" else 0.0, + "mergeable_clean": 1.0 if state == "clean" else 0.0, + "files_small": 1.0 if files and files <= FILES_SMALL else 0.0, + "tests_present": 1.0 if pr.get("tests") else 0.0, + "security_fix": 1.0 if pr.get("security") else 0.0, + "docs_only": 1.0 if pr.get("docs_only") else 0.0, + "stale_base": 1.0 if pr.get("stale_base") else 0.0, + "mergeable_dirty": 1.0 if state == "dirty" else 0.0, + "mergeable_unstable": 1.0 if state == "unstable" else 0.0, + "changed_files_over_40": 1.0 if files > FILES_OVER_40 else 0.0, + "minesweeper_overlap": 1.0 if pr.get("minesweeper") else 0.0, + "ml_wholesale": 1.0 if pr.get("ml_wholesale") else 0.0, + "jules_bot": 1.0 if pr.get("author") == "google-labs-jules[bot]" else 0.0, + "master_staging_base": 1.0 if pr.get("master_staging_base") else 0.0, + "stacked_feature_base": 1.0 if pr.get("stacked_feature_base") else 0.0, + "hitl_risk": 1.0 if pr.get("hitl_risk") else 0.0, + "draft": 1.0 if pr.get("draft") else 0.0, + "replay_landed": 1.0 if pr.get("replay_landed") else 0.0, + } diff --git a/ml/pipelines/moneyball/scorer.py b/ml/pipelines/moneyball/scorer.py new file mode 100644 index 000000000..d8d23ed76 --- /dev/null +++ b/ml/pipelines/moneyball/scorer.py @@ -0,0 +1,37 @@ +"""scorer: Transparent weighted lane classifier.""" +from __future__ import annotations + +from typing import Any, Mapping + +from ml.pipelines.lib.types import Lane +from ml.pipelines.lib.weights import FILES_WHOLESALE, THRESHOLDS, WEIGHTS +from ml.pipelines.moneyball.features import features as _features + + +def score(pr: Mapping[str, Any]) -> float: + feats = _features(pr) + return round(sum(WEIGHTS[key] * value for key, value in feats.items()), 4) + + +def classify(points: float, pr: Mapping[str, Any]) -> Lane: + state = str(pr.get("mergeable_state") or "") + files = int(pr.get("changed_files") or 0) + if pr.get("ml_wholesale") or (files > FILES_WHOLESALE and state == "dirty"): + return Lane.EXTRACT + if pr.get("master_staging_base") or pr.get("stacked_feature_base"): + return Lane.HOLD + if pr.get("hitl_risk"): + return Lane.HOLD + if pr.get("draft"): + return Lane.OBSERVE + if pr.get("supersede"): + return Lane.SUPERSEDE + if state == "dirty" or pr.get("minesweeper"): + return Lane.HOLD if files <= FILES_WHOLESALE else Lane.EXTRACT + if points >= THRESHOLDS["promote"] and pr.get("dual_gate") == "green" and state == "clean": + return Lane.PROMOTE + if state == "unstable" or points >= THRESHOLDS["wait"]: + return Lane.WAIT + if points >= THRESHOLDS["hold"]: + return Lane.OBSERVE + return Lane.HOLD diff --git a/ml/pipelines/moneyball/test_explain.py b/ml/pipelines/moneyball/test_explain.py new file mode 100644 index 000000000..9262878bd --- /dev/null +++ b/ml/pipelines/moneyball/test_explain.py @@ -0,0 +1,10 @@ +import unittest + +from ml.pipelines.moneyball.explain import explain + + +class ExplainTests(unittest.TestCase): + def test_terms(self) -> None: + row = explain({"number": 1, "changed_files": 2, "dual_gate": "green", "mergeable_state": "clean", "tests": True}) + self.assertIn("dual_gate_green", row["terms"]) + self.assertEqual(row["lane"], "promote") diff --git a/ml/pipelines/moneyball/test_features.py b/ml/pipelines/moneyball/test_features.py new file mode 100644 index 000000000..61f4b9b60 --- /dev/null +++ b/ml/pipelines/moneyball/test_features.py @@ -0,0 +1,10 @@ +import unittest + +from ml.pipelines.moneyball.features import features + + +class FeatureTests(unittest.TestCase): + def test_flags(self) -> None: + feats = features({"changed_files": 3, "dual_gate": "green", "mergeable_state": "clean", "tests": True}) + self.assertEqual(feats["files_small"], 1.0) + self.assertEqual(feats["dual_gate_green"], 1.0) diff --git a/ml/pipelines/moneyball/test_scorer.py b/ml/pipelines/moneyball/test_scorer.py new file mode 100644 index 000000000..16158ae1e --- /dev/null +++ b/ml/pipelines/moneyball/test_scorer.py @@ -0,0 +1,60 @@ +import unittest + +from ml.pipelines.lib.types import Lane +from ml.pipelines.moneyball.scorer import classify, score + + +class ScorerTests(unittest.TestCase): + def test_wholesale_is_extract(self) -> None: + pr = {"number": 432, "changed_files": 120, "mergeable_state": "dirty", "ml_wholesale": True} + self.assertEqual(classify(score(pr), pr), Lane.EXTRACT) + + def test_green_small_is_promote(self) -> None: + pr = { + "number": 707, + "changed_files": 3, + "mergeable_state": "clean", + "dual_gate": "green", + "tests": True, + } + self.assertEqual(classify(score(pr), pr), Lane.PROMOTE) + + def test_minesweeper_extract(self) -> None: + pr = {"number": 630, "changed_files": 89, "mergeable_state": "dirty", "minesweeper": True} + self.assertEqual(classify(score(pr), pr), Lane.EXTRACT) + + def test_staging_hold(self) -> None: + pr = { + "number": 48, + "changed_files": 73, + "mergeable_state": "dirty", + "master_staging_base": True, + } + self.assertEqual(classify(score(pr), pr), Lane.HOLD) + + def test_stacked_hold(self) -> None: + pr = { + "number": 69, + "changed_files": 20, + "mergeable_state": "dirty", + "stacked_feature_base": True, + } + self.assertEqual(classify(score(pr), pr), Lane.HOLD) + + def test_unstable_wait(self) -> None: + pr = { + "number": 724, + "changed_files": 18, + "mergeable_state": "unstable", + "dual_gate": "green", + "tests": True, + } + self.assertEqual(classify(score(pr), pr), Lane.WAIT) + + def test_hitl_hold(self) -> None: + pr = {"number": 739, "changed_files": 12, "mergeable_state": "unstable", "hitl_risk": True} + self.assertEqual(classify(score(pr), pr), Lane.HOLD) + + def test_draft_observe(self) -> None: + pr = {"number": 740, "changed_files": 10, "mergeable_state": "unstable", "draft": True} + self.assertEqual(classify(score(pr), pr), Lane.OBSERVE) diff --git a/ml/pipelines/moneyball/test_thresholds.py b/ml/pipelines/moneyball/test_thresholds.py new file mode 100644 index 000000000..52e931f46 --- /dev/null +++ b/ml/pipelines/moneyball/test_thresholds.py @@ -0,0 +1,9 @@ +import unittest + +from ml.pipelines.moneyball.thresholds import THRESHOLDS + + +class ThresholdImportTests(unittest.TestCase): + def test_keys(self) -> None: + self.assertIn("promote", THRESHOLDS) + self.assertIn("wait", THRESHOLDS) diff --git a/ml/pipelines/moneyball/thresholds.py b/ml/pipelines/moneyball/thresholds.py new file mode 100644 index 000000000..4b62df17e --- /dev/null +++ b/ml/pipelines/moneyball/thresholds.py @@ -0,0 +1,3 @@ +from ml.pipelines.lib.weights import THRESHOLDS + +__all__ = ['THRESHOLDS'] diff --git a/ml/pipelines/py.typed b/ml/pipelines/py.typed new file mode 100644 index 000000000..e69de29bb diff --git a/ml/pipelines/replay/__init__.py b/ml/pipelines/replay/__init__.py new file mode 100644 index 000000000..ab1232b05 --- /dev/null +++ b/ml/pipelines/replay/__init__.py @@ -0,0 +1,3 @@ +from .lineage import lineage_fields, sufficient_gain +from .action_effect import import_events +from .corpus import project diff --git a/ml/pipelines/replay/action_effect.py b/ml/pipelines/replay/action_effect.py new file mode 100644 index 000000000..872903025 --- /dev/null +++ b/ml/pipelines/replay/action_effect.py @@ -0,0 +1,17 @@ +"""Import Action Effectiveness events into replay history (#742).""" +from __future__ import annotations +from typing import Any, Iterable, Mapping + +def import_events(events: Iterable[Mapping[str, Any]]) -> list[dict[str, Any]]: + out: list[dict[str, Any]] = [] + for event in events: + kind = str(event.get("kind") or event.get("type") or "action_effect") + out.append( + { + "kind": kind, + "sha": event.get("sha"), + "outcome": event.get("outcome") or "UNKNOWN", + "notes": event.get("notes") or "", + } + ) + return out diff --git a/ml/pipelines/replay/corpus.py b/ml/pipelines/replay/corpus.py new file mode 100644 index 000000000..46ff36162 --- /dev/null +++ b/ml/pipelines/replay/corpus.py @@ -0,0 +1,11 @@ +"""Project replay evidence into corpus experiments (#741).""" +from __future__ import annotations +from typing import Any, Iterable, Mapping +from .lineage import lineage_fields + +def project(events: Iterable[Mapping[str, Any]]) -> list[dict[str, Any]]: + rows = [] + for event in events: + fields = lineage_fields(event) + rows.append({**fields, "corpus": "evolutionary-replay"}) + return rows diff --git a/ml/pipelines/replay/lineage.py b/ml/pipelines/replay/lineage.py new file mode 100644 index 000000000..284ee8d1f --- /dev/null +++ b/ml/pipelines/replay/lineage.py @@ -0,0 +1,15 @@ +"""Project replay evidence (#741/#742) into pipeline context. Isolated adapter.""" +from __future__ import annotations +from typing import Any, Mapping + +def lineage_fields(event: Mapping[str, Any]) -> dict[str, Any]: + return { + "experiment_id": event.get("experiment_id") or event.get("id"), + "parent_sha": event.get("parent_sha"), + "child_sha": event.get("child_sha") or event.get("sha"), + "gain": float(event.get("gain") or 0.0), + "promoted": bool(event.get("promoted")), + } + +def sufficient_gain(event: Mapping[str, Any], floor: float = 0.0) -> bool: + return lineage_fields(event)["gain"] > floor diff --git a/ml/pipelines/replay/test_action_effect.py b/ml/pipelines/replay/test_action_effect.py new file mode 100644 index 000000000..f4c9627c3 --- /dev/null +++ b/ml/pipelines/replay/test_action_effect.py @@ -0,0 +1,7 @@ +import unittest +from ml.pipelines.replay.action_effect import import_events + +class ActionEffectTests(unittest.TestCase): + def test_import(self) -> None: + rows = import_events([{"sha": "d10a7a54", "outcome": "PASS", "kind": "help-wanted"}]) + self.assertEqual(rows[0]["outcome"], "PASS") diff --git a/ml/pipelines/replay/test_corpus.py b/ml/pipelines/replay/test_corpus.py new file mode 100644 index 000000000..a1b4c1259 --- /dev/null +++ b/ml/pipelines/replay/test_corpus.py @@ -0,0 +1,7 @@ +import unittest +from ml.pipelines.replay.corpus import project + +class CorpusTests(unittest.TestCase): + def test_project(self) -> None: + rows = project([{"experiment_id": "x", "gain": 1, "sha": "aa"}]) + self.assertEqual(rows[0]["corpus"], "evolutionary-replay") diff --git a/ml/pipelines/replay/test_lineage.py b/ml/pipelines/replay/test_lineage.py new file mode 100644 index 000000000..a800f6968 --- /dev/null +++ b/ml/pipelines/replay/test_lineage.py @@ -0,0 +1,11 @@ +import unittest +from ml.pipelines.replay.lineage import lineage_fields, sufficient_gain + +class LineageTests(unittest.TestCase): + def test_gain(self) -> None: + event = {"experiment_id": "e1", "gain": 0.2, "sha": "abc"} + self.assertTrue(sufficient_gain(event)) + self.assertEqual(lineage_fields(event)["child_sha"], "abc") + + def test_lock_incumbent(self) -> None: + self.assertFalse(sufficient_gain({"gain": 0.0}, floor=0.1)) diff --git a/ml/pipelines/stages/__init__.py b/ml/pipelines/stages/__init__.py new file mode 100644 index 000000000..3ac40cb2d --- /dev/null +++ b/ml/pipelines/stages/__init__.py @@ -0,0 +1,11 @@ +from . import deploy, evaluate, features, ingest, monitor, recon, train + +STAGES = [ + ("00_recon", recon.run), + ("10_ingest", ingest.run), + ("20_features", features.run), + ("30_train", train.run), + ("40_evaluate", evaluate.run), + ("50_deploy", deploy.run), + ("60_monitor", monitor.run), +] diff --git a/ml/pipelines/stages/deploy.py b/ml/pipelines/stages/deploy.py new file mode 100644 index 000000000..e74b66088 --- /dev/null +++ b/ml/pipelines/stages/deploy.py @@ -0,0 +1,17 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.contracts.gate import assert_promotable +from ml.pipelines.lib.context import prs +from ml.pipelines.lib.errors import GateBlocked +from ml.pipelines.lib.types import StageResult, StageStatus +from ml.pipelines.moneyball.scorer import classify, score + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + blocked = 0 + for pr in prs(ctx): + lane = classify(score(pr), pr) + try: + assert_promotable(pr, lane) + except GateBlocked: + blocked += 1 + return StageResult(stage_id="50_deploy", status=StageStatus.OK, artifacts={"blocked": blocked}) diff --git a/ml/pipelines/stages/evaluate.py b/ml/pipelines/stages/evaluate.py new file mode 100644 index 000000000..eed68d1f2 --- /dev/null +++ b/ml/pipelines/stages/evaluate.py @@ -0,0 +1,13 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.lib.context import prs, put +from ml.pipelines.lib.types import StageResult, StageStatus +from ml.pipelines.moneyball.scorer import classify, score + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + lanes = [ + {"number": pr.get("number"), "lane": classify(score(pr), pr).value, "score": score(pr)} + for pr in prs(ctx) + ] + put(ctx, "lanes", lanes) + return StageResult(stage_id="40_evaluate", status=StageStatus.OK, artifacts={"lanes": lanes}) diff --git a/ml/pipelines/stages/features.py b/ml/pipelines/stages/features.py new file mode 100644 index 000000000..f63bd884f --- /dev/null +++ b/ml/pipelines/stages/features.py @@ -0,0 +1,10 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.lib.context import prs, put +from ml.pipelines.lib.types import StageResult, StageStatus +from ml.pipelines.moneyball.scorer import score + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + scored = [{"number": pr.get("number"), "score": score(pr)} for pr in prs(ctx)] + put(ctx, "scored", scored) + return StageResult(stage_id="20_features", status=StageStatus.OK, artifacts={"n": len(scored)}) diff --git a/ml/pipelines/stages/ingest.py b/ml/pipelines/stages/ingest.py new file mode 100644 index 000000000..d74e70122 --- /dev/null +++ b/ml/pipelines/stages/ingest.py @@ -0,0 +1,8 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.lib.context import prs +from ml.pipelines.lib.types import StageResult, StageStatus + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + items = prs(ctx) + return StageResult(stage_id="10_ingest", status=StageStatus.OK, artifacts={"n": len(items)}) diff --git a/ml/pipelines/stages/monitor.py b/ml/pipelines/stages/monitor.py new file mode 100644 index 000000000..05e851de2 --- /dev/null +++ b/ml/pipelines/stages/monitor.py @@ -0,0 +1,10 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.lib.types import StageResult, StageStatus + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + return StageResult( + stage_id="60_monitor", + status=StageStatus.OK, + artifacts={"wait": "dual-gate on extract", "operator": "ACTIVE"}, + ) diff --git a/ml/pipelines/stages/recon.py b/ml/pipelines/stages/recon.py new file mode 100644 index 000000000..a9b2acd3b --- /dev/null +++ b/ml/pipelines/stages/recon.py @@ -0,0 +1,12 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.lib.context import snapshot +from ml.pipelines.lib.types import StageResult, StageStatus + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + snap = snapshot(ctx) + return StageResult( + stage_id="00_recon", + status=StageStatus.OK, + artifacts={"master_sha": snap.get("master_sha"), "issue": snap.get("issue", 175)}, + ) diff --git a/ml/pipelines/stages/test_stages.py b/ml/pipelines/stages/test_stages.py new file mode 100644 index 000000000..a0fba8df0 --- /dev/null +++ b/ml/pipelines/stages/test_stages.py @@ -0,0 +1,22 @@ +import unittest +from ml.pipelines.lib.engine import run_dag, summarize +from ml.pipelines.stages import STAGES + +SNAP = { + "master_sha": "d10a7a54", + "issue": 175, + "prs": [ + {"number": 707, "changed_files": 3, "mergeable_state": "clean", "dual_gate": "green", "tests": True}, + {"number": 48, "changed_files": 73, "mergeable_state": "dirty", "master_staging_base": True}, + ], +} + +class StageTests(unittest.TestCase): + def test_dag(self) -> None: + ctx: dict = {"snapshot": SNAP} + results = run_dag(STAGES, ctx) + self.assertEqual(len(results), 7) + self.assertTrue(all(v == "ok" for v in summarize(results).values())) + lanes = {row["number"]: row["lane"] for row in ctx["lanes"]} + self.assertEqual(lanes[707], "promote") + self.assertEqual(lanes[48], "hold") diff --git a/ml/pipelines/stages/train.py b/ml/pipelines/stages/train.py new file mode 100644 index 000000000..d20ca245a --- /dev/null +++ b/ml/pipelines/stages/train.py @@ -0,0 +1,11 @@ +from __future__ import annotations +from typing import Any, MutableMapping +from ml.pipelines.lib.types import StageResult, StageStatus +from ml.pipelines.lib.weights import WEIGHTS + +def run(ctx: MutableMapping[str, Any]) -> StageResult: + return StageResult( + stage_id="30_train", + status=StageStatus.OK, + artifacts={"note": "weights static", "n_weights": len(WEIGHTS)}, + ) diff --git a/ml/pipelines/test_cli.py b/ml/pipelines/test_cli.py new file mode 100644 index 000000000..861385dd6 --- /dev/null +++ b/ml/pipelines/test_cli.py @@ -0,0 +1,13 @@ +import json +import unittest +from ml.pipelines.cli import main + +class CliTests(unittest.TestCase): + def test_status(self) -> None: + self.assertEqual(main(["status"]), 0) + + def test_lanes(self) -> None: + self.assertEqual(main(["lanes"]), 0) + + def test_gate_blocks_48(self) -> None: + self.assertEqual(main(["gate", "48"]), 2) diff --git a/ml/pipelines/test_end_to_end.py b/ml/pipelines/test_end_to_end.py new file mode 100644 index 000000000..e51eaca36 --- /dev/null +++ b/ml/pipelines/test_end_to_end.py @@ -0,0 +1,20 @@ +import unittest +from ml.pipelines.cli import _fixture +from ml.pipelines.lanes.classify import classify_pr +from ml.pipelines.lib.types import Lane + + +class EndToEndTests(unittest.TestCase): + def test_fixture_classifications(self) -> None: + snap = _fixture() + got = {pr["number"]: classify_pr(pr) for pr in snap["prs"]} + self.assertEqual(got[48], Lane.HOLD) + self.assertEqual(got[788], Lane.HOLD) + self.assertEqual(got[785], Lane.WAIT) + self.assertEqual(got[682], Lane.EXTRACT) + self.assertEqual(got[432], Lane.EXTRACT) + self.assertEqual(got[630], Lane.EXTRACT) + self.assertEqual(got[764], Lane.HOLD) + self.assertEqual(got[762], Lane.OBSERVE) + self.assertEqual(got[787], Lane.SUPERSEDE) + self.assertEqual(got[707], Lane.PROMOTE) diff --git a/ml/pipelines/test_session_20260923.py b/ml/pipelines/test_session_20260923.py new file mode 100644 index 000000000..c93b2707a --- /dev/null +++ b/ml/pipelines/test_session_20260923.py @@ -0,0 +1,39 @@ +import json +import unittest +from pathlib import Path + +from ml.pipelines.lanes.classify import classify_pr +from ml.pipelines.moneyball.scorer import score + + +class Session20260923(unittest.TestCase): + def setUp(self): + path = Path(__file__).resolve().parent / "fixtures" / "session_20260923.json" + self.payload = json.loads(path.read_text()) + + def test_tip(self): + self.assertTrue(self.payload["master_sha"].startswith("ba5f6b6d")) + self.assertEqual(self.payload["issue"], 175) + + def test_operator_active(self): + self.assertEqual(self.payload["operator"], "ACTIVE") + + def test_hold_48(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 48) + self.assertEqual(classify_pr(pr).value, "hold") + + def test_extract_682(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 682) + self.assertEqual(classify_pr(pr).value, "extract") + + def test_wait_785(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 785) + self.assertEqual(classify_pr(pr).value, "wait") + + def test_scores_finite(self): + for pr in self.payload["prs"]: + self.assertIsInstance(score(pr), float) + + +if __name__ == "__main__": + unittest.main() diff --git a/ml/pipelines/test_session_20260924.py b/ml/pipelines/test_session_20260924.py new file mode 100644 index 000000000..3aa0a8dea --- /dev/null +++ b/ml/pipelines/test_session_20260924.py @@ -0,0 +1,64 @@ +import json +import unittest +from pathlib import Path + +from ml.pipelines.lanes.classify import classify_pr +from ml.pipelines.lib.latest import latest_session_path +from ml.pipelines.moneyball.scorer import score + + +class Session20260924(unittest.TestCase): + def setUp(self): + path = Path(__file__).resolve().parent / "fixtures" / "session_20260924.json" + self.payload = json.loads(path.read_text()) + + def test_tip(self): + self.assertTrue(self.payload["master_sha"].startswith("03ffb33b")) + self.assertEqual(self.payload["issue"], 175) + + def test_latest_points_here(self): + self.assertEqual(latest_session_path().name, "session_20260924.json") + + def test_operator_active(self): + self.assertEqual(self.payload["operator"], "ACTIVE") + + def test_master_dual_gate_green(self): + dg = self.payload["dual_gate_on_master"] + self.assertEqual(dg["repo_gate"], "success") + self.assertEqual(dg["termux_smoke"], "success") + + def test_hold_48(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 48) + self.assertEqual(classify_pr(pr).value, "hold") + + def test_extract_682(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 682) + self.assertEqual(classify_pr(pr).value, "extract") + + def test_extract_432(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 432) + self.assertEqual(classify_pr(pr).value, "extract") + + def test_supersede_787(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 787) + self.assertEqual(classify_pr(pr).value, "supersede") + + def test_extract_630_minesweeper(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 630) + self.assertEqual(classify_pr(pr).value, "extract") + + def test_promote_707(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 707) + self.assertEqual(classify_pr(pr).value, "promote") + + def test_hold_788(self): + pr = next(p for p in self.payload["prs"] if p["number"] == 788) + self.assertEqual(classify_pr(pr).value, "hold") + + def test_scores_finite(self): + for pr in self.payload["prs"]: + self.assertIsInstance(score(pr), float) + + +if __name__ == "__main__": + unittest.main() diff --git a/ml/pipelines/viz/__init__.py b/ml/pipelines/viz/__init__.py new file mode 100644 index 000000000..d38391b53 --- /dev/null +++ b/ml/pipelines/viz/__init__.py @@ -0,0 +1,2 @@ +from .cctv import cctv +from .projection import project diff --git a/ml/pipelines/viz/cctv.py b/ml/pipelines/viz/cctv.py new file mode 100644 index 000000000..1e6324bb0 --- /dev/null +++ b/ml/pipelines/viz/cctv.py @@ -0,0 +1,17 @@ +"""ICM-CCTV projection for ops dashboards (Stepie 2087 step 10023).""" +from __future__ import annotations +from collections import Counter +from typing import Any, Iterable, Mapping + +def cctv(rows: Iterable[Mapping[str, Any]]) -> dict[str, Any]: + items = list(rows) + counts = Counter(str(row.get("lane") or "observe") for row in items) + return { + "surface": "icm-cctv", + "n": len(items), + "lanes": dict(counts), + "operator": "ACTIVE", + "primary_goal": 2149, + "adjacent_goal": 2087, + "step": 10023, + } diff --git a/ml/pipelines/viz/matrix_view.py b/ml/pipelines/viz/matrix_view.py new file mode 100644 index 000000000..2fa60211a --- /dev/null +++ b/ml/pipelines/viz/matrix_view.py @@ -0,0 +1,15 @@ +from __future__ import annotations +from typing import Any, Iterable, Mapping + +def matrix_rows(rows: Iterable[Mapping[str, Any]]) -> list[dict[str, Any]]: + out = [] + for row in rows: + out.append( + { + "pr": row.get("number"), + "lane": row.get("lane"), + "score": row.get("score"), + "why": row.get("why") or row.get("lane"), + } + ) + return out diff --git a/ml/pipelines/viz/projection.py b/ml/pipelines/viz/projection.py new file mode 100644 index 000000000..32f24f68c --- /dev/null +++ b/ml/pipelines/viz/projection.py @@ -0,0 +1,14 @@ +"""JSON projection consumed by operator dashboards. No secrets.""" +from __future__ import annotations +from typing import Any, Mapping +from .cctv import cctv +from .matrix_view import matrix_rows + +def project(snapshot: Mapping[str, Any], lanes: list[dict[str, Any]]) -> dict[str, Any]: + return { + "master_sha": snapshot.get("master_sha"), + "issue": snapshot.get("issue", 175), + "captured_at": snapshot.get("captured_at"), + "cctv": cctv(lanes), + "rows": matrix_rows(lanes), + } diff --git a/ml/pipelines/viz/test_cctv.py b/ml/pipelines/viz/test_cctv.py new file mode 100644 index 000000000..17534802a --- /dev/null +++ b/ml/pipelines/viz/test_cctv.py @@ -0,0 +1,9 @@ +import unittest +from ml.pipelines.viz.cctv import cctv + +class CctvTests(unittest.TestCase): + def test_counts(self) -> None: + view = cctv([{"lane": "hold"}, {"lane": "hold"}, {"lane": "wait"}]) + self.assertEqual(view["n"], 3) + self.assertEqual(view["lanes"]["hold"], 2) + self.assertEqual(view["step"], 10023) diff --git a/ml/pipelines/viz/test_projection.py b/ml/pipelines/viz/test_projection.py new file mode 100644 index 000000000..6bd3fef7b --- /dev/null +++ b/ml/pipelines/viz/test_projection.py @@ -0,0 +1,8 @@ +import unittest +from ml.pipelines.viz.projection import project + +class ProjectionTests(unittest.TestCase): + def test_project(self) -> None: + payload = project({"master_sha": "d10a7a54", "issue": 175}, [{"number": 724, "lane": "wait", "score": 1.2}]) + self.assertEqual(payload["rows"][0]["pr"], 724) + self.assertEqual(payload["cctv"]["surface"], "icm-cctv")