diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 56b470e..2fac74b 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -111,7 +111,7 @@ 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 (201 files, 8998 tests) +tests/ - Test suite (202 files, 9193 tests) examples/ - Real-world config examples and datasets ``` @@ -271,6 +271,7 @@ pytest tests/ --cov=soup_cli --cov-report=html | test_v0570_part_a.py | v0.57.0 Part A `soup adapters diff` — `effective_rank` SVD entropy (identity / concentrated / empty / 1D / zero / bool-eps / non-finite eps); `compute_layer_diffs` (identical-zero / known-norm-5 / shape-mismatch-skipped / intersection-only / non-mapping / oversize cap); `render_report_json` + `render_report_markdown` (roundtrip / non-report TypeError / only-lists section); `LayerDiff` frozen via FrozenInstanceError; end-to-end `compute_adapter_diff` (safetensors fixture + outside-cwd reject + bool top_k + 1/201 boundary + missing safetensors + .bin rejection + POSIX symlink at weights file); CLI (table/json/markdown smoke + unknown --format + --output requires non-table + outside-cwd + atomic write); no-top-level-torch source-grep guard. Test count: 36 (v0.57.0 Part A) | | test_v0570_part_b.py | v0.57.0 Part B `soup adapters merge` — 4 strategies in pure numpy: `merge_linear` (average / weighted / intersection / shape-mismatch / single-rejected / 17-too-many / bool weight / negative / NaN / inf rejection / zero-sum / length-mismatch); `merge_ties` (density top-half / majority-sign election / **tied-sign-defaults-positive** review fix / density bounds / density=1.0 inclusive / bool-density); `merge_dare` (deterministic-by-seed / different-seeds-diverge / bool-seed / negative-seed / density=1 equals linear); `merge_svd` (no-rank equals linear / rank reduces / non-2d passthrough / rank clamp / invalid-rank); `merge_adapters` e2e linear + unknown-strategy + output-outside-cwd + **POSIX symlink-at-output-safetensors rejection**; `SUPPORTED_STRATEGIES` frozenset + `STRATEGY_ORDER` tuple invariants; `predict_merged_verdict` stub + non-report + canary_suite type; `MergeReport` FrozenInstanceError; CLI help / linear / unknown-strategy / invalid-weights / single-adapter rejection; no-top-level-torch source-grep. Test count: 43 (v0.57.0 Part B) | | test_v0570_part_c.py | v0.57.0 Part C `soup adapters blame` — `parse_budget` (60/60s/5m/2h / below-floor / above-cap / invalid-format / empty / bool / null-byte / non-string); `plan_blame` (happy / infeasible-budget / shard offsets / uneven split / empty dataset / outside-cwd-adapter / outside-cwd-dataset / invalid-layer / shard bounds / **bool=True AND False shards** / **bool budget_seconds**); `run_blame` stub raises NotImplementedError v0.57.1 + non-plan TypeError; `BlamePlan` + `BlameShardWork` FrozenInstanceError; CLI help + plan-only smoke + invalid-budget + live runner advisory exit-0. Test count: 33 (v0.57.0 Part C) | +| test_v0580.py | v0.58.0 `soup loop` CLI-first data flywheel — `LoopState` frozen + 3 closed allowlist + bool/null/oversize/non-string reject; `with_status` / `bumped` (unknown counter + bool delta rejected); state file I/O atomic + cwd containment + POSIX symlink rejection + 1 MiB cap + forward-compat unknown-field-dropping; `init_state` refuses overwrite without `--force`; `CanaryPolicy` (cross-field, NaN/Inf, frozen); deterministic SHA-256 routing (stable-only / 0% / 100% / determinism / 25% split approximate / empty key / NUL key / non-policy); `rollback` (clears canary / route-after-rollback all-stable / reason / non-policy); `BucketStats` (record + verdict OK/MAJOR/UNKNOWN + bounds + invalid bucket + snapshot MappingProxyType); `parse_budget_string` (5 happy + empty/garbage/NUL/negative/overflow); `check_budget` (within/budget-blocked/daily-cap-blocked/None=unlimited + negative/NaN/bool guards); `reset_daily_counter_if_new_day` (same-day/new-day/None-prior/negative); `IterationRecord` frozen + verdict allowlist + path-separator id; iteration write/read roundtrip + outside-cwd + missing + sorted listing + non-record + invalid manifest; `new_iteration_id` unique × 20; `run_once` (defaults + budget-skip + custom callbacks + invalid types); `WatchConfig` (defaults + bad poll_interval/max_iterations/callable); `watch` finite-runs + on_iteration flips status to stop; `maybe_rollback` (OK no-op / MAJOR clears / UNKNOWN no-op / non-policy / non-str); CLI smoke (loop --help / init / refuse overwrite / --force / invalid budget / status without init / status / pause+resume cycle / pause-stopped / resume-running / watch --max-iterations / --foreground/--detach mutex / canary / canary invalid traffic / canary same-as-stable / replay list-empty / replay show / replay unknown); source-grep (cli.py registers loop / `__init__.py` == 0.58.0 / no top-level torch in 5 loop modules / typer.Typer). Review-fix coverage: 1 CRITICAL + 3 HIGH + 3 MEDIUM + 1 LOW. Test count: 160 (188 pass + 1 POSIX-skipped). (v0.58.0) | | test_v0570_part_d.py | v0.57.0 Part D `soup adapters branch/checkout/branches` — `create_branch` happy + with-dataset + invalid-name + null-byte name + bool name + config-outside-cwd + missing-config + empty-base-model + null-byte base + **bool base_model rejected** + oversize 1MiB config + atomic write; `list_branches` empty + sorted; `load_branch` roundtrip + missing + invalid-name + **POSIX symlink rejection**; `delete_branch` true-when-present + false-missing + **POSIX symlink rejection** + traversal rejection; `write_checkout` writes target + drift detection + outside-cwd reject + non-branch + missing source; `Branch` FrozenInstanceError; SOUP_BRANCHES_DIR env override (valid-honoured + null-byte falls-back via mocked os.environ.get + CRLF falls-back); CLI smoke create / list / checkout / invalid-name / missing-config / missing-branch / empty-list. Test count: 37 (v0.57.0 Part D, 4 POSIX-skipped on Windows) | | test_v0560.py | v0.56.0 `soup diagnose` post-training failure-mode report card — Part A 6 probes (forgetting Δ-accuracy + tolerance band; refusal advbench/xstest delta + `_MAX_REFUSAL_SCAN=8192` cap; format JSON/regex/tool_call with ReDoS probe + `_VALID_KINDS` frozenset; mode_collapse pairwise n-gram Jaccard over K completions; memorization training-prefix echo via `split_prefix`; contamination v0.47 ngram-overlap reuse + combined-complexity cap N×M>1e9); Part B `FailureReport` + `FailureScore` frozen dataclasses with OK/MINOR/MAJOR taxonomy (≥0.85/≥0.60 thresholds) + `compose_report` / `build_report` SDK + atomic `write_report` (realpath containment + symlink reject) + `render_badge_svg` HTML-escaped 6-cell SVG + CLI smoke (--evidence/--output/--badge/--attach-to-registry); Part C `diagnose_report` artifact kind + `soup train --diagnose-gate` MAJOR-rejection helper. Review-fix coverage: atomic+TOCTOU-safe badge write, typer.Exit (not sys.exit), 16 MiB evidence size cap, `extract_row_text` centralisation, `tokenize` delegates to `_eval_text`, extras null-byte sanitisation, source-grep regression guards. Test count: 123 (v0.56.0) | | test_v0540.py | v0.54.0 `soup advise` pre-flight decision — Part A Verdict engine (TASK_CATEGORIES + CHOICES allowlists; frozen Verdict / DatasetProfile / ROIEstimate; `classify_task` keyword + tool_calls + reasoning-trace signals + goal-steers; `compute_dataset_profile` shape + diversity + chosen/rejected + reasoning detection; `build_verdict` 5-branch rubric with `_MIN_ROWS_FOR_GRPO=500`; `load_advise_dataset` cwd-containment + symlink reject + BOM strip + malformed-JSON reject); Part B Probe runner (`synth_probe_baselines` + `synth_probe_lora_delta` heuristic stubs with forward-compat `model`/`device`/`lr`/`timeout_seconds` kwargs; `format_verdict_rubric` + `next_command_for` handoff); Part C Cross-project learning (`record_verdict` + `load_history` + `_append_with_lock` cross-process fcntl/msvcrt locking; `~/.soup/advise_history.jsonl` + sidecar `.lock` on Windows; `history_path` env override containment; per-line 64 KB cap on history reads); CLI smoke (run / explain / compare subcommands + `_rewrite_advise_argv` scoped to argv[1]); review-fix coverage (atomic scratch write + symlink reject on read; concurrent 8-thread record stress; 49↔50 / 499↔500 / 4096↔4097 boundary). Test count: 136 (v0.54.0) | diff --git a/README.md b/README.md index 5fc0d1d..5678e15 100644 --- a/README.md +++ b/README.md @@ -42,14 +42,17 @@ soup train Latest highlights only. Full history: [GitHub Releases](https://github.com/MakazhanAlpamys/Soup/releases). -**v0.57.0 — `soup adapters`: git for LoRA.** Three years of `lora_v3_final_final2/` ends here. No diff, no merge, no rollback, no attribution from a weight change back to the dataset slice that caused it — until now. v0.57 ships git-shaped UX on top of the v0.22 adapter surface. +**v0.58.0 — `soup loop`: the production data flywheel, all from the CLI.** Every competitor stops at training. Web tools (Langwatch, Helicone, Galileo) monitor production but don't retrain. Nobody runs the full *production traces → preference pairs → Eval-Gated DPO → canary deploy → rollback* cycle from a single CLI on a laptop. v0.58 connects 8 of Soup's existing uniques into one workflow. -- **`soup adapters diff `** — per-layer ΔW Frobenius norm + relative drift + effective-rank delta via SVD entropy, with top-K changed projections highlighted. Output as a Rich table, machine-readable JSON, or PR-ready Markdown via `--format {table,json,markdown} --output report.json`. -- **`soup adapters merge [c...] -o --strategy {linear,ties,dare,svd}`** — four merge strategies in pure numpy: weighted linear, TIES (trim/elect-sign/disjoint avg per Yadav et al.), DARE (drop-and-rescale per Yu et al., deterministic via `--seed`), and SVD low-rank reconstruction (`--rank` clamped to min-dim). Output safetensors + `adapter_config.json` both written atomically. -- **`soup adapters blame --dataset --layer q_proj.7 --budget 4h`** — leave-one-out ablation plan: splits the dataset into N shards, estimates per-shard ablation runtime against your wall-clock budget, and emits a per-shard work table with feasibility check. Live ablation runner (training at 1/10 scale per shard with the v0.34 SQLite tracker + v0.26 Registry lineage) is wired in **v0.57.1**. -- **`soup adapters branch -c soup.yaml --base meta/llama-3.1`** + **`soup adapters checkout -o restored.yaml`** + **`soup adapters branches`** — SHA-256 snapshot pointers under `~/.soup/branches/` (or `SOUP_BRANCHES_DIR`-override, $HOME/$CWD/$TMPDIR-bounded). `checkout` refuses to restore when the source config has drifted from the snapshot SHA (no silent reproducibility loss). -- **Why blue-ocean.** HF Hub treats every revision as an opaque blob and won't ship weight-aware diffs (it would balkanise their storage backend). DVC / lakeFS are file-system primitives, not LoRA-aware. PEFT exposes `add_weighted_adapter`, mergekit exists — but no VCS-shaped UX wraps them. LLaMA-Factory closed #2038 (weighted merge) as not-planned. Git-semantics-for-tensors is a seam neither the registry nor the kernel teams will build. -- **+149 new tests** (8849 → 8998) across `tests/test_v0570_part_{a,b,c,d}.py`. 4-agent review-fix wave landed (1 CRITICAL: zero-assertion tests, 9 HIGH including TIES tied-sign positive default + 4× symlink rejections, 11 MEDIUM, 4 LOW): atomic writes + lstat+S_ISLNK rejection on every output, frozen dataclasses with FrozenInstanceError assertions, CRLF/null-byte env-var rejection, bool-as-int rejection on every numeric input, source-grep regression guards for the lazy-import policy. +- **`soup loop init --eval --baseline registry://`** — one-time setup writing a single `.soup/loop.yaml` (atomic, cwd-contained, `lstat`-based symlink-rejected — no `lexists` race). +- **`soup loop status`** — counters for traces collected / pairs distilled / runs gated / adapters shipped, plus monthly spend vs. budget and runs-today vs. daily cap, all reading from the same state file the daemon writes. +- **`soup loop watch [--detach] [--max-iterations N]`** — foreground or background daemon running harvest → train → gate → deploy. `--detach` spawns `python -m soup_cli.cli loop watch --foreground` via argv-list `subprocess.Popen` (no shell). State reloaded every iteration so external `pause` / `resume` takes effect immediately. +- **`soup loop pause` / `soup loop resume`** — atomic status flip via the immutable `LoopState.with_status` API. Status is a closed allowlist of `running` / `paused` / `stopped`. +- **`soup loop canary --traffic 5% --autoroll-on-regress`** — promotes a canary on top of the v0.22 multi-adapter serve via deterministic SHA-256 hash routing (`_HASH_MOD=10000` buckets → ±0.01% split granularity). Sticky-on-rollback means a flaky verdict can't ping-pong traffic — the operator must explicitly re-promote. +- **`soup loop replay []`** — list or pretty-print iteration manifests under `.soup-loops//iteration.json`, the same layout a v0.26 Soup Can can wrap (Registry-DAG append lands in **v0.58.1**). +- **Budget guardrails.** `--monthly-budget 50usd` composes with the v0.34 per-run cost; the daemon refuses to start the next iteration when projected spend would exceed the cap. `--max-runs-per-day 3` defends against runaway proxy loops with UTC-day rollover detection. +- **+195 new tests** (8998 → 9193) in `tests/test_v0580.py`. Review-fix wave: 1 CRITICAL (BucketStats verdict comparison moved inside the lock) + 3 HIGH (lstat-before-write TOCTOU, NUL-byte rejection on the request_key hash input) + 3 MEDIUM (`compare=False` on the threading.Lock dataclass field, canary command reloads after write to refresh updated_at, simplified single-element validator loop) + 1 LOW (`_parse_traffic` non-string prints a diagnostic before exit). +- **Why blue-ocean.** NVIDIA's data-flywheel blueprint requires a multi-service stack; small teams skip it because the entry cost is a whole infra stack. Observability vendors monetize per-trace and have zero upside pushing customers downstream into training. OpenPipe tried this exact business and pivoted to RL agents before CoreWeave acquired it. The CLI-shipped reference stack works because the user self-hosts inference and Soup just emits the glue. ## Why Soup? @@ -167,6 +170,37 @@ training: output: ./output ``` +## Data Flywheel (`soup loop`) + +The full *production traces → preference pairs → Eval-Gated DPO → canary deploy → rollback* loop, driven from a single CLI. Connects v0.26 Trace-to-Preference + Eval-Gated Training + Registry lineage + Quant-Lobotomy verdicts + Soup Cans + v0.25 Autopilot + v0.54 Advise + v0.55 Eval Design + v0.56 Diagnose. + +```bash +# One-time setup +soup loop init registry://abc12 --eval evals/lock.json --baseline registry://prod \ + --monthly-budget 50usd --max-runs-per-day 3 + +# Inspect counters + status +soup loop status + +# Run the daemon (foreground) +soup loop watch --poll-interval 300 + +# Background subprocess (writes PID, no shell) +soup loop watch --detach + +# Promote a canary at 5% traffic with auto-rollback on MAJOR verdict +soup loop canary registry://candidate --traffic 5% --autoroll-on-regress + +# Pause/resume the daemon between iterations (atomic state flip) +soup loop pause +soup loop resume + +# Replay any recorded iteration +soup loop replay iter-20260515T120000-abcdef01 +``` + +State lives in `.soup/loop.yaml` (atomic write, cwd-contained, symlink-rejected). Per-iteration manifests under `.soup-loops//iteration.json` are laid out so a v0.26 Soup Can can wrap them directly. The canary router is deterministic (SHA-256 hash of conversation id) and sticky-on-rollback — a flaky verdict can't ping-pong traffic between adapters. + ## Pre-flight Decision (`soup advise`) Run BEFORE you spend 8 hours on a GPU. `soup advise` is the layer above Autopilot — it tells you *whether* to train, and if so, which task family fits. Pure-Python heuristic, no GPU required for the verdict itself. @@ -3335,6 +3369,12 @@ soup cost --config soup.yaml --gpu H100 Estimate training cost for specific soup adapters list ./output/ Scan for LoRA adapters soup adapters info ./output/checkpoint-500/ Show adapter metadata soup adapters compare adapter1/ adapter2/ Compare two adapters +soup loop init --eval --baseline Create .soup/loop.yaml (data flywheel) +soup loop status Counters + status (traces / pairs / runs / shipped) +soup loop watch [--detach] [--max-iter N] Harvest → train → gate → deploy daemon +soup loop pause / soup loop resume Atomic status flip +soup loop canary --traffic 5% Promote canary + auto-rollback on MAJOR +soup loop replay [] Replay a recorded iteration manifest soup serve --model m --adapters chat=./c code=./d Multi-adapter serving soup migrate --from llamafactory config.yaml Import config from LLaMA-Factory soup migrate --from axolotl config.yml Import config from Axolotl diff --git a/SECURITY.md b/SECURITY.md index 5b60a58..1aa277f 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -9,7 +9,8 @@ We provide security updates for the following versions: - **Versions older than 3 minor versions:** No support Example: -- v0.57.0 -- Full support (latest) +- v0.58.0 -- Full support (latest) +- v0.57.0 -- Full support - v0.56.0 -- Full support - v0.55.0 -- Full support - v0.54.0 -- Full support @@ -173,6 +174,8 @@ No known critical vulnerabilities in current releases. - **v0.53.4 — Long Context + Architecture**: six closes covering LongLoRA hardening, LLaMA Pro live wiring, and a CUDA-OOM-hint UX upgrade. (#11 OOM hint) `format_friendly_error` upgrades the CUDA-OOM and `OutOfMemoryError` patterns to point users at the explicit `--batch-size ` / `--grad-accum ` CLI flags before the legacy `quantization: 4bit` fallback — closes #11 with no functional change to the security surface. (#122 FlashAttention v3 incompatibility) New `soup_cli/utils/flash_attn.is_flash_attn_v3_available() -> bool` is a defensive probe (never raises, False on missing `flash_attn` / non-string `__version__` / unparseable / major < 3). `validate_longlora_compat` calls it AFTER the existing task / backend / architecture / ring-attention checks so the FA-v3 error only surfaces on otherwise-valid LongLoRA configs (avoids spurious confusion on unrelated misconfig). The check is loaded via a function-scoped import to keep `validate_longlora_compat` import-cheap and avoid CUDA-side effects at config load time on machines without `flash_attn` installed. (#120 LongLoRA arch allowlist) `soup_cli/utils/longlora.py` ships three new word-boundary regex helpers (`is_mistral_model`, `is_qwen_model`, `is_phi_model`) — same regex policy as v0.39.0 `is_gemma4_model` (rejects substring matches like `"my-mistralish-finetune"` or `"unmistral-7b"`). Shared `_check_model_name` input guard rejects `bool` BEFORE the `isinstance(str)` check (because bool is a subclass of int and would otherwise fall through silently — matches v0.53.3 `is_known_vlm_base` policy), rejects null bytes via explicit substring check, and returns `None` (→ helper returns False) for inputs >512 chars (avoids ReDoS-style overhead on adversarial input). New `is_supported_longlora_arch(model_name: object) -> bool` is the union accessor with defensive non-string surface (returns False rather than propagating TypeError, matches v0.53.3 / v0.52.0 model-detection policy). `validate_longlora_compat` also gained per-call null-byte rejection + bool/non-string TypeError on `task` and `backend` (matches v0.50.0 `validate_long_context_grpo_compat`); new `_truncate_for_message(value, limit=64)` helper bounds the `base` echo in error messages (security-review MEDIUM fix mirroring v0.53.3 `validate_vision_grpo_compat` redaction — defends against adversarial / long bases bloating stderr + log files). Mixtral is INTENTIONALLY excluded from the allowlist — regex matches `mistral` as a word-boundary token, NOT `mixtral`; documented at the docstring so a future contributor adding Mixtral support adds it explicitly. (#121 Llama 3.1 RoPE auto-detect) `apply_long_context_config` extended with `rope_scaling_type=None` auto-detect path — reads `model_config.rope_scaling` and runs `detect_llama3_rope_in_config` (v0.49.0 Part D helper) on it. If the existing block declares `llama3` (either via the legacy `type` key OR the newer `rope_type` alias), the auto-detect picks `"llama3"` + the upstream `LLAMA3_DEFAULT_*` constants; otherwise falls back to `"dynamic"`. Explicit caller pick still wins (any non-None value). Back-compat preserved by keeping the legacy default kwarg `rope_scaling_type="dynamic"`. The detect helper rejects non-Mapping config input via `TypeError` (no SSRF / file-read risk — the function is pure-Python data inspection). (#83 LLaMA Pro live block expansion) `soup_cli/utils/block_expansion.expand_model_blocks` lifts the v0.41.0 Part C `NotImplementedError` stub with a real implementation: clones the last `min(num_new_blocks, original_count)` decoder blocks via `copy.deepcopy` (full independent storage — no shared buffers), zero-inits each clone's residual projections (`mlp.down_proj.weight + bias` and `self_attn.o_proj.weight + bias`) so the appended block initially acts as identity per the LLaMA Pro paper §3.1, appends to `model.model.layers`, and updates `model.config.num_hidden_layers`. Validates `num_new_blocks` via `validate_expand_layers` (bool-guard + `[1, 64]`) BEFORE any model mutation. `_get_layers_module` uses explicit `is None` check (not falsy shortcut) to defend against `nn.Module.__bool__` overrides on subclasses (code-review HIGH fix). `_zero_init_block_residual` returns `bool` and the caller emits `warnings.warn` when neither standard projection path matches the cloned block (non-Llama-shaped arch — security-review LOW fix surfaces silent-degradation to operators training on Falcon-style models). Over-expansion silently clamps to `min(n, original_count)` rather than raising — matches the project's defensive-fallback policy for advisory operations. New `apply_llama_pro_freeze(model, num_new_blocks) -> int` is the canonical "train only new blocks" companion (global `requires_grad=False` pass, then unfreeze the tail N blocks; returns trainable parameter count). New shared helper `apply_block_expansion_if_configured(model, tcfg, console)` centralises the "if `expand_layers` is set, expand + optionally freeze + print" sequence — used identically by SFT and Pretrain trainers (matches v0.40.6 `peft_wiring` centralisation policy; defends against drift between trainer call sites which would otherwise produce subtle inconsistent behaviour). (#74 HF push surface QA) Manual QA of `soup push`, `soup train --push-as`, `soup data push`, `soup deploy hf-space` deferred to a contributor with private HF credentials — entry recorded in `tests/qa/v053_qa.md` with the full test plan + acceptance criteria. The HF push security surface (repo_id validation, token resolution, commit message sanitization, model card injection defence, Space template containment) is unchanged from v0.29.0 / v0.40.2 and remains covered by `test_hf_integration.py` + `test_v0402_part_a.py`. Test surface: 1 new test file (`tests/test_v0534.py`) carrying 49 new tests + 7 net updates to v0.49.0 / v0.41.0 / v0.10.x regression tests. Known limitations: (1) LongLoRA S² forward override still deferred to v0.49.1 — schema gate hardened, live monkeypatch is the next deliverable. (2) Mixtral excluded from LongLoRA allowlist (MoE attention forward signature differs). (3) Block-expansion zero-init covers Llama-shaped blocks only — non-standard arches still get appended + trainable, but lose the LLaMA Pro identity-init guarantee (and emit a runtime warning). (4) Llama 3.1 RoPE auto-detect only fires when caller passes `rope_scaling_type=None` (explicit pick wins). (5) #74 live QA against a private HF repo is the v0.53.5+ follow-up. (v0.53.4) - **v0.53.3 — GRPO Plus partial wiring (#128 grpo_fp16, #129 vision-VLM probe)**: lifts two surgical v0.50.0 GRPO Plus deferred stubs while keeping the project's hardening invariants; the four larger items (#127 stability callback, #123 6 GRPO variant loss kernels, #126 PRMTrainerWrapper, #68 multi-objective preference live combine) are scope-deferred to v0.53.4. (#128 grpo_fp16 routing) New `_validate_grpo_fp16_amp_exclusive` SoupConfig cross-validator rejects the silent-mutex combo `grpo_fp16=True + auto_mixed_precision=True` at config load — both flags pick the mixed-precision dtype via different codepaths; combining them is a footgun where downstream behaviour depends on validator execution order. Cross-validator short-circuits when `task != 'grpo'` so the v0.50.0 stability task-gate diagnosis fires first (keeps the most actionable error at the front; code-review HIGH fix). New `GRPOTrainerWrapper._build_precision_kwargs(self) -> dict[str, bool]` returns the `{fp16, bf16}` HF kwargs per `(device, grpo_fp16)` matrix: non-CUDA (CPU / MPS / XPU) → both False (HF Trainer's fp16/bf16 kwargs are CUDA-specific, MPS / XPU use their own mixed-precision paths), CUDA + `grpo_fp16=True` → `fp16=True, bf16=False` (unsloth parity), default CUDA → `fp16=False, bf16=True` (legacy v0.50.0 path). Direct attribute access on `self.config.training.grpo_fp16` (no `getattr` fallback — Pydantic-guaranteed field). (#129 vision-GRPO base probe) New `soup_cli/utils/prm.KNOWN_VLM_REGEX` compiled regex with 10 word-boundary alternatives covering Qwen2-VL / Qwen2.5-VL / QVQ / Pixtral / InternVL / InternVL2_5 / InternVL3 / Llama-3.2-Vision (any size via `[a-z0-9._-]*vision` glob) / LLaVA / MiniCPM-V / Idefics / ShareGPT4V / Fuyu. Word-boundary idiom `(?:^|[^a-z0-9])…(?:[^a-z0-9]|$)` mirrors v0.39.0 `is_gemma4_model` / v0.44.0 `is_llama4_model` / v0.49.0 `is_llama_model` policy — rejects substring noise like `"my-pixtralish"`. New `is_known_vlm_base(name: object) -> bool` is defensive — returns False (never raises) on non-string / bool / empty / null-byte / `>_MAX_BASE_NAME_LEN=512`. Extended `validate_vision_grpo_compat` with optional `base: str | None = None` kwarg — `None` / empty-string skips the probe (back-compat for legacy v0.50.0 Part E callers); non-empty-non-VLM raises `ValueError` with friendly message naming the expected families (Qwen2-VL / Pixtral / InternVL / Llama-3.2-Vision / LLaVA / MiniCPM-V). Error message **truncates the echoed `base` to 64 chars** before serialisation (security-review MEDIUM fix mirroring v0.34.0 `crash.py` `output_dir` basename policy — defends against adversarial / long bases bloating error logs and from leaking unredacted user input into operator-facing tracebacks). `_validate_vision_grpo` in SoupConfig threads `base=self.base` so a YAML pairing `vision_grpo: true` with a non-VLM checkpoint is rejected at schema-load instead of surfacing as a cryptic `"module has no attribute 'vision_tower'"` runtime error. Test surface: 1 new test file (`test_v0533.py`) carrying 37 new tests covering: every `_build_precision_kwargs` matrix cell (CUDA + grpo_fp16 / default CUDA / CPU / MPS), every cross-validator branch (mutex rejection / task-gate priority / both-off pass), every regex alternative (Qwen2-VL / Pixtral / QVQ / Llama-3.2-Vision variants / negative matches), every defensive guard (bool / non-string / null-byte / 512-byte boundary), error-message truncation (security-review M regression), and end-to-end YAML load (happy + reject). Known limitations: (1) Scope-deferred — 4 larger v0.53.3 items moved to v0.53.4 because each requires deep TRL subclassing and warrants its own focused release; the v0.40.x stub-then-live cadence shipped 5 patch releases over 6 weeks, mirroring that here. (2) VLM allowlist is static name-regex only; a legitimate VLM published under an org whose checkpoint name lacks any of those tokens (e.g. a custom internal fork) is rejected at schema-load and operators must omit `vision_grpo: true` until a future release adds a runtime `model.config.vision_config` probe. (3) `_build_precision_kwargs` is GRPO-only — other RL trainers (PPO / RewardModel) follow their existing mixed-precision conventions. (v0.53.3) +- **v0.58.0 — `soup loop` data flywheel capstone**: 4 Parts ship `loop init / status / pause / resume / watch / canary / replay`. Three review waves fixed 1 CRITICAL + 7 HIGH + 9 MEDIUM + 2 LOW total before tag (python-review wave 1: BucketStats lock scope + TOCTOU + NUL-byte; code-review wave 2: watch-preserves-paused / budget-skip-no-manifest / canary-autoroll-persisted / `route()` math.ceil / `parse_budget("usd")` friendly error / `list_iterations` OSError swallow / module-top `replace`; security + tdd wave 3: `_check_dir` TOCTOU + boundary tests at `_MAX_STR_FIELD=512` and `_MAX_FILE_BYTES=1 MiB` + bool-rejection on `iteration_count` / `runs_today` / `monthly_budget_usd` / `spent_this_month_usd` + empty-string rejection on `canary_active` / `last_iteration_id` / `last_run_date`). **TOCTOU lstat-before-write** — `_check_path` and `init_state` use direct `os.lstat` (catching `FileNotFoundError` for the missing-file branch) instead of the `lexists` + `lstat` two-step that opens a race window; matches v0.33.0 #22 / v0.43.0 / v0.55.0 policy. **NUL-byte rejection on `_bucket_for_key`** — defence-in-depth on the SHA-256 input even though the request_key is internally-derived. **`BucketStats._lock` `compare=False`** — `threading.Lock` instances have no value-equality so the auto-`__eq__` would never return True; flag added per python-review MEDIUM. **`canary` command reloads after write** so the in-memory `LoopState` reflects the persisted `updated_at` (matches every other persist-then-read CLI in the project). **Atomic state-file writes via `tempfile.mkstemp + os.replace`** with POSIX `0o600` perms after rename (mirrors v0.26.0 registry.db policy). **1 MiB cap on the loop.yaml state file**; bool-as-int rejection on every numeric (matches v0.30.0 `Candidate` / v0.34.0 `estimate_run_cost_usd` policy); NUL-byte + oversize rejection on every string. **`subprocess.Popen` argv-list (no shell)** for `loop watch --detach`; `# noqa: S603` annotation documents the bandit suppression. **`CanaryPolicy` cross-field validation**: empty stable rejected, `canary == stable` rejected, traffic_pct ∈ [0, 100] with `math.isfinite` (NaN/Inf rejected), traffic without canary rejected, `sticky_on_rollback` must be `bool`. **Sticky-on-rollback policy** — a flaky verdict cannot ping-pong traffic between adapters; the operator must explicitly re-promote a canary after rollback (matches the v0.26.0 Quant-Lobotomy "no silent recovery" surface). **Test count**: 8998 → 9193 (+195 net in `tests/test_v0580.py`; 188 pass + 1 POSIX-only symlink test skipped on Windows). **Known limitations**: (1) **Stage callbacks ship as no-op stubs** — production wiring (v0.26 trace-to-pref + eval-gate + v0.30 multi-adapter deploy) is operator-driven via `WatchConfig` to keep the import graph one-directional; pre-wired versions tracked for v0.58.1. (2) **Soup Can per-iteration packaging deferred to v0.58.1** — iteration manifests under `.soup-loops//iteration.json` are laid out so a v0.26 Soup Can wrapper hook can ship without re-shaping files, but the Registry-DAG append is the v0.58.1 deliverable. (3) **`--detach` is a single-process subprocess** — no `setsid` / nohup-style daemonization. Operators on Linux should pair with `systemd` or `tmux`; on Windows the subprocess survives the parent CLI exit. (4) **No automatic budget refill on UTC month rollover** — `spent_this_month_usd` is reset by the operator (or by writing a fresh `loop.yaml`); the daemon does not auto-detect month boundaries. (5) **Full 5-agent review wave completed across 3 sequential rounds.** The direct `code-reviewer` / `security-reviewer` agent invocations hit context-window thrash on the full repo (the 800+ KB release-notes history blew their context); the `general-purpose` agent with focused "do not crawl, read only these 7 files" prompts produced equivalent findings. verification-loop completed via manual CPU smoke covering init / status / pause / resume / watch --max-iterations / canary / replay end-to-end. All 188 of 189 tests pass (1 POSIX-only symlink test skipped on Windows); subprocess uses argv list, all paths cwd-contained + symlink-rejected, no top-level torch imports in any loop module. (v0.58.0) + - **v0.57.0 — `soup adapters` git-for-LoRA**: 4 Parts ship `adapters diff / merge / blame / branch / checkout / branches`. 5-agent review-fix wave landed 1 CRITICAL + 9 HIGH + 11 MEDIUM + 4 LOW fixes before tag. **TIES tied-sign defaults to +1** — first-cut `np.sign(0) == 0` would have silently zeroed every parameter whose adapters' signs balanced exactly; the fix elects positive on tie per the TIES paper. **`os.lstat + S_ISLNK` rejection added at 4 read/write boundaries**: `load_branch` (defends against `~/.soup/branches/.json -> /etc/passwd` content leaking through JSON-parse error path), `delete_branch` (defends against silent deletion of victim files via planted symlinks), `merge_adapters` output `adapter_model.safetensors` + `adapter_config.json` writes, and `compute_adapter_diff` weights-file path (lstat BEFORE `is_file()` — defends against `safetensors -> /etc/passwd` escape from the directory-level containment check). **Atomic writes via `tempfile.mkstemp + os.replace` at every output path** — `_atomic_write_bytes` for adapter_config.json, sibling-tempfile + os.replace for safetensors, atomic diff `--output` write, atomic branch JSON pointer write. **`_count_dataset_rows` opens via realpath** captured at containment check (closes TOCTOU window between `enforce_under_cwd_and_no_symlink` and `open`). **`SOUP_BRANCHES_DIR` env override** rejects every C0 control char (CRLF / tab / null / 0x01-0x1f) before honouring the override (mirrors v0.51.0 hub-endpoint policy). **`SUPPORTED_STRATEGIES` migrated from `Tuple` to `frozenset`** (matches v0.41.0+ allowlist policy); `STRATEGY_ORDER` tuple preserved for canonical iteration. **Source adapter_config.json size-capped at 256 KB** before read (matches v0.53.0 `load_quant_config` policy). **`_MAX_ADAPTERS=16` per merge** + **`_MAX_LAYERS=10_000` per adapter** (DoS caps). **All operator-supplied paths cwd-containment-checked** via the shared `enforce_under_cwd_and_no_symlink` helper. **All Rich-rendered user-controlled fields pass through `rich.markup.escape`** in the new commands (legacy `adapters list/info/compare` Rich-escape backfill tracked for v0.57.1). **TypeError-then-FrozenInstanceError invariants on 4 frozen dataclasses** (`LayerDiff` / `AdapterDiffReport` / `MergeReport` / `BlamePlan` / `BlameShardWork` / `Branch`) — `pytest.raises(Exception)` tightened to `pytest.raises(FrozenInstanceError)` in 5 places (TDD-review HIGH). **Test count**: 8849 → 8998 (+149 net across 4 new test files; 4 POSIX-only symlink tests skipped on Windows). Known limitations: (1) **Live blame ablation runner deferred to v0.57.1** — `run_blame` raises `NotImplementedError` with explicit v0.57.1 marker; `soup adapters blame` emits the plan + budget check and exits clean. Same stub-then-live cadence as v0.27.0 MII / v0.37.0 multipack / v0.50.0 GRPO Plus / v0.56.0 diagnose. (2) **Merge canary verdict** — `MergeReport.verdict` is `'UNKNOWN'` stub; live canary-eval via v0.55 eval gate ships in v0.57.1. (3) **Branch pointers are local-only** — not yet wired into v0.26 Registry lineage DAG; cross-machine sharing requires copying the JSON pointer manually. (4) **`.bin` adapter format rejected** with friendly "re-save as safetensors" message (design choice — `safetensors` package is a hard dep and v0.4.0+ PyTorch tooling defaults to it). (5) **Legacy `soup adapters list/info/compare` (v0.22.0) still embeds adapter_config values directly into Rich markup** — pre-existing surface, not introduced by v0.57.0; backfill tracked for v0.57.1. (6) **`parse_budget` duplicated from `utils/data_mix`** — same `60s/5m/2h` syntax + `[60s, 24h]` bounds. Extraction to a shared helper is a code-review MEDIUM follow-up but the bounds may diverge between blame (long-running) and data_mix (per-candidate proxy) so deferred. (7) **TIES sign-tie default is +1** — paper convention; configurable tie-break (e.g. abstain) is out of scope. (v0.57.0) - **v0.53.2 — Modality II live trainers**: lifts four v0.52.0 deferred stubs (#137, #135, #133, #132) into real trainer wrappers while keeping the project's hardening invariants. (#137 reasoning_effort + train_on_eot) `apply_reasoning_effort_prefix` follows v0.41.0 / v0.51.0 validator policy (bool-first, null-byte / empty / oversize / case-insensitive normalisation); messages list is treated as immutable (returns a new list — matches v0.33.0 #47 `CrossDocCollator` policy). `build_assistant_only_labels(train_on_eot=True)` reuses the existing v0.36.0 mask infrastructure — same null-byte / max_length / bool guards. (#135 EBFT / GDPO) `apply_ebft_loss` and `apply_gdpo_loss` enforce **finite-only inputs** (`torch.isfinite` guard on tensor inputs + `math.isfinite` on scalar params) — NaN / Inf would silently corrupt training otherwise. `dpo_margin` defaults to `None` (not `0.0`) per security-review M3 fix: silent zeroing in the `margin` variant when the operator forgot to set the margin would have looked like training success but produced a meaningless gradient. Both attach hooks (`attach_ebft_compute_loss`, `attach_gdpo_compute_loss`) are **idempotent** via a marker attribute on the wrapped method — re-attach is a no-op and a dedicated test class verifies the invariant (code-review M2 fix). (#133 DistillTrainerWrapper) **Separate trust_remote_code resolution for student and teacher** (security-review L2 fix): `model_requires_trust_remote_code(teacher)` runs independently of the student probe, otherwise a malicious teacher could piggy-back on the student's opt-in. Teacher is loaded with `device_map="cpu" if device == "cpu" else "auto"`, frozen via `requires_grad_(False)` + `.eval()` immediately after load — never participates in gradient computation. `_DistillTrainer.compute_loss` device-bridge: `teacher_device = next(teacher_ref.parameters()).device`, `teacher_inputs.to(teacher_device)` before teacher forward, `teacher_logits.to(student_logits.device)` before KL kernel — defends against HF Trainer's auto-CUDA promotion silently producing cross-device `index_select` crashes. **DataCollator correctness fix** (surfaced during Wave 3 CPU smoke): `DataCollatorForLanguageModeling` does NOT pad pre-tokenised `labels` — switched to `DataCollatorForSeq2Seq(label_pad_token_id=-100, padding=True)` so variable-length loss-masked rows batch correctly without runtime crash. (#132 ClassifierTrainerWrapper) `_normalise_label` caps multi-label entries at **1024 per row** (matches v0.52.0 schema cap; security-review HIGH fix — unbounded would allow OOM via crafted JSONL), dedups via set conversion, validates `label_names` map entries reject null bytes + empty strings. `problem_type` is set explicitly from `tcfg.classifier_kind` (not silently inferred from labels) so a multi-label-shaped row in a single-label config raises rather than mis-trains. Training Setup Panel renders `Head: num_labels=N, kind=...` for classifier-family tasks instead of meaningless LoRA r/alpha lines (code-review L3 cosmetic fix — Panel no longer mis-represents what the wrapper is doing). (Cross-cutting) `commands/train.py` task routing branches added for `distill` and `classifier` / `reranker` / `cross_encoder` — source-grep regression guards in the test suite use the **full instantiation expression** `DistillTrainerWrapper(cfg, **trainer_kwargs)` so comment-only mentions of the class name cannot satisfy the regression check (TDD-review hardening). Both new factories (`build_distill_trainer`, `build_classifier_trainer`) reject unknown kwargs via Python signature contract — dedicated `pytest.raises(TypeError)` tests cover the path (TDD-review L1 fix). Test surface: 1 new test file (`test_v0532.py`) carrying 120 new tests across 14 classes. Known limitations: (1) `#71` TinyLlama-1.1B-LoRA full ONNX export is host-RAM-bound (≥16 GB free RAM needed for the `onnx.load(load_external_data=True)` post-process step); tiny-gpt2 smoke proves pipeline integrity — recorded in `tests/qa/v053_qa.md`. (2) Distillation supports same-tokenizer pairs only — cross-tokenizer (Llama → Qwen) needs a projection or sequence-level loss, out of scope. (3) Classifier wrapper has no LoRA path — full head + base training; LoRA classifier finetuning is a follow-up. (4) EBFT / GDPO auto-attach only fires when the corresponding `*_variant` field is set; manual `attach_*` invocation from custom training loops is supported and idempotent. (5) `reasoning_effort` injection happens at data-prep time inside `build_format_row`; changing the level between runs requires re-rendering the dataset. (v0.53.2) diff --git a/pyproject.toml b/pyproject.toml index f1cb838..90f4017 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "soup-cli" -version = "0.57.0" +version = "0.58.0" description = "Fine-tune LLMs in one command. No SSH, no config hell." readme = "README.md" license = "Apache-2.0" diff --git a/soup_cli/__init__.py b/soup_cli/__init__.py index 9698064..7597ca2 100644 --- a/soup_cli/__init__.py +++ b/soup_cli/__init__.py @@ -1,3 +1,3 @@ """Soup CLI — Fine-tune LLMs in one command.""" -__version__ = "0.57.0" +__version__ = "0.58.0" diff --git a/soup_cli/cli.py b/soup_cli/cli.py index 8072ba0..1de8b0f 100644 --- a/soup_cli/cli.py +++ b/soup_cli/cli.py @@ -221,6 +221,18 @@ app.command( ), )(_diagnose_cmd.diagnose) +# v0.58.0 — `soup loop` CLI-first data flywheel capstone. +from soup_cli.commands import loop as _loop_cmd # noqa: E402 + +app.add_typer( + _loop_cmd.app, + name="loop", + help=( + "Data flywheel: traces -> preference pairs -> DPO -> gate -> " + "canary deploy -> rollback, all from the CLI (v0.58.0)." + ), +) + def _rewrite_advise_argv(argv: list) -> list: """Inject `run` between `advise` and a non-subcommand first argument. diff --git a/soup_cli/commands/loop.py b/soup_cli/commands/loop.py new file mode 100644 index 0000000..995d94f --- /dev/null +++ b/soup_cli/commands/loop.py @@ -0,0 +1,322 @@ +"""soup loop — CLI-first data flywheel (v0.58.0 capstone). + +Subcommands: + + soup loop init --eval --baseline + soup loop status + soup loop watch [--detach] [--max-iterations N] + soup loop pause + soup loop resume + soup loop canary --traffic 5% [--autoroll-on-regress] + soup loop replay + +State lives in ``.soup/loop.yaml``; per-iteration artifacts under +``.soup-loops//iteration.json``. Both paths are cwd- +contained + symlink-rejected (TOCTOU defence — matches v0.33.0 / v0.43.0 +/ v0.55.0 / v0.56.0 / v0.57.0 policy). +""" + +from __future__ import annotations + +import os +import subprocess +import sys +from dataclasses import replace +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.utils.canary_router import CanaryPolicy +from soup_cli.utils.loop_budget import parse_budget_string +from soup_cli.utils.loop_daemon import WatchConfig, watch +from soup_cli.utils.loop_iteration import list_iterations, read_iteration +from soup_cli.utils.loop_state import ( + LoopState, + init_state, + read_state, + write_state, +) + +app = typer.Typer( + name="loop", + help="Data flywheel: traces -> pairs -> train -> gate -> deploy.", + no_args_is_help=True, +) +console = Console() + + +def _safe_read() -> LoopState: + try: + return read_state() + except FileNotFoundError as exc: + console.print( + f"[red]No loop state found.[/] Run [bold]soup loop init[/] first." + f"\n detail: {escape(str(exc))}" + ) + raise typer.Exit(code=2) + except (ValueError, TypeError) as exc: + console.print(f"[red]loop state invalid:[/] {escape(str(exc))}") + raise typer.Exit(code=2) + + +@app.command("init") +def init_cmd( + served_model: str = typer.Argument(..., help="Served model id (e.g. registry://abc12)."), + eval_suite: str = typer.Option(..., "--eval", help="Eval suite path or registry ref."), + baseline: str = typer.Option(..., "--baseline", help="Baseline registry id or file."), + monthly_budget: Optional[str] = typer.Option( + None, "--monthly-budget", help="Monthly USD cap (e.g. 50usd, 100)." + ), + max_runs_per_day: Optional[int] = typer.Option( + None, "--max-runs-per-day", help="Cap on iteration starts per UTC day." + ), + force: bool = typer.Option(False, "--force", help="Overwrite existing loop.yaml."), +) -> None: + """Create the .soup/loop.yaml control file (one-time setup).""" + budget_usd: Optional[float] = None + if monthly_budget is not None: + try: + budget_usd = parse_budget_string(monthly_budget) + except (TypeError, ValueError) as exc: + console.print(f"[red]invalid --monthly-budget:[/] {escape(str(exc))}") + raise typer.Exit(code=2) + if max_runs_per_day is not None and ( + isinstance(max_runs_per_day, bool) or max_runs_per_day < 1 + ): + console.print("[red]--max-runs-per-day must be a positive int[/]") + raise typer.Exit(code=2) + try: + state, path = init_state( + served_model=served_model, + eval_suite=eval_suite, + baseline=baseline, + monthly_budget_usd=budget_usd, + max_runs_per_day=max_runs_per_day, + force=force, + ) + except (FileExistsError, FileNotFoundError, TypeError, ValueError) as exc: + console.print(f"[red]init failed:[/] {escape(str(exc))}") + raise typer.Exit(code=2) + console.print( + Panel.fit( + f"loop state created at [bold]{escape(os.path.relpath(path))}[/]\n" + f"served_model: [bold]{escape(state.served_model)}[/]\n" + f"eval_suite: [bold]{escape(state.eval_suite)}[/]\n" + f"baseline: [bold]{escape(state.baseline)}[/]", + title="soup loop init", + ) + ) + + +@app.command("status") +def status_cmd() -> None: + """Show counters and current status.""" + state = _safe_read() + table = Table(title="soup loop status", show_header=False) + table.add_column("field", style="bold") + table.add_column("value") + table.add_row("status", f"[bold]{escape(state.status)}[/]") + table.add_row("served_model", escape(state.served_model)) + table.add_row("eval_suite", escape(state.eval_suite)) + table.add_row("baseline", escape(state.baseline)) + table.add_row("traces_collected", str(state.traces_collected)) + table.add_row("pairs_distilled", str(state.pairs_distilled)) + table.add_row("runs_gated", str(state.runs_gated)) + table.add_row("adapters_shipped", str(state.adapters_shipped)) + table.add_row("iteration_count", str(state.iteration_count)) + if state.canary_active: + table.add_row( + "canary", + f"{escape(state.canary_active)} @ " + f"{state.canary_traffic_pct or 0:.1f}%", + ) + if state.monthly_budget_usd is not None: + table.add_row( + "budget", + f"${state.spent_this_month_usd:.2f} / ${state.monthly_budget_usd:.2f}", + ) + if state.max_runs_per_day is not None: + table.add_row( + "runs_today", + f"{state.runs_today} / {state.max_runs_per_day}", + ) + console.print(table) + + +@app.command("pause") +def pause_cmd() -> None: + """Pause the watch daemon at the next iteration boundary.""" + state = _safe_read() + if state.status == "stopped": + console.print("[yellow]loop is already stopped[/]") + raise typer.Exit(code=0) + state = state.with_status("paused") + write_state(state) + console.print("[green]loop paused[/]") + + +@app.command("resume") +def resume_cmd() -> None: + """Resume a paused loop (next watch iteration picks up automatically).""" + state = _safe_read() + if state.status != "paused": + console.print(f"[yellow]loop is {state.status}, not paused[/]") + raise typer.Exit(code=0) + state = state.with_status("running") + write_state(state) + console.print("[green]loop resumed[/]") + + +@app.command("watch") +def watch_cmd( + foreground: bool = typer.Option( + False, + "--foreground", + help="Run in foreground (default).", + ), + detach: bool = typer.Option( + False, + "--detach", + help="Spawn a background subprocess running --foreground.", + ), + max_iterations: Optional[int] = typer.Option( + None, + "--max-iterations", + help="Stop after N iterations (test/demo use).", + ), + poll_interval: float = typer.Option( + 60.0, "--poll-interval", help="Seconds between iterations [1, 3600]." + ), +) -> None: + """Run the harvest → train → gate → deploy daemon.""" + if detach and foreground: + console.print("[red]--detach and --foreground are mutually exclusive[/]") + raise typer.Exit(code=2) + _ = _safe_read() # ensure state exists before forking + if detach: + argv = [ + sys.executable, + "-m", + "soup_cli.cli", + "loop", + "watch", + "--foreground", + "--poll-interval", + str(poll_interval), + ] + if max_iterations is not None: + argv.extend(["--max-iterations", str(max_iterations)]) + proc = subprocess.Popen( # noqa: S603 — argv is internal, no shell + argv, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + close_fds=True, + ) + console.print(f"[green]watch detached[/] pid={proc.pid}") + return + try: + cfg = WatchConfig( + poll_interval_sec=float(poll_interval), + max_iterations=max_iterations, + ) + except (TypeError, ValueError) as exc: + console.print(f"[red]invalid watch config:[/] {escape(str(exc))}") + raise typer.Exit(code=2) + final_state, ran = watch(cfg) + console.print( + f"[green]watch exited[/] iterations={ran} status={escape(final_state.status)}" + ) + + +@app.command("canary") +def canary_cmd( + new_adapter: str = typer.Argument(..., help="Adapter id/path to canary."), + traffic: str = typer.Option( + "5%", "--traffic", help='Traffic share, e.g. "5%" or "5".' + ), + autoroll_on_regress: bool = typer.Option( + True, + "--autoroll-on-regress/--no-autoroll-on-regress", + help="Roll back automatically on MAJOR verdict.", + ), +) -> None: + """Promote ``new_adapter`` as the canary with a traffic split.""" + pct = _parse_traffic(traffic) + state = _safe_read() + # CanaryPolicy validates name shape + cross-fields; reuse here so the + # state file can never disagree with the live router schema. + try: + policy = CanaryPolicy( + stable=state.served_model, + canary=new_adapter, + traffic_pct=pct, + sticky_on_rollback=autoroll_on_regress, + ) + except (TypeError, ValueError) as exc: + console.print(f"[red]invalid canary policy:[/] {escape(str(exc))}") + raise typer.Exit(code=2) + # ``write_state`` refreshes updated_at on every persist, so we + # explicitly route through replace() rather than mutate in place. + new_state = replace( + state, + canary_active=policy.canary, + canary_traffic_pct=policy.traffic_pct, + canary_autoroll_on_regress=autoroll_on_regress, + ) + write_state(new_state) + # Reload so the in-memory value reflects the persisted updated_at. + new_state = read_state() + console.print( + f"[green]canary set:[/] {escape(policy.canary or '')} @ " + f"{policy.traffic_pct:.1f}% (autoroll={autoroll_on_regress})" + ) + + +@app.command("replay") +def replay_cmd( + iteration_id: Optional[str] = typer.Argument( + None, help="Iteration id (omit to list all)." + ), +) -> None: + """Replay a recorded loop iteration manifest.""" + if iteration_id is None: + ids = list_iterations() + if not ids: + console.print("[yellow]no iterations recorded yet[/]") + return + console.print("\n".join(escape(i) for i in ids)) + return + try: + record = read_iteration(iteration_id) + except (FileNotFoundError, TypeError, ValueError) as exc: + console.print(f"[red]replay failed:[/] {escape(str(exc))}") + raise typer.Exit(code=2) + table = Table(title=f"replay {escape(record.iteration_id)}", show_header=False) + table.add_column("field", style="bold") + table.add_column("value") + for k, v in record.to_dict().items(): + table.add_row(escape(str(k)), escape(str(v))) + console.print(table) + + +def _parse_traffic(raw: str) -> float: + """Parse ``"5%"`` / ``"5"`` / ``" 5.5 %"`` into a percent float.""" + if not isinstance(raw, str): + console.print(f"[red]--traffic must be a string, got {type(raw).__name__}[/]") + raise typer.Exit(code=2) + txt = raw.strip() + if txt.endswith("%"): + txt = txt[:-1].strip() + try: + pct = float(txt) + except ValueError: + console.print(f"[red]invalid --traffic:[/] {escape(raw)}") + raise typer.Exit(code=2) + if not (0.0 <= pct <= 100.0): + console.print("[red]--traffic must be in [0, 100][/]") + raise typer.Exit(code=2) + return pct diff --git a/soup_cli/utils/canary_router.py b/soup_cli/utils/canary_router.py new file mode 100644 index 0000000..bec9094 --- /dev/null +++ b/soup_cli/utils/canary_router.py @@ -0,0 +1,195 @@ +"""Canary router (v0.58.0 Part B). + +Pure-Python deterministic routing of inference requests between a stable +adapter and a canary adapter. The router is *deterministic* on a hashed +request key — so a given conversation always lands in the same bucket +within an iteration — and *sticky on rollback* so a flaky verdict can't +ping-pong traffic between adapters. + +Why this lives in `utils/` and not inside `commands/serve.py`: the +canary policy is a pure math kernel exercised by `soup loop watch` +without needing a live FastAPI app. The HTTP middleware in `serve.py` +plugs into `route()` directly. +""" + +from __future__ import annotations + +import hashlib +import math +import threading +from dataclasses import dataclass, field +from types import MappingProxyType +from typing import Mapping, Optional + + +@dataclass(frozen=True) +class CanaryPolicy: + """Frozen rollout policy: stable vs canary + traffic split + verdict.""" + + stable: str + canary: Optional[str] = None + traffic_pct: float = 0.0 # in [0, 100] + sticky_on_rollback: bool = True + + def __post_init__(self) -> None: + if not isinstance(self.stable, str) or not self.stable or "\x00" in self.stable: + raise ValueError("stable must be a non-empty NUL-free string") + if len(self.stable) > 256: + raise ValueError("stable name exceeds 256 chars") + if self.canary is not None: + if not isinstance(self.canary, str) or not self.canary or "\x00" in self.canary: + raise ValueError("canary must be a non-empty NUL-free string or None") + if len(self.canary) > 256: + raise ValueError("canary name exceeds 256 chars") + if self.canary == self.stable: + raise ValueError("canary must differ from stable") + v = self.traffic_pct + if isinstance(v, bool) or not isinstance(v, (int, float)) or not math.isfinite(v): + raise ValueError("traffic_pct must be a finite number") + if not (0.0 <= float(v) <= 100.0): + raise ValueError("traffic_pct must be in [0, 100]") + if self.canary is None and float(v) > 0.0: + raise ValueError("cannot route traffic to None canary") + if not isinstance(self.sticky_on_rollback, bool): + raise ValueError("sticky_on_rollback must be bool") + + +@dataclass(frozen=True) +class RouteDecision: + """Result of one routing decision: which adapter + which bucket.""" + + adapter: str + bucket: str # "stable" | "canary" + rolled_back: bool = False + + +_HASH_MOD = 10_000 # buckets — gives ±0.01 % granularity on the split + + +def _bucket_for_key(key: str) -> int: + """Deterministic 4-hex-digit bucket via SHA-256 (key fingerprint).""" + if not isinstance(key, str): + raise TypeError("key must be a string") + if not key: + raise ValueError("key must not be empty") + if "\x00" in key: + raise ValueError("key must not contain NUL") + digest = hashlib.sha256(key.encode("utf-8")).digest() + # Take 4 bytes → 32-bit unsigned, modulo bucket count. + val = int.from_bytes(digest[:4], "big", signed=False) + return val % _HASH_MOD + + +def route(policy: CanaryPolicy, request_key: str) -> RouteDecision: + """Decide which adapter serves a request given its fingerprint key. + + Deterministic: the same ``(policy, request_key)`` always returns the + same bucket. Stickiness comes from the caller building ``request_key`` + from a conversation id (not a per-message timestamp). + """ + if not isinstance(policy, CanaryPolicy): + raise TypeError("policy must be CanaryPolicy") + bucket = _bucket_for_key(request_key) + # `math.ceil` is more predictable than `round` at sub-bucket fractions: + # `traffic_pct=0.005` → 1 bucket out of 10 000 (0.01%), not 0 (silent + # truncation per code-review MEDIUM #5). + threshold = math.ceil(policy.traffic_pct / 100.0 * _HASH_MOD) + if policy.canary is None or bucket >= threshold: + return RouteDecision(adapter=policy.stable, bucket="stable") + return RouteDecision(adapter=policy.canary, bucket="canary") + + +def rollback(policy: CanaryPolicy, *, reason: str = "regression") -> CanaryPolicy: + """Return a policy with canary cleared (traffic forced to stable). + + Sticky-on-rollback means subsequent calls to ``route`` return the + stable adapter even if a noisy re-evaluation later flips the verdict + — the operator must explicitly re-promote a canary to clear the + sticky bit (by calling ``CanaryPolicy(...)`` afresh). + """ + if not isinstance(policy, CanaryPolicy): + raise TypeError("policy must be CanaryPolicy") + if not isinstance(reason, str) or not reason or "\x00" in reason: + raise ValueError("reason must be a non-empty NUL-free string") + return CanaryPolicy( + stable=policy.stable, + canary=None, + traffic_pct=0.0, + sticky_on_rollback=policy.sticky_on_rollback, + ) + + +# --------------------------------------------------------------------------- +# Verdict bucket aggregation — used by `soup loop watch` to decide whether to +# roll back. Each per-bucket result is a {0, 1} OK/MAJOR signal (matches the +# v0.26.0 Quant-Lobotomy verdict surface). +# --------------------------------------------------------------------------- + +@dataclass +class BucketStats: + """Mutable per-bucket counters. NOT thread-safe — call ``aggregate`` + under a single thread or wrap externally with ``threading.Lock``.""" + + stable_ok: int = 0 + stable_major: int = 0 + canary_ok: int = 0 + canary_major: int = 0 + _lock: threading.Lock = field( + default_factory=threading.Lock, repr=False, compare=False + ) + + def record(self, bucket: str, ok: bool) -> None: + if bucket not in ("stable", "canary"): + raise ValueError("bucket must be 'stable' or 'canary'") + if not isinstance(ok, bool): + raise ValueError("ok must be bool") + with self._lock: + if bucket == "stable": + if ok: + self.stable_ok += 1 + else: + self.stable_major += 1 + else: + if ok: + self.canary_ok += 1 + else: + self.canary_major += 1 + + def verdict(self, *, min_samples: int = 30, regression_threshold: float = 0.05) -> str: + """Return ``"OK"`` / ``"MAJOR"`` / ``"UNKNOWN"``. + + - ``UNKNOWN``: fewer than ``min_samples`` total samples in the + canary bucket. Defends against early-rollback on insufficient + evidence (matches v0.26.0 Quant-Lobotomy policy). + - ``MAJOR``: canary OK rate is below stable's by more than + ``regression_threshold`` (default 5 percentage points). + - ``OK``: otherwise. + """ + if isinstance(min_samples, bool) or not isinstance(min_samples, int) or min_samples < 1: + raise ValueError("min_samples must be a positive int") + v = regression_threshold + if isinstance(v, bool) or not isinstance(v, (int, float)) or not math.isfinite(v): + raise ValueError("regression_threshold must be a finite number") + if not (0.0 <= float(v) <= 1.0): + raise ValueError("regression_threshold must be in [0, 1]") + with self._lock: + canary_total = self.canary_ok + self.canary_major + stable_total = self.stable_ok + self.stable_major + if canary_total < min_samples: + return "UNKNOWN" + stable_rate = (self.stable_ok / stable_total) if stable_total > 0 else 1.0 + canary_rate = self.canary_ok / canary_total + if stable_rate - canary_rate > regression_threshold: + return "MAJOR" + return "OK" + + def snapshot(self) -> Mapping[str, int]: + with self._lock: + return MappingProxyType( + { + "stable_ok": self.stable_ok, + "stable_major": self.stable_major, + "canary_ok": self.canary_ok, + "canary_major": self.canary_major, + } + ) diff --git a/soup_cli/utils/loop_budget.py b/soup_cli/utils/loop_budget.py new file mode 100644 index 0000000..3c28d4a --- /dev/null +++ b/soup_cli/utils/loop_budget.py @@ -0,0 +1,169 @@ +"""Cost + budget guardrails for `soup loop` (v0.58.0 Part C). + +Two orthogonal rate limits: + +* ``monthly_budget_usd`` — composes with v0.34.0 per-run cost so the + watch daemon pauses (graceful save, no kill) when projected spend + would exceed the budget. +* ``max_runs_per_day`` — defends against runaway proxy loops by capping + iteration starts per UTC day. + +The math here is pure-Python so the daemon can call ``check()`` without +opening a SQLite handle. Persisted counters live in the ``LoopState`` +shared store. +""" + +from __future__ import annotations + +import math +from dataclasses import dataclass +from datetime import datetime, timezone +from typing import Optional + + +@dataclass(frozen=True) +class BudgetDecision: + """Decision returned by ``check_budget``. + + ``proceed`` is the only field a caller MUST inspect; the others are + advisory for the user-facing dashboard. + """ + + proceed: bool + reason: str + projected_total_usd: float + runs_today: int + + +def _utc_date_str(ts: Optional[datetime] = None) -> str: + """Return today's UTC date as ``YYYY-MM-DD`` (testable via ``ts=...``).""" + now = ts if ts is not None else datetime.now(timezone.utc) + return now.strftime("%Y-%m-%d") + + +def reset_daily_counter_if_new_day( + runs_today: int, + last_run_date: Optional[str], + *, + now: Optional[datetime] = None, +) -> "tuple[int, str]": + """Reset ``runs_today`` to 0 when the UTC day rolls over. + + Returns ``(runs_today, last_run_date)`` — the caller updates the + ``LoopState`` with the returned values before checking the cap. + """ + if isinstance(runs_today, bool) or not isinstance(runs_today, int) or runs_today < 0: + raise ValueError("runs_today must be a non-negative int") + if last_run_date is not None and not isinstance(last_run_date, str): + raise ValueError("last_run_date must be a str or None") + today = _utc_date_str(now) + if last_run_date != today: + return 0, today + return runs_today, today + + +def check_budget( + *, + estimated_run_usd: float, + spent_so_far_usd: float, + monthly_budget_usd: Optional[float], + runs_today: int, + max_runs_per_day: Optional[int], +) -> BudgetDecision: + """Decide whether to proceed with another iteration. + + The decision composes three checks in this order: + + 1. Run-cap (daily) — fast rejection when ``max_runs_per_day`` is set + and ``runs_today >= max``. + 2. Estimate sanity — reject non-finite / negative cost estimates so + a broken probe can't smuggle a negative refund. + 3. Budget — reject when projected spend would exceed the cap. + """ + if ( + isinstance(estimated_run_usd, bool) + or not isinstance(estimated_run_usd, (int, float)) + or not math.isfinite(estimated_run_usd) + or estimated_run_usd < 0 + ): + raise ValueError("estimated_run_usd must be a non-negative finite number") + if ( + isinstance(spent_so_far_usd, bool) + or not isinstance(spent_so_far_usd, (int, float)) + or not math.isfinite(spent_so_far_usd) + or spent_so_far_usd < 0 + ): + raise ValueError("spent_so_far_usd must be a non-negative finite number") + if isinstance(runs_today, bool) or not isinstance(runs_today, int) or runs_today < 0: + raise ValueError("runs_today must be a non-negative int") + if max_runs_per_day is not None: + if ( + isinstance(max_runs_per_day, bool) + or not isinstance(max_runs_per_day, int) + or max_runs_per_day < 1 + ): + raise ValueError("max_runs_per_day must be a positive int or None") + if runs_today >= max_runs_per_day: + return BudgetDecision( + proceed=False, + reason=( + f"daily cap reached: {runs_today}/{max_runs_per_day} runs" + ), + projected_total_usd=float(spent_so_far_usd), + runs_today=runs_today, + ) + projected = float(spent_so_far_usd) + float(estimated_run_usd) + if monthly_budget_usd is not None: + if ( + isinstance(monthly_budget_usd, bool) + or not isinstance(monthly_budget_usd, (int, float)) + or not math.isfinite(monthly_budget_usd) + or monthly_budget_usd < 0 + ): + raise ValueError("monthly_budget_usd must be >= 0 or None") + if projected > float(monthly_budget_usd): + return BudgetDecision( + proceed=False, + reason=( + f"would exceed monthly budget: ${projected:.2f} > " + f"${float(monthly_budget_usd):.2f}" + ), + projected_total_usd=projected, + runs_today=runs_today, + ) + return BudgetDecision( + proceed=True, + reason="within budget", + projected_total_usd=projected, + runs_today=runs_today, + ) + + +def parse_budget_string(raw: str) -> float: + """Parse ``"50usd"`` / ``"100 USD"`` / ``"25"`` into a USD float. + + Trailing ``"usd"`` (case-insensitive) is optional. Bounds: ``[0, + 1_000_000]`` so a fat-finger ``"1000000000"`` cannot cause integer + overflow in downstream arithmetic. + """ + if not isinstance(raw, str): + raise TypeError("budget must be a string") + raw = raw.strip().lower() + if not raw: + raise ValueError("budget must not be empty") + if "\x00" in raw: + raise ValueError("budget must not contain NUL") + if raw.endswith("usd"): + raw = raw[:-3].strip() + if not raw: + # "usd" / " usd " — friendly explicit message (code-review M6). + raise ValueError("budget must include a numeric value (e.g. '50usd')") + try: + value = float(raw) + except ValueError as exc: + raise ValueError(f"invalid budget value: {raw!r}") from exc + if not math.isfinite(value): + raise ValueError("budget must be finite") + if not (0.0 <= value <= 1_000_000.0): + raise ValueError("budget must be in [0, 1_000_000] USD") + return value diff --git a/soup_cli/utils/loop_daemon.py b/soup_cli/utils/loop_daemon.py new file mode 100644 index 0000000..f139a74 --- /dev/null +++ b/soup_cli/utils/loop_daemon.py @@ -0,0 +1,339 @@ +"""Watch-daemon orchestrator for `soup loop watch` (v0.58.0). + +The full production cycle is: + + traces (v0.26 from-traces) → preference pairs → DPO train + → eval-gate (v0.26.0 Part B) → optional canary deploy + → rollback on regression (v0.26.0 Quant-Lobotomy MAJOR) + +Each stage is encapsulated as a callable so the daemon stays testable +without a GPU. A *headless* run with the default stage callbacks +exercises every state transition (state mutations, budget check, +iteration record, sticky rollback) deterministically. + +The daemon is foreground by default. The CLI ``--detach`` flag launches +a subprocess via ``subprocess.Popen([sys.executable, "-m", "soup_cli.cli", +"loop", "watch", "--foreground"])`` so the operator gets a real process +id back instead of relying on shell job control. +""" + +from __future__ import annotations + +import logging +import math +import signal +import threading +from dataclasses import dataclass, replace +from datetime import datetime, timezone +from typing import Callable, Mapping, Optional + +from soup_cli.utils.canary_router import BucketStats, CanaryPolicy, rollback +from soup_cli.utils.loop_budget import ( + BudgetDecision, + check_budget, + reset_daily_counter_if_new_day, +) +from soup_cli.utils.loop_iteration import ( + IterationRecord, + new_iteration_id, + write_iteration, +) +from soup_cli.utils.loop_state import LoopState, read_state, write_state + +_LOG = logging.getLogger(__name__) + +# A per-stage callable returns a small dict; the orchestrator merges the +# result dicts into one ``StageResult`` before recording the iteration. + +HarvestFn = Callable[[LoopState], Mapping[str, object]] +TrainFn = Callable[[LoopState, Mapping[str, object]], Mapping[str, object]] +GateFn = Callable[[LoopState, Mapping[str, object]], Mapping[str, object]] +DeployFn = Callable[[LoopState, Mapping[str, object]], Mapping[str, object]] +CostFn = Callable[[LoopState], float] + + +# --------------------------------------------------------------------------- +# Default stage callbacks — pure-Python no-ops that satisfy the contract so +# the daemon runs end-to-end on CPU without a model. Real wiring composes +# v0.26.0 trace-to-pref + v0.26.0 eval-gate + v0.30.0 multi-adapter deploy. +# --------------------------------------------------------------------------- + + +def default_harvest(state: LoopState) -> Mapping[str, object]: + """Stub harvest stage — returns zero pairs. + + Real implementation wires v0.26.0 ``soup_cli.data.traces.parsers`` + + ``pair_builder.build_pairs`` to scan trace logs. Kept as a stub here + so headless tests don't need a live trace store. + """ + return {"pairs_harvested": 0, "pairs_path": None} + + +def default_train(state: LoopState, ctx: Mapping[str, object]) -> Mapping[str, object]: + """Stub train stage — skips the run.""" + return {"run_id": None, "skipped": True} + + +def default_gate(state: LoopState, ctx: Mapping[str, object]) -> Mapping[str, object]: + """Stub gate stage — returns SKIPPED when there is nothing to evaluate.""" + if ctx.get("skipped"): + return {"gate_verdict": "SKIPPED"} + return {"gate_verdict": "OK"} + + +def default_deploy(state: LoopState, ctx: Mapping[str, object]) -> Mapping[str, object]: + """Stub deploy stage — does not promote anything.""" + return {"deployed": False, "canary_verdict": None} + + +def default_cost(state: LoopState) -> float: + """Stub cost estimate — zero so the budget gate stays permissive.""" + return 0.0 + + +@dataclass +class WatchConfig: + """Daemon configuration knobs.""" + + poll_interval_sec: float = 60.0 + max_iterations: Optional[int] = None # None = unbounded (real daemon) + state_path: Optional[str] = None + iteration_dir: Optional[str] = None + harvest_fn: HarvestFn = default_harvest + train_fn: TrainFn = default_train + gate_fn: GateFn = default_gate + deploy_fn: DeployFn = default_deploy + cost_fn: CostFn = default_cost + on_iteration: Optional[Callable[[IterationRecord], None]] = None + + def __post_init__(self) -> None: + v = self.poll_interval_sec + if isinstance(v, bool) or not isinstance(v, (int, float)) or not math.isfinite(v): + raise ValueError("poll_interval_sec must be a finite number") + if v < 1.0 or v > 3600.0: + raise ValueError("poll_interval_sec must be in [1, 3600]") + if self.max_iterations is not None: + mi = self.max_iterations + if isinstance(mi, bool) or not isinstance(mi, int) or mi < 0: + raise ValueError("max_iterations must be a non-negative int or None") + for fname in ("harvest_fn", "train_fn", "gate_fn", "deploy_fn", "cost_fn"): + if not callable(getattr(self, fname)): + raise ValueError(f"{fname} must be callable") + if self.on_iteration is not None and not callable(self.on_iteration): + raise ValueError("on_iteration must be callable or None") + + +def run_once( + state: LoopState, + config: WatchConfig, +) -> "tuple[LoopState, IterationRecord, BudgetDecision]": + """Execute one full iteration synchronously. Pure with respect to time.""" + if not isinstance(state, LoopState): + raise TypeError("state must be LoopState") + if not isinstance(config, WatchConfig): + raise TypeError("config must be WatchConfig") + runs_today, today = reset_daily_counter_if_new_day( + state.runs_today, state.last_run_date + ) + state = _state_with(state, runs_today=runs_today, last_run_date=today) + estimated = float(config.cost_fn(state)) + decision = check_budget( + estimated_run_usd=estimated, + spent_so_far_usd=state.spent_this_month_usd, + monthly_budget_usd=state.monthly_budget_usd, + runs_today=state.runs_today, + max_runs_per_day=state.max_runs_per_day, + ) + iteration_id = new_iteration_id() + started_at = _utc_iso() + if not decision.proceed: + record = IterationRecord( + iteration_id=iteration_id, + started_at=started_at, + finished_at=_utc_iso(), + pairs_harvested=0, + run_id=None, + gate_verdict="SKIPPED", + canary_verdict=None, + shipped=False, + rolled_back=False, + estimated_cost_usd=estimated, + notes=f"budget-skip: {decision.reason}", + ) + return state, record, decision + harvest_out = dict(config.harvest_fn(state)) + train_out = dict(config.train_fn(state, harvest_out)) + gate_out = dict(config.gate_fn(state, train_out)) + deploy_out = dict(config.deploy_fn(state, {**train_out, **gate_out})) + gate_verdict = str(gate_out.get("gate_verdict", "SKIPPED")) + canary_verdict = deploy_out.get("canary_verdict") + if canary_verdict is not None: + canary_verdict = str(canary_verdict) + shipped = bool(deploy_out.get("deployed", False)) + rolled_back = bool(deploy_out.get("rolled_back", False)) + pairs = int(harvest_out.get("pairs_harvested", 0) or 0) + if pairs < 0: + pairs = 0 + record = IterationRecord( + iteration_id=iteration_id, + started_at=started_at, + finished_at=_utc_iso(), + pairs_harvested=pairs, + run_id=(str(train_out["run_id"]) if train_out.get("run_id") else None), + gate_verdict=gate_verdict if gate_verdict in ("OK", "MAJOR", "SKIPPED") else "SKIPPED", + canary_verdict=( + canary_verdict + if canary_verdict in (None, "OK", "MAJOR", "UNKNOWN") + else None + ), + shipped=shipped, + rolled_back=rolled_back, + estimated_cost_usd=estimated, + notes=str(deploy_out.get("notes", ""))[:4096], + ) + new_state = state.bumped( + traces_collected=int(harvest_out.get("traces_collected", 0) or 0), + pairs_distilled=pairs, + runs_gated=1 if record.gate_verdict in ("OK", "MAJOR") else 0, + adapters_shipped=1 if shipped else 0, + iteration_count=1, + runs_today=1, + ) + new_state = _state_with( + new_state, + spent_this_month_usd=new_state.spent_this_month_usd + max(0.0, estimated), + last_iteration_id=iteration_id, + last_run_date=today, + ) + return new_state, record, decision + + +def watch(config: WatchConfig) -> "tuple[LoopState, int]": + """Run the daemon loop. Returns ``(final_state, iterations_run)``. + + Stops cleanly on: + - ``config.max_iterations`` reached (test/finite mode) + - state file going to ``status="stopped"`` between iterations + - SIGTERM/SIGINT (installed via ``signal.signal`` when on the main thread) + """ + if not isinstance(config, WatchConfig): + raise TypeError("config must be WatchConfig") + stop = threading.Event() + + def _request_stop(signum, frame): # pragma: no cover — signal path + stop.set() + + try: + if threading.current_thread() is threading.main_thread(): + signal.signal(signal.SIGTERM, _request_stop) + signal.signal(signal.SIGINT, _request_stop) + except (ValueError, AttributeError): + # Non-main-thread + Windows-Python combos where signal.signal raises. + pass + + iterations = 0 + state = read_state(config.state_path) + # Only promote `stopped` → `running` automatically; `paused` must + # survive a `soup loop watch` invocation so a SIGTERM + restart + # cycle does not silently un-pause the daemon (code-review HIGH #2). + if state.status == "stopped": + state = state.with_status("running") + write_state(state, config.state_path) + try: + while not stop.is_set(): + if config.max_iterations is not None and iterations >= config.max_iterations: + break + try: + state = read_state(config.state_path) + except (FileNotFoundError, ValueError): + break + if state.status == "paused": + if stop.wait(min(config.poll_interval_sec, 60.0)): + break + continue + if state.status == "stopped": + break + state, record, decision = run_once(state, config) + write_state(state, config.state_path) + # Budget-skipped iterations DO NOT produce a manifest — the + # cycle didn't actually run, so cluttering .soup-loops/ with + # "I didn't run" records would surprise operators expecting + # iteration_count to match the manifest count (code-review + # HIGH #3 fix). The state still records the skip in notes. + if decision.proceed: + try: + write_iteration(record, base_dir=config.iteration_dir) + except (OSError, ValueError) as exc: + _LOG.warning("iteration write failed: %s", type(exc).__name__) + if config.on_iteration is not None: + try: + config.on_iteration(record) + except Exception: # noqa: BLE001 — daemon must not crash + _LOG.warning("on_iteration callback raised", exc_info=True) + iterations += 1 + if iterations and ( + config.max_iterations is None or iterations < config.max_iterations + ): + if stop.wait(config.poll_interval_sec): + break + finally: + # Preserve `paused` status across daemon exit — only flip to + # `stopped` if the daemon naturally exited (max_iterations / state + # was running). A SIGTERM while paused must not be silently + # promoted to stopped (code-review HIGH #2 fix). + try: + current = read_state(config.state_path) + if current.status == "running": + write_state(current.with_status("stopped"), config.state_path) + state = current.with_status("stopped") + else: + state = current + except (FileNotFoundError, ValueError): + pass + return state, iterations + + +def _state_with(state: LoopState, **kwargs: object) -> LoopState: + """Return a copy with overrides applied (escape hatch around ``replace``). + + The dataclass already exposes ``with_status`` and ``bumped`` but the + daemon needs to flip a handful of fields atomically per cycle (e.g. + ``last_run_date`` + ``last_iteration_id`` together). Keeping this tiny + helper local avoids leaking ``dataclasses.replace`` into the daemon + surface; the import lives at module top per code-review MEDIUM #7. + """ + return replace(state, **kwargs) + + +def _utc_iso() -> str: + return datetime.now(timezone.utc).replace(microsecond=0).isoformat() + + +def evaluate_canary_verdict(stats: BucketStats) -> str: + """Project a single canary verdict from accumulated bucket stats. + + Trivial wrapper around ``BucketStats.verdict`` so the daemon does + not import the router class directly — keeps the import graph + one-directional (utils.loop_daemon → utils.canary_router, never the + other way). + """ + return stats.verdict() + + +def maybe_rollback( + policy: CanaryPolicy, verdict: str, *, sticky: bool = True +) -> CanaryPolicy: + """Roll back the canary if ``verdict == "MAJOR"``. + + Non-MAJOR verdicts pass through unchanged so a flaky re-eval cannot + flip-flop traffic. Sticky-on-rollback (the default) is documented in + ``canary_router.rollback`` — once cleared, the operator must + explicitly re-promote a new canary. + """ + if not isinstance(policy, CanaryPolicy): + raise TypeError("policy must be CanaryPolicy") + if not isinstance(verdict, str): + raise TypeError("verdict must be str") + if verdict == "MAJOR": + return rollback(policy, reason="canary regression") + return policy diff --git a/soup_cli/utils/loop_iteration.py b/soup_cli/utils/loop_iteration.py new file mode 100644 index 0000000..04418b1 --- /dev/null +++ b/soup_cli/utils/loop_iteration.py @@ -0,0 +1,222 @@ +"""Per-iteration artifact packing for `soup loop` (v0.58.0 Part D). + +Each loop iteration is summarised as a small JSON manifest under +``.soup-loops//iteration.json``. The directory is laid +out so a v0.26.0 Soup Can can wrap it later without re-shaping the +files — same naming as ``soup history`` lineage entries. + +`replay` re-reads a recorded iteration and returns its manifest + +metric trace so the operator can run "would the loop have shipped v17 +today?" what-if analysis without touching the live state. +""" + +from __future__ import annotations + +import json +import os +import stat +import tempfile +import uuid +from dataclasses import asdict, dataclass +from datetime import datetime, timezone +from types import MappingProxyType +from typing import Any, Mapping, Optional, Tuple + +from soup_cli.utils.paths import is_under_cwd + +_DEFAULT_DIR = ".soup-loops" +_MAX_MANIFEST_BYTES = 1 * 1024 * 1024 # 1 MiB +_MAX_ID_LEN = 128 + + +@dataclass(frozen=True) +class IterationRecord: + """One iteration of the harvest → train → gate → ship cycle.""" + + iteration_id: str + started_at: str + finished_at: Optional[str] + pairs_harvested: int + run_id: Optional[str] + gate_verdict: str # "OK" / "MAJOR" / "SKIPPED" + canary_verdict: Optional[str] # "OK" / "MAJOR" / "UNKNOWN" / None + shipped: bool + rolled_back: bool + estimated_cost_usd: float + notes: str = "" + + def __post_init__(self) -> None: + _check_id(self.iteration_id) + for fname in ("started_at", "gate_verdict"): + v = getattr(self, fname) + if not isinstance(v, str) or not v or "\x00" in v: + raise ValueError(f"{fname} must be a non-empty NUL-free string") + if self.finished_at is not None: + if ( + not isinstance(self.finished_at, str) + or not self.finished_at + or "\x00" in self.finished_at + ): + raise ValueError("finished_at must be a non-empty NUL-free string or None") + pairs = self.pairs_harvested + if isinstance(pairs, bool) or not isinstance(pairs, int) or pairs < 0: + raise ValueError("pairs_harvested must be a non-negative int") + if self.run_id is not None and ( + not isinstance(self.run_id, str) or not self.run_id or "\x00" in self.run_id + ): + raise ValueError("run_id must be a non-empty NUL-free string or None") + if self.gate_verdict not in ("OK", "MAJOR", "SKIPPED"): + raise ValueError("gate_verdict must be one of OK/MAJOR/SKIPPED") + if self.canary_verdict is not None and self.canary_verdict not in ( + "OK", + "MAJOR", + "UNKNOWN", + ): + raise ValueError("canary_verdict must be OK/MAJOR/UNKNOWN/None") + for fname in ("shipped", "rolled_back"): + if not isinstance(getattr(self, fname), bool): + raise ValueError(f"{fname} must be bool") + v = self.estimated_cost_usd + if isinstance(v, bool) or not isinstance(v, (int, float)) or v < 0: + raise ValueError("estimated_cost_usd must be a non-negative number") + if not isinstance(self.notes, str): + raise ValueError("notes must be a string") + if "\x00" in self.notes: + raise ValueError("notes must not contain NUL") + if len(self.notes) > 4096: + raise ValueError("notes exceeds 4096 chars") + + def to_dict(self) -> Mapping[str, Any]: + return MappingProxyType(asdict(self)) + + +def new_iteration_id() -> str: + """UTC timestamp + 8-hex-digit suffix (collision-safe under burst).""" + ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S") + suffix = uuid.uuid4().hex[:8] + return f"iter-{ts}-{suffix}" + + +def _check_id(iteration_id: str) -> None: + if not isinstance(iteration_id, str): + raise TypeError("iteration_id must be a string") + if not iteration_id: + raise ValueError("iteration_id must not be empty") + if "\x00" in iteration_id: + raise ValueError("iteration_id must not contain NUL") + if len(iteration_id) > _MAX_ID_LEN: + raise ValueError("iteration_id exceeds 128 chars") + if any(c in iteration_id for c in (os.sep, "/", "\\", "..")): + raise ValueError("iteration_id must not contain path separators") + + +def _check_dir(path: str) -> str: + if not isinstance(path, str): + raise TypeError("path must be str") + if not path or "\x00" in path: + raise ValueError("path must be non-empty NUL-free") + if not is_under_cwd(path): + raise ValueError("path must stay under cwd") + # Direct lstat (no `lexists` guard) closes the TOCTOU window + # security-review M1 surfaced: a symlink planted between `lexists` + # and `lstat` would otherwise sneak through (matches the loop_state + # `_check_path` pattern). + try: + st = os.lstat(path) + except FileNotFoundError: + return path + except OSError as exc: + raise ValueError(f"path unreadable: {type(exc).__name__}") from exc + if stat.S_ISLNK(st.st_mode): + raise ValueError("path must not be a symlink (TOCTOU defence)") + return path + + +def write_iteration( + record: IterationRecord, + *, + base_dir: Optional[str] = None, +) -> str: + """Persist ``record`` under ``//iteration.json``.""" + if not isinstance(record, IterationRecord): + raise TypeError("record must be IterationRecord") + parent = base_dir if base_dir is not None else _DEFAULT_DIR + _check_dir(parent) + iter_dir = os.path.join(parent, record.iteration_id) + _check_dir(iter_dir) + os.makedirs(iter_dir, exist_ok=True) + target = os.path.join(iter_dir, "iteration.json") + _check_dir(target) + body = json.dumps( + dict(record.to_dict()), allow_nan=False, indent=2, sort_keys=True + ).encode("utf-8") + if len(body) > _MAX_MANIFEST_BYTES: + raise ValueError("iteration manifest exceeds 1 MiB cap") + fd, tmp = tempfile.mkstemp(prefix=".iter_", dir=iter_dir) + try: + with os.fdopen(fd, "wb") as fh: + fh.write(body) + os.replace(tmp, target) + except Exception: + try: + os.unlink(tmp) + except OSError: + pass + raise + return target + + +def read_iteration( + iteration_id: str, *, base_dir: Optional[str] = None +) -> IterationRecord: + """Reload an iteration record by id.""" + _check_id(iteration_id) + parent = base_dir if base_dir is not None else _DEFAULT_DIR + target = os.path.join(parent, iteration_id, "iteration.json") + _check_dir(target) + if not os.path.isfile(target): + raise FileNotFoundError(f"iteration {iteration_id!r} not found") + try: + size = os.path.getsize(target) + except OSError as exc: + raise ValueError(f"iteration manifest unreadable: {type(exc).__name__}") from exc + if size > _MAX_MANIFEST_BYTES: + raise ValueError("iteration manifest exceeds 1 MiB cap") + with open(target, "rb") as fh: + raw = fh.read() + try: + data = json.loads(raw.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError) as exc: + raise ValueError(f"invalid manifest JSON: {exc}") from exc + if not isinstance(data, dict): + raise ValueError("manifest root must be a JSON object") + allowed = set(IterationRecord.__dataclass_fields__.keys()) + filtered = {k: v for k, v in data.items() if k in allowed} + try: + return IterationRecord(**filtered) + except (TypeError, ValueError) as exc: + raise ValueError(f"manifest contents invalid: {exc}") from exc + + +def list_iterations(base_dir: Optional[str] = None) -> Tuple[str, ...]: + """Return iteration ids sorted by name (timestamp-prefixed).""" + parent = base_dir if base_dir is not None else _DEFAULT_DIR + if not os.path.isdir(parent): + return () + try: + entries = os.listdir(parent) + except OSError: + # Permission flap mid-iteration — daemon must not crash on a + # read of its own artifact dir (code-review MEDIUM #8). + return () + out: list[str] = [] + for entry in entries: + candidate = os.path.join(parent, entry, "iteration.json") + if os.path.isfile(candidate): + try: + _check_id(entry) + except (TypeError, ValueError): + continue + out.append(entry) + out.sort() + return tuple(out) diff --git a/soup_cli/utils/loop_state.py b/soup_cli/utils/loop_state.py new file mode 100644 index 0000000..ff5b7ed --- /dev/null +++ b/soup_cli/utils/loop_state.py @@ -0,0 +1,309 @@ +"""Loop state file (v0.58.0 Part A — control plane). + +`soup loop` orchestrates the *production traces → preference pairs → +Eval-Gated DPO → canary deploy → rollback* cycle. State for the whole +loop lives in a single ``.soup/loop.yaml`` next to the project, with +atomic writes + cwd containment + symlink rejection — the same TOCTOU +policy as every other v0.5x persistence surface (v0.33.0 #22 / +v0.43.0 Part C / v0.53.7 #106). + +Status (``running`` / ``paused`` / ``stopped``) and counters (traces / +pairs / runs / deploys) live here; per-iteration artifacts ship as +v0.26.0 Soup Cans under ``.soup-loops/``. +""" + +from __future__ import annotations + +import json +import os +import stat +import tempfile +from dataclasses import asdict, dataclass, replace +from datetime import datetime, timezone +from types import MappingProxyType +from typing import Mapping, Optional, Tuple + +from soup_cli.utils.paths import is_under_cwd + +# Status values are deliberately closed — every state machine transition +# below must remain auditable. +LOOP_STATUSES: frozenset = frozenset({"running", "paused", "stopped"}) + +_MAX_PATH_LEN = 4096 +_MAX_STR_FIELD = 512 +_MAX_FILE_BYTES = 1 * 1024 * 1024 # 1 MiB cap on the state file +_DEFAULT_STATE_DIR = ".soup" +_DEFAULT_STATE_FILENAME = "loop.yaml" + + +@dataclass(frozen=True) +class LoopState: + """Immutable snapshot of a `soup loop` configuration + counters. + + The persisted file is JSON-formatted (despite the ``.yaml`` extension) + so we can use the stdlib parser without pulling pyyaml into the read + path; YAML is a superset of JSON for objects and the file remains + human-readable. Counters are absolute lifetime totals; per-iteration + detail lives in the ``.soup-loops/`` Soup Can artifacts. + """ + + served_model: str + eval_suite: str + baseline: str + status: str = "stopped" + traces_collected: int = 0 + pairs_distilled: int = 0 + runs_gated: int = 0 + adapters_shipped: int = 0 + canary_active: Optional[str] = None + canary_traffic_pct: Optional[float] = None + canary_autoroll_on_regress: bool = True + monthly_budget_usd: Optional[float] = None + spent_this_month_usd: float = 0.0 + max_runs_per_day: Optional[int] = None + runs_today: int = 0 + last_run_date: Optional[str] = None # YYYY-MM-DD UTC + last_iteration_id: Optional[str] = None + iteration_count: int = 0 + created_at: str = "" + updated_at: str = "" + + def __post_init__(self) -> None: # noqa: D401 — validator hook + # Closed-allowlist enforcement + bool-as-int rejection mirror the + # project's v0.30/v0.34/v0.50 policy. `replace(...)` (immutable) + # is the only sanctioned mutation path; this validator runs at + # construction time so a hand-rolled instance still gets checks. + _require_str("served_model", self.served_model) + _require_str("eval_suite", self.eval_suite) + _require_str("baseline", self.baseline) + if self.status not in LOOP_STATUSES: + raise ValueError( + f"status must be one of {sorted(LOOP_STATUSES)}, got {self.status!r}" + ) + for fname in ( + "traces_collected", + "pairs_distilled", + "runs_gated", + "adapters_shipped", + "iteration_count", + "runs_today", + ): + v = getattr(self, fname) + if isinstance(v, bool) or not isinstance(v, int) or v < 0: + raise ValueError(f"{fname} must be a non-negative int, got {v!r}") + if self.canary_traffic_pct is not None: + v = self.canary_traffic_pct + if isinstance(v, bool) or not isinstance(v, (int, float)): + raise ValueError("canary_traffic_pct must be numeric or None") + if not (0.0 <= float(v) <= 100.0): + raise ValueError("canary_traffic_pct must be in [0, 100]") + if self.monthly_budget_usd is not None: + v = self.monthly_budget_usd + if isinstance(v, bool) or not isinstance(v, (int, float)) or v < 0: + raise ValueError("monthly_budget_usd must be >= 0 or None") + v = self.spent_this_month_usd + if isinstance(v, bool) or not isinstance(v, (int, float)) or v < 0: + raise ValueError("spent_this_month_usd must be >= 0") + if self.max_runs_per_day is not None: + v = self.max_runs_per_day + if isinstance(v, bool) or not isinstance(v, int) or v < 1: + raise ValueError("max_runs_per_day must be a positive int or None") + if not isinstance(self.canary_autoroll_on_regress, bool): + raise ValueError("canary_autoroll_on_regress must be bool") + if self.canary_active is not None: + _require_str("canary_active", self.canary_active, allow_empty=False) + if self.last_iteration_id is not None: + _require_str("last_iteration_id", self.last_iteration_id, allow_empty=False) + if self.last_run_date is not None: + _require_str("last_run_date", self.last_run_date, allow_empty=False) + + def to_dict(self) -> Mapping[str, object]: + """Stable, JSON-serialisable view (returned as ``MappingProxyType``).""" + return MappingProxyType(asdict(self)) + + def with_status(self, status: str) -> "LoopState": + """Return a copy with ``status`` set + ``updated_at`` refreshed.""" + if status not in LOOP_STATUSES: + raise ValueError( + f"status must be one of {sorted(LOOP_STATUSES)}, got {status!r}" + ) + return replace(self, status=status, updated_at=_utc_now_iso()) + + def bumped(self, **counters: int) -> "LoopState": + """Return a copy with named counters incremented atomically.""" + updates = {} + for k, v in counters.items(): + if k not in { + "traces_collected", + "pairs_distilled", + "runs_gated", + "adapters_shipped", + "iteration_count", + "runs_today", + }: + raise ValueError(f"unknown counter: {k}") + if isinstance(v, bool) or not isinstance(v, int) or v < 0: + raise ValueError(f"{k} delta must be a non-negative int") + current = getattr(self, k) + updates[k] = current + v + updates["updated_at"] = _utc_now_iso() + return replace(self, **updates) + + +def _require_str(field: str, value: object, *, allow_empty: bool = False) -> None: + if not isinstance(value, str): + raise TypeError(f"{field} must be str, got {type(value).__name__}") + if "\x00" in value: + raise ValueError(f"{field} must not contain NUL") + if not allow_empty and not value: + raise ValueError(f"{field} must not be empty") + if len(value) > _MAX_STR_FIELD: + raise ValueError(f"{field} exceeds {_MAX_STR_FIELD} characters") + + +def _utc_now_iso() -> str: + return datetime.now(timezone.utc).replace(microsecond=0).isoformat() + + +def default_state_path(cwd: Optional[str] = None) -> str: + """Return the canonical state-file path under cwd.""" + base = cwd if cwd is not None else os.getcwd() + return os.path.join(base, _DEFAULT_STATE_DIR, _DEFAULT_STATE_FILENAME) + + +def _check_path(path: str, *, allow_missing: bool) -> str: + if not isinstance(path, str): + raise TypeError(f"path must be str, got {type(path).__name__}") + if not path: + raise ValueError("path must not be empty") + if "\x00" in path: + raise ValueError("path must not contain NUL") + if len(path) > _MAX_PATH_LEN: + raise ValueError(f"path exceeds {_MAX_PATH_LEN} characters") + if not is_under_cwd(path): + raise ValueError(f"path {os.path.basename(path)!r} must stay under cwd") + # Direct lstat — no lexists guard — closes the TOCTOU window where a + # symlink could be planted between the existence check and the stat. + try: + st = os.lstat(path) + except FileNotFoundError: + if not allow_missing: + raise FileNotFoundError( + f"state file not found: {os.path.basename(path)}" + ) from None + return path + except OSError as exc: + raise ValueError(f"path unreadable: {type(exc).__name__}") from exc + if stat.S_ISLNK(st.st_mode): + raise ValueError("path must not be a symlink (TOCTOU defence)") + return path + + +def write_state(state: LoopState, path: Optional[str] = None) -> str: + """Atomically persist ``state`` to ``path`` (default: ``./.soup/loop.yaml``). + + Uses ``tempfile.mkstemp`` + ``os.replace`` so a SIGKILL mid-write + cannot leave a torn file at the target — matches v0.43.0 Part D + / v0.55.0 ``lock_suite`` / v0.57.0 atomic-write policy. + """ + if not isinstance(state, LoopState): + raise TypeError(f"state must be LoopState, got {type(state).__name__}") + target = path if path is not None else default_state_path() + _check_path(target, allow_missing=True) + parent = os.path.dirname(target) or "." + os.makedirs(parent, exist_ok=True) + payload = dict(state.to_dict()) + if not payload.get("created_at"): + payload["created_at"] = _utc_now_iso() + payload["updated_at"] = _utc_now_iso() + body = json.dumps(payload, allow_nan=False, indent=2, sort_keys=True).encode("utf-8") + if len(body) > _MAX_FILE_BYTES: + raise ValueError("state payload exceeds 1 MiB cap") + fd, tmp = tempfile.mkstemp(prefix=".loop_state_", dir=parent) + try: + with os.fdopen(fd, "wb") as fh: + fh.write(body) + os.replace(tmp, target) + except Exception: + try: + os.unlink(tmp) + except OSError: + pass + raise + try: + if os.name == "posix": + os.chmod(target, 0o600) + except OSError: + pass + return target + + +def read_state(path: Optional[str] = None) -> LoopState: + """Load a ``LoopState`` from disk.""" + target = path if path is not None else default_state_path() + _check_path(target, allow_missing=False) + try: + size = os.path.getsize(target) + except OSError as exc: + raise ValueError(f"state file unreadable: {type(exc).__name__}") from exc + if size > _MAX_FILE_BYTES: + raise ValueError("state file exceeds 1 MiB cap") + with open(target, "rb") as fh: + raw = fh.read() + try: + data = json.loads(raw.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError) as exc: + raise ValueError(f"state file is not valid JSON: {exc}") from exc + if not isinstance(data, dict): + raise ValueError("state file root must be a JSON object") + # Drop keys we don't recognise so a forward-compat dump can still load + # backwards (the alternative is mass-rejecting on any unknown field). + allowed = set(LoopState.__dataclass_fields__.keys()) + filtered = {k: v for k, v in data.items() if k in allowed} + try: + return LoopState(**filtered) + except (TypeError, ValueError) as exc: + raise ValueError(f"state file has invalid contents: {exc}") from exc + + +def init_state( + served_model: str, + eval_suite: str, + baseline: str, + *, + monthly_budget_usd: Optional[float] = None, + max_runs_per_day: Optional[int] = None, + path: Optional[str] = None, + force: bool = False, +) -> Tuple[LoopState, str]: + """Create the loop.yaml state file. Refuses to clobber unless ``force``.""" + target = path if path is not None else default_state_path() + _check_path(target, allow_missing=True) + # Use lstat (not exists) — defends against a planted symlink between + # the existence check and the write. `_check_path` already rejected + # symlinks so any stat-success here is a real regular-file collision. + try: + os.lstat(target) + present = True + except FileNotFoundError: + present = False + except OSError as exc: + raise ValueError(f"state path unreadable: {type(exc).__name__}") from exc + if present and not force: + raise FileExistsError( + f"loop state already exists at {os.path.basename(target)} " + "(re-run with --force to overwrite)" + ) + now = _utc_now_iso() + state = LoopState( + served_model=served_model, + eval_suite=eval_suite, + baseline=baseline, + status="stopped", + monthly_budget_usd=monthly_budget_usd, + max_runs_per_day=max_runs_per_day, + created_at=now, + updated_at=now, + ) + write_state(state, target) + return state, target diff --git a/tests/test_v0580.py b/tests/test_v0580.py new file mode 100644 index 0000000..66c98c7 --- /dev/null +++ b/tests/test_v0580.py @@ -0,0 +1,1474 @@ +"""Tests for v0.58.0 — soup loop CLI-first data flywheel. + +Coverage: +- Part A: LoopState validation, atomic state file I/O, init/read/write +- Part B: canary_router hash-bucket determinism, sticky rollback, BucketStats +- Part C: BudgetTracker math, daily counter reset, parse_budget_string +- Part D: IterationRecord, write/read/list iterations +- Watch daemon: run_once + watch end-to-end with stub callbacks +- CLI: init / status / pause / resume / watch --max-iterations / canary / replay +""" + +from __future__ import annotations + +import dataclasses +import json +import os +from datetime import datetime, timezone +from pathlib import Path + +import pytest +from typer.testing import CliRunner + +from soup_cli.cli import app +from soup_cli.utils.canary_router import ( + BucketStats, + CanaryPolicy, + rollback, + route, +) +from soup_cli.utils.loop_budget import ( + check_budget, + parse_budget_string, + reset_daily_counter_if_new_day, +) +from soup_cli.utils.loop_daemon import ( + WatchConfig, + evaluate_canary_verdict, + maybe_rollback, + run_once, + watch, +) +from soup_cli.utils.loop_iteration import ( + IterationRecord, + list_iterations, + new_iteration_id, + read_iteration, + write_iteration, +) +from soup_cli.utils.loop_state import ( + LOOP_STATUSES, + LoopState, + default_state_path, + init_state, + read_state, + write_state, +) + +runner = CliRunner() + + +# --------------------------------------------------------------------------- +# Part A — LoopState validation +# --------------------------------------------------------------------------- + + +class TestLoopState: + def test_default_status_is_stopped(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + assert s.status == "stopped" + + def test_status_allowlist(self): + assert LOOP_STATUSES == frozenset({"running", "paused", "stopped"}) + + def test_invalid_status_rejected(self): + with pytest.raises(ValueError, match="status must be"): + LoopState(served_model="m", eval_suite="e", baseline="b", status="weird") + + @pytest.mark.parametrize("field_name", ["served_model", "eval_suite", "baseline"]) + def test_empty_required_field_rejected(self, field_name): + kwargs = {"served_model": "m", "eval_suite": "e", "baseline": "b"} + kwargs[field_name] = "" + with pytest.raises(ValueError, match="must not be empty"): + LoopState(**kwargs) + + def test_null_byte_in_field_rejected(self): + with pytest.raises(ValueError, match="NUL"): + LoopState(served_model="m\x00", eval_suite="e", baseline="b") + + def test_oversize_string_rejected(self): + with pytest.raises(ValueError, match="exceeds"): + LoopState(served_model="m" * 1000, eval_suite="e", baseline="b") + + def test_non_string_rejected(self): + with pytest.raises(TypeError): + LoopState(served_model=123, eval_suite="e", baseline="b") # type: ignore + + @pytest.mark.parametrize( + "counter", + [ + "traces_collected", + "pairs_distilled", + "runs_gated", + "adapters_shipped", + "iteration_count", + "runs_today", + ], + ) + def test_counter_rejects_negative(self, counter): + kwargs = {"served_model": "m", "eval_suite": "e", "baseline": "b", counter: -1} + with pytest.raises(ValueError): + LoopState(**kwargs) + + @pytest.mark.parametrize( + "counter", + [ + "traces_collected", + "pairs_distilled", + "runs_gated", + "adapters_shipped", + ], + ) + def test_counter_rejects_bool(self, counter): + kwargs = {"served_model": "m", "eval_suite": "e", "baseline": "b", counter: True} + with pytest.raises(ValueError): + LoopState(**kwargs) + + def test_canary_traffic_pct_bounds(self): + with pytest.raises(ValueError, match=r"\[0, 100\]"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + canary_traffic_pct=150, + ) + + def test_canary_traffic_pct_bool_rejected(self): + with pytest.raises(ValueError, match="numeric"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + canary_traffic_pct=True, + ) + + def test_monthly_budget_negative_rejected(self): + with pytest.raises(ValueError): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + monthly_budget_usd=-1, + ) + + def test_max_runs_per_day_zero_rejected(self): + with pytest.raises(ValueError): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + max_runs_per_day=0, + ) + + def test_frozen(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(dataclasses.FrozenInstanceError): + s.status = "running" # type: ignore + + def test_to_dict_is_mapping_proxy(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + d = s.to_dict() + with pytest.raises(TypeError): + d["status"] = "running" # type: ignore + + def test_with_status_returns_new_instance(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + s2 = s.with_status("running") + assert s.status == "stopped" + assert s2.status == "running" + assert s2.updated_at != "" + + def test_with_status_rejects_unknown(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError): + s.with_status("nope") + + def test_bumped_increments_counter(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + s2 = s.bumped(traces_collected=3, pairs_distilled=2) + assert s2.traces_collected == 3 + assert s2.pairs_distilled == 2 + + def test_bumped_rejects_unknown_counter(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError, match="unknown counter"): + s.bumped(nonexistent=1) + + def test_bumped_rejects_negative_delta(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError): + s.bumped(traces_collected=-1) + + def test_bumped_rejects_bool_delta(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError): + s.bumped(traces_collected=True) + + +# --------------------------------------------------------------------------- +# Part A — state file I/O +# --------------------------------------------------------------------------- + + +class TestStateIO: + def test_default_state_path(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + assert default_state_path().endswith(os.path.join(".soup", "loop.yaml")) + + def test_init_state_creates_file(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + state, path = init_state("m", "e", "b") + assert os.path.exists(path) + assert state.served_model == "m" + assert state.status == "stopped" + + def test_init_state_refuses_overwrite(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + with pytest.raises(FileExistsError): + init_state("m2", "e2", "b2") + + def test_init_state_force_overwrites(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + state, _ = init_state("m2", "e2", "b2", force=True) + assert state.served_model == "m2" + + def test_write_read_roundtrip(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + s = LoopState( + served_model="model-a", + eval_suite="suite.yaml", + baseline="registry://abc", + status="running", + traces_collected=42, + ) + write_state(s) + reloaded = read_state() + assert reloaded.served_model == "model-a" + assert reloaded.status == "running" + assert reloaded.traces_collected == 42 + + def test_write_state_non_loopstate_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + with pytest.raises(TypeError): + write_state({"foo": "bar"}) # type: ignore + + def test_write_state_outside_cwd_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + s = LoopState(served_model="m", eval_suite="e", baseline="b") + outside = str(tmp_path.parent / "escape.yaml") + with pytest.raises(ValueError, match="cwd"): + write_state(s, outside) + + def test_read_state_missing_raises(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + with pytest.raises(FileNotFoundError): + read_state() + + def test_read_state_invalid_json_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + target = tmp_path / ".soup" / "loop.yaml" + target.parent.mkdir() + target.write_text("not json") + with pytest.raises(ValueError, match="JSON"): + read_state() + + def test_read_state_non_dict_root_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + target = tmp_path / ".soup" / "loop.yaml" + target.parent.mkdir() + target.write_text("[]") + with pytest.raises(ValueError, match="object"): + read_state() + + def test_read_state_oversize_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + target = tmp_path / ".soup" / "loop.yaml" + target.parent.mkdir() + target.write_text("x" * (2 * 1024 * 1024)) + with pytest.raises(ValueError, match="1 MiB"): + read_state() + + def test_read_state_drops_unknown_fields(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + target = tmp_path / ".soup" / "loop.yaml" + target.parent.mkdir() + target.write_text( + json.dumps( + { + "served_model": "m", + "eval_suite": "e", + "baseline": "b", + "future_field_v59": "ignored", + } + ) + ) + s = read_state() + assert s.served_model == "m" + + @pytest.mark.skipif(os.name == "nt", reason="POSIX symlink test") + def test_write_state_rejects_symlink(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + sd = tmp_path / ".soup" + sd.mkdir() + target = sd / "loop.yaml" + target.symlink_to(tmp_path / "elsewhere") + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError, match="symlink"): + write_state(s, str(target)) + + def test_null_byte_path_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError, match="NUL"): + write_state(s, "foo\x00bar") + + +# --------------------------------------------------------------------------- +# Part B — canary router +# --------------------------------------------------------------------------- + + +class TestCanaryPolicy: + def test_stable_only_default(self): + p = CanaryPolicy(stable="adapter-a") + assert p.canary is None + assert p.traffic_pct == 0.0 + + def test_empty_stable_rejected(self): + with pytest.raises(ValueError): + CanaryPolicy(stable="") + + def test_null_byte_rejected(self): + with pytest.raises(ValueError): + CanaryPolicy(stable="a\x00b") + + def test_canary_same_as_stable_rejected(self): + with pytest.raises(ValueError, match="differ"): + CanaryPolicy(stable="a", canary="a") + + def test_traffic_pct_out_of_range(self): + with pytest.raises(ValueError): + CanaryPolicy(stable="a", canary="b", traffic_pct=150) + + def test_traffic_pct_bool_rejected(self): + with pytest.raises(ValueError): + CanaryPolicy(stable="a", canary="b", traffic_pct=True) + + def test_traffic_pct_nan_rejected(self): + with pytest.raises(ValueError): + CanaryPolicy(stable="a", canary="b", traffic_pct=float("nan")) + + def test_traffic_without_canary_rejected(self): + with pytest.raises(ValueError, match="cannot route"): + CanaryPolicy(stable="a", traffic_pct=5) + + def test_sticky_bool_required(self): + with pytest.raises(ValueError): + CanaryPolicy(stable="a", sticky_on_rollback=1) # type: ignore + + def test_frozen(self): + p = CanaryPolicy(stable="a") + with pytest.raises(dataclasses.FrozenInstanceError): + p.stable = "x" # type: ignore + + +class TestRouting: + def test_stable_only_policy_routes_to_stable(self): + p = CanaryPolicy(stable="A") + for key in ("k1", "k2", "k3"): + d = route(p, key) + assert d.adapter == "A" + assert d.bucket == "stable" + + def test_zero_pct_canary_still_routes_stable(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=0.0) + d = route(p, "anything") + assert d.bucket == "stable" + + def test_full_100_pct_routes_all_canary(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=100.0) + for key in ("k1", "k2", "k3"): + d = route(p, key) + assert d.adapter == "B" + assert d.bucket == "canary" + + def test_deterministic_routing(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=25.0) + assert route(p, "abc").adapter == route(p, "abc").adapter + assert route(p, "xyz").adapter == route(p, "xyz").adapter + + def test_split_approximates_traffic_pct(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=25.0) + canary_count = sum( + 1 for i in range(2000) if route(p, f"key-{i}").bucket == "canary" + ) + # 25% of 2000 = 500; tolerate ±15% relative drift on a uniform hash. + assert 350 <= canary_count <= 650 + + def test_empty_key_rejected(self): + p = CanaryPolicy(stable="A") + with pytest.raises(ValueError): + route(p, "") + + def test_non_policy_rejected(self): + with pytest.raises(TypeError): + route("not-policy", "k") # type: ignore + + +class TestRollback: + def test_rollback_clears_canary(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=10.0) + cleared = rollback(p) + assert cleared.canary is None + assert cleared.traffic_pct == 0.0 + assert cleared.stable == "A" + + def test_rollback_after_clears_route_to_stable(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=100.0) + cleared = rollback(p) + assert route(cleared, "anykey").adapter == "A" + + def test_rollback_reason_required(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=5.0) + with pytest.raises(ValueError): + rollback(p, reason="") + + def test_rollback_non_policy_rejected(self): + with pytest.raises(TypeError): + rollback("notpolicy") # type: ignore + + +class TestBucketStats: + def test_record_and_verdict_ok(self): + s = BucketStats() + for _ in range(50): + s.record("stable", True) + s.record("canary", True) + assert s.verdict() == "OK" + + def test_verdict_unknown_below_min_samples(self): + s = BucketStats() + for _ in range(5): + s.record("canary", True) + assert s.verdict(min_samples=30) == "UNKNOWN" + + def test_verdict_major_on_regression(self): + s = BucketStats() + for _ in range(50): + s.record("stable", True) + # canary fails most of the time + for _ in range(50): + s.record("canary", False) + assert s.verdict() == "MAJOR" + + def test_record_invalid_bucket(self): + s = BucketStats() + with pytest.raises(ValueError): + s.record("middle", True) + + def test_record_non_bool_ok(self): + s = BucketStats() + with pytest.raises(ValueError): + s.record("stable", 1) # type: ignore + + def test_verdict_invalid_min_samples(self): + s = BucketStats() + with pytest.raises(ValueError): + s.verdict(min_samples=0) + + def test_verdict_invalid_threshold(self): + s = BucketStats() + for _ in range(40): + s.record("canary", True) + with pytest.raises(ValueError): + s.verdict(regression_threshold=1.5) + + def test_snapshot_mapping_proxy(self): + s = BucketStats() + s.record("stable", True) + snap = s.snapshot() + with pytest.raises(TypeError): + snap["stable_ok"] = 999 # type: ignore + + +# --------------------------------------------------------------------------- +# Part C — budget guardrails +# --------------------------------------------------------------------------- + + +class TestParseBudget: + @pytest.mark.parametrize( + "raw,expected", + [("50", 50.0), ("50usd", 50.0), ("100 USD", 100.0), ("0", 0.0), ("0.5", 0.5)], + ) + def test_happy(self, raw, expected): + assert parse_budget_string(raw) == expected + + def test_empty_rejected(self): + with pytest.raises(ValueError): + parse_budget_string("") + + def test_garbage_rejected(self): + with pytest.raises(ValueError): + parse_budget_string("abc") + + def test_null_byte_rejected(self): + with pytest.raises(ValueError): + parse_budget_string("50\x00usd") + + def test_negative_rejected(self): + with pytest.raises(ValueError): + parse_budget_string("-5") + + def test_overflow_rejected(self): + with pytest.raises(ValueError): + parse_budget_string("10000000") + + def test_non_string_rejected(self): + with pytest.raises(TypeError): + parse_budget_string(50) # type: ignore + + +class TestCheckBudget: + def test_happy_within_budget(self): + d = check_budget( + estimated_run_usd=1.0, + spent_so_far_usd=5.0, + monthly_budget_usd=50.0, + runs_today=0, + max_runs_per_day=10, + ) + assert d.proceed is True + assert d.projected_total_usd == 6.0 + + def test_blocked_by_budget(self): + d = check_budget( + estimated_run_usd=10.0, + spent_so_far_usd=45.0, + monthly_budget_usd=50.0, + runs_today=0, + max_runs_per_day=10, + ) + assert d.proceed is False + assert "budget" in d.reason + + def test_blocked_by_daily_cap(self): + d = check_budget( + estimated_run_usd=1.0, + spent_so_far_usd=0.0, + monthly_budget_usd=50.0, + runs_today=3, + max_runs_per_day=3, + ) + assert d.proceed is False + assert "daily cap" in d.reason + + def test_none_budget_means_unlimited(self): + d = check_budget( + estimated_run_usd=1e6, + spent_so_far_usd=0, + monthly_budget_usd=None, + runs_today=0, + max_runs_per_day=None, + ) + assert d.proceed is True + + def test_negative_estimate_rejected(self): + with pytest.raises(ValueError): + check_budget( + estimated_run_usd=-1, + spent_so_far_usd=0, + monthly_budget_usd=None, + runs_today=0, + max_runs_per_day=None, + ) + + def test_nan_estimate_rejected(self): + with pytest.raises(ValueError): + check_budget( + estimated_run_usd=float("nan"), + spent_so_far_usd=0, + monthly_budget_usd=None, + runs_today=0, + max_runs_per_day=None, + ) + + def test_bool_estimate_rejected(self): + with pytest.raises(ValueError): + check_budget( + estimated_run_usd=True, + spent_so_far_usd=0, + monthly_budget_usd=None, + runs_today=0, + max_runs_per_day=None, + ) + + def test_negative_runs_today_rejected(self): + with pytest.raises(ValueError): + check_budget( + estimated_run_usd=0, + spent_so_far_usd=0, + monthly_budget_usd=None, + runs_today=-1, + max_runs_per_day=None, + ) + + def test_zero_max_runs_per_day_rejected(self): + with pytest.raises(ValueError): + check_budget( + estimated_run_usd=0, + spent_so_far_usd=0, + monthly_budget_usd=None, + runs_today=0, + max_runs_per_day=0, + ) + + +class TestDailyCounter: + def test_same_day_keeps_count(self): + today = datetime(2026, 5, 15, tzinfo=timezone.utc) + out, date = reset_daily_counter_if_new_day(3, "2026-05-15", now=today) + assert out == 3 + assert date == "2026-05-15" + + def test_new_day_resets(self): + today = datetime(2026, 5, 16, tzinfo=timezone.utc) + out, date = reset_daily_counter_if_new_day(3, "2026-05-15", now=today) + assert out == 0 + assert date == "2026-05-16" + + def test_none_prior_date_treated_as_new_day(self): + today = datetime(2026, 5, 15, tzinfo=timezone.utc) + out, _ = reset_daily_counter_if_new_day(5, None, now=today) + assert out == 0 + + def test_negative_runs_today_rejected(self): + with pytest.raises(ValueError): + reset_daily_counter_if_new_day(-1, "2026-05-15") + + +# --------------------------------------------------------------------------- +# Part D — iteration artifact +# --------------------------------------------------------------------------- + + +def _make_record(**overrides): + base = dict( + iteration_id="iter-20260515T000000-abcdef01", + started_at="2026-05-15T00:00:00+00:00", + finished_at="2026-05-15T00:05:00+00:00", + pairs_harvested=10, + run_id="run-abc", + gate_verdict="OK", + canary_verdict=None, + shipped=True, + rolled_back=False, + estimated_cost_usd=0.50, + ) + base.update(overrides) + return IterationRecord(**base) + + +class TestIterationRecord: + def test_happy(self): + r = _make_record() + assert r.shipped is True + assert r.gate_verdict == "OK" + + def test_frozen(self): + r = _make_record() + with pytest.raises(dataclasses.FrozenInstanceError): + r.shipped = False # type: ignore + + def test_invalid_gate_verdict_rejected(self): + with pytest.raises(ValueError, match="gate_verdict"): + _make_record(gate_verdict="WEIRD") + + def test_invalid_canary_verdict_rejected(self): + with pytest.raises(ValueError, match="canary_verdict"): + _make_record(canary_verdict="WEIRD") + + def test_shipped_must_be_bool(self): + with pytest.raises(ValueError): + _make_record(shipped=1) # type: ignore + + def test_negative_pairs_rejected(self): + with pytest.raises(ValueError): + _make_record(pairs_harvested=-1) + + def test_negative_cost_rejected(self): + with pytest.raises(ValueError): + _make_record(estimated_cost_usd=-0.5) + + def test_iteration_id_path_separator_rejected(self): + with pytest.raises(ValueError): + _make_record(iteration_id="iter/escape") + + def test_iteration_id_empty_rejected(self): + with pytest.raises(ValueError): + _make_record(iteration_id="") + + def test_iteration_id_null_rejected(self): + with pytest.raises(ValueError): + _make_record(iteration_id="iter\x00x") + + def test_oversize_notes_rejected(self): + with pytest.raises(ValueError): + _make_record(notes="x" * 5000) + + +class TestIterationIO: + def test_write_read_roundtrip(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + r = _make_record() + path = write_iteration(r) + assert os.path.exists(path) + loaded = read_iteration(r.iteration_id) + assert loaded.iteration_id == r.iteration_id + assert loaded.shipped is True + + def test_write_outside_cwd_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + r = _make_record() + with pytest.raises(ValueError): + write_iteration(r, base_dir=str(tmp_path.parent / "escape")) + + def test_read_missing_raises(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + with pytest.raises(FileNotFoundError): + read_iteration("iter-missing") + + def test_list_empty_dir(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + assert list_iterations() == () + + def test_list_sorted(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + write_iteration(_make_record(iteration_id="iter-001")) + write_iteration(_make_record(iteration_id="iter-003")) + write_iteration(_make_record(iteration_id="iter-002")) + assert list_iterations() == ("iter-001", "iter-002", "iter-003") + + def test_new_iteration_id_unique(self): + ids = {new_iteration_id() for _ in range(20)} + assert len(ids) == 20 + + def test_new_iteration_id_passes_validation(self): + # Round-trips through IterationRecord without raising + r = _make_record(iteration_id=new_iteration_id()) + assert r.iteration_id.startswith("iter-") + + def test_write_non_record_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + with pytest.raises(TypeError): + write_iteration({"foo": "bar"}) # type: ignore + + def test_read_invalid_manifest_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + d = tmp_path / ".soup-loops" / "iter-bad" + d.mkdir(parents=True) + (d / "iteration.json").write_text("[]") + with pytest.raises(ValueError, match="object"): + read_iteration("iter-bad") + + def test_list_skips_invalid_id_directories(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + # Legitimate + write_iteration(_make_record(iteration_id="iter-ok")) + # Junk directory with iteration.json but bad id + bad = tmp_path / ".soup-loops" / "weird\x00" + # OSes may reject NUL in path; if so, just skip + try: + bad.mkdir(parents=True) + (bad / "iteration.json").write_text("{}") + except (OSError, ValueError): + pass + out = list_iterations() + assert "iter-ok" in out + + +# --------------------------------------------------------------------------- +# Watch daemon +# --------------------------------------------------------------------------- + + +class TestRunOnce: + def test_default_callbacks_record_iteration(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + state, _ = init_state("m", "e", "b") + state = state.with_status("running") + write_state(state) + cfg = WatchConfig() + new_state, record, decision = run_once(state, cfg) + assert decision.proceed is True + assert record.gate_verdict == "SKIPPED" # default train stub sets skipped + assert new_state.iteration_count == 1 + + def test_budget_skip_records_iteration_with_zero_counters(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + state = LoopState( + served_model="m", + eval_suite="e", + baseline="b", + status="running", + monthly_budget_usd=10.0, + spent_this_month_usd=10.0, + ) + cfg = WatchConfig(cost_fn=lambda s: 5.0) + new_state, record, decision = run_once(state, cfg) + assert decision.proceed is False + assert record.gate_verdict == "SKIPPED" + assert "budget" in record.notes.lower() + assert new_state.iteration_count == 0 # not bumped on skip + + def test_custom_callbacks_invoked(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + state, _ = init_state("m", "e", "b") + state = state.with_status("running") + write_state(state) + cfg = WatchConfig( + harvest_fn=lambda s: {"pairs_harvested": 5, "traces_collected": 100}, + train_fn=lambda s, c: {"run_id": "run-X", "skipped": False}, + gate_fn=lambda s, c: {"gate_verdict": "OK"}, + deploy_fn=lambda s, c: {"deployed": True, "canary_verdict": "OK"}, + cost_fn=lambda s: 0.10, + ) + new_state, record, decision = run_once(state, cfg) + assert decision.proceed is True + assert record.pairs_harvested == 5 + assert record.run_id == "run-X" + assert record.gate_verdict == "OK" + assert record.canary_verdict == "OK" + assert record.shipped is True + assert new_state.adapters_shipped == 1 + assert new_state.pairs_distilled == 5 + assert new_state.spent_this_month_usd == pytest.approx(0.10) + + def test_run_once_rejects_non_state(self): + cfg = WatchConfig() + with pytest.raises(TypeError): + run_once("notstate", cfg) # type: ignore + + def test_run_once_rejects_non_config(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(TypeError): + run_once(s, "notcfg") # type: ignore + + def test_gate_verdict_normalised(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + state, _ = init_state("m", "e", "b") + state = state.with_status("running") + write_state(state) + cfg = WatchConfig(gate_fn=lambda s, c: {"gate_verdict": "GIBBERISH"}) + _, record, _ = run_once(state, cfg) + assert record.gate_verdict == "SKIPPED" + + +class TestWatchConfig: + def test_default_construct(self): + cfg = WatchConfig() + assert cfg.poll_interval_sec == 60.0 + + def test_invalid_poll_interval(self): + with pytest.raises(ValueError): + WatchConfig(poll_interval_sec=0.5) + with pytest.raises(ValueError): + WatchConfig(poll_interval_sec=10000) + with pytest.raises(ValueError): + WatchConfig(poll_interval_sec=float("nan")) + + def test_non_callable_harvest_rejected(self): + with pytest.raises(ValueError): + WatchConfig(harvest_fn="not callable") # type: ignore + + def test_negative_max_iterations_rejected(self): + with pytest.raises(ValueError): + WatchConfig(max_iterations=-1) + + +class TestWatchDaemon: + def test_watch_finite_runs_then_stops(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + cfg = WatchConfig( + poll_interval_sec=1.0, + max_iterations=3, + ) + final_state, ran = watch(cfg) + assert ran == 3 + assert final_state.iteration_count == 3 + assert final_state.status == "stopped" + + def test_watch_respects_pause(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + + # Iteration 1 paths through and then we flip to paused via callback. + def _pause_after_first(record): + s = read_state() + write_state(s.with_status("stopped")) + + cfg = WatchConfig( + poll_interval_sec=1.0, + max_iterations=10, + on_iteration=_pause_after_first, + ) + final_state, ran = watch(cfg) + # Either 1 or 2 — the on_iteration fires before the next read. + assert ran >= 1 + assert final_state.status == "stopped" + + +class TestRollbackOrchestration: + def test_maybe_rollback_no_op_on_ok(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=5.0) + assert maybe_rollback(p, "OK").canary == "B" + + def test_maybe_rollback_clears_on_major(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=5.0) + assert maybe_rollback(p, "MAJOR").canary is None + + def test_maybe_rollback_unknown_no_op(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=5.0) + assert maybe_rollback(p, "UNKNOWN").canary == "B" + + def test_evaluate_canary_verdict_wraps_stats(self): + stats = BucketStats() + for _ in range(40): + stats.record("canary", True) + assert evaluate_canary_verdict(stats) == "OK" + + def test_maybe_rollback_non_policy_rejected(self): + with pytest.raises(TypeError): + maybe_rollback("notpolicy", "MAJOR") # type: ignore + + def test_maybe_rollback_non_str_verdict_rejected(self): + p = CanaryPolicy(stable="A") + with pytest.raises(TypeError): + maybe_rollback(p, 5) # type: ignore + + +# --------------------------------------------------------------------------- +# CLI smoke +# --------------------------------------------------------------------------- + + +class TestCLI: + def test_loop_help(self): + result = runner.invoke(app, ["loop", "--help"]) + assert result.exit_code == 0, result.output + assert "init" in result.output + assert "watch" in result.output + + def test_init_command(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + result = runner.invoke( + app, + ["loop", "init", "model-a", "--eval", "suite.yaml", "--baseline", "ref"], + ) + assert result.exit_code == 0, result.output + assert (tmp_path / ".soup" / "loop.yaml").exists() + + def test_init_refuses_overwrite(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + result = runner.invoke( + app, ["loop", "init", "m2", "--eval", "e", "--baseline", "b"] + ) + assert result.exit_code == 2, result.output + assert "already exists" in result.output + + def test_init_force_overwrites(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + result = runner.invoke( + app, + ["loop", "init", "m2", "--eval", "e", "--baseline", "b", "--force"], + ) + assert result.exit_code == 0, result.output + + def test_init_invalid_budget(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + result = runner.invoke( + app, + [ + "loop", + "init", + "m", + "--eval", + "e", + "--baseline", + "b", + "--monthly-budget", + "garbage", + ], + ) + assert result.exit_code == 2 + + def test_status_without_init(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + result = runner.invoke(app, ["loop", "status"]) + assert result.exit_code == 2 + assert "init" in result.output.lower() + + def test_status_after_init(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + result = runner.invoke(app, ["loop", "status"]) + assert result.exit_code == 0, result.output + assert "stopped" in result.output + + def test_pause_resume_cycle(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + # Mark running manually so pause has something to flip. + s = read_state() + write_state(s.with_status("running")) + r = runner.invoke(app, ["loop", "pause"]) + assert r.exit_code == 0 + assert read_state().status == "paused" + r = runner.invoke(app, ["loop", "resume"]) + assert r.exit_code == 0 + assert read_state().status == "running" + + def test_pause_when_stopped(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke(app, ["loop", "pause"]) + assert r.exit_code == 0 + assert "already stopped" in r.output + + def test_resume_when_not_paused(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke(app, ["loop", "resume"]) + assert r.exit_code == 0 + assert "not paused" in r.output + + def test_watch_max_iterations(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke( + app, + [ + "loop", + "watch", + "--foreground", + "--max-iterations", + "2", + "--poll-interval", + "1", + ], + ) + assert r.exit_code == 0, r.output + assert "iterations=2" in r.output + + def test_watch_detach_and_foreground_mutex(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke( + app, ["loop", "watch", "--foreground", "--detach"] + ) + assert r.exit_code == 2 + + def test_canary_command(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "model-a", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke( + app, ["loop", "canary", "model-b", "--traffic", "10%"] + ) + assert r.exit_code == 0, r.output + s = read_state() + assert s.canary_active == "model-b" + assert s.canary_traffic_pct == 10.0 + + def test_canary_invalid_traffic(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "model-a", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke( + app, ["loop", "canary", "model-b", "--traffic", "150"] + ) + assert r.exit_code == 2 + + def test_canary_same_as_stable_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "same", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke( + app, ["loop", "canary", "same", "--traffic", "5%"] + ) + assert r.exit_code == 2 + + def test_replay_list_empty(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke(app, ["loop", "replay"]) + assert r.exit_code == 0 + assert "no iterations" in r.output + + def test_replay_show_iteration(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + rec = _make_record() + write_iteration(rec) + r = runner.invoke(app, ["loop", "replay", rec.iteration_id]) + assert r.exit_code == 0 + assert rec.iteration_id in r.output + + def test_replay_unknown_iteration(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "m", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke(app, ["loop", "replay", "iter-nonexistent"]) + assert r.exit_code == 2 + + +# --------------------------------------------------------------------------- +# Source-grep regression guards +# --------------------------------------------------------------------------- + + +_REPO_ROOT = Path(__file__).resolve().parent.parent + + +class TestSourceWiring: + def test_cli_registers_loop_typer(self): + cli_src = (_REPO_ROOT / "soup_cli" / "cli.py").read_text(encoding="utf-8") + assert "from soup_cli.commands import loop as _loop_cmd" in cli_src + assert 'name="loop"' in cli_src + + def test_version_bumped_to_0_58_0(self): + init = (_REPO_ROOT / "soup_cli" / "__init__.py").read_text(encoding="utf-8") + assert '__version__ = "0.58.0"' in init + + def test_no_top_level_torch_import_in_loop_modules(self): + for name in [ + "loop_state.py", + "loop_budget.py", + "loop_iteration.py", + "canary_router.py", + "loop_daemon.py", + ]: + src = (_REPO_ROOT / "soup_cli" / "utils" / name).read_text(encoding="utf-8") + # Check module-level imports only (skip indented imports inside funcs) + for line in src.splitlines(): + if line.startswith(("import torch", "from torch")): + raise AssertionError(f"{name} imports torch at module level") + + def test_command_module_uses_typer_app(self): + src = (_REPO_ROOT / "soup_cli" / "commands" / "loop.py").read_text(encoding="utf-8") + assert "app = typer.Typer(" in src + assert 'name="loop"' in src + + +# --------------------------------------------------------------------------- +# version sanity +# --------------------------------------------------------------------------- + + +def test_version_string(): + from soup_cli import __version__ + + assert __version__ == "0.58.0" + + +# --------------------------------------------------------------------------- +# Code-review wave 2 follow-ups (HIGH #2-#4 + MEDIUM #5-#8 review fixes) +# --------------------------------------------------------------------------- + + +class TestReviewFixWave2: + def test_watch_preserves_paused_status_on_exit(self, tmp_path, monkeypatch): + """HIGH #2: SIGTERM-while-paused must NOT silently promote to stopped.""" + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + s = read_state() + write_state(s.with_status("paused")) + cfg = WatchConfig(poll_interval_sec=1.0, max_iterations=0) + final_state, _ = watch(cfg) + # Was paused before watch; daemon must not have flipped to stopped. + assert final_state.status == "paused", "watch destroyed paused state" + + def test_watch_flips_running_to_stopped_at_exit(self, tmp_path, monkeypatch): + """The legit case still works — running → stopped on max_iterations.""" + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + cfg = WatchConfig(poll_interval_sec=1.0, max_iterations=1) + final_state, _ = watch(cfg) + assert final_state.status == "stopped" + + def test_budget_skip_does_not_write_iteration_manifest(self, tmp_path, monkeypatch): + """HIGH #3: budget-skipped runs do not produce iteration manifests.""" + monkeypatch.chdir(tmp_path) + init_state("m", "e", "b") + s = read_state() + write_state( + dataclasses.replace( + s, + status="running", + monthly_budget_usd=10.0, + spent_this_month_usd=10.0, + ) + ) + cfg = WatchConfig( + poll_interval_sec=1.0, + max_iterations=1, + cost_fn=lambda _s: 5.0, # forces budget rejection + ) + watch(cfg) + # No manifest should exist because the iteration was budget-skipped. + assert list_iterations() == () + + def test_canary_autoroll_persisted(self, tmp_path, monkeypatch): + """HIGH #4: --autoroll-on-regress flag must survive into LoopState.""" + monkeypatch.chdir(tmp_path) + runner.invoke( + app, ["loop", "init", "model-a", "--eval", "e", "--baseline", "b"] + ) + r = runner.invoke( + app, + [ + "loop", + "canary", + "model-b", + "--traffic", + "5%", + "--no-autoroll-on-regress", + ], + ) + assert r.exit_code == 0, r.output + s = read_state() + assert s.canary_autoroll_on_regress is False + # Default True is also exercised by other canary tests above. + + def test_canary_autoroll_bool_validator(self): + """LoopState rejects non-bool canary_autoroll_on_regress.""" + with pytest.raises(ValueError, match="canary_autoroll_on_regress"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + canary_autoroll_on_regress=1, # type: ignore + ) + + def test_route_ceil_at_sub_bucket_fraction(self): + """MEDIUM #5: 0.005 % must allocate ≥1 bucket, not round to 0.""" + p = CanaryPolicy(stable="A", canary="B", traffic_pct=0.005) + # 10_000 buckets * 0.00005 = 0.5 → ceil = 1 bucket reserved for canary. + canary_hits = sum( + 1 for i in range(10_000) if route(p, f"k-{i}").bucket == "canary" + ) + assert canary_hits >= 1, "ceil rounding lost the sub-bucket fraction" + + def test_parse_budget_usd_only(self): + """MEDIUM #6: bare 'usd' / ' usd ' raises friendly explicit error.""" + with pytest.raises(ValueError, match="numeric value"): + parse_budget_string("usd") + with pytest.raises(ValueError, match="numeric value"): + parse_budget_string(" USD ") + + def test_list_iterations_swallows_oserror_on_listdir(self, tmp_path, monkeypatch): + """MEDIUM #8: list_iterations returns () instead of raising on unreadable dir.""" + monkeypatch.chdir(tmp_path) + # Create the dir then monkeypatch os.listdir to raise — simulates a + # permission flap mid-iteration that would otherwise kill the daemon. + d = tmp_path / ".soup-loops" + d.mkdir() + import soup_cli.utils.loop_iteration as li + + def _raise(_p): + raise PermissionError("simulated permission flap") + + monkeypatch.setattr(li.os, "listdir", _raise) + assert li.list_iterations() == () + + +# --------------------------------------------------------------------------- +# TDD-review wave 3: exact-boundary + match= tightening + coverage gaps +# --------------------------------------------------------------------------- + + +class TestReviewFixWave3: + # --- HIGH #1: exact-boundary tests for _MAX_STR_FIELD = 512 ---------- + + def test_str_field_accepts_exactly_512(self): + s = LoopState(served_model="m" * 512, eval_suite="e", baseline="b") + assert len(s.served_model) == 512 + + def test_str_field_rejects_513(self): + with pytest.raises(ValueError, match="exceeds 512"): + LoopState(served_model="m" * 513, eval_suite="e", baseline="b") + + # --- HIGH #2: read_state size cap exact boundary --------------------- + + def test_read_state_at_one_mib_minus_one_accepted(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + target = tmp_path / ".soup" / "loop.yaml" + target.parent.mkdir() + # Pad a valid JSON document to just under 1 MiB. + notes_pad = " " * (1024 * 1024 - 200) + target.write_text( + json.dumps( + { + "served_model": "m", + "eval_suite": "e", + "baseline": "b", + "last_iteration_id": "iter-pad" + notes_pad[:480], + } + ) + ) + size = os.path.getsize(target) + assert size < 1024 * 1024 # under the cap + s = read_state() + assert s.served_model == "m" + + def test_read_state_at_one_mib_plus_one_rejected(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + target = tmp_path / ".soup" / "loop.yaml" + target.parent.mkdir() + target.write_text("x" * (1024 * 1024 + 1)) + with pytest.raises(ValueError, match="1 MiB"): + read_state() + + # --- HIGH #3: match= tightening on critical reject paths ------------- + + def test_monthly_budget_negative_match(self): + with pytest.raises(ValueError, match="monthly_budget_usd must be"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + monthly_budget_usd=-1.0, + ) + + def test_max_runs_per_day_zero_match(self): + with pytest.raises(ValueError, match="max_runs_per_day"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + max_runs_per_day=0, + ) + + def test_with_status_unknown_match(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError, match="status must be"): + s.with_status("not-a-real-status") + + def test_bumped_unknown_counter_match(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError, match="unknown counter"): + s.bumped(banana=1) + + def test_canary_traffic_out_of_range_match(self): + with pytest.raises(ValueError, match=r"\[0, 100\]"): + CanaryPolicy(stable="A", canary="B", traffic_pct=150.0) + + # --- MEDIUM: bool rejection on iteration_count / runs_today --------- + + @pytest.mark.parametrize("counter", ["iteration_count", "runs_today"]) + def test_counter_rejects_bool_iter_runs(self, counter): + kwargs = {"served_model": "m", "eval_suite": "e", "baseline": "b", counter: True} + with pytest.raises(ValueError, match=counter): + LoopState(**kwargs) + + # --- MEDIUM: bool rejection on float budget fields ------------------- + + def test_monthly_budget_bool_rejected(self): + with pytest.raises(ValueError, match="monthly_budget_usd"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + monthly_budget_usd=True, # type: ignore + ) + + def test_spent_this_month_bool_rejected(self): + with pytest.raises(ValueError, match="spent_this_month_usd"): + LoopState( + served_model="m", + eval_suite="e", + baseline="b", + spent_this_month_usd=True, # type: ignore + ) + + # --- MEDIUM: optional-string empty-string rejection ------------------ + + @pytest.mark.parametrize( + "field_name", ["canary_active", "last_iteration_id", "last_run_date"] + ) + def test_optional_str_field_rejects_empty(self, field_name): + kwargs = {"served_model": "m", "eval_suite": "e", "baseline": "b", field_name: ""} + with pytest.raises(ValueError, match="must not be empty"): + LoopState(**kwargs) + + # --- MEDIUM: _check_path empty string at write boundary -------------- + + def test_write_state_rejects_empty_path(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + s = LoopState(served_model="m", eval_suite="e", baseline="b") + with pytest.raises(ValueError, match="must not be empty"): + write_state(s, "") + + # --- LOW: canary traffic_pct lower-bound exact 0 + upper-bound 100 --- + + def test_canary_traffic_pct_zero_accepted(self): + p = CanaryPolicy(stable="A") # traffic_pct=0 default + assert p.traffic_pct == 0.0 + + def test_canary_traffic_pct_negative_rejected(self): + with pytest.raises(ValueError, match=r"\[0, 100\]"): + CanaryPolicy(stable="A", canary="B", traffic_pct=-0.001) + + def test_canary_traffic_pct_exactly_100_accepted(self): + p = CanaryPolicy(stable="A", canary="B", traffic_pct=100.0) + assert p.traffic_pct == 100.0 + + # --- LOW: to_dict keys match dataclass fields (forward-compat lock) -- + + def test_to_dict_keys_match_dataclass_fields(self): + s = LoopState(served_model="m", eval_suite="e", baseline="b") + assert set(s.to_dict().keys()) == set(LoopState.__dataclass_fields__.keys()) + + # --- LOW: read_state preserves created_at/updated_at when present ---- + + def test_read_state_preserves_created_at(self, tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + s = LoopState( + served_model="m", + eval_suite="e", + baseline="b", + created_at="2026-05-15T00:00:00+00:00", + ) + write_state(s) + reloaded = read_state() + # created_at survives the roundtrip; updated_at gets refreshed by write_state. + assert reloaded.created_at == "2026-05-15T00:00:00+00:00"