Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 4 additions & 12 deletions .github/workflows/cd.yml
Original file line number Diff line number Diff line change
Expand Up @@ -60,18 +60,10 @@ jobs:
- run: python -m pip install -e ".[dev]"
- name: Verify the serving configuration matches the catalog
run: python -m llm_router.serving > /tmp/ray-serve.yaml && diff -u config/ray-serve.yaml /tmp/ray-serve.yaml
- name: Render the canary and rollback plan
run: |
python - <<'PY' > canary-plan.json
import json
from llm_router.registry import load_registry
from llm_router.serving import canary_config
registry = load_registry("config/registry.yaml")
production = [item for item in registry.deployments if item.stage.value == "production"]
plan = canary_config(registry, production[0].id)
plan["release_ref"] = "${{ env.RELEASE_REF }}"
print(json.dumps(plan, indent=2))
PY
- name: Render the canary and rollback plans
# 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: Validate the Kubernetes manifests
run: |
curl -sSLo kubeconform.tar.gz https://github.com/yannh/kubeconform/releases/download/v0.7.0/kubeconform-linux-amd64.tar.gz
Expand Down
50 changes: 49 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,7 @@ in-process and correct for a single replica only. Install the client with the ex
python -m pip install -e ".[redis]"
```

CD renders the canary plan (with its rollback target and triggers), verifies
CD renders the canary plans (one per track, each with its rollback target), verifies
`config/ray-serve.yaml` against the catalog, and validates the manifests with kubeconform.
Applying to a cluster stays disabled until a deployment destination is configured.

Expand All @@ -211,6 +211,7 @@ be served. A request can never introduce a model path, revision, or adapter.
| `GET /v1/registry/adapters` | Promoted LoRA and QLoRA adapters. |
| `GET /v1/registry/deployments` | Deployment revisions and rollback targets. |
| `GET /v1/registry/variants` | Optimization variants with their measured deltas. |
| `GET /v1/registry/canaries` | Canary plans by track, with live adapter state. |

Send `routing.domain` to request a domain adapter; the router applies the promoted adapter
with the largest measured quality gain for that base revision and task, or none at all.
Expand Down Expand Up @@ -238,6 +239,53 @@ warm replica never pay it on the request path. The high-capability tier scales t
first request after an idle period waits for a full model load: budget for it, or raise its
`min_replicas`.

## Canaries and rollback

Models, adapters, and router policies are canaried on separate tracks, each with its own plan
naming exactly what it rolls back to:

```bash
python -m llm_router.canary
```

| Track | Canary | Rolls back to |
|---|---|---|
| `model` | A deployment revision | The revision in `previous_revision_id`. |
| `adapter` | A `staging` adapter | The `production` adapter for the same base revision and domain, or the base model alone. |
| `policy` | The policy `version` | `previous_version`. |

One rule decides every track. Failed readiness rolls back at once. Once 50 requests have been
observed, the canary rolls back if its error rate exceeds 1%, its p95 latency exceeds the
strictest tier objective among the models it serves, or its observed quality falls more than
0.05 below their benchmarked quality. Error rate and latency only count against the canary
when the stable baseline does not share the problem, so an engine outage that degrades both is
not blamed on the change. A canary that stays clean for 500 requests is reported ready.

**The adapter track runs inside the gateway.** A staged adapter is offered
`canary_traffic_percent` of eligible requests, bucketed by tenant and prompt so a retry cannot
flip between adapters. It is suspended automatically the moment it fails, and the production
adapter serves everything again. A canary's responses are never cached, so nothing it produced
outlives a rollback. Before this, a staged adapter with the larger measured gain took all of the
traffic.

`GET /v1/registry/canaries` lists every plan with live state for adapters (`in-progress`,
`ready-to-promote`, or `rolled-back`, with the reasons). `router_canary_requests_total` and
`router_canary_rollbacks_total` report it to Prometheus. State is per gateway replica: each
reaches the same verdict from its own share of traffic.

**Model and policy canaries are rollouts** of the serving pool or the gateway, so whatever
controls that rollout evaluates the plan through the same rule:

```bash
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.

Promotion is never automatic on any track. It is a catalog change and goes through review.

## Tenants

What a caller may use is governance and lives in the catalog; the credential that proves which
Expand Down
63 changes: 60 additions & 3 deletions src/llm_router/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
exact_cache_eligible,
semantic_cache_eligible,
)
from llm_router.canary import CanaryMonitor, canary_plans
from llm_router.classifier import TaskClassifier, load_classifier
from llm_router.config import Settings, get_settings
from llm_router.engine_stats import ColdStartTracker, EngineStatsCollector
Expand Down Expand Up @@ -125,12 +126,15 @@ def create_app(
catalog.policy.version if catalog is not None else runtime_settings.routing_policy_version
)
load = LoadTracker()
plans = canary_plans(catalog) if catalog is not None else ()
canary = CanaryMonitor(plans)
router = Router(
profiles=profiles,
external_fallback_enabled=runtime_settings.external_fallback_enabled,
registry=catalog,
classifier=_load_task_classifier(runtime_settings.task_classifier_path),
load=load,
canary=canary,
)
admission = AdmissionController(
runtime_settings.max_concurrency,
Expand Down Expand Up @@ -400,6 +404,23 @@ async def variants() -> dict[str, object]:
)
return {"object": "list", "data": data}

@app.get("/v1/registry/canaries", dependencies=[Depends(authenticate)])
async def canaries() -> dict[str, object]:
"""Every canary plan by track, with live state for the adapter track."""

return {
"object": "list",
"data": [
{
**plan.model_dump(mode="json"),
"live": (
canary.status(plan.subject) if plan.subject in canary.subjects else None
),
}
for plan in plans
],
}

@app.get("/v1/registry/deployments", dependencies=[Depends(authenticate)])
async def deployments() -> dict[str, object]:
if catalog is None:
Expand Down Expand Up @@ -496,7 +517,9 @@ async def _store_cache(
decision: RouteDecision,
result: BackendResult,
) -> None:
if not runtime_settings.cache_enabled:
# A canary's responses are never cached: if it is rolled back, nothing
# it produced may keep being served from the cache afterwards.
if not runtime_settings.cache_enabled or decision.canary_arm == "canary":
return
entry = CachedCompletion(
text=result.text,
Expand All @@ -513,6 +536,25 @@ async def _store_cache(
scope = semantic_cache.scope(payload, tenant=tenant, model_revision=catalog_version)
semantic_cache.store(scope, prompt, entry)

def _record_canary(
decision: RouteDecision, *, ok: bool, started: float, quality: float | None = None
) -> None:
"""Feed a canaried route's outcome to the monitor, suspending on failure."""

if decision.canary_subject is None or decision.canary_arm is None:
return
on_canary = decision.canary_arm == "canary"
telemetry.record_canary(decision.canary_subject, arm=decision.canary_arm, ok=ok)
verdict = canary.record(
decision.canary_subject,
canary=on_canary,
ok=ok,
latency_ms=(time.perf_counter() - started) * 1000,
quality=quality,
)
if verdict is not None and verdict.action == "rollback":
telemetry.record_canary_rollback(decision.canary_subject)

async def _consume_quota(subject: str, tenant: TenantRecord | None) -> None:
# A tenant may carry its own limit; absent one the platform default
# applies, which is why None means "defer" rather than "unlimited".
Expand Down Expand Up @@ -608,8 +650,12 @@ async def _stream_completion(
prompt_tokens = max(1, len(prompt) // 4)
completion_tokens = max(1, len(text) // 4)
span.set_usage(prompt_tokens=prompt_tokens, completion_tokens=completion_tokens)
validity: float | None = None
if payload.routing.structured:
telemetry.record_structured_output(decision, valid=structured_output_valid(text))
valid = structured_output_valid(text)
validity = 1.0 if valid else 0.0
telemetry.record_structured_output(decision, valid=valid)
_record_canary(decision, ok=True, started=started, quality=validity)
telemetry.record_completion(
decision,
latency_seconds=time.perf_counter() - started,
Expand All @@ -629,6 +675,8 @@ async def _stream_completion(
)
except Exception as error:
span.fail(error)
if isinstance(error, BackendUnavailableError):
_record_canary(decision, ok=False, started=started)
raise
finally:
lease.close()
Expand Down Expand Up @@ -731,6 +779,7 @@ async def _complete(
tenant.allow_external_fallback is not False if tenant is not None else True
),
privacy_raised_from=privacy_raised_from,
canary_key=f"{subject}|{prompt}",
)
span.set_route(decision)
if runtime_settings.cache_enabled:
Expand Down Expand Up @@ -779,15 +828,22 @@ async def _complete(

try:
result = await inference_backend.generate(payload, decision)
except BackendUnavailableError:
_record_canary(decision, ok=False, started=started)
raise
finally:
telemetry.inflight_requests.dec()
admission.release()

span.set_usage(
prompt_tokens=result.prompt_tokens, completion_tokens=result.completion_tokens
)
validity: float | None = None
if payload.routing.structured:
telemetry.record_structured_output(decision, valid=structured_output_valid(result.text))
valid = structured_output_valid(result.text)
validity = 1.0 if valid else 0.0
telemetry.record_structured_output(decision, valid=valid)
_record_canary(decision, ok=True, started=started, quality=validity)
telemetry.record_completion(
decision,
latency_seconds=time.perf_counter() - started,
Expand Down Expand Up @@ -827,6 +883,7 @@ async def _complete(
"task_source": decision.task_source,
"task_confidence": decision.task_confidence,
"complexity": decision.complexity,
"canary_arm": decision.canary_arm,
"reason": decision.reason,
"score": decision.score,
"candidate_count": decision.candidate_count,
Expand Down
Loading
Loading