"""Image update checker. Compares the locally-pulled manifest digest (from RepoDigests) against the current manifest digest in the registry. Supports Docker Hub, ghcr.io, lscr.io and other token-auth v2 registries. """ from __future__ import annotations import asyncio import logging import time from dataclasses import asdict, dataclass from typing import Optional import httpx from config import settings from docker_client import DockerError, get_client, safe_call from models.setting import EVENT_UPDATE_AVAILABLE from services import notify_service, settings_service logger = logging.getLogger("stackpilot.update") _MANIFEST_ACCEPT = ", ".join( [ "application/vnd.docker.distribution.manifest.v2+json", "application/vnd.docker.distribution.manifest.list.v2+json", "application/vnd.oci.image.manifest.v1+json", "application/vnd.oci.image.index.v1+json", ] ) @dataclass class UpdateStatus: image: str update_available: bool current_digest: Optional[str] remote_digest: Optional[str] checked_at: float error: Optional[str] = None def to_dict(self) -> dict: return asdict(self) # image ref -> UpdateStatus _CACHE: dict[str, UpdateStatus] = {} # images we've already sent an "update available" notification for, so the # background loop doesn't re-notify on every cycle. _NOTIFIED: set[str] = set() # --------------------------------------------------------------------------- # # Image reference parsing # --------------------------------------------------------------------------- # def parse_ref(image: str) -> tuple[str, str, str]: """Return (registry_host, repository, tag/digest).""" ref = image tag = "latest" # split tag (but not the registry port colon) if "@" in ref: ref, tag = ref.split("@", 1) else: # find last colon after last slash slash = ref.rfind("/") colon = ref.rfind(":") if colon > slash: tag = ref[colon + 1 :] ref = ref[:colon] parts = ref.split("/", 1) if len(parts) == 2 and ("." in parts[0] or ":" in parts[0] or parts[0] == "localhost"): registry = parts[0] repo = parts[1] else: registry = "registry-1.docker.io" repo = ref if "/" not in repo: repo = "library/" + repo return registry, repo, tag def _local_digest(image: str) -> Optional[str]: try: client = get_client() img = safe_call(client.images.get, image) except DockerError: return None repo_digests = img.attrs.get("RepoDigests") or [] for rd in repo_digests: if "@" in rd: return rd.split("@", 1)[1] return None # --------------------------------------------------------------------------- # # Registry manifest digest # --------------------------------------------------------------------------- # async def _get_token(client: httpx.AsyncClient, www_auth: str) -> Optional[str]: # Parse: Bearer realm="...",service="...",scope="..." params = {} if not www_auth.lower().startswith("bearer"): return None for part in www_auth[len("Bearer ") :].split(","): if "=" in part: k, v = part.split("=", 1) params[k.strip()] = v.strip().strip('"') realm = params.pop("realm", None) if not realm: return None try: resp = await client.get(realm, params=params, timeout=10) resp.raise_for_status() data = resp.json() return data.get("token") or data.get("access_token") except (httpx.HTTPError, ValueError): return None async def remote_digest(image: str) -> Optional[str]: registry, repo, tag = parse_ref(image) if tag.startswith("sha256:"): return tag scheme = "https" url = f"{scheme}://{registry}/v2/{repo}/manifests/{tag}" headers = {"Accept": _MANIFEST_ACCEPT} async with httpx.AsyncClient(follow_redirects=True) as client: try: resp = await client.head(url, headers=headers, timeout=10) if resp.status_code == 401: token = await _get_token(client, resp.headers.get("WWW-Authenticate", "")) if not token: return None headers["Authorization"] = f"Bearer {token}" resp = await client.head(url, headers=headers, timeout=10) if resp.status_code == 405 or "Docker-Content-Digest" not in resp.headers: # Some registries don't support HEAD; fall back to GET. resp = await client.get(url, headers=headers, timeout=10) digest = resp.headers.get("Docker-Content-Digest") return digest except httpx.HTTPError as exc: logger.debug("remote_digest failed for %s: %s", image, exc) return None # --------------------------------------------------------------------------- # # Public API # --------------------------------------------------------------------------- # async def check_image(image: str) -> UpdateStatus: local = _local_digest(image) remote = await remote_digest(image) error = None if remote is None: error = "could not reach registry" update_available = bool(local and remote and local != remote) status = UpdateStatus( image=image, update_available=update_available, current_digest=local, remote_digest=remote, checked_at=time.time(), error=error, ) _CACHE[image] = status if update_available and image not in _NOTIFIED: _NOTIFIED.add(image) try: await notify_service.notify( EVENT_UPDATE_AVAILABLE, "Image update available", f"A newer image is available for {image}.", ) except Exception as exc: # noqa: BLE001 - notifications are best-effort logger.debug("update notify failed for %s: %s", image, exc) elif not update_available: _NOTIFIED.discard(image) return status def _all_running_images() -> set[str]: images: set[str] = set() try: client = get_client() for c in safe_call(client.containers.list, all=True): cfg_image = c.attrs.get("Config", {}).get("Image") if cfg_image: images.add(cfg_image) except DockerError: pass return images def stack_images(stack_id: str) -> set[str]: """Images used by the containers of one compose project (== stack id).""" images: set[str] = set() try: client = get_client() for c in safe_call(client.containers.list, all=True): labels = c.labels or {} if labels.get("com.docker.compose.project") != stack_id: continue cfg_image = c.attrs.get("Config", {}).get("Image") if cfg_image: images.add(cfg_image) except DockerError: pass return images def stacks_update_summary() -> dict[str, dict]: """Per-stack image-update status for every running compose project, read from the digest cache the background loop maintains — no registry calls, so it's cheap enough for the stacks list to poll. Stacks with no cached image yet are simply absent (treated as "no update" by the UI).""" by_stack: dict[str, set[str]] = {} try: client = get_client() for c in safe_call(client.containers.list, all=True): project = (c.labels or {}).get("com.docker.compose.project") if not project: continue cfg_image = c.attrs.get("Config", {}).get("Image") if cfg_image: by_stack.setdefault(project, set()).add(cfg_image) except DockerError: return {} summary: dict[str, dict] = {} for stack_id, images in by_stack.items(): stale = [ img for img in images if (st := _CACHE.get(img)) is not None and st.update_available ] summary[stack_id] = { "update_available": bool(stale), "stale_images": stale, } return summary async def stack_updates(stack_id: str, refresh: bool = True) -> dict: """Update status for one stack's images. ``refresh=True`` queries the registry now; ``False`` reads the cache the background loop already populated (so the auto-update pass adds no extra registry round-trips). DB-free, so the agent can reuse it verbatim. """ images = stack_images(stack_id) result: dict[str, dict] = {} for image in images: status = await check_image(image) if refresh else _CACHE.get(image) if status is not None: result[image] = status.to_dict() stale = [img for img, st in result.items() if st.get("update_available")] return { "stack_id": stack_id, "update_available": bool(stale), "stale_images": stale, "images": result, } async def check_all() -> dict[str, dict]: images = _all_running_images() for image in images: await check_image(image) return {k: v.to_dict() for k, v in _CACHE.items()} def get_cache() -> dict[str, dict]: return {k: v.to_dict() for k, v in _CACHE.items()} async def background_loop(): # initial delay so startup isn't blocked await asyncio.sleep(30) while True: try: await check_all() logger.info("Image update check complete (%d images)", len(_CACHE)) except Exception as exc: # noqa: BLE001 logger.warning("Image update check failed: %s", exc) # Apply auto-update policies using the digest cache we just refreshed. # Lazy import avoids a circular import (auto_update_service imports us). try: from services import auto_update_service await auto_update_service.run_due() except Exception as exc: # noqa: BLE001 logger.warning("Auto-update pass failed: %s", exc) # Re-read the interval each cycle so Settings changes take effect. interval = max(settings_service.get_update_interval(), 5) * 60 await asyncio.sleep(interval)