"""Registry ops — k3d mode (local dev cluster). Uses k3d's built-in registry or falls back to a plain Docker registry container. """ from __future__ import annotations from pathlib import Path from ._services_common import ( _LogFn, _docker, _k3d, _log, _to_bool, ) def _patch_k3d_registry_mirrors( registry_name: str, port: str, cluster_name: str = "knoe-dev-cluster", *, env: dict | None = None, log: _LogFn | None = None, ) -> None: """Write containerd mirror config for the k3d registry onto every cluster node. k3d nodes resolve `k3d-NAME.localhost:PORT` as 127.0.0.1 (loopback) which has nothing listening. The registries.yaml mirror tells containerd to use the Docker-network DNS name of the registry container (`http://k3d-NAME:PORT`) instead, which works because all k3d containers share the same Docker network. After patching, each affected node is restarted so containerd picks up the new config. This is idempotent — if the mirror entry is already present the node is left untouched. """ import subprocess as _sp docker_name = f"k3d-{registry_name}" host_alias = f"{docker_name}.localhost:{port}" endpoint = f"http://{docker_name}:{port}" registries_yaml = ( "mirrors:\n" f' "{host_alias}":\n' " endpoint:\n" f' - "{endpoint}"\n' ) node_result = _sp.run( ["k3d", "node", "list", "-o", "json"], capture_output=True, text=True, timeout=15, ) import json try: nodes = json.loads(node_result.stdout or "[]") except Exception: nodes = [] node_containers = [ n.get("name") or "" for n in nodes if isinstance(n, dict) and n.get("role", "") not in ("loadbalancer", "registry") and (n.get("runtimeLabels") or {}).get("k3d.cluster") == cluster_name ] if not node_containers: # Fallback: get containers by Docker label (two separate queries for server+agent) for role in ("server", "agent"): r = _sp.run( ["docker", "ps", "--filter", f"label=k3d.cluster={cluster_name}", "--filter", f"label=k3d.role={role}", "--format", "{{.Names}}"], capture_output=True, text=True, timeout=15, ) node_containers += [n.strip() for n in r.stdout.splitlines() if n.strip()] cfg_path = "/etc/rancher/k3s/registries.yaml" restarted = [] for container in node_containers: if not container: continue # Check if mirror is already configured check = _sp.run( ["docker", "exec", container, "cat", cfg_path], capture_output=True, text=True, timeout=10, ) if check.returncode == 0 and host_alias in (check.stdout or ""): _log(log, f"[REGISTRY] Mirror already configured on {container}") continue _log(log, f"[REGISTRY] Patching containerd mirror config on {container}") _sp.run( ["docker", "exec", container, "mkdir", "-p", "/etc/rancher/k3s"], capture_output=True, timeout=10, ) _sp.run( ["docker", "exec", "-i", container, "sh", "-c", f"cat > {cfg_path}"], input=registries_yaml, capture_output=True, text=True, timeout=15, ) restarted.append(container) for container in restarted: _log(log, f"[REGISTRY] Restarting {container} to apply mirror config") _sp.run(["docker", "restart", container], capture_output=True, text=True, timeout=60) if restarted: import time _log(log, "[REGISTRY] Waiting for cluster nodes to recover after restart...") time.sleep(10) def _k3d_registry_running(raw_list: str, name: str) -> bool: """Return True only if a line whose NAME ends with `name` has STATUS=running.""" for line in raw_list.splitlines(): parts = line.split() if not parts: continue # k3d may prefix the name with "k3d-"; match the suffix. if parts[0].endswith(name) and parts[-1].lower() == "running": return True return False def _ensure_k3d_registry(*, env: dict | None = None, log: _LogFn | None = None) -> None: name = str((env or {}).get("K3D_REGISTRY_NAME") or "knoe-registry") port = str((env or {}).get("REGISTRY_PORT") or "5000") legacy_name = "prole-registry" listed = _k3d(["registry", "list"], env=env, timeout=60) raw = listed.stdout if listed.returncode == 0 else "" # Remove legacy prole-registry that may be squatting on the port. if legacy_name in raw: _log(log, f"[REGISTRY] Removing legacy k3d registry {legacy_name}") _k3d(["registry", "delete", legacy_name], env=env, timeout=60) listed = _k3d(["registry", "list"], env=env, timeout=60) raw = listed.stdout if listed.returncode == 0 else "" cluster_name = str((env or {}).get("K3D_CLUSTER_NAME") or "knoe-dev-cluster") if name in raw: if _k3d_registry_running(raw, name): _log(log, f"[REGISTRY] k3d registry {name} already running") _patch_k3d_registry_mirrors(name, port, cluster_name, env=env, log=log) return # Exists but not running (e.g. "created" — port was unavailable). _log(log, f"[REGISTRY] k3d registry {name} exists but not running; recreating") _k3d(["registry", "delete", name], env=env, timeout=60) _log(log, f"[REGISTRY] Creating k3d registry {name} on port {port}") created = _k3d(["registry", "create", name, "--port", f"{port}:{port}"], env=env, timeout=180) if created.returncode == 0: _patch_k3d_registry_mirrors(name, port, cluster_name, env=env, log=log) return # fallback docker registry _log(log, "[REGISTRY] k3d registry creation failed; trying docker registry fallback") running = _docker(["ps", "--format", "{{.Names}}"], env=env, timeout=30) if running.returncode == 0 and name in running.stdout.splitlines(): return _docker(["rm", "-f", name], env=env, timeout=30) _docker( [ "run", "-d", "--restart=always", "-p", f"{port}:5000", "--name", name, "registry:2", ], env=env, timeout=120, check=True, ) def initialize( *, namespace: str | None = None, env: dict | None = None, project_root: str | Path = ".", log: _LogFn | None = None, mode: str | None = None, ) -> None: update(namespace=namespace, env=env, project_root=project_root, log=log, mode=mode) def start( *, namespace: str | None = None, env: dict | None = None, project_root: str | Path = ".", log: _LogFn | None = None, mode: str | None = None, ) -> None: update(namespace=namespace, env=env, project_root=project_root, log=log, mode=mode) def update( *, namespace: str | None = None, env: dict | None = None, project_root: str | Path = ".", log: _LogFn | None = None, mode: str | None = None, ) -> None: _ensure_k3d_registry(env=env, log=log) def stop( *, namespace: str | None = None, env: dict | None = None, mode: str | None = None, log: _LogFn | None = None, ) -> None: name = str((env or {}).get("K3D_REGISTRY_NAME") or "knoe-registry") _log(log, f"[REGISTRY] Removing k3d registry {name}") _k3d(["registry", "delete", name], env=env, timeout=120) _docker(["rm", "-f", name], env=env, timeout=30) def restart( *, namespace: str | None = None, env: dict | None = None, mode: str | None = None, log: _LogFn | None = None, ) -> None: update(namespace=namespace, env=env, mode=mode, log=log) def status( *, namespace: str | None = None, env: dict | None = None, mode: str | None = None, ) -> bool: name = str((env or {}).get("K3D_REGISTRY_NAME") or "knoe-registry") reg = _k3d(["registry", "list"], env=env, timeout=30) return _to_bool(reg.returncode == 0 and name in reg.stdout)