diff --git a/.github/workflows/cd.yml b/.github/workflows/cd.yml index 90568b8..1857704 100644 --- a/.github/workflows/cd.yml +++ b/.github/workflows/cd.yml @@ -72,6 +72,10 @@ jobs: # One plan per track (model, adapter, policy), each naming what it # rolls back to and the criteria that trigger it. run: python -m llm_router.canary > canary-plan.json + - name: Render the atomic deploy command + # What the deploy step will run once a destination is configured: + # an atomic upgrade, which Helm rolls back itself on failed readiness. + run: python -m llm_router.rollout deploy > deploy-command.json - name: Render the governance plan # What a sync would record in MLflow for this catalog: every model and # adapter revision, its stage, and its benchmark evidence. Planned @@ -92,6 +96,7 @@ jobs: path: | canary-plan.json governance-plan.json + deploy-command.json config/ray-serve.yaml chart rendered-base.yaml diff --git a/README.md b/README.md index 543ba04..64a01a1 100644 --- a/README.md +++ b/README.md @@ -479,9 +479,28 @@ controls that rollout evaluates the plan through the same rule: python -m llm_router.canary --plan model:deploy-0002 --observation observed.json --baseline stable.json ``` -It exits `0` to promote, `2` to hold, and `3` to roll back, and prints the rollback target. No -rollout controller is wired to it yet, so on those two tracks the decision is automatic and the -action is not. +It exits `0` to promote, `2` to hold, and `3` to roll back, and prints the rollback target. + +`llm_router.rollout` acts on that verdict, so neither kind of rollback waits for a person: + +```bash +python -m llm_router.rollout deploy --set serving.mode=ray --execute +python -m llm_router.rollout evaluate --plan model:deploy-0002 \ + --observation observed.json --baseline stable.json --execute +``` + +- **Failed readiness.** `deploy` runs `helm upgrade --install --atomic`: if a workload does not + become ready within the timeout, Helm restores the previous release. +- **Failed canary criteria.** `evaluate` judges the plan and, on a rollback verdict for the model + or policy track, runs `helm rollback` to the release that was live before. It exits `4` if that + command fails. +- A plan with no recorded rollback target runs nothing and stops: restoring whatever Helm happens + to hold is not a rollback to a known state. +- Without `--execute` both print the command and change nothing. + +The observation is a file you supply: nothing here collects a model or policy canary's metrics +from Prometheus, and nothing schedules `evaluate`. The Helm commands have been checked against +`helm` for their flags, not run against a cluster. Promotion is never automatic on any track. It is a catalog change and goes through review. diff --git a/src/llm_router/rollout.py b/src/llm_router/rollout.py new file mode 100644 index 0000000..26f6a2f --- /dev/null +++ b/src/llm_router/rollout.py @@ -0,0 +1,198 @@ +"""Acting on a rollout: deploy atomically, and roll back when a canary fails. + +Section 16 asks for automatic rollback after failed readiness or failed canary +criteria. Readiness is covered by deploying with Helm's ``--atomic``: a release +whose pods never become ready is rolled back by Helm itself. Canary criteria +are judged by ``llm_router.canary``; this module turns that verdict into the +command that carries it out, so the model and policy tracks roll back without +a person reading a decision and typing one. + +Adapters are not handled here. The gateway withdraws a failing adapter canary +on its own, with no redeploy. +""" + +import json +import subprocess +from collections.abc import Callable, Sequence +from pathlib import Path + +from pydantic import BaseModel + +from llm_router.canary import ( + CanaryDecision, + CanaryObservation, + CanaryPlan, + CanaryTrack, + canary_plans, + evaluate, +) + +DEFAULT_RELEASE = "llm-routing" +DEFAULT_NAMESPACE = "llm-routing" +DEFAULT_CHART = "deploy/helm/llm-routing" +DEFAULT_TIMEOUT = "10m" +EXIT_CODES = {"promote": 0, "hold": 2, "rollback": 3} +# The decision was to roll back and the rollback itself did not succeed. +EXIT_ROLLBACK_FAILED = 4 + +Runner = Callable[[Sequence[str]], int] + + +class RolloutStep(BaseModel): + """What a canary verdict means for the release, and the command that does it.""" + + plan: str + action: str + reasons: tuple[str, ...] + rollback_to: str | None + command: tuple[str, ...] = () + note: str = "" + + +def deploy_command( + *, + release: str = DEFAULT_RELEASE, + namespace: str = DEFAULT_NAMESPACE, + chart: str = DEFAULT_CHART, + values: Sequence[str] = (), + timeout: str = DEFAULT_TIMEOUT, +) -> tuple[str, ...]: + """Install or upgrade the release so that failed readiness undoes it. + + ``--atomic`` waits for every workload to become ready and, if one does not + within the timeout, restores the previous release. + """ + + command = [ + "helm", + "upgrade", + "--install", + release, + chart, + "--namespace", + namespace, + "--create-namespace", + "--atomic", + "--timeout", + timeout, + ] + for value in values: + command += ["--set", value] + return tuple(command) + + +def rollback_command( + *, + release: str = DEFAULT_RELEASE, + namespace: str = DEFAULT_NAMESPACE, + timeout: str = DEFAULT_TIMEOUT, +) -> tuple[str, ...]: + """Restore the release that was live before the current one.""" + + return ("helm", "rollback", release, "--namespace", namespace, "--wait", "--timeout", timeout) + + +def next_step( + plan: CanaryPlan, + decision: CanaryDecision, + *, + release: str = DEFAULT_RELEASE, + namespace: str = DEFAULT_NAMESPACE, +) -> RolloutStep: + """Translate a verdict on one plan into the step that carries it out.""" + + step = RolloutStep( + plan=plan.id, + action=decision.action, + reasons=decision.reasons, + rollback_to=plan.rollback_to, + ) + if decision.action != "rollback": + return step + if plan.track is CanaryTrack.ADAPTER: + return step.model_copy( + update={"note": "the gateway withdraws a failing adapter itself; nothing to redeploy"} + ) + if plan.rollback_to is None: + # With nothing recorded to return to, a rollback would restore + # whatever Helm happens to hold. Stopping is the only safe step. + return step.model_copy( + update={"note": f"{plan.rollback_action}; no command is run without a recorded target"} + ) + return step.model_copy( + update={ + "command": rollback_command(release=release, namespace=namespace), + "note": plan.rollback_action, + } + ) + + +def run_command(command: Sequence[str]) -> int: # pragma: no cover - runs helm + return subprocess.run(list(command), check=False).returncode + + +def main(argv: Sequence[str] | None = None, runner: Runner = run_command) -> int: + """Deploy the release, or act on a canary observation. + + ``deploy`` prints the atomic install command. ``evaluate`` judges one plan + against an observation and prints the step that follows. Either runs its + command only with --execute; without it nothing is changed. + + ``evaluate`` exits 0 to promote, 2 to hold, 3 after a rollback decision, + and 4 if the rollback command itself failed. + """ + + import argparse + + from llm_router.registry import load_registry + + parser = argparse.ArgumentParser(description=main.__doc__) + parser.add_argument("command", choices=["deploy", "evaluate"]) + parser.add_argument("--release", default=DEFAULT_RELEASE) + parser.add_argument("--namespace", default=DEFAULT_NAMESPACE) + parser.add_argument("--chart", default=DEFAULT_CHART) + parser.add_argument("--set", action="append", default=[], dest="values") + parser.add_argument("--catalog", default="config/registry.yaml") + parser.add_argument("--plan", help="plan id to evaluate, for example model:deploy-0002") + parser.add_argument("--observation", help="JSON file holding the observed canary metrics") + parser.add_argument("--baseline", help="JSON file holding the stable baseline metrics") + parser.add_argument("--execute", action="store_true", help="run the command, not just print it") + arguments = parser.parse_args(argv) + + if arguments.command == "deploy": + command = deploy_command( + release=arguments.release, + namespace=arguments.namespace, + chart=arguments.chart, + values=arguments.values, + ) + print(json.dumps({"command": list(command), "executed": arguments.execute}, indent=2)) + return runner(command) if arguments.execute else 0 + + plans = {plan.id: plan for plan in canary_plans(load_registry(arguments.catalog))} + plan = plans.get(arguments.plan or "") + if plan is None or arguments.observation is None: + parser.error(f"evaluate needs --observation and --plan, one of: {', '.join(plans)}") + + def read(path: str) -> CanaryObservation: + return CanaryObservation.model_validate_json(Path(path).read_text(encoding="utf-8")) + + decision = evaluate( + plan.criteria, + read(arguments.observation), + read(arguments.baseline) if arguments.baseline else None, + ) + step = next_step(plan, decision, release=arguments.release, namespace=arguments.namespace) + executed = bool(step.command) and arguments.execute + failed = executed and runner(step.command) != 0 + print( + json.dumps( + {**step.model_dump(mode="json"), "executed": executed, "succeeded": not failed}, + indent=2, + ) + ) + return EXIT_ROLLBACK_FAILED if failed else EXIT_CODES[decision.action] + + +if __name__ == "__main__": # pragma: no cover - command-line entry point + raise SystemExit(main()) diff --git a/tests/unit/test_rollout.py b/tests/unit/test_rollout.py new file mode 100644 index 0000000..b573818 --- /dev/null +++ b/tests/unit/test_rollout.py @@ -0,0 +1,181 @@ +import json +import shutil +import subprocess +from collections.abc import Sequence +from pathlib import Path + +import pytest + +from llm_router.canary import CanaryDecision, canary_plans +from llm_router.registry import load_registry +from llm_router.rollout import deploy_command, main, next_step, rollback_command + +PLANS = {plan.id: plan for plan in canary_plans(load_registry("config/registry.yaml"))} +MODEL, ADAPTER, POLICY = "model:deploy-0002", "adapter:claims-extraction-lora-next", "policy:v1" +HEALTHY = {"requests": 800, "errors": 1, "p95_latency_ms": 400, "quality": 0.9} +FAILING = {"requests": 800, "errors": 200, "p95_latency_ms": 400, "quality": 0.9} +EARLY = {"requests": 10, "errors": 0} +ROLLBACK = CanaryDecision(action="rollback", reasons=("error rate 0.250 exceeds 0.010",)) + + +class Recorder: + def __init__(self, result: int = 0) -> None: + self.commands: list[tuple[str, ...]] = [] + self.result = result + + def __call__(self, command: Sequence[str]) -> int: + self.commands.append(tuple(command)) + return self.result + + +def observation(tmp_path: Path, values: dict[str, object]) -> str: + path = tmp_path / f"observation-{len(list(tmp_path.iterdir()))}.json" + path.write_text(json.dumps(values), encoding="utf-8") + return str(path) + + +def test_a_deploy_is_atomic_so_failed_readiness_undoes_it() -> None: + command = deploy_command(values=["serving.mode=ray"]) + + assert command[:5] == ("helm", "upgrade", "--install", "llm-routing", "deploy/helm/llm-routing") + assert "--atomic" in command + assert command[command.index("--timeout") + 1] == "10m" + assert command[-2:] == ("--set", "serving.mode=ray") + + +def test_a_failed_model_canary_rolls_the_release_back_to_its_recorded_target() -> None: + step = next_step(PLANS[MODEL], ROLLBACK, release="prod", namespace="inference") + + assert step.command == rollback_command(release="prod", namespace="inference") + assert step.command[:3] == ("helm", "rollback", "prod") and "--wait" in step.command + assert step.rollback_to == "deploy-0001" + assert "redeploy deploy-0001" in step.note + + +def test_a_failed_adapter_canary_needs_no_redeploy() -> None: + step = next_step(PLANS[ADAPTER], ROLLBACK) + + assert step.command == () + assert "gateway withdraws a failing adapter itself" in step.note + + +def test_a_rollback_with_no_recorded_target_runs_nothing() -> None: + step = next_step(PLANS[POLICY], ROLLBACK) + + assert step.rollback_to is None and step.command == () + assert "no command is run without a recorded target" in step.note + + +@pytest.mark.parametrize("action", ["promote", "hold"]) +def test_only_a_rollback_decision_produces_a_command(action: str) -> None: + step = next_step(PLANS[MODEL], CanaryDecision(action=action)) + + assert step.command == () and step.note == "" + + +@pytest.mark.parametrize( + ("values", "code", "action"), + [(HEALTHY, 0, "promote"), (EARLY, 2, "hold"), (FAILING, 3, "rollback")], +) +def test_the_exit_code_tells_a_pipeline_what_was_decided( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + values: dict[str, object], + code: int, + action: str, +) -> None: + runner = Recorder() + + result = main( + ["evaluate", "--plan", MODEL, "--observation", observation(tmp_path, values)], runner + ) + + printed = json.loads(capsys.readouterr().out) + assert result == code and printed["action"] == action + # Without --execute nothing is ever run. + assert runner.commands == [] and printed["executed"] is False + + +def test_execute_runs_the_rollback_and_reports_when_it_fails( + tmp_path: Path, capsys: pytest.CaptureFixture[str] +) -> None: + arguments = [ + "evaluate", + "--plan", + MODEL, + "--observation", + observation(tmp_path, FAILING), + "--execute", + ] + succeeded, failed = Recorder(), Recorder(result=1) + + assert main(arguments, succeeded) == 3 + assert succeeded.commands == [rollback_command()] + assert json.loads(capsys.readouterr().out)["succeeded"] is True + + assert main(arguments, failed) == 4 + assert json.loads(capsys.readouterr().out)["succeeded"] is False + + +def test_execute_runs_nothing_for_a_healthy_canary_or_an_adapter(tmp_path: Path) -> None: + runner = Recorder() + + healthy = ["evaluate", "--plan", MODEL, "--observation", observation(tmp_path, HEALTHY)] + adapter = ["evaluate", "--plan", ADAPTER, "--observation", observation(tmp_path, FAILING)] + assert main([*healthy, "--execute"], runner) == 0 + assert main([*adapter, "--execute"], runner) == 3 + assert runner.commands == [] + + +def test_a_baseline_keeps_a_shared_outage_from_being_blamed_on_the_canary(tmp_path: Path) -> None: + baseline = tmp_path / "baseline.json" + baseline.write_text(json.dumps(FAILING), encoding="utf-8") + runner = Recorder() + + result = main( + [ + "evaluate", + "--plan", + MODEL, + "--observation", + observation(tmp_path, FAILING), + "--baseline", + str(baseline), + "--execute", + ], + runner, + ) + + assert result != 3 and runner.commands == [] + + +def test_deploy_prints_its_command_and_runs_it_only_when_asked( + capsys: pytest.CaptureFixture[str], +) -> None: + runner = Recorder() + + assert main(["deploy", "--set", "serving.mode=ray"], runner) == 0 + assert runner.commands == [] + printed = json.loads(capsys.readouterr().out) + assert "--atomic" in printed["command"] and printed["executed"] is False + + assert main(["deploy", "--execute"], runner) == 0 + assert runner.commands == [deploy_command()] + + +@pytest.mark.parametrize("arguments", [["evaluate"], ["evaluate", "--plan", "model:unknown"]]) +def test_evaluate_says_what_it_is_missing(arguments: list[str]) -> None: + with pytest.raises(SystemExit) as raised: + main(arguments, Recorder()) + assert raised.value.code == 2 + + +@pytest.mark.skipif(shutil.which("helm") is None, reason="helm is not installed") +def test_helm_accepts_the_flags_the_commands_use() -> None: + for command in (deploy_command(), rollback_command()): + flags = [item for item in command if item.startswith("--")] + described = subprocess.run( + [*command[:2], "--help"], capture_output=True, text=True, check=True + ).stdout + for flag in flags: + assert flag in described, f"{flag} is not a flag of {' '.join(command[:2])}"