#!/usr/bin/env python3 """Pico reference CLI for the Open Engineering Academy. Implements the subcommands taught by the labs: pico parse --out pico compose --out pico runtime inspect [--lab ] [--namespace ] pico runtime emit --channel (--from-file | --value ) [--namespace ] [--from ] [--event-type ] pico runtime observe --pico [--namespace ] [--timeout ] [--expect ] This is a small, dependency-free reference implementation. The `runtime` subcommand group is a scriptable ControlSurface (Phase 7 vocabulary) that acts against the same declared runtime truth the Manifold labs deploy: the InteractionTopology ConfigMap, the Channel Service, and the hosted Pico Pods. It shells out to `kubectl`; it does not replace or bypass Kubernetes, Manifold, or Wrangler. """ from __future__ import annotations import argparse import json import os import stat import subprocess import sys from datetime import datetime, timezone from pathlib import Path PARSER_VERSION = "0.1.0" COMPOSER_VERSION = "0.1.0" RUNTIME_VERSION = "0.1.0" # Known labs the runtime ControlSurface can inspect against. Each entry # names the on-cluster carriers a lab applies: the topology ConfigMap, # the Channel Service(s), and the Pico Pods. The CLI reads these via # `kubectl` — it does not maintain its own state. _KNOWN_LABS: dict[str, dict] = { "hello-pico-on-manifold": { "namespace": "manifold", "topology_configmap": "hello-world-pico-on-manifold-topology", "channels": ["hello"], "picos": ["pico-engine"], }, "hello-two-picos": { "namespace": "manifold", "topology_configmap": "hello-two-picos-topology", "channels": ["greeting"], "picos": ["hello-consumer-pico", "hello-producer-pico"], }, } _SUPPORTED_KEYS = {"id", "kind", "value"} def _strip_quotes(text: str) -> str: if len(text) >= 2 and text[0] == text[-1] and text[0] in ('"', "'"): return text[1:-1] return text def parse_rule_yaml(path: Path) -> dict: """Parse the flat Rule YAML shape used by Hello Pico. Recognizes ``key: value`` lines for the supported keys and ignores blank lines and ``#`` comments. Rejects anything else so the learner gets a clear error instead of silent data loss. """ fields: dict[str, str] = {} for lineno, raw in enumerate(path.read_text().splitlines(), start=1): line = raw.split("#", 1)[0].rstrip() if not line.strip(): continue if ":" not in line: raise SystemExit(f"pico parse: {path}:{lineno}: expected 'key: value'") key, _, value = line.partition(":") key = key.strip() value = _strip_quotes(value.strip()) if key not in _SUPPORTED_KEYS: raise SystemExit( f"pico parse: {path}:{lineno}: unsupported key {key!r} " f"(expected one of {sorted(_SUPPORTED_KEYS)})" ) fields[key] = value missing = sorted(_SUPPORTED_KEYS - fields.keys()) if missing: raise SystemExit(f"pico parse: {path}: missing required keys: {missing}") return { "id": fields["id"], "kind": fields["kind"], "value": fields["value"], "_source": str(path), "_parsed_at": datetime.now(timezone.utc).isoformat(timespec="seconds"), "_parser_version": PARSER_VERSION, } def compose_artifact(parsed: dict, out_path: Path) -> None: """Write an executable shell artifact that prints the Rule's value.""" if parsed.get("kind") != "greeting": raise SystemExit( f"pico compose: unsupported kind {parsed.get('kind')!r} " "(this reference Composer only handles 'greeting')" ) value = parsed["value"].replace("'", "'\"'\"'") script = ( "#!/usr/bin/env bash\n" f"# Generated by pico compose (composer v{COMPOSER_VERSION})\n" f"# Source rule: {parsed.get('_source', 'unknown')}\n" f"printf '%s\\n' '{value}'\n" ) out_path.parent.mkdir(parents=True, exist_ok=True) out_path.write_text(script) mode = out_path.stat().st_mode out_path.chmod(mode | stat.S_IXUSR | stat.S_IXGRP | stat.S_IXOTH) def _cmd_parse(args: argparse.Namespace) -> None: parsed = parse_rule_yaml(Path(args.rule)) out = Path(args.out) out.parent.mkdir(parents=True, exist_ok=True) out.write_text(json.dumps(parsed, indent=2, sort_keys=True) + "\n") print(f"pico parse: wrote {out}") def _cmd_compose(args: argparse.Namespace) -> None: parsed = json.loads(Path(args.parsed).read_text()) out = Path(args.out) compose_artifact(parsed, out) print(f"pico compose: wrote {out}") # --- runtime ControlSurface helpers --------------------------------------- # # These helpers keep the shell-out layer thin and testable. The public # functions below (build_emit_command, build_observe_wait_command, # build_inspect_commands, load_event_payload, match_observation) are pure # and covered by pytest tests. The `_run_kubectl` wrapper is the only # place that actually invokes `kubectl`. def _run_kubectl(argv: list[str], *, stdin: bytes | None = None, check: bool = True) -> subprocess.CompletedProcess: """Run `kubectl` with the given argv. Thin wrapper for testability.""" return subprocess.run( ["kubectl", *argv], input=stdin, check=check, capture_output=True, ) def load_event_payload(*, from_file: str | None, value: str | None, channel: str, event_type: str, from_pico: str | None) -> bytes: """Return the raw JSON payload bytes to send onto a Channel. Either reads an existing event file (authored declarative truth) or synthesizes the minimal shape used by the labs. Synthesis carries the Channel name and EventType so the payload names its own contract. """ if from_file: return Path(from_file).read_bytes() if value is None: raise SystemExit( "pico runtime emit: need --from-file or --value") payload: dict[str, str] = { "eventType": event_type, "channel": channel, "value": value, } if from_pico: payload["from"] = from_pico return (json.dumps(payload) + "\n").encode("utf-8") def build_emit_command(*, namespace: str, channel: str, sender_name: str = "pico-cli-event-sender", image: str = "busybox:1.36", nc_timeout: int = 2) -> list[str]: """Return the kubectl argv that runs an in-cluster nc client. The channel host is the Kubernetes Service name for the Channel in the Manifold namespace, matching the labs' 03-*-service.yaml files. """ host = f"{channel}.{namespace}.svc.cluster.local" port = 8080 return [ "-n", namespace, "run", sender_name, f"--image={image}", "--restart=Never", "--rm", "-i", "--quiet", "--command", "--", "sh", "-c", f"nc -w {nc_timeout} {host} {port}", ] def build_observe_wait_command(*, namespace: str, pico: str, timeout_seconds: int) -> list[str]: """Return kubectl argv that waits for a Pico Pod to Succeed.""" return [ "-n", namespace, "wait", "--for=jsonpath={.status.phase}=Succeeded", f"--timeout={timeout_seconds}s", f"pod/{pico}", ] def build_inspect_commands(*, namespace: str, topology_configmap: str, picos: list[str]) -> list[list[str]]: """Return the kubectl argv sequence used by `runtime inspect`. Reads (in order) the topology ConfigMap YAML (the Wrangler-declared InteractionTopology carried on-cluster) and the Pico Pods' status. """ cmds: list[list[str]] = [ ["-n", namespace, "get", "configmap", topology_configmap, "-o", "yaml"], ] for pico in picos: cmds.append(["-n", namespace, "get", "pod", pico, "-o", "jsonpath={.metadata.name}\t{.status.phase}\n"]) return cmds def match_observation(logs: str, expected_line: str) -> bool: """Return True iff `logs` contains `expected_line` as a full line.""" return any(line == expected_line for line in logs.splitlines()) def _resolve_lab_defaults(lab: str | None) -> dict: if lab is None: return {} if lab not in _KNOWN_LABS: known = ", ".join(sorted(_KNOWN_LABS)) raise SystemExit(f"pico runtime: unknown --lab {lab!r} (known: {known})") return _KNOWN_LABS[lab] def _cmd_runtime_inspect(args: argparse.Namespace) -> None: defaults = _resolve_lab_defaults(args.lab) namespace = args.namespace or defaults.get("namespace") or "manifold" configmap = args.topology_configmap or defaults.get("topology_configmap") picos = args.pico or defaults.get("picos") or [] if not configmap: raise SystemExit( "pico runtime inspect: need --lab or --topology-configmap") cmds = build_inspect_commands( namespace=namespace, topology_configmap=configmap, picos=list(picos), ) print(f"pico runtime inspect: namespace={namespace} configmap={configmap}") for cmd in cmds: result = _run_kubectl(cmd, check=False) sys.stdout.write(result.stdout.decode("utf-8", errors="replace")) if result.returncode != 0: sys.stderr.write(result.stderr.decode("utf-8", errors="replace")) def _cmd_runtime_emit(args: argparse.Namespace) -> None: defaults = _resolve_lab_defaults(args.lab) namespace = args.namespace or defaults.get("namespace") or "manifold" payload = load_event_payload( from_file=args.from_file, value=args.value, channel=args.channel, event_type=args.event_type, from_pico=getattr(args, "from"), ) cmd = build_emit_command(namespace=namespace, channel=args.channel) print(f"pico runtime emit: namespace={namespace} channel={args.channel} " f"bytes={len(payload)}") result = _run_kubectl(cmd, stdin=payload, check=False) sys.stdout.write(result.stdout.decode("utf-8", errors="replace")) if result.returncode != 0: sys.stderr.write(result.stderr.decode("utf-8", errors="replace")) raise SystemExit(result.returncode) def _cmd_runtime_observe(args: argparse.Namespace) -> None: defaults = _resolve_lab_defaults(args.lab) namespace = args.namespace or defaults.get("namespace") or "manifold" wait_cmd = build_observe_wait_command( namespace=namespace, pico=args.pico_name, timeout_seconds=args.timeout, ) print(f"pico runtime observe: namespace={namespace} pico={args.pico_name} " f"timeout={args.timeout}s") wait = _run_kubectl(wait_cmd, check=False) sys.stderr.write(wait.stderr.decode("utf-8", errors="replace")) logs = _run_kubectl( ["-n", namespace, "logs", f"pod/{args.pico_name}"], check=False, ) text = logs.stdout.decode("utf-8", errors="replace") sys.stdout.write(text) if args.expect is not None: if match_observation(text, args.expect): print(f"pico runtime observe: OK — observed {args.expect!r}") return raise SystemExit( f"pico runtime observe: FAIL — expected line not found: {args.expect!r}") def _add_runtime_parser(subs: argparse._SubParsersAction) -> None: p_runtime = subs.add_parser( "runtime", help="ControlSurface subcommands over the Manifold runtime.", ) r_subs = p_runtime.add_subparsers(dest="runtime_command", required=True) p_inspect = r_subs.add_parser( "inspect", help="Show declared topology (ConfigMap) and Pico Pod status.", ) p_inspect.add_argument("--lab", choices=sorted(_KNOWN_LABS)) p_inspect.add_argument("--namespace") p_inspect.add_argument("--topology-configmap") p_inspect.add_argument("--pico", action="append", help="Repeatable. Pico Pod name to include.") p_inspect.set_defaults(func=_cmd_runtime_inspect) p_emit = r_subs.add_parser( "emit", help="Inject one event onto a Wrangler-declared Channel.", ) p_emit.add_argument("--lab", choices=sorted(_KNOWN_LABS)) p_emit.add_argument("--namespace") p_emit.add_argument("--channel", required=True) p_emit.add_argument("--event-type", default="hello.request") p_emit.add_argument("--from", dest="from", help="Optional 'from' pico name to include in payload.") src = p_emit.add_mutually_exclusive_group(required=True) src.add_argument("--from-file") src.add_argument("--value") p_emit.set_defaults(func=_cmd_runtime_emit) p_observe = r_subs.add_parser( "observe", help="Wait for a Pico Pod to Succeed and print its logs.", ) p_observe.add_argument("--lab", choices=sorted(_KNOWN_LABS)) p_observe.add_argument("--namespace") p_observe.add_argument("--pico", dest="pico_name", required=True, help="Pico Pod name whose Observation to capture.") p_observe.add_argument("--timeout", type=int, default=60) p_observe.add_argument("--expect", help="Optional exact log line to assert.") p_observe.set_defaults(func=_cmd_runtime_observe) def main(argv: list[str] | None = None) -> int: root = argparse.ArgumentParser(prog="pico", description="Pico reference CLI.") subs = root.add_subparsers(dest="command", required=True) p_parse = subs.add_parser("parse", help="Parse a Rule YAML file to JSON.") p_parse.add_argument("rule", help="Path to the Rule YAML file.") p_parse.add_argument("--out", required=True, help="Path to write parsed JSON.") p_parse.set_defaults(func=_cmd_parse) p_compose = subs.add_parser("compose", help="Compose a Pico artifact from parsed JSON.") p_compose.add_argument("parsed", help="Path to the parsed JSON file.") p_compose.add_argument("--out", required=True, help="Path to write the artifact.") p_compose.set_defaults(func=_cmd_compose) _add_runtime_parser(subs) args = root.parse_args(argv) args.func(args) return 0 if __name__ == "__main__": sys.exit(main())