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
11 changes: 9 additions & 2 deletions deploy/otel-demo/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,11 +59,18 @@ Then generate the infra-fault cases (Chaos-Mesh manifests are in
# dry-run first (no cluster): logs the apply/delete + emits the 4 case dirs
python scripts/eval/otel_corpus.py --chaos --dry-run --out-dir /tmp/otel-chaos-dry

# for real (needs kubectl context on the cluster + the Collector export at --otlp-dir)
python scripts/eval/otel_corpus.py --chaos \
# for real: name the (non-prod) target context explicitly + the Collector export at --otlp-dir
python scripts/eval/otel_corpus.py --chaos --kube-context kind-kind \
--otlp-dir /path/to/collector/export --out-dir data/eval-cases/otel
```

**Target safety.** A real `--chaos` run never uses the kubeconfig's current context.
`--kube-context` is required, and every `kubectl` call passes `--context <it>`. The run
refuses a context that isn't in the kubeconfig. It also refuses one whose name, cluster
or API server contains `prod`, `prd` or `live`. If a non-prod context trips that check
(e.g. `delivery-dev`), repeat its exact name in `--allow-context delivery-dev` to let
it through. Use a dedicated, disposable cluster.

The driver `kubectl apply`s each manifest, waits the incident window, captures, and
`kubectl delete`s it. **Before a real run, edit the manifests** in
`deploy/otel-demo/chaos/` so the namespace + label selectors (and `IOChaos`
Expand Down
88 changes: 80 additions & 8 deletions scripts/eval/otel_corpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,20 +13,26 @@
with flagd reachable at --flagd-url and its Collector writing OTLP-JSON to
--otlp-dir (logs.json / traces.json / metrics.json).

`--chaos` mutates a cluster (`kubectl apply/delete`), so it never uses the kubeconfig's
current context: a real chaos run requires an explicit `--kube-context`, refuses any
context whose name, cluster or API server looks like production (prod / prd / live),
and passes `--context` to every kubectl call. A non-prod context that merely contains
one of those markers can be let through with `--allow-context <exact same name>`.

Real runs are serial and slow — each case waits `baseline + post` seconds for its
windows to accrue (≈15 min/case at the 300/600 defaults). Use --dry-run first to
validate wiring and case emission without a cluster (no HTTP, no waiting).
"""
from __future__ import annotations

import argparse
import json
import subprocess
import sys
import time
from datetime import datetime, timezone
from pathlib import Path

import subprocess

from src.eval.otel_demo import (
CHAOS_SCENARIOS,
FLAG_SCENARIOS,
Expand Down Expand Up @@ -83,13 +89,67 @@ def reset() -> None:
return reset


def _kubectl_chaos_hooks(chaos_dir: Path):
"""apply/delete a Chaos Mesh experiment via `kubectl` from a per-scenario
manifest (deploy/otel-demo/chaos/<scenario.name>.yaml)."""
# Substrings that mark a kube context, its cluster or its API server as production. Fail closed:
# a non-prod context that merely contains one (e.g. "delivery-dev") needs --allow-context.
_PROD_MARKERS = ("prod", "prd", "live")


class KubeContextRefused(Exception):
"""The capture will not run kubectl against this target."""


def _kube_context_facts(context: str) -> tuple[str, str] | None:
"""(cluster name, API server) that ``context`` points at in the local kubeconfig, or None
if no such context exists. Read-only: ``kubectl config view`` never contacts a cluster."""
try:
out = subprocess.run(["kubectl", "config", "view", "-o", "json"],
capture_output=True, text=True, check=True).stdout
except (OSError, subprocess.CalledProcessError) as exc:
raise KubeContextRefused(f"cannot read kubeconfig: {exc}") from exc
cfg = json.loads(out or "{}")
ctx = next((c for c in cfg.get("contexts") or [] if c.get("name") == context), None)
if ctx is None:
return None
cluster = (ctx.get("context") or {}).get("cluster") or ""
server = next(((c.get("cluster") or {}).get("server") or ""
for c in cfg.get("clusters") or [] if c.get("name") == cluster), "")
return cluster, server


def resolve_kube_context(context: str | None, allow: str | None, facts=_kube_context_facts) -> str:
"""The kube context a mutating capture may target, or raise KubeContextRefused.

There is no implicit fallback to the kubeconfig's current context: the target must be named.
It must exist, and neither its name, its cluster nor its API server may carry a production
marker — unless ``allow`` repeats the exact same context name (an affirmative override for a
non-prod context whose name happens to match)."""
if not context:
raise KubeContextRefused("--kube-context is required: a chaos capture never uses the current context")
if allow is not None and allow != context:
raise KubeContextRefused(f"--allow-context {allow!r} does not match --kube-context {context!r}")
found = facts(context)
if found is None:
raise KubeContextRefused(f"kube context {context!r} is not in the kubeconfig")
cluster, server = found
hits = sorted({m for name in (context, cluster, server) for m in _PROD_MARKERS if m in name.lower()})
if hits and allow != context:
raise KubeContextRefused(
f"kube context {context!r} (cluster {cluster!r}, server {server!r}) looks like production "
f"({', '.join(hits)}); if it is not, pass --allow-context {context}")
print(f"kube target: context={context} cluster={cluster} server={server}"
+ (" (allowed by --allow-context)" if hits else ""), flush=True)
return context


def _kubectl_chaos_hooks(chaos_dir: Path, context: str):
"""apply/delete a Chaos Mesh experiment via `kubectl --context <context>` from a per-scenario
manifest (deploy/otel-demo/chaos/<scenario.name>.yaml). ``context`` comes from
``resolve_kube_context``; every call names it, so the current context is never used."""
def _run(verb: str, sc, extra=()):
manifest = chaos_dir / f"{sc.name}.yaml"
subprocess.run(["kubectl", verb, "-f", str(manifest), *extra], check=(verb == "apply"))
print(f" kubectl {verb}: {sc.kind}={sc.name}", flush=True)
subprocess.run(["kubectl", "--context", context, verb, "-f", str(manifest), *extra],
check=(verb == "apply"))
print(f" kubectl --context {context} {verb}: {sc.kind}={sc.name}", flush=True)

return (lambda sc: _run("apply", sc)), (lambda sc: _run("delete", sc, ("--ignore-not-found",)))

Expand Down Expand Up @@ -118,12 +178,23 @@ def main() -> int:
help="generate Chaos-Mesh infra-fault cases (kubectl apply/delete); needs k8s")
ap.add_argument("--chaos-dir", type=Path, default=Path("deploy/otel-demo/chaos"),
help="dir of Chaos-Mesh manifests, one per scenario name")
ap.add_argument("--kube-context", default=None,
help="kube context a real --chaos run targets (required; the current context is never used)")
ap.add_argument("--allow-context", default=None,
help="exact --kube-context name to allow even though it matches a prod/live marker")
ap.add_argument("--dry-run", action="store_true", help="no cluster: validate the loop + case emission")
args = ap.parse_args()

if not args.dry_run and args.otlp_dir is None:
print("--otlp-dir is required unless --dry-run", file=sys.stderr)
return 2
context = None
if args.chaos and not args.dry_run:
try:
context = resolve_kube_context(args.kube_context, args.allow_context)
except KubeContextRefused as exc:
print(f"refusing to run chaos: {exc}", file=sys.stderr)
return 2

capture = _dry_capture if args.dry_run else capture_from_otlp_dir(args.otlp_dir)
sleep = (lambda _s: None) if args.dry_run else time.sleep
Expand All @@ -136,7 +207,8 @@ def main() -> int:
# Chaos-Mesh infra faults (k8s + Chaos Mesh) — a separate mode from the
# flag/negative corpus. Reads deploy/otel-demo/chaos/<scenario>.yaml.
if args.chaos:
apply_chaos, delete_chaos = _dry_chaos_hooks() if args.dry_run else _kubectl_chaos_hooks(args.chaos_dir)
apply_chaos, delete_chaos = (_dry_chaos_hooks() if args.dry_run
else _kubectl_chaos_hooks(args.chaos_dir, context))
for sc in CHAOS_SCENARIOS:
case_id = f"otel_chaos_{sc.name}"
print(f"[{len(written)+1}] chaos {sc.kind}={sc.name} ({sc.service})", flush=True)
Expand Down
92 changes: 92 additions & 0 deletions tests/unit/test_otel_corpus.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,98 @@ def test_chaos_dry_run_emits_loadable_corpus(tmp_path, monkeypatch):
assert {c.trigger.type for c in cases} == {s.trigger_type for s in CHAOS_SCENARIOS}


class _FakeKubectl:
"""Stands in for ``subprocess.run``: serves a synthetic kubeconfig to ``kubectl config view``
and records every other call. No test here can reach a real kubectl or cluster."""

KUBECONFIG = {
"current-context": "cfbc-live-prod",
"contexts": [
{"name": "cfbc-live-prod", "context": {"cluster": "cfbc-live-prod"}},
{"name": "kind-otel", "context": {"cluster": "kind-otel"}},
# innocuous alias for a production cluster
{"name": "sandbox", "context": {"cluster": "arn:aws:eks:us-east-1:1:cluster/cfbc-live-prod"}},
{"name": "delivery-dev", "context": {"cluster": "delivery-dev"}},
],
"clusters": [
{"name": "cfbc-live-prod", "cluster": {"server": "https://ABC.gr7.us-east-1.eks.amazonaws.com"}},
{"name": "kind-otel", "cluster": {"server": "https://127.0.0.1:6443"}},
{"name": "arn:aws:eks:us-east-1:1:cluster/cfbc-live-prod",
"cluster": {"server": "https://DEF.gr7.us-east-1.eks.amazonaws.com"}},
{"name": "delivery-dev", "cluster": {"server": "https://10.0.0.5:6443"}},
],
}

def __init__(self):
self.calls: list[list[str]] = []

def __call__(self, argv, **_kw):
import json
import subprocess

if argv[:3] == ["kubectl", "config", "view"]:
return subprocess.CompletedProcess(argv, 0, stdout=json.dumps(self.KUBECONFIG), stderr="")
self.calls.append(list(argv))
return subprocess.CompletedProcess(argv, 0, stdout="", stderr="")


def _run_chaos(tmp_path, monkeypatch, *extra):
fake = _FakeKubectl()
monkeypatch.setattr("subprocess.run", fake)
code = _run(["--chaos", "--otlp-dir", str(tmp_path / "otlp"), "--out-dir", str(tmp_path / "out"),
"--baseline", "0", "--post", "0", *extra], monkeypatch)
return code, fake


class TestKubeContextGuard:
def test_real_chaos_without_kube_context_refuses_and_runs_no_kubectl(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch)
assert code == 2 and fake.calls == [] # never falls back to the current (prod) context

def test_prod_named_context_is_refused(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "cfbc-live-prod")
assert code == 2 and fake.calls == []

def test_alias_pointing_at_a_prod_cluster_is_refused(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "sandbox")
assert code == 2 and fake.calls == []

def test_unknown_context_is_refused(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "kind-typo")
assert code == 2 and fake.calls == []

def test_allow_context_must_repeat_the_exact_name(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "cfbc-live-prod",
"--allow-context", "delivery-dev")
assert code == 2 and fake.calls == []

def test_allow_context_does_not_admit_a_missing_context(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "kind-typo",
"--allow-context", "kind-typo")
assert code == 2 and fake.calls == []

def test_marker_false_positive_needs_allow_context(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "delivery-dev") # "live" in "delivery"
assert code == 2 and fake.calls == []
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "delivery-dev",
"--allow-context", "delivery-dev")
assert code == 0 and fake.calls

def test_every_kubectl_call_names_the_context(self, tmp_path, monkeypatch):
code, fake = _run_chaos(tmp_path, monkeypatch, "--kube-context", "kind-otel")
assert code == 0
assert len(fake.calls) == 2 * len(CHAOS_SCENARIOS) # one apply + one delete per scenario
for argv in fake.calls:
assert argv[:3] == ["kubectl", "--context", "kind-otel"], argv
assert [a[3] for a in fake.calls] == ["apply", "delete"] * len(CHAOS_SCENARIOS)

def test_dry_run_never_touches_kubectl(self, tmp_path, monkeypatch):
fake = _FakeKubectl()
monkeypatch.setattr("subprocess.run", fake)
assert _run(["--chaos", "--dry-run", "--out-dir", str(tmp_path)], monkeypatch) == 0
assert fake.calls == []


def test_every_chaos_scenario_has_a_committed_manifest():
chaos_dir = _REPO / "deploy" / "otel-demo" / "chaos"
for s in CHAOS_SCENARIOS:
Expand Down
Loading