From 3e4187d93e61934e17b10e6d1d1d227f78575812 Mon Sep 17 00:00:00 2001 From: bmcquilkin Date: Wed, 29 Jul 2026 16:26:59 -0700 Subject: [PATCH] feat(operator): pod management --- .../sentry_streams_k8s/operator/__init__.py | 6 +- .../sentry_streams_k8s/operator/operator.py | 92 +++- .../sentry_streams_k8s/operator/pod_health.py | 53 +- .../sentry_streams_k8s/operator/pod_status.py | 4 - .../sentry_streams_k8s/operator/reconcile.py | 476 ++++++++++++++---- .../operator/streaming_pipeline.py | 10 +- sentry_streams_k8s/tests/k8s_fixtures.py | 135 +++++ sentry_streams_k8s/tests/test_operator.py | 363 +++++++------ sentry_streams_k8s/tests/test_pods.py | 251 +++++++-- .../tests/test_streaming_pipeline.py | 12 +- 10 files changed, 1051 insertions(+), 351 deletions(-) diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/__init__.py b/sentry_streams_k8s/sentry_streams_k8s/operator/__init__.py index 0da8da2b..96ca5016 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/__init__.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/__init__.py @@ -1,13 +1,15 @@ from sentry_streams_k8s.operator.streaming_pipeline import ( StreamingPipelineSpec, from_crd_spec, - render, + render_deployments, + render_pods, validate, ) __all__ = [ "StreamingPipelineSpec", "from_crd_spec", - "render", + "render_deployments", + "render_pods", "validate", ] diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py b/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py index e10aa96c..851ac554 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/operator.py @@ -2,26 +2,36 @@ import asyncio import copy +import json import os from collections.abc import Mapping +from datetime import datetime, timezone +from types import SimpleNamespace from typing import Any, cast import kopf from kubernetes import client +from kubernetes.client import V1Pod +from kubernetes.client.exceptions import ApiException from sentry_streams_k8s.k8s_types import V1ConditionDict from sentry_streams_k8s.operator.constants import ( + FIELD_MANAGER, GROUP, HEALTH_SCAN_INTERVAL_SECONDS, + MANAGED_BY_LABEL, MAX_CONCURRENT_RECONCILES, + OWNER_UID_LABEL, PLURAL, VERSION, WORKLOAD_NAMESPACE_ENV, Logger, ) +from sentry_streams_k8s.operator.pod_health import pod_health +from sentry_streams_k8s.operator.pod_resources import delete_owned_pods from sentry_streams_k8s.operator.reconcile import ( PipelineStatusPatch, - _prune_stale_resources, + prune_stale_configmaps, reconcile_pipeline, ) @@ -85,16 +95,39 @@ def _workload_namespace() -> str: return namespace -def _previous_conditions(body: kopf.Body) -> list[V1ConditionDict] | None: - status = body.get("status") - if not isinstance(status, Mapping): - return None +def _published_conditions(status: Mapping[str, object]) -> list[V1ConditionDict] | None: conditions = status.get("conditions") if not isinstance(conditions, list): return None return cast(list[V1ConditionDict], conditions) +def _deserialize_pod(body: kopf.Body) -> V1Pod: + # kopf provides the event's raw JSON body so we need to deserialize. + # Use a simple namespace since ApiClient expects a RESTResponse: + json_response = SimpleNamespace(data=json.dumps(dict(body))) + return cast(V1Pod, client.ApiClient().deserialize(json_response, "V1Pod")) + + +def _get_pipeline_status(name: str, namespace: str) -> Mapping[str, object]: + api = client.CustomObjectsApi() + try: + obj = api.get_namespaced_custom_object( + group=GROUP, + version=VERSION, + namespace=namespace, + plural=PLURAL, + name=name, + ) + except ApiException as e: + if e.status == 404: + return {} + raise + obj = cast(Mapping[str, object], obj) + status = obj.get("status") + return cast(Mapping[str, object], status) if isinstance(status, Mapping) else {} + + def _patch_pipeline_status(name: str, namespace: str, status: PipelineStatusPatch) -> None: api = client.CustomObjectsApi() api.patch_namespaced_custom_object_status( @@ -147,7 +180,6 @@ async def _reconcile_once( logger: Logger, scheduler: ReconcileScheduler, stopped: kopf.DaemonStopped, - previous_conditions: list[V1ConditionDict] | None = None, ) -> float | None: status_patch: PipelineStatusPatch = {} try: @@ -156,6 +188,7 @@ async def _reconcile_once( if stopped: return None try: + published = await asyncio.to_thread(_get_pipeline_status, name, namespace) await asyncio.to_thread( reconcile_pipeline, spec=copy.deepcopy(dict(spec)), @@ -165,7 +198,8 @@ async def _reconcile_once( workload_namespace=_workload_namespace(), logger=logger, status=status_patch, - previous_conditions=previous_conditions, + previous_conditions=_published_conditions(published), + previous_generations=published.get("generations"), ) except kopf.PermanentError as e: if status_patch: @@ -188,7 +222,6 @@ async def _reconcile_once( async def reconcile_pipeline_daemon( stopped: kopf.DaemonStopped, spec: kopf.Spec, - body: kopf.Body, name: str, namespace: str | None, uid: str, @@ -212,7 +245,6 @@ async def reconcile_pipeline_daemon( logger=logger, scheduler=scheduler, stopped=stopped, - previous_conditions=_previous_conditions(body), ) if stopped: break @@ -221,16 +253,52 @@ async def reconcile_pipeline_daemon( scheduler.unregister(uid, event) +@kopf.on.event("", "v1", "pods", labels={MANAGED_BY_LABEL: FIELD_MANAGER}) +async def handle_pipeline_pod_event( + type: str | None, + body: kopf.Body, + meta: kopf.Meta, + labels: kopf.Labels, + name: str | None, + namespace: str | None, + memo: kopf.Memo, + logger: Logger, + **_: Any, +) -> None: + if type not in {"DELETED", "MODIFIED"}: + return + + if type == "MODIFIED" and meta.deletion_timestamp is None: + health = pod_health(_deserialize_pod(body), datetime.now(timezone.utc)) + if not health.delete: + return + + owner_uid = labels.get(OWNER_UID_LABEL) + if not owner_uid: + logger.warning( + "managed Pod %s/%s is missing its owner UID label; cannot reconcile", + namespace, + name, + ) + return + + if _scheduler(memo).notify(owner_uid): + logger.info("requested reconciliation after Pod %s event=%s", name, type) + + @kopf.on.delete(GROUP, VERSION, PLURAL) async def cleanup(uid: str, memo: kopf.Memo, logger: Logger, **_: Any) -> None: scheduler = _scheduler(memo) async with scheduler.limit: async with scheduler.lock(uid): + workload_namespace = _workload_namespace() + core = client.CoreV1Api() + await asyncio.to_thread(delete_owned_pods, core, workload_namespace, uid, logger) await asyncio.to_thread( - _prune_stale_resources, - workload_namespace=_workload_namespace(), + prune_stale_configmaps, + core=core, + workload_namespace=workload_namespace, owner_uid=uid, - desired_deployments=set(), desired_configmaps=set(), logger=logger, ) diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/pod_health.py b/sentry_streams_k8s/sentry_streams_k8s/operator/pod_health.py index c3dad270..d0849a38 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/pod_health.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/pod_health.py @@ -82,14 +82,6 @@ def _first_waiting_reason( return None -def _first_unhealthy_waiting(statuses: list[V1ContainerStatus] | None) -> str | None: - return _first_waiting_reason(statuses, UNHEALTHY_WAITING_REASONS) - - -def _first_permanent_waiting(statuses: list[V1ContainerStatus] | None) -> str | None: - return _first_waiting_reason(statuses, PERMANENT_WAITING_REASONS) - - def _container_failed_terminated_reason(status: V1ContainerStatus) -> str | None: terminated = status.state.terminated if status.state else None if terminated is None: @@ -140,27 +132,6 @@ class PodHealth: permanent: bool = False -def _verdict( - pod_name: str, - *, - ready: bool = False, - reason: str | None = None, - unhealthy: bool = False, - delete: bool = False, - force: bool = False, - permanent: bool = False, -) -> PodHealth: - return PodHealth( - name=pod_name, - ready=ready, - reason=reason, - unhealthy=unhealthy, - delete=delete, - force=force, - permanent=permanent, - ) - - def _container_statuses_verdict( pod_name: str, statuses: list[V1ContainerStatus] | None, @@ -169,9 +140,9 @@ def _container_statuses_verdict( *, reason_prefix: str = "", ) -> PodHealth | None: - permanent_waiting = _first_permanent_waiting(statuses) + permanent_waiting = _first_waiting_reason(statuses, PERMANENT_WAITING_REASONS) if permanent_waiting is not None: - return _verdict( + return PodHealth( pod_name, reason=f"{reason_prefix}{permanent_waiting}", unhealthy=True, @@ -180,16 +151,16 @@ def _container_statuses_verdict( terminated_reason = _first_failed_terminated_reason(statuses) if terminated_reason is not None: - return _verdict( + return PodHealth( pod_name, reason=f"{reason_prefix}{terminated_reason}", unhealthy=True, delete=True, ) - unhealthy_waiting = _first_unhealthy_waiting(statuses) + unhealthy_waiting = _first_waiting_reason(statuses, UNHEALTHY_WAITING_REASONS) if unhealthy_waiting is not None: - return _verdict( + return PodHealth( pod_name, reason=f"{reason_prefix}{unhealthy_waiting}", unhealthy=True, @@ -198,7 +169,7 @@ def _container_statuses_verdict( waiting_reason = _first_waiting_reason(statuses) if waiting_reason is not None: - return _verdict( + return PodHealth( pod_name, reason=f"{reason_prefix}{waiting_reason}", ) @@ -228,7 +199,7 @@ def pod_health(pod: V1Pod, now: datetime) -> PodHealth: if terminating_age is not None: stuck = terminating_age >= POD_TERMINATING_GRACE_SECONDS reason = "StuckTerminating" if stuck else "Terminating" - return _verdict( + return PodHealth( pod_name, reason=reason, unhealthy=stuck, @@ -238,10 +209,10 @@ def pod_health(pod: V1Pod, now: datetime) -> PodHealth: status = pod.status if status is None: - return _verdict(pod_name) + return PodHealth(pod_name) if _pod_unschedulable(status.conditions): - return _verdict(pod_name, reason="Unschedulable", unhealthy=True) + return PodHealth(pod_name, reason="Unschedulable", unhealthy=True) init_verdict = _container_statuses_verdict( pod_name, @@ -256,7 +227,7 @@ def pod_health(pod: V1Pod, now: datetime) -> PodHealth: phase = status.phase if phase == "Succeeded": - return _verdict(pod_name, reason="Succeeded", delete=True) + return PodHealth(pod_name, reason="Succeeded", delete=True) app_verdict = _container_statuses_verdict( pod_name, @@ -269,6 +240,6 @@ def pod_health(pod: V1Pod, now: datetime) -> PodHealth: if phase == "Failed": reason = status.reason or phase or "Terminated" - return _verdict(pod_name, reason=reason, unhealthy=True, delete=True) + return PodHealth(pod_name, reason=reason, unhealthy=True, delete=True) - return _verdict(pod_name, ready=pod_is_ready(pod)) + return PodHealth(pod_name, ready=pod_is_ready(pod)) diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/pod_status.py b/sentry_streams_k8s/sentry_streams_k8s/operator/pod_status.py index 6c91035e..579b8604 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/pod_status.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/pod_status.py @@ -35,10 +35,6 @@ class ReportedPodStatus: reason: str | None = None permanent: bool = False - @property - def is_unhealthy(self) -> bool: - return self.unhealthy - def to_status_dict(self) -> PodStatusEntry: entry: PodStatusEntry = { "name": self.name, diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/reconcile.py b/sentry_streams_k8s/sentry_streams_k8s/operator/reconcile.py index d4132f24..db3c08e7 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/reconcile.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/reconcile.py @@ -1,42 +1,83 @@ from __future__ import annotations +from collections.abc import Mapping from datetime import datetime, timezone from typing import Any, TypedDict, cast import kopf from kubernetes import client -from kubernetes.client import V1Condition +from kubernetes.client import V1Condition, V1Pod -from sentry_streams_k8s.consumer_builder import compute_config_version +from sentry_streams_k8s.constants import ALL_WORKLOAD_SETS +from sentry_streams_k8s.consumer_builder import WorkloadSet, compute_config_version from sentry_streams_k8s.k8s_types import ( V1ConditionDict, V1ConfigMapDict, - V1DeploymentDict, + V1PodDict, ) from sentry_streams_k8s.operator.constants import ( APPLY_PATCH_CONTENT_TYPE, FIELD_MANAGER, + MAX_BASE_NAME_LENGTH, + MAX_GENERATION, + MAX_REPLICAS, OWNER_NAME_ANNOTATION, OWNER_NAMESPACE_ANNOTATION, OWNER_UID_LABEL, Logger, ) +from sentry_streams_k8s.operator.pod_health import ( + PodHealth, + is_deleting, + pod_health, + pod_spec_changed, +) +from sentry_streams_k8s.operator.pod_resources import ( + apply_pod, + build_pipeline_pod, + consumer_pod_name, + delete_pod, + list_owned_pods, + pod_generation, + pod_keep_key, + pod_name, + pod_ordinal, + pod_workload_set, +) +from sentry_streams_k8s.operator.pod_status import ( + PodStatusEntry, + ReportedPodStatus, + reported_pod_status, +) from sentry_streams_k8s.operator.streaming_pipeline import ( from_crd_spec, - render, + render_pods, validate, ) +class PodSetResult(TypedDict): + childPods: list[str] + desiredReplicas: int + readyReplicas: int + unhealthyPods: list[PodStatusEntry] + permanentErrors: list[PodStatusEntry] + + +class CombinedPodResult(PodSetResult): + sets: dict[str, PodSetResult | None] + + class PipelineStatusPatch(TypedDict, total=False): conditions: list[V1ConditionDict] config_version: str - replicas: dict[str, int] + pods: CombinedPodResult workload_namespace: str + generations: dict[str, dict[str, int] | None] def _prepare_manifest( - manifest: V1ConfigMapDict | V1DeploymentDict, + manifest: V1ConfigMapDict, *, workload_namespace: str, owner_uid: str, @@ -72,50 +113,16 @@ def _apply_configmap( ) -def _apply_deployment( - apps: client.AppsV1Api, - manifest: V1DeploymentDict, - *, - workload_namespace: str, -) -> None: - apps.patch_namespaced_deployment( - name=manifest["metadata"]["name"], - namespace=workload_namespace, - body=manifest, - field_manager=FIELD_MANAGER, - force=True, - _content_type=APPLY_PATCH_CONTENT_TYPE, - ) - - -def _prune_stale_resources( +def prune_stale_configmaps( *, + core: client.CoreV1Api, workload_namespace: str, owner_uid: str, - desired_deployments: set[str], desired_configmaps: set[str], logger: Logger, ) -> None: selector = f"{OWNER_UID_LABEL}={owner_uid}" - apps = client.AppsV1Api() - deployments = apps.list_namespaced_deployment( - namespace=workload_namespace, - label_selector=selector, - ) - for deployment in deployments.items: - if deployment.metadata.name not in desired_deployments: - logger.info( - "Pruning stale deployment %s/%s", - workload_namespace, - deployment.metadata.name, - ) - apps.delete_namespaced_deployment( - name=deployment.metadata.name, - namespace=workload_namespace, - ) - - core = client.CoreV1Api() configmaps = core.list_namespaced_config_map( namespace=workload_namespace, label_selector=selector, @@ -152,6 +159,276 @@ def _merge_conditions( return merged +def _parse_generations(data: object) -> dict[int, int]: + """ + Returns a dict mapping ordinal -> highest generation. + Empty when there is no ledger yet (e.g. first reconcile). + """ + + if not isinstance(data, Mapping): + return {} + return {int(ordinal): generation for ordinal, generation in data.items()} + + +def _serialize_generations(generations: dict[int, int]) -> dict[str, int]: + """ + Pods for a replica are named {name}-{ordinal}-{generation}. The generation + increments every time the operator replaces a Pod, so a newly created Pod + never shares a name with the old one that is still terminating. This is what + lets us recreate instantly instead of waiting for the old Pod to be deleted. + We save the highest generation per ordinal in the CR's status.generations. + """ + + return {str(ordinal): generation for ordinal, generation in sorted(generations.items())} + + +def _delete_current_pod( + core: client.CoreV1Api, + pod: V1Pod, + namespace: str, + logger: Logger, + health: PodHealth, + reason: str, +) -> None: + name = pod_name(pod) + if is_deleting(pod) and not health.delete: + logger.info("pipeline Pod %s/%s is already deleting reason=%s", namespace, name, reason) + return + delete_pod(core, name, namespace, force=health.force) + logger.info("deleted pipeline Pod %s/%s reason=%s", namespace, name, reason) + + +def delete_obsolete_pod_sets( + core: client.CoreV1Api, + namespace: str, + owner_uid: str, + desired_sets: set[str], + logger: Logger, +) -> None: + now = datetime.now(timezone.utc) + for pod in list_owned_pods(core, namespace, owner_uid): + workload_set = pod_workload_set(pod) + if workload_set in desired_sets: + continue + health = pod_health(pod, now) + _delete_current_pod( + core, + pod, + namespace, + logger, + health, + f"StaleWorkloadSet:{workload_set or 'missing'}", + ) + + +def _allocate_generation(generations: dict[int, int], ordinal: int, pods: list[V1Pod]) -> int: + live_max = max((pod_generation(pod) for pod in pods), default=-1) + generation = max(generations.get(ordinal, -1), live_max) + 1 + + if generation > MAX_GENERATION: + # TODO: We should see if it is safe to wrap back to generation 0 at this point. + raise kopf.PermanentError(f"Replica {ordinal} exceeds {MAX_GENERATION} generations.") + + generations[ordinal] = generation + return generation + + +def reconcile_pipeline_pods( + *, + core: client.CoreV1Api, + workload_namespace: str, + owner_uid: str, + owner_name: str, + owner_namespace: str, + base_name: str, + template_metadata: Mapping[str, Any], + template_spec: Mapping[str, Any], + replicas: int, + generations: dict[int, int], + logger: Logger, + workload_set: str, +) -> PodSetResult: + desired_ordinals = set(range(max(replicas, 0))) + current = list_owned_pods( + core, + workload_namespace, + owner_uid, + workload_set=workload_set, + ) + now = datetime.now(timezone.utc) + + pods_by_ordinal: dict[int, list[V1Pod]] = {} + health_by_name: dict[str, PodHealth] = {} + reported_statuses: list[ReportedPodStatus] = [] + active_pod_names: list[str] = [] + + def _build(ordinal: int, generation: int) -> V1PodDict: + return build_pipeline_pod( + base_name=base_name, + template_metadata=template_metadata, + template_spec=template_spec, + ordinal=ordinal, + generation=generation, + owner_uid=owner_uid, + owner_name=owner_name, + owner_namespace=owner_namespace, + workload_set=workload_set, + ) + + for pod in current: + name = pod_name(pod) + health = pod_health(pod, now) + health_by_name[name] = health + ordinal = pod_ordinal(pod) + if ordinal is None or ordinal not in desired_ordinals: + _delete_current_pod(core, pod, workload_namespace, logger, health, "Stale") + continue + pods_by_ordinal.setdefault(ordinal, []).append(pod) + reported_statuses.append(reported_pod_status(pod, health)) + + for ordinal in sorted(desired_ordinals): + pods = pods_by_ordinal.get(ordinal, []) + desired_template = _build(ordinal, 0) + candidates = [ + pod + for pod in pods + if ( + not is_deleting(pod) + and not pod_spec_changed(pod, desired_template) + and not health_by_name[pod_name(pod)].delete + ) + ] + + keep: V1Pod | None + if candidates: + # Prefer a ready Pod, otherwise keep the newest generation: + keep = max(candidates, key=lambda pod: pod_keep_key(pod, health_by_name[pod_name(pod)])) + active_pod_names.append(pod_name(keep)) + generations[ordinal] = max(generations.get(ordinal, -1), pod_generation(keep)) + else: + # Create the replacement before deleting the old Pod: + generation = _allocate_generation(generations, ordinal, pods) + keep_manifest = _build(ordinal, generation) + apply_pod(core, keep_manifest, workload_namespace) + keep = None + keep_name = consumer_pod_name(base_name, ordinal, generation) + active_pod_names.append(keep_name) + logger.info( + "applied replacement pipeline Pod %s/%s", + workload_namespace, + keep_name, + ) + + for pod in pods: + if pod is keep: + continue + health = health_by_name[pod_name(pod)] + if pod_spec_changed(pod, desired_template): + _delete_current_pod(core, pod, workload_namespace, logger, health, "Outdated") + elif health.delete: + _delete_current_pod( + core, pod, workload_namespace, logger, health, health.reason or "Unhealthy" + ) + elif pod in candidates: + _delete_current_pod(core, pod, workload_namespace, logger, health, "Duplicate") + + active = set(active_pod_names) + ready_ordinals = { + ordinal + for ordinal, pods in pods_by_ordinal.items() + if any( + not is_deleting(pod) and pod_name(pod) in active and health_by_name[pod_name(pod)].ready + for pod in pods + ) + } + unhealthy_pods = [status.to_status_dict() for status in reported_statuses if status.unhealthy] + permanent_errors = [ + status.to_status_dict() + for status in reported_statuses + if status.unhealthy and status.permanent + ] + return { + "childPods": sorted(active_pod_names), + "desiredReplicas": len(desired_ordinals), + "readyReplicas": len(ready_ordinals), + "unhealthyPods": unhealthy_pods, + "permanentErrors": permanent_errors, + } + + +def _reconcile_pod_set( + *, + core: client.CoreV1Api, + workload: WorkloadSet, + workload_set: str, + workload_namespace: str, + owner_uid: str, + owner_name: str, + owner_namespace: str, + logger: Logger, + previous_generations: dict[int, int], +) -> tuple[PodSetResult, dict[int, int]]: + base_name = workload.name + + if len(base_name) > MAX_BASE_NAME_LENGTH: + raise kopf.PermanentError( + f"{workload_set} name cannot exceed {MAX_BASE_NAME_LENGTH} characters." + ) + + replicas = workload.replicas + + if type(replicas) is not int or replicas < 0: + raise kopf.PermanentError(f"{workload_set} replica count must be a non-negative integer.") + + if replicas > MAX_REPLICAS: + raise kopf.PermanentError(f"{workload_set} replica count cannot exceed {MAX_REPLICAS}.") + + generations = dict(previous_generations) + pod_result = reconcile_pipeline_pods( + core=core, + workload_namespace=workload_namespace, + owner_uid=owner_uid, + owner_name=owner_name, + owner_namespace=owner_namespace, + base_name=base_name, + template_metadata=workload.pod_template["metadata"], + template_spec=workload.pod_template["spec"], + replicas=replicas, + generations=generations, + logger=logger, + workload_set=workload_set, + ) + return pod_result, generations + + +def _with_workload_set(entry: PodStatusEntry, workload_set: str) -> PodStatusEntry: + return {**entry, "workloadSet": workload_set} + + +def _combine_pod_results(results: dict[str, PodSetResult]) -> CombinedPodResult: + child_pods: list[str] = [] + unhealthy_pods: list[PodStatusEntry] = [] + permanent_errors: list[PodStatusEntry] = [] + for workload_set, result in results.items(): + child_pods.extend(result["childPods"]) + unhealthy_pods.extend( + _with_workload_set(entry, workload_set) for entry in result["unhealthyPods"] + ) + permanent_errors.extend( + _with_workload_set(entry, workload_set) for entry in result["permanentErrors"] + ) + return { + "childPods": sorted(child_pods), + "desiredReplicas": sum(result["desiredReplicas"] for result in results.values()), + "readyReplicas": sum(result["readyReplicas"] for result in results.values()), + "unhealthyPods": unhealthy_pods, + "permanentErrors": permanent_errors, + # Status is updated as a JSON merge patch, so a set that is not included keeps its old + # value. Explicitly include all sets and null out the removed ones to clear them: + "sets": {workload_set: results.get(workload_set) for workload_set in ALL_WORKLOAD_SETS}, + } + + def reconcile_pipeline( *, spec: Any, @@ -163,7 +440,8 @@ def reconcile_pipeline( patch: kopf.Patch | None = None, status: PipelineStatusPatch | None = None, previous_conditions: list[V1ConditionDict] | None = None, -) -> None: + previous_generations: object = None, +) -> CombinedPodResult: status_patch = ( status if status is not None @@ -173,12 +451,12 @@ def reconcile_pipeline( now = datetime.now(timezone.utc).replace(microsecond=0) core = client.CoreV1Api() - apps = client.AppsV1Api() + ledger = previous_generations if isinstance(previous_generations, Mapping) else {} consumer = from_crd_spec(dict(spec), name=name) try: validate(consumer) - result = render(consumer) + result = render_pods(consumer) except Exception as e: if status_patch is not None: failed = [ @@ -199,56 +477,46 @@ def reconcile_pipeline( ) raise kopf.PermanentError(f"StreamingPipeline {namespace}/{name} failed to render: {e}") - manifests: list[V1ConfigMapDict | V1DeploymentDict] = [ - result["configmap"], - result["deployment"], - ] - if "canary_deployment" in result: - manifests.append(result["canary_deployment"]) + configmap = result["configmap"] + _prepare_manifest( + configmap, + workload_namespace=workload_namespace, + owner_uid=uid, + owner_name=name, + owner_namespace=namespace, + ) + _apply_configmap(core, configmap, workload_namespace=workload_namespace) - for manifest in manifests: - _prepare_manifest( - manifest, + pod_set_results: dict[str, PodSetResult] = {} + generations_by_set: dict[str, dict[int, int]] = {} + for workload_set, workload in result["sets"].items(): + set_result, generations = _reconcile_pod_set( + core=core, + workload=workload, + workload_set=workload_set, workload_namespace=workload_namespace, owner_uid=uid, owner_name=name, owner_namespace=namespace, + logger=logger, + previous_generations=_parse_generations(ledger.get(workload_set)), ) - kind = manifest["kind"] - if kind == "ConfigMap": - _apply_configmap( - core, - cast(V1ConfigMapDict, manifest), - workload_namespace=workload_namespace, - ) - elif kind == "Deployment": - _apply_deployment( - apps, - cast(V1DeploymentDict, manifest), - workload_namespace=workload_namespace, - ) - else: - raise kopf.PermanentError(f"Cannot apply unsupported manifest kind {kind}.") + pod_set_results[workload_set] = set_result + generations_by_set[workload_set] = generations + + # Remove Pods left behind by a workload set that is no longer rendered: + delete_obsolete_pod_sets(core, workload_namespace, uid, set(result["sets"]), logger) - _prune_stale_resources( + pod_result = _combine_pod_results(pod_set_results) + prune_stale_configmaps( + core=core, workload_namespace=workload_namespace, owner_uid=uid, - desired_deployments={ - manifest["metadata"]["name"] - for manifest in manifests - if manifest["kind"] == "Deployment" - }, - desired_configmaps={ - manifest["metadata"]["name"] - for manifest in manifests - if manifest["kind"] == "ConfigMap" - }, + desired_configmaps={configmap["metadata"]["name"]}, logger=logger, ) if status_patch is not None: - replicas = consumer.get("replicas", 1) - canary = 1 if "canary_deployment" in result else 0 applied = [ V1Condition( type="Rendered", @@ -257,14 +525,32 @@ def reconcile_pipeline( message="", last_transition_time=now, ), - V1Condition( - type="Applied", - status="True", - reason="Applied", - message="", - last_transition_time=now, - ), ] + permanent_errors = pod_result["permanentErrors"] + if permanent_errors: + error = permanent_errors[0] + applied.append( + V1Condition( + type="Applied", + status="False", + reason="PermanentPodFailure", + message=( + f"Pod {error['name']} is permanently unhealthy: " + f"{error.get('reason', 'Unknown')}. Update the StreamingPipeline spec." + ), + last_transition_time=now, + ) + ) + else: + applied.append( + V1Condition( + type="Applied", + status="True", + reason="Applied", + message="", + last_transition_time=now, + ) + ) status_patch["conditions"] = _merge_conditions( previous_conditions, cast( @@ -273,5 +559,15 @@ def reconcile_pipeline( ), ) status_patch["config_version"] = compute_config_version(consumer["pipeline_config"]) - status_patch["replicas"] = {"primary": replicas - canary, "canary": canary} + status_patch["pods"] = pod_result status_patch["workload_namespace"] = workload_namespace + status_patch["generations"] = { + workload_set: ( + _serialize_generations(generations_by_set[workload_set]) + if workload_set in generations_by_set + else None + ) + for workload_set in ALL_WORKLOAD_SETS + } + + return pod_result diff --git a/sentry_streams_k8s/sentry_streams_k8s/operator/streaming_pipeline.py b/sentry_streams_k8s/sentry_streams_k8s/operator/streaming_pipeline.py index d5029c58..b6faa3fc 100644 --- a/sentry_streams_k8s/sentry_streams_k8s/operator/streaming_pipeline.py +++ b/sentry_streams_k8s/sentry_streams_k8s/operator/streaming_pipeline.py @@ -6,6 +6,7 @@ ConsumerBuilder, ConsumerSpec, RenderedDeployments, + RenderedPods, ) @@ -74,8 +75,15 @@ def to_consumer_spec(spec: StreamingPipelineSpec) -> ConsumerSpec: ) -def render(spec: StreamingPipelineSpec) -> RenderedDeployments: +def render_deployments(spec: StreamingPipelineSpec) -> RenderedDeployments: builder = ConsumerBuilder(spec["deployment_template"], spec["container_template"]) consumer = to_consumer_spec(spec) builder.validate(consumer, spec["pipeline_config"]) return builder.build_deployments(consumer, spec["pipeline_config"]) + + +def render_pods(spec: StreamingPipelineSpec) -> RenderedPods: + builder = ConsumerBuilder(spec["deployment_template"], spec["container_template"]) + consumer = to_consumer_spec(spec) + builder.validate(consumer, spec["pipeline_config"]) + return builder.build_pods(consumer, spec["pipeline_config"]) diff --git a/sentry_streams_k8s/tests/k8s_fixtures.py b/sentry_streams_k8s/tests/k8s_fixtures.py index 2a4d6d80..2cc9912d 100644 --- a/sentry_streams_k8s/tests/k8s_fixtures.py +++ b/sentry_streams_k8s/tests/k8s_fixtures.py @@ -1,9 +1,14 @@ from __future__ import annotations +import copy from collections.abc import Mapping +from dataclasses import dataclass, field from datetime import datetime +from typing import Any, cast from kubernetes.client import ( + V1ConfigMap, + V1ConfigMapList, V1ContainerState, V1ContainerStateTerminated, V1ContainerStateWaiting, @@ -11,9 +16,139 @@ V1ObjectMeta, V1Pod, V1PodCondition, + V1PodList, V1PodStatus, ) +from sentry_streams_k8s.k8s_types import V1ConfigMapDict, V1PodDict + + +def _matches_selector(labels: Mapping[str, str] | None, selector: str | None) -> bool: + if not selector: + return True + labels = labels or {} + return all( + labels.get(key) == value + for requirement in selector.split(",") + for key, value in [requirement.split("=", maxsplit=1)] + ) + + +@dataclass +class FakeCoreV1Api: + """In-memory CoreV1Api subset used by operator tests.""" + + pods: list[V1Pod] = field(default_factory=list) + configmaps: list[V1ConfigMap] = field(default_factory=list) + + applied_pods: list[V1PodDict] = field(default_factory=list, init=False) + deleted_pods: list[tuple[str, bool]] = field(default_factory=list, init=False) + + applied_configmaps: list[V1ConfigMapDict] = field(default_factory=list, init=False) + deleted_configmaps: list[str] = field(default_factory=list, init=False) + + pod_list_calls: list[tuple[str, str | None]] = field(default_factory=list, init=False) + configmap_list_calls: list[tuple[str, str | None]] = field(default_factory=list, init=False) + + pod_patch_calls: list[dict[str, Any]] = field(default_factory=list, init=False) + configmap_patch_calls: list[dict[str, Any]] = field(default_factory=list, init=False) + + operations: list[str] = field(default_factory=list, init=False) + + def list_namespaced_pod( + self, + *, + namespace: str, + label_selector: str | None = None, + ) -> V1PodList: + self.pod_list_calls.append((namespace, label_selector)) + items = [ + copy.deepcopy(pod) + for pod in self.pods + if _matches_selector(pod.metadata.labels if pod.metadata else None, label_selector) + ] + return V1PodList(items=items) + + def patch_namespaced_pod( + self, + *, + name: str, + namespace: str, + body: object, + field_manager: str, + force: bool, + _content_type: str, + ) -> None: + self.applied_pods.append(copy.deepcopy(cast(V1PodDict, body))) + self.pod_patch_calls.append( + { + "name": name, + "namespace": namespace, + "body": copy.deepcopy(body), + "field_manager": field_manager, + "force": force, + "_content_type": _content_type, + } + ) + self.operations.append(f"apply:{name}") + + def delete_namespaced_pod( + self, + *, + name: str, + namespace: str, + body: object | None = None, + ) -> None: + del namespace + force = getattr(body, "grace_period_seconds", None) == 0 + self.deleted_pods.append((name, force)) + self.operations.append(f"delete:{name}") + + def list_namespaced_config_map( + self, + *, + namespace: str, + label_selector: str | None = None, + ) -> V1ConfigMapList: + self.configmap_list_calls.append((namespace, label_selector)) + items = [ + copy.deepcopy(configmap) + for configmap in self.configmaps + if _matches_selector( + configmap.metadata.labels if configmap.metadata else None, + label_selector, + ) + ] + return V1ConfigMapList(items=items) + + def patch_namespaced_config_map( + self, + *, + name: str, + namespace: str, + body: object, + field_manager: str, + force: bool, + _content_type: str, + ) -> None: + self.applied_configmaps.append(copy.deepcopy(cast(V1ConfigMapDict, body))) + self.configmap_patch_calls.append( + { + "name": name, + "namespace": namespace, + "body": copy.deepcopy(body), + "field_manager": field_manager, + "force": force, + "_content_type": _content_type, + } + ) + self.operations.append(f"apply-configmap:{name}") + + def delete_namespaced_config_map(self, *, name: str, namespace: str) -> None: + del namespace + self.deleted_configmaps.append(name) + self.operations.append(f"delete-configmap:{name}") + def make_condition(type_: str, status: str) -> V1PodCondition: return V1PodCondition(type=type_, status=status) diff --git a/sentry_streams_k8s/tests/test_operator.py b/sentry_streams_k8s/tests/test_operator.py index 350ffa33..35f77cbd 100644 --- a/sentry_streams_k8s/tests/test_operator.py +++ b/sentry_streams_k8s/tests/test_operator.py @@ -12,6 +12,9 @@ import pytest from kubernetes import client +from sentry_streams_k8s.constants import CANARY_WORKLOAD_SET, PRIMARY_WORKLOAD_SET +from sentry_streams_k8s.consumer_builder import WorkloadSet +from sentry_streams_k8s.k8s_types import V1ConditionDict from sentry_streams_k8s.operator import operator as operator_module from sentry_streams_k8s.operator.constants import ( HEALTH_SCAN_INTERVAL_SECONDS, @@ -25,6 +28,7 @@ _reconcile_once, _wait_for_reconcile, cleanup, + handle_pipeline_pod_event, reconcile_pipeline_daemon, request_pipeline_reconcile, ) @@ -32,19 +36,20 @@ APPLY_PATCH_CONTENT_TYPE, PipelineStatusPatch, _apply_configmap, - _apply_deployment, + _merge_conditions, _prepare_manifest, - _prune_stale_resources, + prune_stale_configmaps, reconcile_pipeline, ) +from tests.k8s_fixtures import FakeCoreV1Api WORKLOAD_NAMESPACE = "test-streaming-pipelines" def test_prepare_manifest_routes_workload_and_records_source_cr() -> None: manifest = { - "apiVersion": "apps/v1", - "kind": "Deployment", + "apiVersion": "v1", + "kind": "ConfigMap", "metadata": { "name": "pipeline", "labels": {"service": "test"}, @@ -84,79 +89,51 @@ def _configmap_manifest() -> dict[str, Any]: } -def _deployment_manifest() -> dict[str, Any]: - return { - "apiVersion": "apps/v1", - "kind": "Deployment", - "metadata": {"name": "pipeline", "namespace": WORKLOAD_NAMESPACE}, - } - - def test_apply_configmap() -> None: - core = MagicMock() + core = FakeCoreV1Api() manifest = _configmap_manifest() _apply_configmap(core, manifest, workload_namespace=WORKLOAD_NAMESPACE) - core.patch_namespaced_config_map.assert_called_once_with( - name="pipeline", - namespace=WORKLOAD_NAMESPACE, - body=manifest, - field_manager="streaming-operator", - force=True, - _content_type=APPLY_PATCH_CONTENT_TYPE, - ) - - -def test_apply_deployment() -> None: - apps = MagicMock() - manifest = _deployment_manifest() + assert core.configmap_patch_calls == [ + { + "name": "pipeline", + "namespace": WORKLOAD_NAMESPACE, + "body": manifest, + "field_manager": "streaming-operator", + "force": True, + "_content_type": APPLY_PATCH_CONTENT_TYPE, + } + ] - _apply_deployment(apps, manifest, workload_namespace=WORKLOAD_NAMESPACE) - apps.patch_namespaced_deployment.assert_called_once_with( - name="pipeline", - namespace=WORKLOAD_NAMESPACE, - body=manifest, - field_manager="streaming-operator", - force=True, - _content_type=APPLY_PATCH_CONTENT_TYPE, +def test_prune_removes_only_stale_configmaps() -> None: + core = FakeCoreV1Api( + configmaps=[ + client.V1ConfigMap( + metadata=client.V1ObjectMeta( + name="desired-configmap", + labels={OWNER_UID_LABEL: "owner-uid"}, + ) + ), + client.V1ConfigMap( + metadata=client.V1ObjectMeta( + name="stale-configmap", + labels={OWNER_UID_LABEL: "owner-uid"}, + ) + ), + ] ) - -@patch("sentry_streams_k8s.operator.reconcile.client.CoreV1Api") -@patch("sentry_streams_k8s.operator.reconcile.client.AppsV1Api") -def test_prune_removes_only_stale_resources( - apps_api: MagicMock, - core_api: MagicMock, -) -> None: - apps = apps_api.return_value - apps.list_namespaced_deployment.return_value.items = [ - SimpleNamespace(metadata=SimpleNamespace(name="desired-deployment")), - SimpleNamespace(metadata=SimpleNamespace(name="stale-deployment")), - ] - core = core_api.return_value - core.list_namespaced_config_map.return_value.items = [ - SimpleNamespace(metadata=SimpleNamespace(name="desired-configmap")), - SimpleNamespace(metadata=SimpleNamespace(name="stale-configmap")), - ] - - _prune_stale_resources( + prune_stale_configmaps( + core=core, workload_namespace=WORKLOAD_NAMESPACE, owner_uid="owner-uid", - desired_deployments={"desired-deployment"}, desired_configmaps={"desired-configmap"}, logger=MagicMock(), ) - apps.delete_namespaced_deployment.assert_called_once_with( - name="stale-deployment", - namespace=WORKLOAD_NAMESPACE, - ) - core.delete_namespaced_config_map.assert_called_once_with( - name="stale-configmap", - namespace=WORKLOAD_NAMESPACE, - ) + assert core.deleted_configmaps == ["stale-configmap"] class FakeStopped: @@ -199,57 +176,57 @@ def _pipeline_spec() -> dict[str, Any]: return {"pipeline_config": {"steps": []}, "replicas": 2, "with_canary": True} -def _stub_render(monkeypatch: pytest.MonkeyPatch) -> tuple[MagicMock, MagicMock]: +def _workload(name: str, replicas: int) -> WorkloadSet: + return WorkloadSet( + name=name, + replicas=replicas, + labels={"app": "consumer"}, + pod_template={ + "metadata": {"labels": {"app": "consumer"}}, + "spec": {"containers": [{"name": "consumer", "image": "example/consumer:v1"}]}, + }, + ) + + +def _stub_render(monkeypatch: pytest.MonkeyPatch) -> MagicMock: configmap = { "apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "pipeline-config"}, } - primary = { - "apiVersion": "apps/v1", - "kind": "Deployment", - "metadata": {"name": "consumer"}, - "spec": {"replicas": 1}, - } - canary = { - "apiVersion": "apps/v1", - "kind": "Deployment", - "metadata": {"name": "consumer-canary"}, - "spec": {"replicas": 1}, - } core = MagicMock() core.list_namespaced_config_map.return_value.items = [] + core.list_namespaced_pod.return_value.items = [] core.api_client = _api_client() - apps = MagicMock() - apps.list_namespaced_deployment.return_value.items = [] monkeypatch.setattr( "sentry_streams_k8s.operator.reconcile.from_crd_spec", lambda spec, name: spec ) monkeypatch.setattr("sentry_streams_k8s.operator.reconcile.validate", lambda _consumer: None) monkeypatch.setattr( - "sentry_streams_k8s.operator.reconcile.render", + "sentry_streams_k8s.operator.reconcile.render_pods", lambda _consumer: { "configmap": configmap, - "deployment": primary, - "canary_deployment": canary, + "sets": { + PRIMARY_WORKLOAD_SET: _workload("consumer", 1), + CANARY_WORKLOAD_SET: _workload("consumer-canary", 1), + }, }, ) monkeypatch.setattr( "sentry_streams_k8s.operator.reconcile.compute_config_version", lambda _config: "version" ) monkeypatch.setattr("sentry_streams_k8s.operator.reconcile.client.CoreV1Api", lambda: core) - monkeypatch.setattr("sentry_streams_k8s.operator.reconcile.client.AppsV1Api", lambda: apps) - return core, apps + return core -def test_reconcile_applies_deployments_and_reports_status_through_a_plain_dict( +def test_reconcile_applies_pods_and_reports_status_through_a_plain_dict( monkeypatch: pytest.MonkeyPatch, ) -> None: - core, apps = _stub_render(monkeypatch) + core = _stub_render(monkeypatch) status: PipelineStatusPatch = {} - reconcile_pipeline( + result = reconcile_pipeline( spec=_pipeline_spec(), name="pipeline", namespace="source", @@ -260,23 +237,64 @@ def test_reconcile_applies_deployments_and_reports_status_through_a_plain_dict( ) assert core.patch_namespaced_config_map.call_count == 1 - assert [call.kwargs["name"] for call in apps.patch_namespaced_deployment.call_args_list] == [ - "consumer", - "consumer-canary", + assert [call.kwargs["name"] for call in core.patch_namespaced_pod.call_args_list] == [ + "consumer-0-0", + "consumer-canary-0-0", ] + assert result["childPods"] == ["consumer-0-0", "consumer-canary-0-0"] + assert result["desiredReplicas"] == 2 + assert result["readyReplicas"] == 0 + conditions = cast(list[dict[str, Any]], status.pop("conditions")) assert _undated(conditions) == [ {"type": "Rendered", "status": "True", "reason": "Rendered", "message": ""}, {"type": "Applied", "status": "True", "reason": "Applied", "message": ""}, ] - assert status == { - "config_version": "version", - "replicas": {"primary": 1, "canary": 1}, - "workload_namespace": WORKLOAD_NAMESPACE, + assert status["config_version"] == "version" + assert status["workload_namespace"] == WORKLOAD_NAMESPACE + assert status["pods"] == result + assert status["generations"] == { + PRIMARY_WORKLOAD_SET: {"0": 0}, + CANARY_WORKLOAD_SET: {"0": 0}, } # Both write paths hand the status to json.dumps, so no model objects # may survive into the payload: - json.dumps(conditions) + json.dumps(status) + + +def test_reconcile_nulls_out_a_workload_set_that_is_no_longer_rendered( + monkeypatch: pytest.MonkeyPatch, +) -> None: + core = _stub_render(monkeypatch) + monkeypatch.setattr( + "sentry_streams_k8s.operator.reconcile.render_pods", + lambda _consumer: { + "configmap": { + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": {"name": "pipeline-config"}, + }, + "sets": {PRIMARY_WORKLOAD_SET: _workload("consumer", 1)}, + }, + ) + status: PipelineStatusPatch = {} + + result = reconcile_pipeline( + spec=_pipeline_spec(), + name="pipeline", + namespace="source", + uid="owner-uid", + workload_namespace=WORKLOAD_NAMESPACE, + logger=MagicMock(), + status=status, + previous_generations={PRIMARY_WORKLOAD_SET: {"0": 4}, CANARY_WORKLOAD_SET: {"0": 2}}, + ) + + assert result["sets"][CANARY_WORKLOAD_SET] is None + assert status["generations"] == {PRIMARY_WORKLOAD_SET: {"0": 5}, CANARY_WORKLOAD_SET: None} + assert [call.kwargs["name"] for call in core.patch_namespaced_pod.call_args_list] == [ + "consumer-0-5" + ] def test_reconcile_records_render_failure_in_the_status_dict( @@ -284,7 +302,7 @@ def test_reconcile_records_render_failure_in_the_status_dict( ) -> None: _stub_render(monkeypatch) monkeypatch.setattr( - "sentry_streams_k8s.operator.reconcile.render", + "sentry_streams_k8s.operator.reconcile.render_pods", MagicMock(side_effect=ValueError("bad spec")), ) status: PipelineStatusPatch = {} @@ -305,73 +323,53 @@ def test_reconcile_records_render_failure_in_the_status_dict( ] -def test_unchanged_conditions_keep_their_transition_timestamp( - monkeypatch: pytest.MonkeyPatch, +@pytest.mark.parametrize( + ("previous_status", "expected_timestamp"), + [ + ("True", "2020-01-01T00:00:00+00:00"), + ("False", "2026-01-01T00:00:00+00:00"), + ], +) +def test_merge_conditions_only_preserves_unchanged_transition_timestamp( + previous_status: str, + expected_timestamp: str, ) -> None: - _stub_render(monkeypatch) - published = [ - { - "type": "Rendered", - "status": "True", - "reason": "Rendered", - "message": "", - "lastTransitionTime": "2020-01-01T00:00:00+00:00", - } - ] - status: PipelineStatusPatch = {} - - reconcile_pipeline( - spec=_pipeline_spec(), - name="pipeline", - namespace="source", - uid="owner-uid", - workload_namespace=WORKLOAD_NAMESPACE, - logger=MagicMock(), - status=status, - previous_conditions=published, + previous = cast( + list[V1ConditionDict], + [ + { + "type": "Rendered", + "status": previous_status, + "reason": "Previous", + "message": "", + "lastTransitionTime": "2020-01-01T00:00:00+00:00", + } + ], + ) + current = cast( + list[V1ConditionDict], + [ + { + "type": "Rendered", + "status": "True", + "reason": "Rendered", + "message": "", + "lastTransitionTime": "2026-01-01T00:00:00+00:00", + } + ], ) - conditions = cast(list[dict[str, Any]], status["conditions"]) - # Rendered was already True, so its timestamp is carried over untouched; - # Applied was not published before, so it is stamped now: - assert conditions[0]["lastTransitionTime"] == "2020-01-01T00:00:00+00:00" - assert conditions[1]["type"] == "Applied" - assert conditions[1]["lastTransitionTime"] != "2020-01-01T00:00:00+00:00" - + merged = _merge_conditions(previous, current) -def test_a_flipped_condition_is_restamped(monkeypatch: pytest.MonkeyPatch) -> None: - _stub_render(monkeypatch) - published = [ + assert merged == [ { "type": "Rendered", - "status": "False", - "reason": "ValueError", - "message": "bad spec", - "lastTransitionTime": "2020-01-01T00:00:00+00:00", + "status": "True", + "reason": "Rendered", + "message": "", + "lastTransitionTime": expected_timestamp, } ] - status: PipelineStatusPatch = {} - - reconcile_pipeline( - spec=_pipeline_spec(), - name="pipeline", - namespace="source", - uid="owner-uid", - workload_namespace=WORKLOAD_NAMESPACE, - logger=MagicMock(), - status=status, - previous_conditions=published, - ) - - conditions = cast(list[dict[str, Any]], status["conditions"]) - assert conditions[0] == { - "type": "Rendered", - "status": "True", - "reason": "Rendered", - "message": "", - "lastTransitionTime": ANY, - } - assert conditions[0]["lastTransitionTime"] != "2020-01-01T00:00:00+00:00" @patch("sentry_streams_k8s.operator.operator.client.CustomObjectsApi") @@ -391,11 +389,18 @@ def test_patch_pipeline_status_targets_the_status_subresource(api: MagicMock) -> def _run_reconcile_once( + monkeypatch: pytest.MonkeyPatch, scheduler: ReconcileScheduler, *, uid: str = "owner-uid", stopped: FakeStopped | None = None, + published: dict[str, Any] | None = None, ) -> float | None: + monkeypatch.setattr( + operator_module, + "_get_pipeline_status", + lambda _name, _namespace: published if published is not None else {}, + ) return asyncio.run( _reconcile_once( spec=cast(kopf.Spec, _pipeline_spec()), @@ -421,7 +426,7 @@ def render_status(*, status: PipelineStatusPatch, **_: Any) -> None: monkeypatch.setattr(operator_module, "reconcile_pipeline", render_status) - timeout = _run_reconcile_once(ReconcileScheduler()) + timeout = _run_reconcile_once(monkeypatch, ReconcileScheduler()) assert timeout == float(HEALTH_SCAN_INTERVAL_SECONDS) patch_status.assert_called_once_with( @@ -443,7 +448,7 @@ def fail(*, status: PipelineStatusPatch, **_: Any) -> None: raise kopf.PermanentError("failed to render") monkeypatch.setattr(operator_module, "reconcile_pipeline", fail) - assert _run_reconcile_once(ReconcileScheduler()) is None + assert _run_reconcile_once(monkeypatch, ReconcileScheduler()) is None assert patch_status.call_count == 1 @@ -456,7 +461,7 @@ def test_reconcile_once_retries_soon_after_an_unexpected_error( operator_module, "reconcile_pipeline", MagicMock(side_effect=RuntimeError("boom")) ) - assert _run_reconcile_once(ReconcileScheduler()) == 5.0 + assert _run_reconcile_once(monkeypatch, ReconcileScheduler()) == 5.0 def test_reconcile_once_skips_the_pass_when_already_stopped( @@ -468,7 +473,7 @@ def test_reconcile_once_skips_the_pass_when_already_stopped( stopped = FakeStopped() stopped.set() - assert _run_reconcile_once(ReconcileScheduler(), stopped=stopped) is None + assert _run_reconcile_once(monkeypatch, ReconcileScheduler(), stopped=stopped) is None reconcile.assert_not_called() @@ -565,12 +570,16 @@ async def reconcile_once(**_: Any) -> float | None: assert not scheduler.notify("uid") -def test_cleanup_prunes_every_resource_and_forgets_the_pipeline( +def test_cleanup_deletes_owned_pods_and_prunes_configmaps( monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.setenv("WORKLOAD_NAMESPACE", WORKLOAD_NAMESPACE) + core = MagicMock() + delete_pods = MagicMock() prune = MagicMock() - monkeypatch.setattr(operator_module, "_prune_stale_resources", prune) + monkeypatch.setattr(operator_module.client, "CoreV1Api", lambda: core) + monkeypatch.setattr(operator_module, "delete_owned_pods", delete_pods) + monkeypatch.setattr(operator_module, "prune_stale_configmaps", prune) scheduler = ReconcileScheduler() memo = SimpleNamespace(reconcile_scheduler=scheduler) @@ -580,11 +589,45 @@ async def scenario() -> None: asyncio.run(scenario()) + delete_pods.assert_called_once_with(core, WORKLOAD_NAMESPACE, "uid", ANY) prune.assert_called_once_with( + core=core, workload_namespace=WORKLOAD_NAMESPACE, owner_uid="uid", - desired_deployments=set(), desired_configmaps=set(), logger=ANY, ) assert not scheduler.notify("uid") + + +def test_pod_event_wakes_for_unhealthy_updates_and_ignores_healthy_updates() -> None: + scheduler = ReconcileScheduler() + memo = SimpleNamespace(reconcile_scheduler=scheduler) + + async def scenario() -> tuple[bool, bool]: + event = scheduler.register("uid") + await handle_pipeline_pod_event( + type="MODIFIED", + body=kopf.Body({"metadata": {"name": "consumer-0-0"}, "status": {"phase": "Running"}}), + meta=kopf.Meta({}), + labels={OWNER_UID_LABEL: "uid"}, + name="consumer-0-0", + namespace=WORKLOAD_NAMESPACE, + memo=memo, + logger=MagicMock(), + ) + healthy_woke = event.is_set() + + await handle_pipeline_pod_event( + type="MODIFIED", + body=kopf.Body({"metadata": {"name": "consumer-0-0"}, "status": {"phase": "Failed"}}), + meta=kopf.Meta({}), + labels={OWNER_UID_LABEL: "uid"}, + name="consumer-0-0", + namespace=WORKLOAD_NAMESPACE, + memo=memo, + logger=MagicMock(), + ) + return healthy_woke, event.is_set() + + assert asyncio.run(scenario()) == (False, True) diff --git a/sentry_streams_k8s/tests/test_pods.py b/sentry_streams_k8s/tests/test_pods.py index a602c630..aa96494e 100644 --- a/sentry_streams_k8s/tests/test_pods.py +++ b/sentry_streams_k8s/tests/test_pods.py @@ -1,13 +1,15 @@ from __future__ import annotations +import copy +import logging from datetime import datetime, timezone from typing import Any -from unittest.mock import MagicMock import pytest from kubernetes.client import V1Pod from sentry_streams_k8s.constants import CANARY_WORKLOAD_SET, PRIMARY_WORKLOAD_SET +from sentry_streams_k8s.consumer_builder import WorkloadSet from sentry_streams_k8s.k8s_types import V1PodDict from sentry_streams_k8s.operator.constants import ( GENERATION_LABEL, @@ -31,10 +33,20 @@ pod_ordinal, pod_workload_set, ) -from tests.k8s_fixtures import make_condition, make_pod +from sentry_streams_k8s.operator.reconcile import ( + PodSetResult, + delete_obsolete_pod_sets, + reconcile_pipeline_pods, +) +from tests.k8s_fixtures import ( + FakeCoreV1Api, + make_condition, + make_pod, +) NAMESPACE = "workloads" OWNER_UID = "owner-uid" +LOGGER = logging.getLogger(__name__) def _template() -> tuple[dict[str, Any], dict[str, Any]]: @@ -54,6 +66,7 @@ def _pod( deletion_timestamp: datetime | None = None, ) -> V1Pod: labels = { + OWNER_UID_LABEL: OWNER_UID, ORDINAL_LABEL: str(ordinal), GENERATION_LABEL: str(generation), WORKLOAD_SET_LABEL: PRIMARY_WORKLOAD_SET, @@ -173,50 +186,32 @@ def test_list_owned_pods_selects_owner_and_optional_workload_set( workload_set: str | None, label_selector: str, ) -> None: - core = MagicMock() - core.list_namespaced_pod.return_value.items = [] + core = FakeCoreV1Api() assert list_owned_pods(core, NAMESPACE, OWNER_UID, workload_set) == [] - core.list_namespaced_pod.assert_called_once_with( - namespace=NAMESPACE, - label_selector=label_selector, - ) + assert core.pod_list_calls == [(NAMESPACE, label_selector)] def test_delete_pod_only_drops_the_grace_period_when_forced() -> None: - core = MagicMock() + core = FakeCoreV1Api() delete_pod(core, "consumer-0-0", NAMESPACE) - - core.delete_namespaced_pod.assert_called_once_with( - name="consumer-0-0", - namespace=NAMESPACE, - ) - - core.reset_mock() delete_pod(core, "consumer-0-0", NAMESPACE, force=True) - kwargs = core.delete_namespaced_pod.call_args.kwargs - assert kwargs["name"] == "consumer-0-0" - assert kwargs["namespace"] == NAMESPACE - assert kwargs["body"].grace_period_seconds == 0 + assert core.deleted_pods == [ + ("consumer-0-0", False), + ("consumer-0-0", True), + ] -def test_delete_owned_pods_skips_pods_already_terminating(monkeypatch: pytest.MonkeyPatch) -> None: +def test_delete_owned_pods_skips_pods_already_terminating() -> None: pods = [_pod(0, 0), _pod(1, 0, deletion_timestamp=datetime(2026, 7, 16, tzinfo=timezone.utc))] - deleted: list[str] = [] - monkeypatch.setattr( - "sentry_streams_k8s.operator.pod_resources.list_owned_pods", lambda *_: pods - ) - monkeypatch.setattr( - "sentry_streams_k8s.operator.pod_resources.delete_pod", - lambda _core, name, _namespace: deleted.append(name), - ) + core = FakeCoreV1Api(pods=pods) - delete_owned_pods(MagicMock(), NAMESPACE, OWNER_UID, MagicMock()) + delete_owned_pods(core, NAMESPACE, OWNER_UID, LOGGER) - assert deleted == ["consumer-0-0"] + assert core.deleted_pods == [("consumer-0-0", False)] def test_metadata_readers_parse_operator_labels() -> None: @@ -226,10 +221,6 @@ def test_metadata_readers_parse_operator_labels() -> None: assert pod_generation(labelled) == 7 assert pod_workload_set(labelled) == PRIMARY_WORKLOAD_SET - garbage = make_pod(labels={ORDINAL_LABEL: "one", GENERATION_LABEL: "latest"}) - assert pod_ordinal(garbage) is None - assert pod_generation(garbage) == 0 - def test_pod_keep_key_prefers_ready_then_the_newest_generation() -> None: ready_old = _pod(0, 1, ready=True) @@ -245,3 +236,193 @@ def test_pod_keep_key_prefers_ready_then_the_newest_generation() -> None: assert pod_keep_key(newer, PodHealth(name=pod_name(newer))) > pod_keep_key( older, PodHealth(name=pod_name(older)) ) + + +def _manifest( + ordinal: int, + generation: int, + *, + base_name: str = "consumer", + workload_set: str = PRIMARY_WORKLOAD_SET, + template_metadata: dict[str, Any] | None = None, +) -> V1PodDict: + metadata, spec = _template() + return build_pipeline_pod( + base_name=base_name, + template_metadata=template_metadata or metadata, + template_spec=spec, + ordinal=ordinal, + generation=generation, + owner_uid=OWNER_UID, + owner_name="pipeline", + owner_namespace="source", + workload_set=workload_set, + ) + + +def _observed( + manifest: V1PodDict, + *, + phase: str = "Pending", + ready: bool = False, + container_statuses: list[Any] | None = None, + start_time: datetime | None = None, + deletion_timestamp: datetime | None = None, +) -> V1Pod: + metadata = manifest["metadata"] + return make_pod( + name=metadata["name"], + labels=metadata.get("labels"), + annotations=metadata.get("annotations"), + phase=phase, + conditions=[make_condition("Ready", "True")] if ready else [], + container_statuses=container_statuses, + start_time=start_time, + deletion_timestamp=deletion_timestamp, + ) + + +def _workload(name: str = "consumer", replicas: int = 1) -> WorkloadSet: + metadata, spec = _template() + return WorkloadSet( + name=name, + replicas=replicas, + labels=dict(metadata["labels"]), + pod_template={"metadata": metadata, "spec": spec}, + ) + + +def _reconcile( + pods: list[V1Pod], + *, + replicas: int = 1, + generations: dict[int, int] | None = None, + workload_set: str = PRIMARY_WORKLOAD_SET, +) -> tuple[list[V1PodDict], list[tuple[str, bool]], PodSetResult, dict[int, int]]: + core = FakeCoreV1Api(pods=pods) + metadata, spec = _template() + ledger = generations if generations is not None else {} + result = reconcile_pipeline_pods( + core=core, + workload_namespace=NAMESPACE, + owner_uid=OWNER_UID, + owner_name="pipeline", + owner_namespace="source", + base_name="consumer", + template_metadata=metadata, + template_spec=spec, + replicas=replicas, + generations=ledger, + logger=LOGGER, + workload_set=workload_set, + ) + return core.applied_pods, core.deleted_pods, result, ledger + + +def test_reconcile_creates_all_missing_ordinals() -> None: + applied, deleted, result, ledger = _reconcile([], replicas=2) + + assert [pod["metadata"]["name"] for pod in applied] == ["consumer-0-0", "consumer-1-0"] + assert deleted == [] + assert ledger == {0: 0, 1: 0} + assert result["childPods"] == ["consumer-0-0", "consumer-1-0"] + assert result["desiredReplicas"] == 2 + assert result["readyReplicas"] == 0 + + current = [ + _observed(_manifest(ordinal, 0), phase="Running", ready=True) for ordinal in range(2) + ] + applied, deleted, result, ledger = _reconcile(current, replicas=2) + + assert applied == [] + assert deleted == [] + assert ledger == {0: 0, 1: 0} + assert result["readyReplicas"] == 2 + + +def test_reconcile_scales_down() -> None: + stale = _pod(2, 0) + + applied, deleted, result, _ledger = _reconcile([stale], replicas=0) + + assert applied == [] + assert deleted == [("consumer-2-0", False)] + assert result["desiredReplicas"] == 0 + + +def test_reconcile_replaces_a_failed_pod_with_the_next_generation() -> None: + current = _pod(0, 3, phase="Failed") + + applied, deleted, result, ledger = _reconcile([current], generations={0: 7}) + + assert [pod["metadata"]["name"] for pod in applied] == ["consumer-0-8"] + assert deleted == [("consumer-0-3", False)] + assert ledger == {0: 8} + assert result["unhealthyPods"] == [ + { + "name": "consumer-0-3", + "ready": False, + "phase": "Failed", + "ordinal": "0", + "reason": "Failed", + } + ] + + +def test_reconcile_applies_the_replacement_before_deleting_outdated() -> None: + outdated = _observed(_manifest(0, 1), phase="Running", ready=True) + assert outdated.metadata is not None + outdated.metadata.annotations[SPEC_HASH_ANNOTATION] = "old" + core = FakeCoreV1Api(pods=[outdated]) + metadata, spec = _template() + + reconcile_pipeline_pods( + core=core, + workload_namespace=NAMESPACE, + owner_uid=OWNER_UID, + owner_name="pipeline", + owner_namespace="source", + base_name="consumer", + template_metadata=metadata, + template_spec=spec, + replicas=1, + generations={}, + logger=LOGGER, + workload_set=PRIMARY_WORKLOAD_SET, + ) + + assert core.operations == ["apply:consumer-0-2", "delete:consumer-0-1"] + + +def test_reconcile_keeps_the_best_duplicate_and_deletes_the_rest() -> None: + outdated = _pod(0, 1, ready=True, phase="Running", annotations={SPEC_HASH_ANNOTATION: "old"}) + duplicate_manifest = copy.deepcopy(_manifest(0, 0)) + duplicate_manifest["metadata"]["name"] = "consumer-0-2" + duplicate_manifest["metadata"]["labels"][GENERATION_LABEL] = "2" + duplicate = _observed(duplicate_manifest, phase="Running", ready=True) + + applied, deleted, result, ledger = _reconcile([outdated, duplicate]) + + assert applied == [] + assert deleted == [("consumer-0-1", False)] + assert result["childPods"] == ["consumer-0-2"] + assert ledger == {0: 2} + + +def test_delete_obsolete_pod_sets_removes_canary_when_disabled() -> None: + primary = _pod(0, 1, ready=True, phase="Running") + canary = _pod(0, 2, ready=True, phase="Running") + assert canary.metadata is not None + canary.metadata.name = "consumer-canary-0-2" + canary.metadata.labels[WORKLOAD_SET_LABEL] = CANARY_WORKLOAD_SET + core = FakeCoreV1Api(pods=[primary, canary]) + + delete_obsolete_pod_sets( + core, + NAMESPACE, + OWNER_UID, + {PRIMARY_WORKLOAD_SET}, + LOGGER, + ) + + assert core.deleted_pods == [("consumer-canary-0-2", False)] diff --git a/sentry_streams_k8s/tests/test_streaming_pipeline.py b/sentry_streams_k8s/tests/test_streaming_pipeline.py index cca4ddbf..8dddeebc 100644 --- a/sentry_streams_k8s/tests/test_streaming_pipeline.py +++ b/sentry_streams_k8s/tests/test_streaming_pipeline.py @@ -12,7 +12,7 @@ REQUIRED_FIELDS, StreamingPipelineSpec, from_crd_spec, - render, + render_deployments, validate, ) @@ -119,7 +119,7 @@ def test_validate_lists_missing_fields() -> None: def test_render_produces_manifests_from_inputs() -> None: - result = render(consumer_spec()) + result = render_deployments(consumer_spec()) deployment = result["deployment"] configmap = result["configmap"] @@ -150,12 +150,12 @@ def test_render_produces_manifests_from_inputs() -> None: def test_render_passes_envvar_tokens_through_config() -> None: config = pipeline_config() config["metrics"] = {"type": "datadog", "host": "${envvar:HOST_IP}", "port": 8128} - result = render(consumer_spec(pipeline_config=config)) + result = render_deployments(consumer_spec(pipeline_config=config)) assert "${envvar:HOST_IP}" in result["configmap"]["data"]["pipeline_config.yaml"] def test_render_canary_split() -> None: - result = render(consumer_spec(with_canary=True, replicas=4)) + result = render_deployments(consumer_spec(with_canary=True, replicas=4)) assert result["deployment"]["spec"]["replicas"] == 3 assert result["canary_deployment"]["spec"]["replicas"] == 1 assert result["deployment"]["spec"]["selector"]["matchLabels"]["env"] == "primary" @@ -166,7 +166,7 @@ def test_render_rejects_liveness_probe_conflict() -> None: template = container_template() template["livenessProbe"] = {"exec": {"command": ["true"]}} with pytest.raises(ValueError, match="livenessProbe"): - render(consumer_spec(container_template=template)) + render_deployments(consumer_spec(container_template=template)) def test_render_matches_macro_adapter() -> None: @@ -190,7 +190,7 @@ def test_render_matches_macro_adapter() -> None: "replicas": spec["replicas"], } ) - operator_result = render(spec) + operator_result = render_deployments(spec) assert operator_result == macro_result