From 1f6339342112ff10df801431afe85d206d06a526 Mon Sep 17 00:00:00 2001 From: Alpamys Date: Tue, 2 Jun 2026 14:34:36 +0500 Subject: [PATCH] feat(v0.71.5): ingest/data/prompt/drift polish MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes #157, #205, #207, #149, #164, #163. Defers #204 (live SaaS pull — paid accounts, infra-blocked, kept open). - #164: get_metric_series falls back to eval_results when metrics is empty - #163: build_verdict confidence biased by advise_history (same project+choice, >=3 precedents); decision never changes - #207: shared utils/webhooks.py (SSRF-hardened) + --slack-url/--discord-url on ingest/prune-prompt/ab/active-sample; ab fires only on a decision - #205: soup prune-prompt --tokenizer (token-prefix detect + decode remainder, boundary-safe) - #149: DynamicCurriculumCallback buckets by loss/perplexity percentile; length keeps round-robin - #157: soup data push/forge --hub modelscope|modelers (data score N/A) 107 new tests in tests/test_v0715.py (12474 -> 12581). ruff clean. --- CHANGELOG.md | 45 + CONTRIBUTING.md | 2 +- README.md | 29 +- docs/backends-and-ops.md | 5 + docs/commands.md | 6 +- docs/data.md | 12 + docs/evaluation.md | 4 + docs/training.md | 3 + pyproject.toml | 2 +- src/soup_cli/__init__.py | 2 +- src/soup_cli/commands/_webhook_cli.py | 71 + src/soup_cli/commands/ab.py | 40 + src/soup_cli/commands/active_sample.py | 28 + src/soup_cli/commands/advise.py | 22 +- src/soup_cli/commands/data.py | 76 +- src/soup_cli/commands/data_forge.py | 46 +- src/soup_cli/commands/ingest.py | 25 + src/soup_cli/commands/prune_prompt.py | 34 + src/soup_cli/experiment/tracker.py | 38 +- .../monitoring/curriculum_callback.py | 80 +- src/soup_cli/utils/advise.py | 196 ++- src/soup_cli/utils/advise_history.py | 10 + src/soup_cli/utils/curriculum_dynamic.py | 43 + src/soup_cli/utils/drift_alarm.py | 112 +- src/soup_cli/utils/peft_wiring.py | 14 +- src/soup_cli/utils/prune_prompt.py | 183 ++- src/soup_cli/utils/webhooks.py | 152 ++ tests/test_v0715.py | 1222 +++++++++++++++++ 28 files changed, 2361 insertions(+), 141 deletions(-) create mode 100644 src/soup_cli/commands/_webhook_cli.py create mode 100644 src/soup_cli/utils/webhooks.py create mode 100644 tests/test_v0715.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 68d6507..9d04713 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,51 @@ reproducing 70+ versions of notes. ## [Unreleased] +## [0.71.5] - 2026-06-02 + +### Added +- **`soup eval against` now reads eval metrics** — `ExperimentTracker.get_metric_series` + falls back to the `eval_results` table when the metric is not a per-step + training column (`loss` / `lr` / `grad_norm` / `speed` / `gpu_mem`). So + `soup eval against --candidate --metric task_accuracy` returns a + real score series (benchmark scores live in `eval_results`, not `metrics`) + instead of "Empty series". Per-step columns still read from `metrics` — no + regression for existing callers. +- **`soup advise` learns from past project outcomes** — `soup advise` now reads + this project's accepted-verdict history (`~/.soup/advise_history.jsonl`) and + biases the rubric: 3+ successful SFT precedents flip a marginal RAG call to + SFT; 3+ negative GRPO outcomes suppress GRPO in favour of SFT-on-traces; an + encouraged choice gets a small confidence nudge. Scoped per-project (one + project's record never biases another). No history → identical to before. +- **Slack/Discord webhooks on four more commands** — `--slack-url` / `--discord-url` + (SSRF-hardened, loopback-only HTTP, RFC1918 rejected, never crashes the + command) now ship on `soup ingest`, `soup prune-prompt`, `soup ab` (fires only + on a `reject_h0` / `accept_h0` decision, not `continue`), and + `soup data active-sample` — not just `soup drift-alarm`. The validator + sender + moved to a shared `soup_cli/utils/webhooks.py`. +- **Tokenizer-aware `soup prune-prompt`** — `--tokenizer ` detects + and strips the shared system-prompt prefix on **token** boundaries instead of + characters, so a multi-byte UTF-8 prefix can never be split mid-code-point. + Default (no `--tokenizer`) keeps the whitespace-character behaviour. +- **Curriculum bucketing by loss percentile** — `DynamicCurriculumCallback` now + buckets samples by the percentile rank of the live loss (or perplexity) + signal within a rolling window when `data.curriculum_metric` is `loss` / + `perplexity`, so a consistently-hard sample is routed to the same difficulty + bucket across recomputes. `length` and warm-up still use round-robin. +- **`--hub` on `soup data push` and `soup data forge`** — `soup data push + --hub modelscope|modelers` uploads a dataset via the matching SDK + (`repo_type=dataset`, commit message sanitised); `soup data forge --hub + --teacher owner/name` pre-fetches the teacher model from that hub + (and warns when the teacher is not a repo id so `--hub` is never silently + ignored). HF stays the default. + +### Notes +- Live SaaS *pull* adapters for `soup ingest` (Langfuse / LangSmith / Helicone / + OpenPipe / OpenAI SDKs, issue #204) remain deferred: they need credentialed + vendor accounts with populated trace data to validate honestly. Tracked as an + open, `infra-blocked` (external-account) item. `soup ingest` continues to parse + the JSONL export you pull from your dashboard. + ## [0.71.4] - 2026-06-02 ### Added diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 15bf211..2770c1b 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -120,7 +120,7 @@ src/soup_cli/ templates/ - 17 built-in soup.yaml templates (YAML + manifest.json) with load_template loader (v0.39.0, +bco v0.40.0) ui/ - Web UI (FastAPI + HTML/JS SPA) -tests/ - Test suite (274 files, 12474 tests) +tests/ - Test suite (275 files, 12581 tests) examples/ - Real-world config examples and datasets ``` diff --git a/README.md b/README.md index 76d4bd9..9b4d16f 100644 --- a/README.md +++ b/README.md @@ -49,21 +49,22 @@ infrastructure instead of improving models. Soup fixes that. ## What's New -**v0.71.4 — Adapter lifecycle + loop wiring.** The merge, PR, and continuous-loop surfaces go live: +**v0.71.5 — Ingest, data & prompt polish.** Sharper production-loop ergonomics: -- **Canary verdict on merge** — `soup adapters merge … --canary suite.json` scores the merged - adapter and reports **OK / MINOR / MAJOR**; `--strict-verdict` exits non-zero on a MAJOR - regression. Works with no model load using a pre-scored canary suite. -- **Evolutionary merge for real** — `soup adapters merge --strategy cmaes --eval suite --budget 1h` - now runs the full CMA-ES search (merge → score → optimise) and writes the best blend, instead of - just printing a plan. -- **Publish an adapter PR** — `soup adapters pr --base-sha <hex> --adapter <path> --push - owner/repo#42` posts the rendered PR straight to a GitHub PR comment. -- **Continuous fine-tuning loop, wired up** — `soup loop watch --pre-wired` runs the real - traces → DPO → eval-gate → canary pipeline; `--pack-cans` snapshots every iteration as a - shareable Soup Can with Registry lineage (`soup loop replay <id> --extract dir`). -- **Branches ↔ Registry** — `soup adapters branch <name> --attach-to-registry <id>` / - `--from-registry <id>` links training-env snapshots into the Registry lineage DAG. +- **Alternative model hubs for data** — `soup data push --hub modelscope|modelers` uploads a + local JSONL to ModelScope / Modelers, and `soup data forge --hub … --teacher owner/name` + pre-fetches the teacher from that hub. +- **Tokenizer-aware prompt pruning** — `soup prune-prompt --tokenizer <id-or-path>` finds the + shared *token* prefix and decodes only the remainder, so BPE multi-byte sequences never get + truncated mid-token the way char-slicing can. +- **Webhooks everywhere** — `--slack-url` / `--discord-url` now work on `soup ingest`, + `soup prune-prompt`, `soup ab`, and `soup data active-sample` (same SSRF-hardened validator as + `soup drift-alarm`). The A/B harness only pings when the sequential test actually decides. +- **Curriculum by difficulty percentile** — dynamic curriculum can bucket by `loss` / + `perplexity` percentile instead of length round-robin. +- **Smarter pre-flight `advise`** — `soup advise` now nudges its confidence using your prior + verdicts for the same project, and `soup runs replay` can plot a benchmark-score curve, not + just the loss curve. Full history: [CHANGELOG.md](CHANGELOG.md) · [GitHub Releases](https://github.com/MakazhanAlpamys/Soup/releases). diff --git a/docs/backends-and-ops.md b/docs/backends-and-ops.md index 3caf16d..d6d0d98 100644 --- a/docs/backends-and-ops.md +++ b/docs/backends-and-ops.md @@ -459,6 +459,11 @@ Every completed run also stores an estimated cost (`$` per run) computed from th captured GPU device name and duration. `soup runs show` renders `—` for CPU / MPS / unknown GPUs (no fabricated zeros). +As of v0.71.5, the metric-series lookup that powers replay (`ExperimentTracker.get_metric_series`) +transparently falls back to the `eval_results` table when a metric has no per-step +rows — so you can plot a benchmark-score curve (e.g. `mmlu`, `gsm8k`) the same way +you plot `loss`, without caring which table holds the series. + ### Tracker integrations (--tracker mlflow / swanlab / trackio) ```bash diff --git a/docs/commands.md b/docs/commands.md index 8d192b5..7d9a915 100644 --- a/docs/commands.md +++ b/docs/commands.md @@ -90,10 +90,12 @@ soup data download user/ds --samples 1000 Stream first 1000 samples soup data register --name my-ds --path d.jsonl --format alpaca Register dataset soup data unregister --name my-ds Remove from registry soup data push --input d.jsonl --hf-dataset user/name Upload local JSONL as HF dataset +soup data push --input d.jsonl --hf-dataset u/n --hub modelscope|modelers Upload to an alternative hub soup data registry List all registered datasets soup data demo List bundled demo JSONL fixtures soup data demo alpaca_demo --output ./d.jsonl Copy a bundled demo JSONL fixture soup data forge --docs ./docs --task sft --target-rows 1000 Synthetic data pipeline + provenance +soup data forge --docs ./docs --hub modelscope --teacher owner/name Pre-fetch the teacher from an alternative hub soup data score --input rows.jsonl Composite quality scorecard (PII + toxicity + lang + edu) soup data decontaminate --input rows.jsonl --benchmarks mmlu,gsm8k Drop benchmark-overlap rows soup data toxicity --input rows.jsonl -o tox.jsonl Flag toxic rows (keyword baseline) @@ -137,7 +139,7 @@ soup can publish r.can --hf-hub user/name Publish .can to HF Hub as dataset soup runs List training runs soup runs show <run_id> Run details + loss graph + cost soup runs compare <run_1> <run_2> Compare two runs -soup runs replay <run_id> Replay summary + loss curve from history +soup runs replay <run_id> Replay summary + loss curve from history (also plots a benchmark-score curve when the metric lives in eval_results) soup why [run_id] Explain training anomalies (heuristic) soup tui Full-screen Textual dashboard (requires [tui] extra) soup train --config soup.yaml --profile Record torch.profiler trace to <output>/profiles/ @@ -169,8 +171,10 @@ soup edit set --base <m> --method rome|memit|alphaedit --subject "..." --target soup edit diff <before-run> <after-run> --probes p.jsonl Knowledge-injection diff visualizer soup ingest --source langfuse|langsmith|helicone|openpipe|otel|openai-stored --logs <jsonl> Universal trace importer (6 SaaS adapters → normalised JSONL) soup prune-prompt --input <jsonl> --output <jsonl> --min-frequency 0.95 Detect + strip shared system-prompt prefix +soup prune-prompt ... --tokenizer <id-or-path> Tokenizer-aware prefix detection (decodes remaining ids, boundary-safe) soup data active-sample --input <jsonl> --output <jsonl> --budget N Top-N uncertain prod traces for human review soup ab --input <jsonl> --metric latency|judge_score|retry_rate mSPRT sequential A/B (decision: continue / reject_h0 / accept_h0) +soup ingest|prune-prompt|ab|data active-sample ... --slack-url <https> | --discord-url <https> Shared SSRF-validated webhook on completion soup drift-alarm --reference <jsonl> --live <jsonl> --threshold 0.2 Rolling-KL drift alarm (exit 3 on drift) soup drift-alarm ... --slack-url <https> | --discord-url <https> Optional SSRF-validated webhook on drift detected soup tunability --list List built-in candidate-base catalogue diff --git a/docs/data.md b/docs/data.md index 9ed61e5..8464cb0 100644 --- a/docs/data.md +++ b/docs/data.md @@ -94,6 +94,14 @@ soup prune-prompt --input traces.jsonl --output pruned.jsonl --min-frequency 0.9 Binary-search over up-to-32 candidate templates finds the longest qualifying prefix (a longer threshold-meeting prefix may exist beyond the universal one — Soup does not early-exit on the 100% match). Two-pass file read with a 100 000-row DoS cap. +**Tokenizer-aware mode (v0.71.5).** Pass `--tokenizer <id-or-path>` (a HuggingFace repo id, a local path, or anything `AutoTokenizer.from_pretrained` accepts) to detect the shared prefix in *token* space and decode only the remaining ids: + +```bash +soup prune-prompt --input traces.jsonl --output pruned.jsonl --tokenizer Qwen/Qwen2.5-0.5B +``` + +Char-level stripping can cut a BPE multi-byte sequence in half when the shared prefix ends mid-token; token-aware pruning finds the longest shared *token-id* prefix and decodes the remainder, so the boundary always lands on a real token. Per-row encoding is capped at 50 000 tokens. Omit `--tokenizer` to keep the original character-level behaviour. + ## Active-Learning Sampler (`soup data active-sample`) @@ -108,6 +116,8 @@ soup data active-sample --input traces.jsonl --output for-review.jsonl --budget The output JSONL is a drop-in prompt set for `soup eval human` (v0.19). Budget is bounded `[1, 100 000]`. +**Webhooks (v0.71.5).** `soup ingest`, `soup prune-prompt`, `soup ab`, and `soup data active-sample` all accept `--slack-url` / `--discord-url` and POST a one-line summary on completion through the same SSRF-hardened validator as `soup drift-alarm` (scheme allowlist, loopback-only HTTP, RFC1918 / link-local / reserved / multicast rejected; the post never raises, so a flaky webhook can't fail the command). `soup ab` only fires when the sequential test actually decides (`reject_h0` / `accept_h0`), not while it's still `continue`-ing. + ## Synthetic Data Generation @@ -479,6 +489,8 @@ Three tasks supported: `sft` (Q&A pairs), `preference` (chosen/rejected), `tool` Document discovery is one level deep over `.txt` / `.md` / `.json` / `.jsonl`; dotfiles + symlinked directories are skipped. All paths are cwd-contained, all writes are atomic via staged-tempfile + `os.replace`, and write targets are rejected if they're symlinks. **Judge providers are live**: `--judge-provider ollama` (localhost-only), `--judge-provider anthropic` (env-only API key), `--judge-provider vllm` (scheme-validated). Per-call judge exceptions logged at DEBUG. +**Alternative teacher hubs (v0.71.5).** `--hub modelscope|modelers` pre-fetches the `--teacher` from that hub when the teacher is a routable repo id (`owner/name`); `--hub hf` (default) is a no-op and leaves the teacher as a provenance label. If `--hub` is non-HF but `--teacher` is not a repo id (e.g. the default `local-judge`), Soup prints a loud yellow warning rather than silently dropping the flag. + ## Data Quality Scorecard diff --git a/docs/evaluation.md b/docs/evaluation.md index b215bfc..56bf17c 100644 --- a/docs/evaluation.md +++ b/docs/evaluation.md @@ -82,6 +82,8 @@ soup advise compare 4. Task is `factual_lookup` with high output variance → **RAG**. 5. Otherwise → **SFT**. +**Cross-project confidence bias (v0.71.5).** When `~/.soup/advise_history.jsonl` holds ≥3 prior verdicts for the *same choice* in the *same project*, `soup advise` nudges its confidence (not its decision) toward what worked before: a net-positive precedent record (you accepted it AND its recorded outcome was good) bumps confidence up by a small constant; a net-negative one bumps it down. The rubric verdict itself never changes — only how sure Soup is. Verdicts must be `--record`ed for the bias to kick in. + **Why this command exists.** "Choose fine-tuning vs RAG vs prompt-engineering" is the most-mis-made decision in the space. Reddit, HN, IBM, and Google Cloud all converge on the same advice (start with prompts, escalate to RAG, fine-tune as last resort) and almost everyone ignores it because nobody has the data to prove their case is the exception. Soup `autopilot` picks hyperparameters AFTER you've decided to train; `soup advise` owns the layer above. No trainer library has an incentive to tell users *not to train* — Unsloth's funnel, Axolotl's hosted business, LLaMA-Factory's Alibaba alignment all monetise the training event. @@ -213,6 +215,8 @@ soup ab --input ab.jsonl --metric judge_score --alpha 0.01 --beta 0.10 --effect- Input rows look like `{"arm": "control", "latency": 1.23}` or `{"arm": "treatment", "judge_score": 0.91}`. Decision is one of `continue` (keep collecting samples), `reject_h0` (real difference detected), `accept_h0` (no significant difference). Composes with `soup loop canary` (v0.58) — promote or roll back as soon as the LLR clears a decision boundary. +`soup ab` accepts `--slack-url` / `--discord-url` (v0.71.5) and pings the webhook **only when the test actually decides** (`reject_h0` / `accept_h0`) — a still-running `continue` stays quiet so you're not paged on every peek. Same SSRF-hardened validator as `soup drift-alarm`. + ## Drift Alarm (`soup drift-alarm`) diff --git a/docs/training.md b/docs/training.md index 708e549..325a51c 100644 --- a/docs/training.md +++ b/docs/training.md @@ -907,12 +907,15 @@ Layer dynamic re-weighting on top of the static `curriculum` bucketer. Every N s training: curriculum: true # static bucketer (v0.23.0) curriculum_buckets: 4 + curriculum_metric: perplexity # length (default) | loss | perplexity curriculum_dynamic: true # NEW — dynamic re-weighting curriculum_dynamic_recompute_steps: 50 # refresh every 50 global steps curriculum_dynamic_floor: 0.05 # min weight per bucket curriculum_dynamic_temperature: 1.0 # softmax temp on uncertainty ``` +**Bucketing by difficulty percentile (v0.71.5).** When `curriculum_metric` is `loss` or `perplexity`, the dynamic callback assigns each step's sample to a bucket by its *rank* within a rolling 512-step window of the difficulty signal (perplexity = `exp(min(loss, 50))`), instead of the round-robin fallback used for `length`. This keeps the buckets calibrated to the live loss distribution rather than a static length sort. `length` (the default) keeps the round-robin assignment. + Visualise the recorded bucket-weight evolution with `soup runs curriculum-curve <run_id>`. DDP / grad-accum safety: multi-rank launches must wire an `all_reduce` hook on per-bucket stats (a cross-validator rejects un-coordinated multi-rank runs upfront). Multi-trainer expansion beyond `sft` / `pretrain` is tracked for v0.48.1. diff --git a/pyproject.toml b/pyproject.toml index cc22aa0..f8c467f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "soup-cli" -version = "0.71.4" +version = "0.71.5" description = "Fine-tune LLMs in one command. No SSH, no config hell." readme = "README.md" license = "Apache-2.0" diff --git a/src/soup_cli/__init__.py b/src/soup_cli/__init__.py index f05ab25..f8832d7 100644 --- a/src/soup_cli/__init__.py +++ b/src/soup_cli/__init__.py @@ -1,3 +1,3 @@ """Soup CLI — Fine-tune LLMs in one command.""" -__version__ = "0.71.4" +__version__ = "0.71.5" diff --git a/src/soup_cli/commands/_webhook_cli.py b/src/soup_cli/commands/_webhook_cli.py new file mode 100644 index 0000000..362a96d --- /dev/null +++ b/src/soup_cli/commands/_webhook_cli.py @@ -0,0 +1,71 @@ +"""Shared CLI glue for the --slack-url / --discord-url webhook flags (v0.71.5 #207). + +Keeps Typer + Rich Console concerns in the commands layer (``utils/webhooks`` +stays import-light + framework-free). Used by ``ingest`` / ``prune-prompt`` / +``ab`` / ``data active-sample`` so the validate-then-deliver pattern is defined +once. +""" + +from __future__ import annotations + +from typing import Optional, Tuple + +import typer +from rich.console import Console +from rich.markup import escape + + +def validate_webhook_flags( + slack_url: Optional[str], + discord_url: Optional[str], + *, + console: Console, +) -> Tuple[Optional[str], Optional[str]]: + """Validate webhook URLs at the CLI boundary (``typer.Exit(2)`` on bad). + + Returns the (canonical) URLs. Mirrors the v0.63.0 ``drift-alarm`` + early-rejection pattern so a typo'd / SSRF-y URL fails fast with a + friendly message instead of being silently swallowed at delivery time. + """ + from soup_cli.utils.webhooks import validate_webhook_url + + out = [] + for label, value in (("--slack-url", slack_url), ("--discord-url", discord_url)): + if value is None: + out.append(None) + continue + try: + out.append(validate_webhook_url(value)) + except (TypeError, ValueError) as exc: + console.print(f"[red]{label}: {escape(str(exc))}[/]") + raise typer.Exit(2) from exc + return out[0], out[1] + + +def emit_webhooks( + slack_url: Optional[str], + discord_url: Optional[str], + *, + payload: dict, + console: Console, +) -> None: + """POST the completion payload to any configured webhooks (best-effort). + + Never raises (delegates to the never-raising + :func:`soup_cli.utils.webhooks.send_webhooks`). Prints a per-target + delivered/failed line so the operator sees whether the alert landed. + """ + if slack_url is None and discord_url is None: + return + from soup_cli.utils.webhooks import send_webhooks + + for label, ok in send_webhooks( + payload, slack_url=slack_url, discord_url=discord_url + ): + colour = "green" if ok else "yellow" + console.print( + f"[{colour}]{label} webhook: {'delivered' if ok else 'failed'}[/]" + ) + + +__all__ = ["emit_webhooks", "validate_webhook_flags"] diff --git a/src/soup_cli/commands/ab.py b/src/soup_cli/commands/ab.py index 8b407d1..9d2cd6d 100644 --- a/src/soup_cli/commands/ab.py +++ b/src/soup_cli/commands/ab.py @@ -2,12 +2,15 @@ from __future__ import annotations +from typing import Optional + import typer from rich.console import Console from rich.markup import escape from rich.panel import Panel from rich.table import Table +from soup_cli.commands._webhook_cli import emit_webhooks, validate_webhook_flags from soup_cli.utils.ab_test import ( MsprtConfig, run_msprt, @@ -35,6 +38,20 @@ def ab( 0.1, "--effect-size", help="Minimum detectable difference in means.", ), + slack_url: Optional[str] = typer.Option( + None, "--slack-url", + help=( + "Optional Slack webhook URL — POSTed on a reject_h0 / accept_h0 " + "decision (not on continue). SSRF-validated." + ), + ), + discord_url: Optional[str] = typer.Option( + None, "--discord-url", + help=( + "Optional Discord webhook URL — POSTed on a reject_h0 / accept_h0 " + "decision (not on continue). SSRF-validated." + ), + ), ) -> None: """Sequential A/B test with early-stop guarantees (mSPRT).""" try: @@ -43,6 +60,10 @@ def ab( console.print(f"[red]{escape(str(exc))}[/]") raise typer.Exit(2) from exc + slack_url, discord_url = validate_webhook_flags( + slack_url, discord_url, console=console + ) + try: cfg = MsprtConfig( metric=canonical, alpha=alpha, beta=beta, effect_size=effect_size, @@ -101,5 +122,24 @@ def ab( ) ) + # Webhook only fires on a terminal decision (reject_h0 / accept_h0) — + # a `continue` verdict carries no actionable signal (issue #207). + if verdict.decision != "continue": + emit_webhooks( + slack_url, + discord_url, + payload={ + "command": "ab", + "metric": canonical, + "decision": verdict.decision, + "log_likelihood_ratio": verdict.log_likelihood_ratio, + "n_control": verdict.n_control, + "n_treatment": verdict.n_treatment, + "mean_control": verdict.mean_control, + "mean_treatment": verdict.mean_treatment, + }, + console=console, + ) + __all__ = ["ab"] diff --git a/src/soup_cli/commands/active_sample.py b/src/soup_cli/commands/active_sample.py index 62e0397..f818e6a 100644 --- a/src/soup_cli/commands/active_sample.py +++ b/src/soup_cli/commands/active_sample.py @@ -2,11 +2,14 @@ from __future__ import annotations +from typing import Optional + import typer from rich.console import Console from rich.markup import escape from rich.panel import Panel +from soup_cli.commands._webhook_cli import emit_webhooks, validate_webhook_flags from soup_cli.utils.active_sampler import sample_uncertain_rows, validate_budget console = Console() @@ -24,6 +27,14 @@ def active_sample( 100, "--budget", help="Max rows to surface for human review (1 - 100_000).", ), + slack_url: Optional[str] = typer.Option( + None, "--slack-url", + help="Optional Slack webhook URL — POSTed on completion. SSRF-validated.", + ), + discord_url: Optional[str] = typer.Option( + None, "--discord-url", + help="Optional Discord webhook URL — POSTed on completion. SSRF-validated.", + ), ) -> None: """Surface the most uncertain prod traces for human review.""" try: @@ -32,6 +43,10 @@ def active_sample( console.print(f"[red]{escape(str(exc))}[/]") raise typer.Exit(2) from exc + slack_url, discord_url = validate_webhook_flags( + slack_url, discord_url, console=console + ) + try: plan = sample_uncertain_rows( input_path, @@ -55,5 +70,18 @@ def active_sample( ) ) + emit_webhooks( + slack_url, + discord_url, + payload={ + "command": "active-sample", + "rows_selected": plan.rows_selected, + "rows_in": plan.rows_in, + "mean_uncertainty": plan.mean_uncertainty, + "budget": plan.budget, + }, + console=console, + ) + __all__ = ["active_sample"] diff --git a/src/soup_cli/commands/advise.py b/src/soup_cli/commands/advise.py index 1e33511..412d589 100644 --- a/src/soup_cli/commands/advise.py +++ b/src/soup_cli/commands/advise.py @@ -39,6 +39,7 @@ from soup_cli.utils.advise import ( synth_probe_lora_delta, ) from soup_cli.utils.advise_history import ( + current_project_name, history_path, load_history, record_verdict, @@ -269,8 +270,27 @@ def advise_run( sft_wall_clock_secs=wall_clock, ) + # v0.71.5 #163 — bias the rubric by this project's past accepted-verdict + # outcomes. Best-effort: a missing / unreadable history must never block + # a verdict, so any failure falls back to the un-biased rubric. + history = None + project = None try: - verdict = build_verdict(profile, task_category, goal=goal, roi=roi) + history = load_history(limit=20) + project = current_project_name() + except (TypeError, ValueError, OSError): + history = None + project = None + + try: + verdict = build_verdict( + profile, + task_category, + goal=goal, + roi=roi, + history=history, + project=project, + ) except (TypeError, ValueError) as exc: console.print(f"[red]Verdict build failed:[/] {escape(str(exc))}") raise typer.Exit(1) from exc diff --git a/src/soup_cli/commands/data.py b/src/soup_cli/commands/data.py index 92e6233..e905829 100644 --- a/src/soup_cli/commands/data.py +++ b/src/soup_cli/commands/data.py @@ -1914,16 +1914,32 @@ def push_dataset_cmd( "--message", help="Commit message for the dataset upload", ), + hub: str = typer.Option( + "hf", + "--hub", + help=( + "Target hub: hf (default) / modelscope / modelers. Non-HF hubs " + "upload via the matching SDK (v0.71.5 #157) — repo_type=dataset, " + "commit message sanitised to first line + 200 chars." + ), + ), ): - """Upload a local JSONL dataset to HuggingFace Hub as a dataset repo.""" + """Upload a local JSONL dataset to a model hub as a dataset repo.""" from soup_cli.utils.hf import ( get_hf_api, resolve_endpoint, resolve_token, validate_repo_id, ) + from soup_cli.utils.hubs import validate_hub_name from soup_cli.utils.paths import is_under_cwd + try: + hub_canonical = validate_hub_name(hub) + except (TypeError, ValueError) as exc: + console.print(f"[red]{exc}[/]") + raise typer.Exit(2) from exc + file_path = Path(input_path) if not file_path.exists(): console.print(f"[red]Dataset file not found: {file_path}[/]") @@ -1940,9 +1956,18 @@ def push_dataset_cmd( try: validate_repo_id(hf_dataset) except ValueError as exc: - console.print(f"[red]Invalid --hf-dataset repo id:[/] {exc}") + console.print(f"[red]Invalid --{hub_canonical}-dataset repo id:[/] {exc}") raise typer.Exit(1) from exc + # v0.71.5 #157 — non-HF hubs route through the shared upload_repo adapter. + # ModelScope / Modelers SDKs upload a folder, so the single JSONL is + # staged into a temp dir and uploaded as a dataset repo. + if hub_canonical != "hf": + _push_dataset_non_hf( + hub_canonical, hf_dataset, file_path, commit_message + ) + return + token = resolve_token() if token is None: console.print( @@ -1987,6 +2012,53 @@ def push_dataset_cmd( ) +def _push_dataset_non_hf( + hub: str, repo_id: str, file_path: Path, commit_message: str +) -> None: + """Upload a single JSONL to a non-HF hub as a dataset (v0.71.5 #157). + + ``upload_repo`` uploads a folder, so the file is staged into a temp dir + first. ``upload_repo`` already sanitises the commit message (first line + + 200 chars) and validates the repo id shape. + """ + import shutil + import tempfile + + from rich.markup import escape + + from soup_cli.utils.hubs import upload_repo + + # Stage UNDER cwd — `upload_repo` enforces cwd-containment on folder_path + # (the system tempdir would be rejected). + staging = tempfile.mkdtemp(prefix=".soup_dataset_push.", dir=os.getcwd()) + try: + shutil.copy2(str(file_path), os.path.join(staging, file_path.name)) + try: + upload_repo( + hub, + repo_id, + folder_path=staging, + commit_message=commit_message, + repo_type="dataset", + ) + except ImportError as exc: + console.print(f"[red]{exc}[/]") + raise typer.Exit(1) from exc + except (TypeError, ValueError) as exc: + console.print(f"[red]Upload failed:[/] {exc}") + raise typer.Exit(1) from exc + except Exception as exc: # noqa: BLE001 — surface SDK errors generically + console.print(f"[red]Upload failed:[/] {exc}") + raise typer.Exit(1) from exc + finally: + shutil.rmtree(staging, ignore_errors=True) + + console.print( + f"[green]Uploaded[/] {escape(file_path.name)} to {escape(hub)} " + f"dataset [bold]{escape(repo_id)}[/]" + ) + + # --- v0.42.0 Part C / F: AOT preprocess + document ingestion --------------- @app.command(name="preprocess") diff --git a/src/soup_cli/commands/data_forge.py b/src/soup_cli/commands/data_forge.py index 4b21310..797cdba 100644 --- a/src/soup_cli/commands/data_forge.py +++ b/src/soup_cli/commands/data_forge.py @@ -77,6 +77,14 @@ def forge( "(scheme allowlist + loopback). Ignored for Anthropic." ), ), + hub: str = typer.Option( + "hf", "--hub", + help=( + "Teacher hub: hf (default) / modelscope / modelers. When non-HF " + "and --teacher is a repo id (owner/name), the teacher is " + "pre-fetched from that hub (v0.71.5 #157)." + ), + ), ): """Run the multi-stage synthetic data pipeline with provenance. @@ -94,9 +102,45 @@ def forge( write_forge_dataset, write_provenance, ) + from soup_cli.utils.hubs import validate_hub_name + + try: + hub_canonical = validate_hub_name(hub) + except (TypeError, ValueError) as exc: + console.print(f"[red]{escape(str(exc))}[/]") + raise typer.Exit(2) from exc + + effective_teacher = teacher + # v0.71.5 #157 — non-HF hub + repo-id teacher → pre-fetch the teacher from + # that hub and record the resolved local path in provenance. HF (default) + # is a no-op: the teacher stays a provenance label. + if hub_canonical != "hf": + if "/" in teacher: + from soup_cli.utils.hubs import prefetch_model_from_hub + + try: + effective_teacher = prefetch_model_from_hub( + teacher, hub_canonical, console=console + ) + except ImportError as exc: + console.print(f"[red]{escape(str(exc))}[/]") + raise typer.Exit(1) from exc + except (TypeError, ValueError) as exc: + console.print( + f"[red]Teacher pre-fetch failed:[/] {escape(str(exc))}" + ) + raise typer.Exit(1) from exc + else: + # Non-HF hub requested but the teacher is not a routable repo id + # (owner/name) — warn loudly instead of silently ignoring --hub + # (code-review MEDIUM fix v0.71.5 #157). + console.print( + f"[yellow]--hub {escape(hub_canonical)} ignored:[/] --teacher " + f"{escape(teacher)!r} is not a repo id (owner/name), so there " + "is nothing to pre-fetch." + ) judge_fn = _default_judge - effective_teacher = teacher if judge_provider is not None: canonical = judge_provider.strip().lower() if canonical not in JUDGE_PROVIDERS: diff --git a/src/soup_cli/commands/ingest.py b/src/soup_cli/commands/ingest.py index a9f3893..525edb3 100644 --- a/src/soup_cli/commands/ingest.py +++ b/src/soup_cli/commands/ingest.py @@ -21,6 +21,7 @@ from rich.console import Console from rich.markup import escape from rich.panel import Panel +from soup_cli.commands._webhook_cli import emit_webhooks, validate_webhook_flags from soup_cli.utils import ingest_sources as _ingest_sources from soup_cli.utils.ingest_sources import ( SUPPORTED_INGEST_SOURCES, @@ -53,6 +54,14 @@ def ingest( "-o", help="Output JSONL (default: traces.jsonl in cwd).", ), + slack_url: Optional[str] = typer.Option( + None, "--slack-url", + help="Optional Slack webhook URL — POSTed on completion. SSRF-validated.", + ), + discord_url: Optional[str] = typer.Option( + None, "--discord-url", + help="Optional Discord webhook URL — POSTed on completion. SSRF-validated.", + ), ) -> None: """Import production traces from a SaaS observability vendor (v0.63.0). @@ -66,6 +75,10 @@ def ingest( console.print(f"[red]{escape(str(exc))}[/]") raise typer.Exit(2) from exc + slack_url, discord_url = validate_webhook_flags( + slack_url, discord_url, console=console + ) + if not is_under_cwd(logs): console.print(f"[red]--logs '{escape(logs)}' is outside cwd — refusing[/]") raise typer.Exit(1) @@ -111,6 +124,18 @@ def ingest( f"{escape(output_path.name)}[/]" ) + emit_webhooks( + slack_url, + discord_url, + payload={ + "command": "ingest", + "source": canonical, + "traces_written": count, + "auth_env_set": auth_value is not None, + }, + console=console, + ) + def _env_label(source: str) -> str: """Return the env-var name that authenticates ``source``. diff --git a/src/soup_cli/commands/prune_prompt.py b/src/soup_cli/commands/prune_prompt.py index 4d4865e..7ed640d 100644 --- a/src/soup_cli/commands/prune_prompt.py +++ b/src/soup_cli/commands/prune_prompt.py @@ -7,11 +7,14 @@ signature trick — shipped OSS for v0.63.0 Part B). from __future__ import annotations +from typing import Optional + import typer from rich.console import Console from rich.markup import escape from rich.panel import Panel +from soup_cli.commands._webhook_cli import emit_webhooks, validate_webhook_flags from soup_cli.utils.prune_prompt import prune_traces, validate_min_frequency console = Console() @@ -35,6 +38,22 @@ def prune_prompt_cmd( "--min-frequency", help="Prefix must appear in >= this fraction of rows to be stripped (0.0 - 1.0).", ), + tokenizer: Optional[str] = typer.Option( + None, + "--tokenizer", + help=( + "Optional HF tokenizer (model id or local path) for token-aware " + "prefix detection. Default: whitespace-character level." + ), + ), + slack_url: Optional[str] = typer.Option( + None, "--slack-url", + help="Optional Slack webhook URL — POSTed on completion. SSRF-validated.", + ), + discord_url: Optional[str] = typer.Option( + None, "--discord-url", + help="Optional Discord webhook URL — POSTed on completion. SSRF-validated.", + ), ) -> None: """Detect + strip a shared system-prompt prefix (v0.63.0 Part B).""" try: @@ -43,11 +62,16 @@ def prune_prompt_cmd( console.print(f"[red]{escape(str(exc))}[/]") raise typer.Exit(2) from exc + slack_url, discord_url = validate_webhook_flags( + slack_url, discord_url, console=console + ) + try: report = prune_traces( input_path, output_path=output_path, min_frequency=min_frequency, + tokenizer=tokenizer, ) except FileNotFoundError: console.print(f"[red]Input not found: {escape(input_path)}[/]") @@ -56,6 +80,14 @@ def prune_prompt_cmd( console.print(f"[red]{escape(str(exc))}[/]") raise typer.Exit(1) from exc + payload = { + "command": "prune-prompt", + "prefix_found": bool(report.prefix), + "prefix_chars": report.prefix_chars, + "rows_pruned": report.rows_pruned, + "rows_total": report.rows_total, + } + if not report.prefix: console.print( Panel( @@ -65,6 +97,7 @@ def prune_prompt_cmd( border_style="yellow", ) ) + emit_webhooks(slack_url, discord_url, payload=payload, console=console) return snippet = report.prefix if len(report.prefix) <= 200 else report.prefix[:200] + "..." @@ -77,6 +110,7 @@ def prune_prompt_cmd( border_style="green", ) ) + emit_webhooks(slack_url, discord_url, payload=payload, console=console) __all__ = ["prune_prompt_cmd"] diff --git a/src/soup_cli/experiment/tracker.py b/src/soup_cli/experiment/tracker.py index c7fb88a..ba7e774 100644 --- a/src/soup_cli/experiment/tracker.py +++ b/src/soup_cli/experiment/tracker.py @@ -318,6 +318,15 @@ class ExperimentTracker: Used by ``soup eval against`` for run-vs-run paired-bootstrap CI. Returns an empty list when the metric does not appear in any row — the caller treats that as "no signal, do not gate". + + v0.71.5 #164: the per-step ``metrics`` table only carries training + columns (``loss`` / ``lr`` / ``grad_norm`` / ``speed`` / ``gpu_mem``). + Eval metrics like ``task_accuracy`` / ``refusal_rate`` live in the + ``eval_results`` table instead. So when the per-step pass yields no + rows we fall back to the per-benchmark scores in ``eval_results``. + Querying ``metrics`` first preserves the established behaviour for + every training-loop column (no regression for existing callers); + the fallback only fires when the column path is empty. """ if not isinstance(run_id, str) or not run_id: raise ValueError("run_id must be a non-empty string") @@ -335,7 +344,34 @@ class ExperimentTracker: # Skip non-numeric cells silently — same-run inconsistency # is not the caller's problem; they get a shorter series. continue - return series + if series: + return series + # Bridge to eval_results (v0.71.5 #164) — benchmark scores for + # `soup eval against`. Empty when neither table has data. + return self._eval_score_series(run_id, metric) + + def _eval_score_series(self, run_id: str, benchmark: str) -> list[float]: + """Return the per-row ``score`` series from ``eval_results``. + + Ordered by insertion (``id``) for deterministic pairing in the + paired-bootstrap CI. Non-numeric cells are skipped silently. + """ + conn = self._get_conn() + rows = conn.execute( + "SELECT score FROM eval_results " + "WHERE run_id = ? AND benchmark = ? ORDER BY id", + (run_id, benchmark), + ).fetchall() + out: list[float] = [] + for row in rows: + value = row["score"] + if value is None: + continue + try: + out.append(float(value)) + except (TypeError, ValueError): + continue + return out def save_eval_result( self, diff --git a/src/soup_cli/monitoring/curriculum_callback.py b/src/soup_cli/monitoring/curriculum_callback.py index 7b0e136..da94b44 100644 --- a/src/soup_cli/monitoring/curriculum_callback.py +++ b/src/soup_cli/monitoring/curriculum_callback.py @@ -31,14 +31,17 @@ from __future__ import annotations import json import logging +import math import os import stat import tempfile -from typing import Any, Dict, Optional, Tuple +from collections import deque +from typing import Any, Deque, Dict, Optional, Tuple from soup_cli.utils.curriculum_dynamic import ( DynamicCurriculumPolicy, compute_bucket_weights, + percentile_bucket, ) from soup_cli.utils.paths import is_under_cwd @@ -46,6 +49,14 @@ logger = logging.getLogger(__name__) _MAX_PATH_LEN = 4096 _HISTORY_FILENAME = "curriculum_history.jsonl" +# Curriculum metrics that drive percentile bucketing from the live loss +# signal (v0.71.5 #149). ``length`` has no per-step signal in HF logs, so it +# falls back to round-robin (length bucketing is the static curriculum's +# data-prep-time job — see utils/curriculum.py). +_PERCENTILE_METRICS = frozenset({"loss", "perplexity"}) +_VALID_CURRICULUM_METRICS = frozenset({"length", "perplexity", "loss"}) +# Rolling-window size for the percentile reference distribution. +_SIGNAL_WINDOW = 512 __all__ = [ "DynamicCurriculumCallback", @@ -149,14 +160,25 @@ class DynamicCurriculumCallback(_try_import_callback_base()): # type: ignore[mi self, policy: DynamicCurriculumPolicy, output_dir: str, + curriculum_metric: str = "length", ) -> None: if not isinstance(policy, DynamicCurriculumPolicy): raise TypeError( "policy must be DynamicCurriculumPolicy, got " f"{type(policy).__name__}" ) + if curriculum_metric not in _VALID_CURRICULUM_METRICS: + raise ValueError( + "curriculum_metric must be one of " + f"{sorted(_VALID_CURRICULUM_METRICS)}, got {curriculum_metric!r}" + ) self._policy = policy + self._curriculum_metric = curriculum_metric self._output_dir = _validate_output_dir(output_dir) + # Rolling difficulty-signal window for percentile bucketing (v0.71.5 + # #149). Persists across recomputes so a consistently-hard sample + # keeps landing in the same bucket. + self._signal_window: Deque[float] = deque(maxlen=_SIGNAL_WINDOW) # Per-bucket accumulator: bucket_id -> {num_samples, loss_sum, grad_norm_sum} self._stats: Dict[int, Dict[str, float]] = {} # Most recently computed weights (read by external sampler hook). @@ -173,6 +195,10 @@ class DynamicCurriculumCallback(_try_import_callback_base()): # type: ignore[mi def policy(self) -> DynamicCurriculumPolicy: return self._policy + @property + def curriculum_metric(self) -> str: + return self._curriculum_metric + @property def output_dir(self) -> str: return self._output_dir @@ -215,13 +241,55 @@ class DynamicCurriculumCallback(_try_import_callback_base()): # type: ignore[mi # HF emits "loss" + (optionally) "grad_norm" in `logs`. loss = logs.get("loss") grad_norm = logs.get("grad_norm") - # Bucket id derived from the global step (round-robin BETA strategy). - try: - bucket_id = _pick_bucket(global_step, self._policy.num_buckets) - except (TypeError, ValueError): - return + nb = self._policy.num_buckets + # v0.71.5 #149: percentile bucketing on the difficulty signal for + # loss / perplexity once the rolling window has warmed up; otherwise + # round-robin (warm-up + the `length` metric, which has no per-step + # signal in HF logs). + signal = self._difficulty_signal(loss) + bucket_id: Optional[int] = None + if ( + self._curriculum_metric in _PERCENTILE_METRICS + and signal is not None + and len(self._signal_window) > 0 + ): + try: + bucket_id = percentile_bucket( + signal, list(self._signal_window), nb + ) + except (TypeError, ValueError): + bucket_id = None + if bucket_id is None: + try: + bucket_id = _pick_bucket(global_step, nb) + except (TypeError, ValueError): + return + if signal is not None: + self._signal_window.append(signal) self._record_sample(bucket_id, loss, grad_norm) + def _difficulty_signal(self, loss: object) -> Optional[float]: + """Map the logged loss to the configured difficulty signal. + + Returns ``None`` when no usable signal is available (metric is + ``length`` — no per-step length in HF logs — or the loss is + missing / non-finite), which routes the step through the + round-robin fallback. + """ + if self._curriculum_metric not in _PERCENTILE_METRICS: + return None + try: + loss_f = float(loss) if loss is not None else None + except (TypeError, ValueError): + return None + if loss_f is None or not math.isfinite(loss_f): + return None + if self._curriculum_metric == "perplexity": + # exp is monotonic in loss, so percentile ranks are identical; + # clamp the exponent to avoid overflow on a stray loss spike. + return math.exp(min(loss_f, 50.0)) + return loss_f + def on_step_end( self, args: Any, diff --git a/src/soup_cli/utils/advise.py b/src/soup_cli/utils/advise.py index 233db93..15308c0 100644 --- a/src/soup_cli/utils/advise.py +++ b/src/soup_cli/utils/advise.py @@ -26,7 +26,10 @@ import re import stat from collections.abc import Iterable from dataclasses import dataclass, field -from typing import Dict, List, Mapping, Optional, Sequence, Tuple +from typing import TYPE_CHECKING, Dict, List, Mapping, Optional, Sequence, Tuple + +if TYPE_CHECKING: # pragma: no cover — annotation only, avoids circular import. + from soup_cli.utils.advise_history import HistoryEntry from soup_cli.utils.paths import is_under_cwd @@ -59,6 +62,17 @@ _MIN_ROWS_FOR_TRAINING = 50 # routed into a GPU-intensive RL run). _MIN_ROWS_FOR_GRPO = 500 +# History-bias thresholds (v0.71.5 #163). A choice is "encouraged" when the +# project has >= _HISTORY_MIN_PRECEDENTS accepted verdicts whose mean measured +# outcome is >= _HISTORY_POSITIVE_OUTCOME; "discouraged" when the mean is +# < _HISTORY_NEGATIVE_OUTCOME over the same minimum count. Encouraged choices +# get a small confidence nudge and can flip a marginal tie; discouraged choices +# are suppressed (their effective row-floor is raised). +_HISTORY_MIN_PRECEDENTS = 3 +_HISTORY_POSITIVE_OUTCOME = 0.3 +_HISTORY_NEGATIVE_OUTCOME = 0.0 +_HISTORY_CONFIDENCE_NUDGE = 0.05 + # Probe defaults — tiny, no GPU required for the heuristic stubs. _PROBE_HOLDOUT_DEFAULT = 100 _PROBE_LORA_STEPS_DEFAULT = 100 @@ -460,7 +474,7 @@ def _confidence_from_signals(*, row_count: int, diversity: float) -> float: return max(0.2, min(0.95, 0.4 + 0.4 * size_score + 0.2 * diversity)) -def build_verdict( +def _base_verdict( profile: DatasetProfile, task_category: str, *, @@ -584,6 +598,184 @@ def build_verdict( ) +# --------------------------------------------------------------------------- +# History bias (v0.71.5 #163) — tune the rubric from past project outcomes +# --------------------------------------------------------------------------- + +def _summarise_history_outcomes( + history: Optional[Sequence["HistoryEntry"]], + *, + project: Optional[str] = None, +) -> Dict[str, Tuple[float, int]]: + """Aggregate accepted-verdict outcomes per choice from prior history. + + Duck-typed (reads ``.choice`` / ``.accepted`` / ``.outcome`` / ``.project`` + attributes) so :mod:`advise` never imports :mod:`advise_history` at runtime + — that import direction is owned by ``advise_history`` and reversing it + would create a cycle. + + Filtering: + - only entries whose ``accepted is True`` (a rejected verdict carries no + endorsement signal), + - only entries with a finite ``outcome`` in ``[-1, 1]`` (bool / None / + out-of-range skipped — defensive even though ``HistoryEntry`` validates + on construction), + - only entries matching ``project`` when supplied (per-project scoping: + one project's SFT wins must not bias another project's verdict). + + Returns ``{choice: (mean_outcome, count)}`` for every choice with >= 1 + qualifying entry. + """ + if history is None: + return {} + if isinstance(history, (str, bytes)) or not isinstance(history, Sequence): + raise TypeError("history must be a non-string Sequence or None") + acc: Dict[str, Tuple[float, int]] = {} + for entry in history: + choice = getattr(entry, "choice", None) + if choice not in CHOICES: + continue + if getattr(entry, "accepted", None) is not True: + continue + if project is not None and getattr(entry, "project", None) != project: + continue + outcome = getattr(entry, "outcome", None) + if outcome is None or isinstance(outcome, bool): + continue + if not isinstance(outcome, (int, float)): + continue + f_out = float(outcome) + if not math.isfinite(f_out) or not (-1.0 <= f_out <= 1.0): + continue + total, count = acc.get(choice, (0.0, 0)) + acc[choice] = (total + f_out, count + 1) + return {ch: (total / count, count) for ch, (total, count) in acc.items()} + + +def _is_encouraged(bias: Mapping[str, Tuple[float, int]], choice: str) -> bool: + mean_count = bias.get(choice) + if mean_count is None: + return False + mean, count = mean_count + return count >= _HISTORY_MIN_PRECEDENTS and mean >= _HISTORY_POSITIVE_OUTCOME + + +def _is_discouraged(bias: Mapping[str, Tuple[float, int]], choice: str) -> bool: + mean_count = bias.get(choice) + if mean_count is None: + return False + mean, count = mean_count + return count >= _HISTORY_MIN_PRECEDENTS and mean < _HISTORY_NEGATIVE_OUTCOME + + +def _precedent_count(bias: Mapping[str, Tuple[float, int]], choice: str) -> int: + mean_count = bias.get(choice) + return mean_count[1] if mean_count is not None else 0 + + +def build_verdict( + profile: DatasetProfile, + task_category: str, + *, + goal: Optional[str] = None, + roi: Optional[ROIEstimate] = None, + history: Optional[Sequence["HistoryEntry"]] = None, + project: Optional[str] = None, +) -> Verdict: + """Combine profile + task into a recommendation, optionally history-biased. + + Without ``history`` this is byte-identical to the v0.54.0 rubric + (regression-guarded). When ``history`` is supplied the per-project + outcome record nudges marginal decisions (v0.71.5 #163): + + - A choice with >= 3 accepted verdicts averaging >= +0.3 outcome is + "encouraged": its confidence is nudged up, and it can flip a marginal + RAG-vs-SFT tie toward SFT. + - A choice with >= 3 accepted verdicts averaging < 0.0 outcome is + "discouraged": it is suppressed (e.g. a project that keeps regressing on + GRPO falls back to SFT-on-traces). + + The confidence FLOOR is unchanged; only the tie-break + per-choice + suppression shift. Per-project scoping is enforced in + :func:`_summarise_history_outcomes`. + """ + base = _base_verdict(profile, task_category, goal=goal, roi=roi) + if history is None: + return base + bias = _summarise_history_outcomes(history, project=project) + if not bias: + return base + return _apply_history_bias(base, bias) + + +def _apply_history_bias( + base: Verdict, + bias: Mapping[str, Tuple[float, int]], +) -> Verdict: + """Adjust a base verdict using per-project history outcomes.""" + roi = base.estimated_roi + task_category = base.task_category + + # Marginal RAG → SFT flip: strong SFT track record, no comparable RAG + # track record. RAG is the marginal call (it fired on a heuristic + # variance threshold), so prior SFT success is decisive. + if ( + base.choice == "RAG" + and _is_encouraged(bias, "SFT") + and not _is_encouraged(bias, "RAG") + ): + n_sft = _precedent_count(bias, "SFT") + return Verdict( + choice="SFT", + confidence=min(0.95, base.confidence + _HISTORY_CONFIDENCE_NUDGE), + reason=( + f"Base rubric leaned RAG, but {n_sft} prior SFT verdicts in " + "this project averaged a positive outcome (precedent) — " + "routing to SFT over RAG." + ), + reverse_when=( + "the answer space is small and stable and RAG's freshness " + "outweighs the historical SFT lift — re-run after measuring." + ), + task_category=task_category, + estimated_roi=roi, + ) + + # Discouraged GRPO: a project that keeps regressing on RL falls back to + # SFT-on-traces (raises GRPO's effective floor for this project). + if base.choice == "GRPO" and _is_discouraged(bias, "GRPO"): + n_grpo = _precedent_count(bias, "GRPO") + return Verdict( + choice="SFT", + confidence=base.confidence, + reason=( + f"Base rubric leaned GRPO, but {n_grpo} prior GRPO verdicts in " + "this project averaged a negative outcome (precedent) — " + "falling back to SFT on the reasoning traces." + ), + reverse_when=( + "a reliable programmatic reward is now available and the " + "earlier GRPO regressions were reward-shaping bugs, not a " + "fundamental mismatch." + ), + task_category=task_category, + estimated_roi=roi, + ) + + # No flip — nudge confidence when the chosen path has a positive track + # record (DPO keeps its own +0.1; we only ever raise, never lower). + if _is_encouraged(bias, base.choice): + return Verdict( + choice=base.choice, + confidence=min(0.95, base.confidence + _HISTORY_CONFIDENCE_NUDGE), + reason=base.reason, + reverse_when=base.reverse_when, + task_category=task_category, + estimated_roi=roi, + ) + return base + + # --------------------------------------------------------------------------- # Probe runner (Part B) — heuristic stubs; real model loading is opt-in # --------------------------------------------------------------------------- diff --git a/src/soup_cli/utils/advise_history.py b/src/soup_cli/utils/advise_history.py index eb9021b..105f8b3 100644 --- a/src/soup_cli/utils/advise_history.py +++ b/src/soup_cli/utils/advise_history.py @@ -100,6 +100,16 @@ def _project_name() -> str: return name[:128] +def current_project_name() -> str: + """Public accessor for the current project label (v0.71.5 #163). + + The CLI passes this to ``build_verdict(..., project=...)`` so the history + bias is scoped to the project the verdict was recorded under. Returns the + same string :func:`record_verdict` stamps into the ``project`` field. + """ + return _project_name() + + def record_verdict( verdict: Verdict, *, diff --git a/src/soup_cli/utils/curriculum_dynamic.py b/src/soup_cli/utils/curriculum_dynamic.py index dca77c2..772f5ec 100644 --- a/src/soup_cli/utils/curriculum_dynamic.py +++ b/src/soup_cli/utils/curriculum_dynamic.py @@ -36,6 +36,7 @@ __all__ = [ "DynamicCurriculumPolicy", "BucketStats", "compute_bucket_weights", + "percentile_bucket", "validate_distributed_curriculum", ] @@ -179,6 +180,48 @@ def _softmax(values: Sequence[float], temperature: float) -> List[float]: return [e / total for e in exps] +def percentile_bucket( + value: float, + window: Sequence[float], + num_buckets: int, +) -> int: + """Bucket ``value`` by its percentile rank within a rolling ``window``. + + v0.71.5 #149 — replaces step-mod round-robin with a difficulty-signal + bucketing for ``curriculum_metric in {loss, perplexity}``. A value at or + above every window member lands in the top (hardest) bucket; a value + below every member lands in bucket 0. Because the bucket is a function of + the value's rank (not the step), a consistently-high-loss sample is + routed to the same bucket on every recompute. + + Args: + value: The current sample's difficulty signal (e.g. loss). + window: Recent difficulty signals (rolling reference distribution). + An empty / ``None`` window returns bucket 0 (warm-up — the caller + should fall back to round-robin until the window fills). + num_buckets: Number of difficulty buckets. + + Returns: + Bucket id in ``[0, num_buckets - 1]``. + """ + nb = _reject_bool_int("num_buckets", num_buckets) + if nb < _MIN_BUCKETS or nb > _MAX_BUCKETS: + raise ValueError( + f"num_buckets must be in [{_MIN_BUCKETS}, {_MAX_BUCKETS}], got {nb}" + ) + fv = _reject_bool_float("value", value) + if nb == 1: + return 0 + if not window: + return 0 + le = sum( + 1 for w in window if _reject_bool_float("window value", w) <= fv + ) + rank = le / len(window) + bucket = int(rank * nb) + return min(nb - 1, max(0, bucket)) + + def compute_bucket_weights( stats: Mapping[int, Mapping[str, float]], policy: DynamicCurriculumPolicy, diff --git a/src/soup_cli/utils/drift_alarm.py b/src/soup_cli/utils/drift_alarm.py index eb41ed7..91ad2d7 100644 --- a/src/soup_cli/utils/drift_alarm.py +++ b/src/soup_cli/utils/drift_alarm.py @@ -23,20 +23,22 @@ quant-check thresholds; operators can tune via --threshold. from __future__ import annotations -import ipaddress import json import math import os from dataclasses import dataclass -from typing import Iterable, Mapping, Optional, Tuple -from urllib.parse import urlparse +from typing import Iterable, Mapping, Tuple from soup_cli.utils.paths import is_under_cwd +# Webhook helpers were lifted into the shared ``utils/webhooks`` module in +# v0.71.5 #207 so every production-trace command can offer --slack-url / +# --discord-url. Re-exported here for back-compat (callers + tests that import +# ``validate_webhook_url`` / ``post_webhook`` from ``drift_alarm``). +from soup_cli.utils.webhooks import post_webhook, validate_webhook_url + _MAX_REFERENCE_ROWS = 1_000_000 -_MAX_WEBHOOK_URL_LEN = 4096 _MAX_TEXT_LEN = 1_000_000 # 1 MB / row -_LOOPBACK_HOSTS = frozenset({"localhost", "127.0.0.1", "::1"}) # Default smoothing constant for the KL kernel (Laplace-style add-epsilon # to defend against `log(0)` when a token in `p` is absent from `q`). _EPS = 1e-9 @@ -84,73 +86,6 @@ def validate_threshold(value: object) -> float: return f_val -def _is_private_or_link_local(host: str) -> bool: - """Return True iff ``host`` resolves to a non-loopback private/reserved IP. - - Explicit parentheses on the final clause (code-review MEDIUM fix - v0.63.0): Python binds `and` tighter than `or`, but the SSRF gate is - safety-critical and a future edit should not need to re-derive the - precedence rules to verify the logic. - """ - try: - ip = ipaddress.ip_address(host) - except ValueError: - return False - return ( - ip.is_private - or ip.is_link_local - or (ip.is_loopback is False and (ip.is_reserved or ip.is_multicast)) - ) - - -def validate_webhook_url(url: object) -> str: - """SSRF-hardened webhook URL validator. - - Mirrors v0.29.0 `HF_ENDPOINT` / v0.30.0 OTLP / v0.51.0 `validate_hub_endpoint` - policy: - - scheme allowlist {http, https} - - null-byte / control-char rejection - - ``0.0.0.0`` rejected - - plain HTTP only permitted for loopback hosts - - private / link-local / cloud-metadata IPs rejected - """ - if isinstance(url, bool): - raise TypeError("webhook URL must be str, not bool") - if not isinstance(url, str): - raise TypeError(f"webhook URL must be str, got {type(url).__name__}") - if not url: - raise ValueError("webhook URL must be non-empty") - if "\x00" in url: - raise ValueError("webhook URL must not contain null bytes") - if any(ord(c) < 0x20 for c in url): - raise ValueError("webhook URL must not contain control characters") - if len(url) > _MAX_WEBHOOK_URL_LEN: - raise ValueError(f"webhook URL must be <= {_MAX_WEBHOOK_URL_LEN} chars") - stripped = url.rstrip("/") - parsed = urlparse(stripped) - if parsed.scheme not in ("http", "https"): - raise ValueError( - f"webhook URL must use http/https scheme, got {parsed.scheme!r}" - ) - if not parsed.netloc: - raise ValueError("webhook URL is missing a host") - host = parsed.hostname or "" - if host == "0.0.0.0": - raise ValueError( - "webhook URL 0.0.0.0 is ambiguous; use 127.0.0.1 or localhost" - ) - if parsed.scheme == "http" and host not in _LOOPBACK_HOSTS: - if _is_private_or_link_local(host): - raise ValueError( - "webhook URL plain HTTP is only allowed for loopback; " - "private/link-local hosts require HTTPS" - ) - raise ValueError( - "webhook URL for remote hosts must use HTTPS" - ) - return stripped - - def compute_token_distribution(rows: Iterable[object]) -> Mapping[str, float]: """Compute a normalised whitespace-token frequency distribution. @@ -305,39 +240,6 @@ def run_drift_check( ) -def post_webhook( - *, - url: Optional[str], - payload: Mapping[str, object], - timeout_seconds: float = 5.0, -) -> bool: - """POST ``payload`` as JSON to ``url``. Returns True on 2xx, False otherwise. - - Never raises — webhook delivery must NOT crash the drift-check run. - Lazy-imports ``httpx`` so the runtime cost is paid only when an alarm - actually fires. - """ - if url is None: - return False - try: - validated = validate_webhook_url(url) - except (TypeError, ValueError): - return False - try: - import httpx # type: ignore[import-untyped] - except ImportError: - return False - try: - response = httpx.post( - validated, - json=dict(payload), - timeout=timeout_seconds, - ) - return 200 <= response.status_code < 300 - except Exception: # noqa: BLE001 — webhook must never crash drift check - return False - - __all__ = [ "DriftReport", "compute_token_distribution", diff --git a/src/soup_cli/utils/peft_wiring.py b/src/soup_cli/utils/peft_wiring.py index a96173d..f08ac5d 100644 --- a/src/soup_cli/utils/peft_wiring.py +++ b/src/soup_cli/utils/peft_wiring.py @@ -113,8 +113,20 @@ def attach_curriculum_callback( getattr(tcfg, "curriculum_dynamic_temperature", 1.0) or 1.0 ), ) + # v0.71.5 #149 — thread curriculum_metric so the callback can bucket by + # loss / perplexity percentile (round-robin fallback for `length`). Any + # value that is not one of the three valid metrics (e.g. a missing field + # or a test MagicMock) falls back to `length` so the callback always + # constructs. + curriculum_metric = getattr(tcfg, "curriculum_metric", "length") + if curriculum_metric not in ("length", "perplexity", "loss"): + curriculum_metric = "length" try: - callback = DynamicCurriculumCallback(policy=policy, output_dir=output_dir) + callback = DynamicCurriculumCallback( + policy=policy, + output_dir=output_dir, + curriculum_metric=curriculum_metric, + ) except (TypeError, ValueError) as exc: logger.debug("attach_curriculum_callback rejected: %s", exc) return False diff --git a/src/soup_cli/utils/prune_prompt.py b/src/soup_cli/utils/prune_prompt.py index 6e83926..8adfd52 100644 --- a/src/soup_cli/utils/prune_prompt.py +++ b/src/soup_cli/utils/prune_prompt.py @@ -27,7 +27,7 @@ import json import math import os from dataclasses import dataclass -from typing import Sequence +from typing import Any, List, Optional, Sequence, Union from soup_cli.utils.paths import is_under_cwd @@ -35,6 +35,7 @@ from soup_cli.utils.paths import is_under_cwd _MAX_SCAN_ROWS = 100_000 _MAX_ROW_CHARS = 1_000_000 # 1 MB / row _MAX_PREFIX_LEN = 100_000 # hard cap on returned prefix length +_MAX_TOKENS_PER_ROW = 50_000 # token-aware mode (v0.71.5 #205) per-row cap # Tunable: a frequency below this is meaningless (we want a *near-universal* # prefix). Operator can pick anything in [0, 1] via --min-frequency. @@ -163,11 +164,112 @@ def detect_common_prefix( return best_prefix +def detect_common_prefix_tokens( + token_rows: Sequence[Sequence[int]], + *, + min_frequency: float, +) -> List[int]: + """Token-level analogue of :func:`detect_common_prefix` (v0.71.5 #205). + + Returns the longest token-id prefix shared by >= ``min_frequency`` of + rows. Operating on token IDs (not characters) guarantees the prefix + always ends on a token boundary — a multi-byte UTF-8 sequence can never + be split mid-code-point. + """ + threshold = validate_min_frequency(min_frequency) + if isinstance(token_rows, (str, bytes)) or not hasattr(token_rows, "__iter__"): + raise TypeError( + f"token_rows must be an iterable of int sequences, " + f"got {type(token_rows).__name__}" + ) + + materialised: List[List[int]] = [] + for idx, row in enumerate(token_rows): + if isinstance(row, (str, bytes)) or not hasattr(row, "__iter__"): + raise TypeError( + f"token_rows[{idx}] must be a sequence of ints, " + f"got {type(row).__name__}" + ) + materialised.append(list(row)[:_MAX_TOKENS_PER_ROW]) + if len(materialised) >= _MAX_SCAN_ROWS: + break + + if not materialised: + return [] + if len(materialised) == 1: + return list(materialised[0]) if threshold >= 1.0 else [] + + need = max(1, int(math.ceil(threshold * len(materialised)))) + best_prefix: List[int] = [] + sample_templates = materialised[: min(32, len(materialised))] + for template in sample_templates: + lo, hi = 0, len(template) + best_len = 0 + while lo <= hi: + mid = (lo + hi) // 2 + if mid == 0: + lo = mid + 1 + continue + pfx = template[:mid] + count = sum(1 for r in materialised if r[:mid] == pfx) + if count >= need: + best_len = mid + lo = mid + 1 + else: + hi = mid - 1 + if best_len > len(best_prefix): + best_prefix = template[:best_len] + return best_prefix + + +def _resolve_tokenizer(tokenizer: Union[str, Any]) -> Any: + """Return a tokenizer object from a name (lazy AutoTokenizer) or object. + + A pre-built tokenizer-like object (duck-typed ``encode`` / ``decode``) + is returned as-is — this is the injectable test seam + lets advanced + callers pass an already-loaded tokenizer. A string is treated as an HF + model id / local path and lazy-loaded via ``transformers.AutoTokenizer`` + (so importing this module never pulls transformers). + """ + if hasattr(tokenizer, "encode") and hasattr(tokenizer, "decode"): + return tokenizer + if not isinstance(tokenizer, str): + raise TypeError( + "tokenizer must be a model id / path string or a tokenizer object" + ) + if not tokenizer: + raise ValueError("tokenizer name must be non-empty") + try: + from transformers import AutoTokenizer # noqa: PLC0415 + except ImportError as exc: + raise ValueError( + "tokenizer-aware prune-prompt needs transformers — " + "install with: pip install 'soup-cli[train]'" + ) from exc + try: + return AutoTokenizer.from_pretrained(tokenizer) + except Exception as exc: # noqa: BLE001 — surface a friendly message. + raise ValueError( + f"could not load tokenizer {tokenizer!r}: {type(exc).__name__}: {exc}" + ) from exc + + +def _encode(tok: Any, text: str) -> List[int]: + """Encode ``text`` to token IDs (no special tokens), capped per-row.""" + try: + ids = tok.encode(text, add_special_tokens=False) + except TypeError: + # Tokenizers / fakes without the kwarg. + ids = tok.encode(text) + return list(ids)[:_MAX_TOKENS_PER_ROW] + + def prune_traces( input_path: str, *, output_path: str, min_frequency: float = _DEFAULT_MIN_FREQUENCY, + tokenizer: Optional[Union[str, Any]] = None, ) -> PrunePromptReport: """Read a JSONL of {prompt, output} rows, strip shared prefix, write. @@ -176,6 +278,12 @@ def prune_traces( ``prompt`` field (other fields untouched). When no prefix clears the threshold, the output is byte-identical to the input plus a ``rows_pruned=0`` report. + + v0.71.5 #205: pass ``tokenizer`` (an HF model id / local path string, or + a pre-built tokenizer object) to detect + strip the prefix on token + boundaries instead of characters — the prefix can then never end + mid-UTF-8-code-point. Default (``None``) keeps the whitespace-character + behaviour. """ threshold = validate_min_frequency(min_frequency) @@ -198,6 +306,10 @@ def prune_traces( if not os.path.isfile(input_path): raise FileNotFoundError(input_path) + # Resolve the tokenizer up front (fails fast on a bad name even on an + # empty input file) — None keeps the legacy character path. + tok = _resolve_tokenizer(tokenizer) if tokenizer is not None else None + # First pass: collect prompts (capped). prompts: list[str] = [] rows_total = 0 @@ -233,9 +345,24 @@ def prune_traces( min_frequency=threshold, ) - prefix = detect_common_prefix(prompts, min_frequency=threshold) + if tok is None: + return _prune_char_level( + input_path, output_path, prompts, rows_total, threshold + ) + return _prune_token_level( + tok, input_path, output_path, prompts, rows_total, threshold + ) - # Second pass: write output with prefix stripped where applicable. + +def _prune_char_level( + input_path: str, + output_path: str, + prompts: List[str], + rows_total: int, + threshold: float, +) -> PrunePromptReport: + """Character-level prefix strip (the v0.63.0 default behaviour).""" + prefix = detect_common_prefix(prompts, min_frequency=threshold) rows_pruned = 0 with open(input_path, encoding="utf-8") as fh_in, \ open(output_path, "w", encoding="utf-8") as fh_out: @@ -253,7 +380,6 @@ def prune_traces( row["prompt"] = row["prompt"][len(prefix):] rows_pruned += 1 fh_out.write(json.dumps(row, ensure_ascii=False) + "\n") - return PrunePromptReport( prefix=prefix, prefix_chars=len(prefix), @@ -263,9 +389,58 @@ def prune_traces( ) +def _prune_token_level( + tok: Any, + input_path: str, + output_path: str, + prompts: List[str], + rows_total: int, + threshold: float, +) -> PrunePromptReport: + """Token-aware prefix strip (v0.71.5 #205). + + The detected prefix is a list of token IDs; stripping a row decodes the + REMAINING token IDs so the boundary is always a clean token break. + """ + token_rows = [_encode(tok, p) for p in prompts] + prefix_ids = detect_common_prefix_tokens(token_rows, min_frequency=threshold) + prefix_text = tok.decode(prefix_ids) if prefix_ids else "" + plen = len(prefix_ids) + + rows_pruned = 0 + with open(input_path, encoding="utf-8") as fh_in, \ + open(output_path, "w", encoding="utf-8") as fh_out: + for line in fh_in: + line = line.strip() + if not line: + continue + try: + row = json.loads(line) + except json.JSONDecodeError: + continue + if not isinstance(row, dict): + continue + prompt = row.get("prompt") + if prefix_ids and isinstance(prompt, str): + ids = _encode(tok, prompt) + if ids[:plen] == prefix_ids: + row["prompt"] = tok.decode(ids[plen:]) + rows_pruned += 1 + fh_out.write(json.dumps(row, ensure_ascii=False) + "\n") + + return PrunePromptReport( + prefix=prefix_text, + prefix_chars=len(prefix_text), + rows_total=rows_total, + rows_pruned=rows_pruned, + min_frequency=threshold, + ) + + __all__ = [ "PrunePromptReport", "detect_common_prefix", + "detect_common_prefix_tokens", "prune_traces", "validate_min_frequency", ] diff --git a/src/soup_cli/utils/webhooks.py b/src/soup_cli/utils/webhooks.py new file mode 100644 index 0000000..164d91b --- /dev/null +++ b/src/soup_cli/utils/webhooks.py @@ -0,0 +1,152 @@ +"""Shared Slack / Discord webhook helpers (v0.71.5 #207). + +Lifts the SSRF-hardened ``validate_webhook_url`` + best-effort ``post_webhook`` +out of ``utils/drift_alarm.py`` (v0.63.0 Part E) so every production-trace +command can offer ``--slack-url`` / ``--discord-url`` without re-implementing +the SSRF gate. ``drift_alarm`` now re-exports these for back-compat. + +SSRF policy — full parity with v0.29.0 ``HF_ENDPOINT`` / v0.30.0 OTLP / +v0.51.0 ``validate_hub_endpoint`` / v0.63.0 drift-alarm: +- scheme allowlist {http, https} +- null-byte / control-char rejection +- ``0.0.0.0`` rejected +- plain HTTP only permitted for loopback hosts +- private / link-local / reserved / multicast IPs rejected + +``post_webhook`` NEVER raises — webhook delivery must not crash the command +that triggered it. ``httpx`` is lazy-imported so the runtime cost is paid only +when an alarm actually fires. +""" + +from __future__ import annotations + +import ipaddress +from typing import List, Mapping, Optional, Tuple +from urllib.parse import urlparse + +_MAX_WEBHOOK_URL_LEN = 4096 +_LOOPBACK_HOSTS = frozenset({"localhost", "127.0.0.1", "::1"}) + + +def _is_private_or_link_local(host: str) -> bool: + """Return True iff ``host`` resolves to a non-loopback private/reserved IP. + + Explicit parentheses on the final clause (mirrors v0.63.0 drift-alarm + code-review MEDIUM fix): Python binds ``and`` tighter than ``or``, but + the SSRF gate is safety-critical and a future edit should not need to + re-derive the precedence rules to verify the logic. + """ + try: + ip = ipaddress.ip_address(host) + except ValueError: + return False + return ( + ip.is_private + or ip.is_link_local + or (ip.is_loopback is False and (ip.is_reserved or ip.is_multicast)) + ) + + +def validate_webhook_url(url: object) -> str: + """SSRF-hardened webhook URL validator (returns the canonical URL).""" + if isinstance(url, bool): + raise TypeError("webhook URL must be str, not bool") + if not isinstance(url, str): + raise TypeError(f"webhook URL must be str, got {type(url).__name__}") + if not url: + raise ValueError("webhook URL must be non-empty") + if "\x00" in url: + raise ValueError("webhook URL must not contain null bytes") + if any(ord(c) < 0x20 for c in url): + raise ValueError("webhook URL must not contain control characters") + if len(url) > _MAX_WEBHOOK_URL_LEN: + raise ValueError(f"webhook URL must be <= {_MAX_WEBHOOK_URL_LEN} chars") + stripped = url.rstrip("/") + parsed = urlparse(stripped) + if parsed.scheme not in ("http", "https"): + raise ValueError( + f"webhook URL must use http/https scheme, got {parsed.scheme!r}" + ) + if not parsed.netloc: + raise ValueError("webhook URL is missing a host") + host = parsed.hostname or "" + if host == "0.0.0.0": + raise ValueError( + "webhook URL 0.0.0.0 is ambiguous; use 127.0.0.1 or localhost" + ) + if parsed.scheme == "http" and host not in _LOOPBACK_HOSTS: + if _is_private_or_link_local(host): + raise ValueError( + "webhook URL plain HTTP is only allowed for loopback; " + "private/link-local hosts require HTTPS" + ) + raise ValueError("webhook URL for remote hosts must use HTTPS") + return stripped + + +def post_webhook( + *, + url: Optional[str], + payload: Mapping[str, object], + timeout_seconds: float = 5.0, +) -> bool: + """POST ``payload`` as JSON to ``url``. Returns True on 2xx, False otherwise. + + Never raises — webhook delivery must NOT crash the calling command. + """ + if url is None: + return False + try: + validated = validate_webhook_url(url) + except (TypeError, ValueError): + return False + try: + import httpx # type: ignore[import-untyped] + except ImportError: + return False + try: + response = httpx.post( + validated, + json=dict(payload), + timeout=timeout_seconds, + ) + return 200 <= response.status_code < 300 + except Exception: # noqa: BLE001 — webhook must never crash the command + return False + + +def send_webhooks( + payload: Mapping[str, object], + *, + slack_url: Optional[str] = None, + discord_url: Optional[str] = None, + timeout_seconds: float = 5.0, +) -> List[Tuple[str, bool]]: + """POST ``payload`` to each provided webhook; return per-target delivery. + + Returns a list of ``(label, delivered)`` for every non-``None`` URL, + in ``slack`` then ``discord`` order. ``None`` URLs are skipped (not + attempted). Never raises (delegates to the never-raising + :func:`post_webhook`). + """ + results: List[Tuple[str, bool]] = [] + for label, url in (("slack", slack_url), ("discord", discord_url)): + if url: + results.append( + ( + label, + post_webhook( + url=url, + payload=payload, + timeout_seconds=timeout_seconds, + ), + ) + ) + return results + + +__all__ = [ + "post_webhook", + "send_webhooks", + "validate_webhook_url", +] diff --git a/tests/test_v0715.py b/tests/test_v0715.py new file mode 100644 index 0000000..800c168 --- /dev/null +++ b/tests/test_v0715.py @@ -0,0 +1,1222 @@ +"""Tests for v0.71.5 — Ingest / data / prompt / drift patch. + +Closes (6 of 7): #164, #163, #207, #205, #149, #157. +Deferred: #204 (live SaaS pull adapters — external-account-gated, infra-blocked). + +One test file per release (project convention). Grouped by issue. +""" + +from __future__ import annotations + +import sys +from pathlib import Path + +import pytest +from typer.testing import CliRunner + +runner = CliRunner() + + +# =========================================================================== +# #164 — extend get_metric_series to the eval_results table +# =========================================================================== + + +class TestGetMetricSeriesEvalResults: + """`tracker.get_metric_series` falls back to eval_results for benchmarks.""" + + def _tracker(self, tmp_path): + from soup_cli.experiment.tracker import ExperimentTracker + + return ExperimentTracker(db_path=tmp_path / "exp.db") + + def test_eval_results_series_returned(self, tmp_path): + tracker = self._tracker(tmp_path) + tracker.save_eval_result("m1", "task_accuracy", 0.81, {}, run_id="run-1") + tracker.save_eval_result("m1", "task_accuracy", 0.79, {}, run_id="run-1") + series = tracker.get_metric_series("run-1", "task_accuracy") + assert series == [0.81, 0.79] + + def test_eval_results_filtered_by_benchmark(self, tmp_path): + tracker = self._tracker(tmp_path) + tracker.save_eval_result("m1", "task_accuracy", 0.81, {}, run_id="run-1") + tracker.save_eval_result("m1", "refusal_rate", 0.02, {}, run_id="run-1") + assert tracker.get_metric_series("run-1", "refusal_rate") == [0.02] + + def test_eval_results_filtered_by_run_id(self, tmp_path): + tracker = self._tracker(tmp_path) + tracker.save_eval_result("m1", "task_accuracy", 0.81, {}, run_id="run-1") + tracker.save_eval_result("m2", "task_accuracy", 0.55, {}, run_id="run-2") + assert tracker.get_metric_series("run-1", "task_accuracy") == [0.81] + assert tracker.get_metric_series("run-2", "task_accuracy") == [0.55] + + def test_per_step_column_still_uses_metrics_table(self, tmp_path): + # No regression: known training-loop columns read from `metrics`. + tracker = self._tracker(tmp_path) + tracker.log_metrics("run-1", step=1, loss=2.0) + tracker.log_metrics("run-1", step=2, loss=1.5) + # Also write a (different) eval result to prove we DON'T read it for loss. + tracker.save_eval_result("m1", "loss", 99.0, {}, run_id="run-1") + assert tracker.get_metric_series("run-1", "loss") == [2.0, 1.5] + + def test_missing_metric_returns_empty_list(self, tmp_path): + tracker = self._tracker(tmp_path) + assert tracker.get_metric_series("nope", "task_accuracy") == [] + + def test_unknown_benchmark_returns_empty(self, tmp_path): + tracker = self._tracker(tmp_path) + tracker.save_eval_result("m1", "task_accuracy", 0.81, {}, run_id="run-1") + assert tracker.get_metric_series("run-1", "bleu") == [] + + def test_eval_series_ordered_by_id(self, tmp_path): + # Deterministic pairing for the paired-bootstrap CI: the series is + # ordered by insertion id, independent of the created_at timestamp. + # (The `value is None` / non-numeric branches in _eval_score_series + # are defensive-only — `eval_results.score` is REAL NOT NULL, so a + # NULL / non-numeric cell cannot be inserted; no test can construct + # one without violating the schema.) + tracker = self._tracker(tmp_path) + conn = tracker._get_conn() + for score, ts in ((0.1, "2026-01-09"), (0.2, "2026-01-01"), (0.3, "2026-01-05")): + conn.execute( + "INSERT INTO eval_results (run_id, model_path, benchmark, score, " + "details_json, created_at) VALUES (?, ?, ?, ?, ?, ?)", + ("run-1", "m", "task_accuracy", score, "{}", ts), + ) + conn.commit() + assert tracker.get_metric_series("run-1", "task_accuracy") == [0.1, 0.2, 0.3] + + def test_empty_run_id_rejected(self, tmp_path): + tracker = self._tracker(tmp_path) + with pytest.raises(ValueError, match="non-empty string"): + tracker.get_metric_series("", "task_accuracy") + + def test_empty_metric_rejected(self, tmp_path): + tracker = self._tracker(tmp_path) + with pytest.raises(ValueError, match="non-empty string"): + tracker.get_metric_series("run-1", "") + + @pytest.mark.parametrize("bad", [True, 123, None]) + def test_non_string_run_id_rejected(self, tmp_path, bad): + tracker = self._tracker(tmp_path) + with pytest.raises(ValueError): + tracker.get_metric_series(bad, "task_accuracy") + + @pytest.mark.parametrize("bad", [True, 123, None]) + def test_non_string_metric_rejected(self, tmp_path, bad): + tracker = self._tracker(tmp_path) + with pytest.raises(ValueError): + tracker.get_metric_series("run-1", bad) + + +# =========================================================================== +# #163 — bias build_verdict from advise_history outcomes +# =========================================================================== + + +class _FakeHistoryEntry: + """Duck-typed HistoryEntry stand-in for build_verdict bias tests.""" + + def __init__(self, *, choice, task_category="factual_lookup", accepted=True, + outcome=0.5, project="proj-a"): + self.choice = choice + self.task_category = task_category + self.accepted = accepted + self.outcome = outcome + self.project = project + + +class TestBuildVerdictHistoryBias: + def _profile(self, **over): + from soup_cli.utils.advise import DatasetProfile + + base = dict( + row_count=2000, + avg_input_chars=120.0, + avg_output_chars=80.0, + type_token_diversity=0.6, + label_variance=0.7, + has_chosen_rejected=False, + has_reasoning_traces=False, + ) + base.update(over) + return DatasetProfile(**base) + + def test_summarise_empty_and_none(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + assert _summarise_history_outcomes(None) == {} + assert _summarise_history_outcomes([]) == {} + + def test_summarise_filters_by_project(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + hist = [ + _FakeHistoryEntry(choice="SFT", outcome=0.5, project="proj-a"), + _FakeHistoryEntry(choice="SFT", outcome=0.5, project="proj-b"), + ] + out = _summarise_history_outcomes(hist, project="proj-a") + assert out["SFT"] == (0.5, 1) # only proj-a counted + + def test_summarise_skips_unaccepted_and_none_outcome(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + hist = [ + _FakeHistoryEntry(choice="SFT", outcome=0.5, accepted=False), + _FakeHistoryEntry(choice="SFT", outcome=None, accepted=True), + _FakeHistoryEntry(choice="SFT", outcome=0.4, accepted=True), + ] + out = _summarise_history_outcomes(hist) + assert out["SFT"] == (0.4, 1) + + def test_summarise_skips_bool_outcome(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + hist = [_FakeHistoryEntry(choice="SFT", outcome=True)] + assert _summarise_history_outcomes(hist) == {} + + def test_summarise_non_sequence_rejected(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + with pytest.raises(TypeError): + _summarise_history_outcomes(123) # type: ignore[arg-type] + + def test_no_history_identical_to_base(self): + from soup_cli.utils.advise import build_verdict + + profile = self._profile() # factual_lookup + high variance → RAG + base = build_verdict(profile, "factual_lookup") + with_none = build_verdict(profile, "factual_lookup", history=None) + assert base.choice == with_none.choice == "RAG" + assert base.confidence == with_none.confidence + + def test_sft_precedents_flip_marginal_rag_to_sft(self): + from soup_cli.utils.advise import build_verdict + + profile = self._profile() # would be RAG by default + hist = [ + _FakeHistoryEntry(choice="SFT", outcome=0.4, project="proj-a") + for _ in range(3) + ] + verdict = build_verdict( + profile, "factual_lookup", history=hist, project="proj-a" + ) + assert verdict.choice == "SFT" + assert "precedent" in verdict.reason.lower() + + def test_sft_precedents_other_project_do_not_flip(self): + from soup_cli.utils.advise import build_verdict + + profile = self._profile() + hist = [ + _FakeHistoryEntry(choice="SFT", outcome=0.4, project="other") + for _ in range(3) + ] + verdict = build_verdict( + profile, "factual_lookup", history=hist, project="proj-a" + ) + assert verdict.choice == "RAG" # unchanged — wrong project + + def test_fewer_than_three_sft_precedents_no_flip(self): + from soup_cli.utils.advise import build_verdict + + profile = self._profile() + hist = [ + _FakeHistoryEntry(choice="SFT", outcome=0.4, project="proj-a") + for _ in range(2) + ] + verdict = build_verdict( + profile, "factual_lookup", history=hist, project="proj-a" + ) + assert verdict.choice == "RAG" + + def test_negative_grpo_precedents_suppress_grpo(self): + from soup_cli.utils.advise import build_verdict + + profile = self._profile( + row_count=800, has_reasoning_traces=True, label_variance=0.3, + ) + # Default (no history) → GRPO. + assert build_verdict(profile, "reasoning").choice == "GRPO" + hist = [ + _FakeHistoryEntry(choice="GRPO", outcome=-0.2, project="proj-a") + for _ in range(3) + ] + verdict = build_verdict( + profile, "reasoning", history=hist, project="proj-a" + ) + assert verdict.choice == "SFT" + assert "precedent" in verdict.reason.lower() + + def test_encouraged_choice_nudges_confidence(self): + from soup_cli.utils.advise import build_verdict + + # SFT task (not factual_lookup) so base is SFT; encourage SFT. + profile = self._profile(label_variance=0.3) + base = build_verdict(profile, "style_shaping") + assert base.choice == "SFT" + hist = [ + _FakeHistoryEntry(choice="SFT", outcome=0.5, project="proj-a") + for _ in range(3) + ] + biased = build_verdict( + profile, "style_shaping", history=hist, project="proj-a" + ) + assert biased.choice == "SFT" + # Exact nudge (+0.05, clamped at 0.95) — a zero-nudge bug would fail. + assert biased.confidence == pytest.approx(min(0.95, base.confidence + 0.05)) + + def test_summarise_outcome_out_of_range_skipped(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + hist = [_FakeHistoryEntry(choice="SFT", outcome=2.0)] + assert _summarise_history_outcomes(hist) == {} + + def test_summarise_accepted_must_be_true_not_truthy(self): + from soup_cli.utils.advise import _summarise_history_outcomes + + # accepted=1 (int) is NOT True → rejected (the guard is `is not True`). + hist = [_FakeHistoryEntry(choice="SFT", accepted=1, outcome=0.5)] + assert _summarise_history_outcomes(hist) == {} + + def test_is_encouraged_threshold_boundaries(self): + from soup_cli.utils.advise import _is_encouraged + + # >= 0.3 over >= 3 precedents. + assert _is_encouraged({"SFT": (0.3, 3)}, "SFT") is True + assert _is_encouraged({"SFT": (0.29, 3)}, "SFT") is False + assert _is_encouraged({"SFT": (0.5, 2)}, "SFT") is False # count < 3 + assert _is_encouraged({}, "SFT") is False + + def test_is_discouraged_threshold_boundaries(self): + from soup_cli.utils.advise import _is_discouraged + + # strict < 0.0 over >= 3 precedents. + assert _is_discouraged({"GRPO": (0.0, 3)}, "GRPO") is False + assert _is_discouraged({"GRPO": (-0.01, 3)}, "GRPO") is True + assert _is_discouraged({"GRPO": (-0.5, 2)}, "GRPO") is False # count < 3 + + def test_real_history_entry_duck_types(self): + from soup_cli.utils.advise import _summarise_history_outcomes + from soup_cli.utils.advise_history import HistoryEntry + + entry = HistoryEntry( + timestamp="2026-01-01T00:00:00+00:00", + project="proj-a", + choice="SFT", + task_category="style_shaping", + confidence=0.8, + reason="r", + reverse_when="w", + accepted=True, + outcome=0.5, + notes="", + ) + out = _summarise_history_outcomes([entry], project="proj-a") + assert out["SFT"] == (0.5, 1) + + +# =========================================================================== +# #207 — shared utils/webhooks.py + --slack-url/--discord-url on 4 commands +# =========================================================================== + + +class TestSharedWebhooks: + def test_webhooks_module_exports(self): + from soup_cli.utils import webhooks + + assert hasattr(webhooks, "validate_webhook_url") + assert hasattr(webhooks, "post_webhook") + assert hasattr(webhooks, "send_webhooks") + + def test_drift_alarm_reexports_same_objects(self): + from soup_cli.utils import drift_alarm, webhooks + + assert drift_alarm.validate_webhook_url is webhooks.validate_webhook_url + assert drift_alarm.post_webhook is webhooks.post_webhook + + def test_validate_webhook_url_https_ok(self): + from soup_cli.utils.webhooks import validate_webhook_url + + assert validate_webhook_url("https://hooks.slack.com/x") is not None + + def test_validate_webhook_url_rejects_rfc1918(self): + from soup_cli.utils.webhooks import validate_webhook_url + + with pytest.raises(ValueError): + validate_webhook_url("http://10.0.0.5/hook") + + def test_validate_webhook_url_rejects_loopback_only_http(self): + from soup_cli.utils.webhooks import validate_webhook_url + + assert validate_webhook_url("http://127.0.0.1:9000/h") is not None + with pytest.raises(ValueError): + validate_webhook_url("http://example.com/h") # remote http + + @pytest.mark.parametrize("bad", [True, 123, None]) + def test_validate_webhook_url_type_rejection(self, bad): + from soup_cli.utils.webhooks import validate_webhook_url + + with pytest.raises(TypeError): + validate_webhook_url(bad) + + @pytest.mark.parametrize( + "bad", + [ + "", + "ftp://example.com/h", + "file:///etc/passwd", + "javascript:alert(1)", + "http://0.0.0.0/h", + "http://169.254.169.254/latest", # link-local cloud metadata + "https://host/h\x00", # null byte + "https://host/h\nX", # control char + "https://" + "a" * 5000, # oversize + ], + ) + def test_validate_webhook_url_value_rejection_matrix(self, bad): + # Re-prove the SSRF gate against the NEW module path (the validator + # moved out of drift_alarm in #207 — don't rely on the re-export). + from soup_cli.utils.webhooks import validate_webhook_url + + with pytest.raises(ValueError): + validate_webhook_url(bad) + + def test_send_webhooks_posts_both(self, monkeypatch): + from soup_cli.utils import webhooks + + calls = [] + monkeypatch.setattr( + webhooks, "post_webhook", + lambda **kw: (calls.append(kw), True)[1], + ) + results = webhooks.send_webhooks( + {"k": 1}, + slack_url="https://hooks.slack.com/x", + discord_url="https://discord.com/api/webhooks/x", + ) + assert results == [("slack", True), ("discord", True)] + assert len(calls) == 2 + assert calls[0]["payload"] == {"k": 1} + + def test_send_webhooks_skips_none(self, monkeypatch): + from soup_cli.utils import webhooks + + calls = [] + monkeypatch.setattr( + webhooks, "post_webhook", + lambda **kw: (calls.append(kw), True)[1], + ) + results = webhooks.send_webhooks({"k": 1}, slack_url=None, discord_url=None) + assert results == [] + assert calls == [] + + def test_send_webhooks_swallows_failure(self, monkeypatch): + from soup_cli.utils import webhooks + + monkeypatch.setattr(webhooks, "post_webhook", lambda **kw: False) + results = webhooks.send_webhooks( + {"k": 1}, slack_url="https://hooks.slack.com/x" + ) + assert results == [("slack", False)] + + # --- CLI flag plumbing on each of the 4 commands ---------------------- + + @pytest.mark.parametrize( + "argv", + [ + ["ingest", "--help"], + ["prune-prompt", "--help"], + ["ab", "--help"], + ["data", "active-sample", "--help"], + ], + ) + def test_webhook_flags_in_help(self, argv): + import re as _re + + from soup_cli.cli import app + + result = runner.invoke(app, argv) + assert result.exit_code == 0, (result.output, repr(result.exception)) + clean = _re.sub(r"\x1b\[[0-9;]*m", "", result.output) + clean = clean.replace("\n", " ") + clean = _re.sub(r"\s+", " ", clean) + assert "--slack-url" in clean + assert "--discord-url" in clean + + def test_ingest_rejects_bad_webhook(self, tmp_path, monkeypatch): + from soup_cli.cli import app + + monkeypatch.chdir(tmp_path) + (tmp_path / "logs.jsonl").write_text( + '{"input": "hi", "output": "yo"}\n', encoding="utf-8" + ) + result = runner.invoke( + app, + ["ingest", "--source", "langfuse", "--logs", "logs.jsonl", + "--slack-url", "http://10.0.0.1/h"], + ) + assert result.exit_code == 2 + + def test_bad_discord_url_labels_correct_flag(self, tmp_path, monkeypatch): + # The second loop iteration in validate_webhook_flags must label + # --discord-url (guards a copy-paste bug printing --slack-url twice). + import re as _re + + from soup_cli.cli import app + + monkeypatch.chdir(tmp_path) + (tmp_path / "logs.jsonl").write_text( + '{"input": "hi", "output": "yo"}\n', encoding="utf-8" + ) + result = runner.invoke( + app, + ["ingest", "--source", "langfuse", "--logs", "logs.jsonl", + "--discord-url", "http://10.0.0.1/h"], + ) + assert result.exit_code == 2 + clean = _re.sub(r"\x1b\[[0-9;]*m", "", result.output).replace("\n", " ") + assert "--discord-url" in clean + + def test_emit_webhooks_early_return_no_urls(self): + # Direct unit test of the CLI helper (only exercised via 4 CLIs). + from rich.console import Console + + from soup_cli.commands._webhook_cli import emit_webhooks + + # Both None → no-op, no exception. + emit_webhooks(None, None, payload={"k": 1}, console=Console()) + + def test_ingest_posts_payload_on_success(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import webhooks + + captured = [] + monkeypatch.setattr( + webhooks, "post_webhook", + lambda **kw: (captured.append(kw), True)[1], + ) + monkeypatch.chdir(tmp_path) + (tmp_path / "logs.jsonl").write_text( + '{"input": "hi", "output": "yo"}\n', encoding="utf-8" + ) + result = runner.invoke( + app, + ["ingest", "--source", "langfuse", "--logs", "logs.jsonl", + "--slack-url", "https://hooks.slack.com/x"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert len(captured) == 1 + payload = captured[0]["payload"] + assert payload["source"] == "langfuse" + assert payload["traces_written"] == 1 + assert "auth_env_set" in payload + + def test_prune_prompt_posts_payload(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import webhooks + + captured = [] + monkeypatch.setattr( + webhooks, "post_webhook", + lambda **kw: (captured.append(kw), True)[1], + ) + monkeypatch.chdir(tmp_path) + rows = "".join( + '{"prompt": "SYS PREAMBLE. ask %d", "output": "a"}\n' % i + for i in range(5) + ) + (tmp_path / "in.jsonl").write_text(rows, encoding="utf-8") + result = runner.invoke( + app, + ["prune-prompt", "--input", "in.jsonl", "--output", "out.jsonl", + "--min-frequency", "0.9", + "--discord-url", "https://discord.com/api/webhooks/x"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert len(captured) == 1 + payload = captured[0]["payload"] + assert "rows_pruned" in payload + assert "prefix_chars" in payload + + def test_ab_posts_only_on_decision(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import webhooks + + captured = [] + monkeypatch.setattr( + webhooks, "post_webhook", + lambda **kw: (captured.append(kw), True)[1], + ) + monkeypatch.chdir(tmp_path) + # Two-row "continue" case → no webhook fired. + (tmp_path / "ab.jsonl").write_text( + '{"arm": "control", "latency": 1.0}\n' + '{"arm": "treatment", "latency": 1.0}\n', + encoding="utf-8", + ) + result = runner.invoke( + app, + ["ab", "--input", "ab.jsonl", "--metric", "latency", + "--slack-url", "https://hooks.slack.com/x"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + # decision == continue → no webhook + assert captured == [] + + def test_active_sample_posts_payload(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import webhooks + + captured = [] + monkeypatch.setattr( + webhooks, "post_webhook", + lambda **kw: (captured.append(kw), True)[1], + ) + monkeypatch.chdir(tmp_path) + (tmp_path / "traces.jsonl").write_text( + '{"prompt": "a", "rm_score": 0.5}\n' + '{"prompt": "b", "rm_score": 0.9}\n', + encoding="utf-8", + ) + result = runner.invoke( + app, + ["data", "active-sample", "--input", "traces.jsonl", + "--output", "sel.jsonl", "--budget", "1", + "--slack-url", "https://hooks.slack.com/x"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert len(captured) == 1 + payload = captured[0]["payload"] + assert "rows_selected" in payload + assert "mean_uncertainty" in payload + + +# =========================================================================== +# #205 — tokenizer-aware prefix detection for prune-prompt +# =========================================================================== + + +class _FakeWordTokenizer: + """Whitespace tokenizer with a stable str->id vocab (offline test seam). + + Mimics the `transformers` tokenizer surface: ``encode(text, + add_special_tokens=...)`` -> list[int]; ``decode(ids)`` -> str. + """ + + def __init__(self): + self._stoi: dict[str, int] = {} + self._itos: dict[int, str] = {} + + def _id(self, tok: str) -> int: + if tok not in self._stoi: + idx = len(self._stoi) + self._stoi[tok] = idx + self._itos[idx] = tok + return self._stoi[tok] + + def encode(self, text, add_special_tokens=True): # noqa: ARG002 + return [self._id(t) for t in text.split(" ") if t != ""] + + def decode(self, ids): + return " ".join(self._itos[i] for i in ids) + + +class TestPrunePromptTokenizer: + def test_detect_common_prefix_tokens_happy(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + rows = [ + [1, 2, 3, 4], + [1, 2, 3, 9], + [1, 2, 3, 7], + ] + assert detect_common_prefix_tokens(rows, min_frequency=1.0) == [1, 2, 3] + + def test_detect_common_prefix_tokens_partial_majority(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + rows = [ + [1, 2, 3], + [1, 2, 3], + [9, 9, 9], + ] + # 2/3 share [1,2,3] → at 0.6 threshold, returns it. + assert detect_common_prefix_tokens(rows, min_frequency=0.6) == [1, 2, 3] + # at 0.9 threshold, only [9..]? no shared prefix across all → [] + assert detect_common_prefix_tokens(rows, min_frequency=0.9) == [] + + def test_detect_common_prefix_tokens_empty(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + assert detect_common_prefix_tokens([], min_frequency=1.0) == [] + + def test_detect_common_prefix_tokens_single_row(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + assert detect_common_prefix_tokens([[1, 2]], min_frequency=1.0) == [1, 2] + assert detect_common_prefix_tokens([[1, 2]], min_frequency=0.5) == [] + + def test_detect_common_prefix_tokens_invalid_min_frequency(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + with pytest.raises(ValueError): + detect_common_prefix_tokens([[1]], min_frequency=1.5) + + def test_detect_common_prefix_tokens_min_freq_lower_and_nan(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + with pytest.raises(ValueError): + detect_common_prefix_tokens([[1]], min_frequency=-0.1) + with pytest.raises(ValueError): + detect_common_prefix_tokens([[1]], min_frequency=float("nan")) + + def test_detect_common_prefix_tokens_non_iterable_rejected(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + with pytest.raises(TypeError): + detect_common_prefix_tokens(123, min_frequency=1.0) # type: ignore[arg-type] + + def test_detect_common_prefix_tokens_non_iterable_row_rejected(self): + from soup_cli.utils.prune_prompt import detect_common_prefix_tokens + + with pytest.raises(TypeError): + detect_common_prefix_tokens([123, [1, 2]], min_frequency=1.0) # type: ignore[list-item] + + def test_prune_traces_with_fake_tokenizer(self, tmp_path, monkeypatch): + from soup_cli.utils.prune_prompt import prune_traces + + monkeypatch.chdir(tmp_path) + # The varying token (ask0/ask1/...) carries no leading space so the + # shared prefix is exactly the 3-token preamble. + rows = "".join( + '{"prompt": "SYS PREAMBLE HERE ask%d", "output": "a"}\n' % i + for i in range(5) + ) + (tmp_path / "in.jsonl").write_text(rows, encoding="utf-8") + report = prune_traces( + "in.jsonl", + output_path="out.jsonl", + min_frequency=0.9, + tokenizer=_FakeWordTokenizer(), + ) + assert report.rows_pruned == 5 + assert report.prefix == "SYS PREAMBLE HERE" + # Output prompts no longer carry the preamble tokens. + import json as _json + + lines = [ + _json.loads(line) + for line in (tmp_path / "out.jsonl").read_text( + encoding="utf-8" + ).splitlines() + if line.strip() + ] + assert all(not row["prompt"].startswith("SYS PREAMBLE") for row in lines) + assert lines[0]["prompt"] == "ask0" + + def test_prune_traces_tokenizer_multibyte_safe(self, tmp_path, monkeypatch): + # A shared prefix containing a multi-byte token is stripped on a + # token boundary — never mid-code-point. + from soup_cli.utils.prune_prompt import prune_traces + + monkeypatch.chdir(tmp_path) + rows = "".join( + '{"prompt": "café ☕ menu item %d", "output": "x"}\n' % i + for i in range(4) + ) + (tmp_path / "in.jsonl").write_text(rows, encoding="utf-8") + report = prune_traces( + "in.jsonl", + output_path="out.jsonl", + min_frequency=0.9, + tokenizer=_FakeWordTokenizer(), + ) + assert "café" in report.prefix and "☕" in report.prefix + assert report.rows_pruned == 4 + + def test_prune_traces_string_tokenizer_lazy_loads(self, tmp_path, monkeypatch): + # `tokenizer` as a string → lazy AutoTokenizer.from_pretrained. + import types + + from soup_cli.utils import prune_prompt as pp + + monkeypatch.chdir(tmp_path) + loaded = {} + + def _fake_from_pretrained(name, **kw): # noqa: ARG001 + loaded["name"] = name + return _FakeWordTokenizer() + + fake_auto = types.SimpleNamespace(from_pretrained=_fake_from_pretrained) + fake_transformers = types.SimpleNamespace(AutoTokenizer=fake_auto) + monkeypatch.setitem(sys.modules, "transformers", fake_transformers) + + rows = "".join( + '{"prompt": "SYS ask %d", "output": "a"}\n' % i for i in range(3) + ) + (tmp_path / "in.jsonl").write_text(rows, encoding="utf-8") + report = pp.prune_traces( + "in.jsonl", + output_path="out.jsonl", + min_frequency=0.9, + tokenizer="my/model", + ) + assert loaded["name"] == "my/model" + assert report.rows_pruned == 3 + + def test_prune_traces_friendly_error_on_missing_tokenizer(self, tmp_path, monkeypatch): + import types + + from soup_cli.utils import prune_prompt as pp + + monkeypatch.chdir(tmp_path) + + def _boom(name, **kw): # noqa: ARG001 + raise OSError("no such model") + + fake_auto = types.SimpleNamespace(from_pretrained=_boom) + monkeypatch.setitem( + sys.modules, "transformers", + types.SimpleNamespace(AutoTokenizer=fake_auto), + ) + (tmp_path / "in.jsonl").write_text( + '{"prompt": "a", "output": "b"}\n', encoding="utf-8" + ) + with pytest.raises(ValueError, match="could not load tokenizer"): + pp.prune_traces( + "in.jsonl", + output_path="out.jsonl", + tokenizer="bad/model", + ) + + def test_char_level_unchanged_when_no_tokenizer(self, tmp_path, monkeypatch): + from soup_cli.utils.prune_prompt import prune_traces + + monkeypatch.chdir(tmp_path) + rows = "".join( + '{"prompt": "PREFIX %d", "output": "a"}\n' % i for i in range(4) + ) + (tmp_path / "in.jsonl").write_text(rows, encoding="utf-8") + report = prune_traces( + "in.jsonl", + output_path="out.jsonl", + min_frequency=0.9, + ) + assert report.prefix.startswith("PREFIX ") + assert report.rows_pruned == 4 + + def test_no_top_level_transformers_import(self): + # Lazy import — `soup prune-prompt --help` must not pull transformers. + # Anchor on __file__ (not cwd) so the guard survives another test's + # monkeypatch.chdir (v0.58.0 source-grep precedent). + repo_root = Path(__file__).resolve().parent.parent + src = ( + repo_root / "src" / "soup_cli" / "utils" / "prune_prompt.py" + ).read_text(encoding="utf-8") + assert "\nimport transformers" not in src + assert "\nfrom transformers" not in src + + +# =========================================================================== +# #149 — DynamicCurriculumCallback bucket selection by curriculum_metric +# =========================================================================== + + +class _FakeState: + def __init__(self, global_step): + self.global_step = global_step + + +class TestPercentileBucket: + def test_max_value_top_bucket(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + window = [0.1, 0.2, 0.3, 0.4] + assert percentile_bucket(0.9, window, 4) == 3 + + def test_min_value_bucket_zero(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + window = [0.5, 0.6, 0.7, 0.8] + assert percentile_bucket(0.1, window, 4) == 0 + + def test_mid_value_mid_bucket(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + window = [0.0, 1.0, 2.0, 3.0] + # value 1.5 → 2/4 le → rank 0.5 → bucket 2 + assert percentile_bucket(1.5, window, 4) == 2 + + def test_single_bucket(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + assert percentile_bucket(5.0, [1.0, 2.0], 1) == 0 + + def test_empty_window_bucket_zero(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + assert percentile_bucket(5.0, [], 4) == 0 + + def test_bool_value_rejected(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + with pytest.raises(ValueError): + percentile_bucket(True, [1.0], 4) + + def test_non_finite_value_rejected(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + with pytest.raises(ValueError): + percentile_bucket(float("nan"), [1.0], 4) + + def test_bool_num_buckets_rejected(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + with pytest.raises(ValueError): + percentile_bucket(1.0, [1.0], True) + + @pytest.mark.parametrize("nb", [0, 21]) + def test_num_buckets_out_of_bounds_rejected(self, nb): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + with pytest.raises(ValueError): + percentile_bucket(1.0, [1.0], nb) + + def test_equal_to_all_window_values_top_bucket(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + # value == every member → le == len → rank 1.0 → clamped to top. + # Pins the `<=` (not `<`) in the rank tally. + assert percentile_bucket(0.5, [0.5, 0.5], 4) == 3 + + def test_consistently_high_value_stable_top_bucket(self): + from soup_cli.utils.curriculum_dynamic import percentile_bucket + + window = [0.1, 0.2, 0.3, 0.4, 0.5] + # Two independent appearances of a high loss → same top bucket. + b1 = percentile_bucket(9.0, window, 4) + b2 = percentile_bucket(9.0, window + [0.15, 0.25], 4) + assert b1 == b2 == 3 + + +class TestCurriculumCallbackMetric: + def _callback(self, tmp_path, metric): + from soup_cli.monitoring.curriculum_callback import ( + DynamicCurriculumCallback, + ) + from soup_cli.utils.curriculum_dynamic import DynamicCurriculumPolicy + + policy = DynamicCurriculumPolicy(num_buckets=4, recompute_every_n_steps=10) + return DynamicCurriculumCallback( + policy=policy, output_dir=str(tmp_path), curriculum_metric=metric + ) + + def test_metric_validated(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + with pytest.raises(ValueError): + self._callback(tmp_path, "bogus") + + def test_metric_property(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + cb = self._callback(tmp_path, "loss") + assert cb.curriculum_metric == "loss" + + def test_loss_metric_routes_high_to_top_bucket(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + cb = self._callback(tmp_path, "loss") + # Warm the window with low losses (round-robin during warm-up). + for step in range(1, 9): + cb.on_log(None, _FakeState(step), None, logs={"loss": 0.1 * step}) + # Now a clearly-high loss → top bucket (index 3). + cb.on_log(None, _FakeState(50), None, logs={"loss": 99.0}) + # The high-loss sample must be in the top bucket. + assert 3 in cb._stats + assert cb._stats[3]["num_samples"] >= 1.0 + + def test_length_metric_falls_back_to_round_robin(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + cb = self._callback(tmp_path, "length") + # No length signal in logs → round-robin (step % num_buckets). + cb.on_log(None, _FakeState(2), None, logs={"loss": 99.0}) + # step 2 % 4 buckets == 2. + assert 2 in cb._stats + + def test_default_metric_is_length_round_robin(self, tmp_path, monkeypatch): + from soup_cli.monitoring.curriculum_callback import ( + DynamicCurriculumCallback, + ) + from soup_cli.utils.curriculum_dynamic import DynamicCurriculumPolicy + + monkeypatch.chdir(tmp_path) + policy = DynamicCurriculumPolicy(num_buckets=4, recompute_every_n_steps=10) + cb = DynamicCurriculumCallback(policy=policy, output_dir=str(tmp_path)) + assert cb.curriculum_metric == "length" + cb.on_log(None, _FakeState(3), None, logs={"loss": 1.0}) + assert 3 in cb._stats # round-robin: step 3 % 4 + + def test_perplexity_metric_buckets_like_loss(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + cb = self._callback(tmp_path, "perplexity") + for step in range(1, 9): + cb.on_log(None, _FakeState(step), None, logs={"loss": 0.1 * step}) + cb.on_log(None, _FakeState(50), None, logs={"loss": 20.0}) + assert 3 in cb._stats + + def test_attach_threads_curriculum_metric(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + from types import SimpleNamespace + + from soup_cli.utils.peft_wiring import attach_curriculum_callback + + captured = {} + + class _FakeTrainer: + def add_callback(self, cb): + captured["cb"] = cb + + tcfg = SimpleNamespace( + curriculum_dynamic=True, + curriculum_buckets=4, + curriculum_metric="loss", + curriculum_dynamic_recompute_steps=10, + curriculum_dynamic_floor=0.05, + curriculum_dynamic_temperature=1.0, + ) + ok = attach_curriculum_callback(_FakeTrainer(), tcfg, str(tmp_path)) + assert ok is True + assert captured["cb"].curriculum_metric == "loss" + + +# =========================================================================== +# #157 — extend --hub to data push / data forge +# =========================================================================== + + +class TestDataPushHub: + def test_hub_flag_in_help(self): + from soup_cli.cli import app + + result = runner.invoke(app, ["data", "push", "--help"]) + assert result.exit_code == 0, (result.output, repr(result.exception)) + import re as _re + + clean = _re.sub(r"\x1b\[[0-9;]*m", "", result.output).replace("\n", " ") + assert "--hub" in clean + + def test_invalid_hub_rejected(self, tmp_path, monkeypatch): + from soup_cli.cli import app + + monkeypatch.chdir(tmp_path) + (tmp_path / "ds.jsonl").write_text('{"a": 1}\n', encoding="utf-8") + result = runner.invoke( + app, + ["data", "push", "--input", "ds.jsonl", + "--hf-dataset", "user/ds", "--hub", "bogus"], + ) + assert result.exit_code == 2 + + def test_modelscope_routes_via_upload_repo(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import hubs + + captured = {} + + def _fake_upload(hub, repo_id, **kw): + captured["hub"] = hub + captured["repo_id"] = repo_id + captured.update(kw) + # The staging dir is cleaned up after this returns, so check the + # file is present at call time (proves the JSONL was staged). + captured["staged_file_present"] = ( + Path(kw["folder_path"]) / "ds.jsonl" + ).exists() + # upload_repo enforces cwd-containment on folder_path; the staging + # dir MUST be under cwd (regression guard for the system-tempdir + # bug found in the v0.71.5 step-6 smoke). + from soup_cli.utils.paths import is_under_cwd + + captured["folder_under_cwd"] = is_under_cwd(kw["folder_path"]) + + monkeypatch.setattr(hubs, "upload_repo", _fake_upload) + monkeypatch.chdir(tmp_path) + (tmp_path / "ds.jsonl").write_text('{"a": 1}\n', encoding="utf-8") + result = runner.invoke( + app, + ["data", "push", "--input", "ds.jsonl", + "--hf-dataset", "user/ds", "--hub", "modelscope"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert captured["hub"] == "modelscope" + assert captured["repo_id"] == "user/ds" + assert captured["repo_type"] == "dataset" + assert captured["staged_file_present"] is True + assert captured["folder_under_cwd"] is True + + def test_modelers_missing_sdk_friendly_error(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import hubs + + def _boom(hub, repo_id, **kw): # noqa: ARG001 + raise ImportError("openmind_hub is not installed. pip install ...") + + monkeypatch.setattr(hubs, "upload_repo", _boom) + monkeypatch.chdir(tmp_path) + (tmp_path / "ds.jsonl").write_text('{"a": 1}\n', encoding="utf-8") + result = runner.invoke( + app, + ["data", "push", "--input", "ds.jsonl", + "--hf-dataset", "user/ds", "--hub", "modelers"], + ) + assert result.exit_code == 1 + assert "openmind_hub" in result.output + + def test_generic_upload_error_exit_1(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import hubs + + def _boom(hub, repo_id, **kw): # noqa: ARG001 + raise RuntimeError("network down") + + monkeypatch.setattr(hubs, "upload_repo", _boom) + monkeypatch.chdir(tmp_path) + (tmp_path / "ds.jsonl").write_text('{"a": 1}\n', encoding="utf-8") + result = runner.invoke( + app, + ["data", "push", "--input", "ds.jsonl", + "--hf-dataset", "user/ds", "--hub", "modelscope"], + ) + assert result.exit_code == 1 + assert "Upload failed" in result.output + + def test_hf_default_unchanged_no_token(self, tmp_path, monkeypatch): + # Default --hub hf with no token → existing "no token" error (regression). + from soup_cli.cli import app + from soup_cli.utils import hf as _hf + + monkeypatch.setattr(_hf, "resolve_token", lambda: None) + monkeypatch.chdir(tmp_path) + (tmp_path / "ds.jsonl").write_text('{"a": 1}\n', encoding="utf-8") + result = runner.invoke( + app, + ["data", "push", "--input", "ds.jsonl", "--hf-dataset", "user/ds"], + ) + assert result.exit_code == 1 + assert "token" in result.output.lower() + + +class TestDataForgeHub: + def _docs(self, tmp_path): + docs = tmp_path / "docs" + docs.mkdir() + (docs / "a.txt").write_text( + "Paragraph one with content.\n\nParagraph two here.\n", + encoding="utf-8", + ) + return docs + + def test_hub_flag_in_help(self): + from soup_cli.cli import app + + result = runner.invoke(app, ["data", "forge", "--help"]) + assert result.exit_code == 0, (result.output, repr(result.exception)) + import re as _re + + clean = _re.sub(r"\x1b\[[0-9;]*m", "", result.output).replace("\n", " ") + assert "--hub" in clean + + def test_invalid_hub_rejected(self, tmp_path, monkeypatch): + from soup_cli.cli import app + + monkeypatch.chdir(tmp_path) + self._docs(tmp_path) + result = runner.invoke( + app, + ["data", "forge", "--docs", "docs", "--hub", "bogus", + "--teacher", "owner/repo"], + ) + assert result.exit_code == 2 + + def test_non_hf_prefetches_teacher(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import hubs + + captured = {} + + def _fake_prefetch(base, hub, **kw): # noqa: ARG001 + captured["base"] = base + captured["hub"] = hub + return str(tmp_path / "cache" / "teacher") + + monkeypatch.setattr(hubs, "prefetch_model_from_hub", _fake_prefetch) + monkeypatch.chdir(tmp_path) + self._docs(tmp_path) + result = runner.invoke( + app, + ["data", "forge", "--docs", "docs", "--hub", "modelers", + "--teacher", "owner/teacher-model", "--target-rows", "2"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert captured["base"] == "owner/teacher-model" + assert captured["hub"] == "modelers" + + def test_hf_does_not_prefetch(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import hubs + + called = {"n": 0} + + def _fake_prefetch(base, hub, **kw): # noqa: ARG001 + called["n"] += 1 + return "x" + + monkeypatch.setattr(hubs, "prefetch_model_from_hub", _fake_prefetch) + monkeypatch.chdir(tmp_path) + self._docs(tmp_path) + result = runner.invoke( + app, + ["data", "forge", "--docs", "docs", "--target-rows", "2", + "--teacher", "owner/teacher-model"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert called["n"] == 0 # --hub hf default → no prefetch + + def test_non_hf_bare_teacher_warns_no_prefetch(self, tmp_path, monkeypatch): + # Non-HF hub but teacher lacks owner/name → warn, do NOT prefetch + # (code-review MEDIUM fix — no silent no-op of --hub). + from soup_cli.cli import app + from soup_cli.utils import hubs + + called = {"n": 0} + monkeypatch.setattr( + hubs, "prefetch_model_from_hub", + lambda *a, **k: (called.__setitem__("n", called["n"] + 1), "x")[1], + ) + monkeypatch.chdir(tmp_path) + self._docs(tmp_path) + result = runner.invoke( + app, + ["data", "forge", "--docs", "docs", "--hub", "modelers", + "--teacher", "barename", "--target-rows", "2"], + ) + assert result.exit_code == 0, (result.output, repr(result.exception)) + assert called["n"] == 0 + assert "ignored" in result.output.lower() + + def test_non_hf_prefetch_import_error_exit_1(self, tmp_path, monkeypatch): + from soup_cli.cli import app + from soup_cli.utils import hubs + + def _boom(base, hub, **kw): # noqa: ARG001 + raise ImportError("openmind_hub is not installed") + + monkeypatch.setattr(hubs, "prefetch_model_from_hub", _boom) + monkeypatch.chdir(tmp_path) + self._docs(tmp_path) + result = runner.invoke( + app, + ["data", "forge", "--docs", "docs", "--hub", "modelers", + "--teacher", "owner/teacher", "--target-rows", "2"], + ) + assert result.exit_code == 1 + assert "openmind_hub" in result.output