diff --git a/docs/operator/.nav.yml b/docs/operator/.nav.yml index 304cf4f2f4..186114da5d 100644 --- a/docs/operator/.nav.yml +++ b/docs/operator/.nav.yml @@ -3,6 +3,7 @@ nav: - Deployment Verification: deployment-verification.md - Health Checks: health-checks.md - Scheduled Health Checks: scheduled-health-checks.md + - Triggered Health Checks: triggered-health-checks.md - Alert Destinations: destinations.md - Configuration: configuration.md - Development Guide: development.md diff --git a/docs/operator/deployment-verification.md b/docs/operator/deployment-verification.md index eb03d3ce0b..c4cbeee7ee 100644 --- a/docs/operator/deployment-verification.md +++ b/docs/operator/deployment-verification.md @@ -1,10 +1,75 @@ # Deployment Verification -A common pattern is deploying a HealthCheck alongside your application to verify the new version is working correctly. Since HealthChecks run immediately when created, you can include one in the same manifest (or CI/CD step) as your deployment and use the result to gate rollout progression. +Verifying that a new version is healthy right after it ships is one of the most common +uses of the operator. There are two ways to do it, and they compose: -## One-Time Verification with HealthCheck +- **Automatically, on every rollout** — declare one [TriggeredHealthCheck](triggered-health-checks.md) + and Holmes investigates *every* future rollout of the service, no matter how it was + triggered (CI, GitOps/Argo sync, `kubectl set image`, or a rollback). This is the + recommended default — declare once, no per-deploy wiring. +- **Inline, to gate a pipeline** — include a one-time [HealthCheck](health-checks.md) in + the deploy manifest (or CI/CD step) and block the pipeline on its result. Use this when + CI must wait synchronously for the verdict before proceeding. -Include a [HealthCheck](health-checks.md) in the same manifest as your deployment. It runs immediately after `kubectl apply` and reports whether the new version started correctly. +A typical setup uses both: a `TriggeredHealthCheck` for hands-off coverage of all +rollouts, plus an inline `HealthCheck` in the specific pipeline stage where you want a +hard gate. + +## Automatic verification with TriggeredHealthCheck + +Apply this once. From then on, any rollout of a Deployment matching the selector +automatically spawns a check — including deploys you didn't make through CI. + +```yaml +# verify-checkout-deploys.yaml — apply once, verifies every future rollout +apiVersion: holmesgpt.dev/v1alpha1 +kind: TriggeredHealthCheck +metadata: + name: verify-checkout-deploys + namespace: production +spec: + deploymentRollout: + selector: + matchLabels: + app: checkout-api + delaySeconds: 300 # wait 5m after the rollout, then check (default) + query: | + checkout-api was just rolled out to {{ .new.image }} (previously {{ .old.image }}). + Is the new version healthy? Compare error rates, latency, restarts, and logs + before vs after the rollout and flag any regressions. + timeout: 120 + mode: alert + destinations: + - type: slack + config: + channel: "#deploy-alerts" +``` + +```bash +kubectl apply -f verify-checkout-deploys.yaml + +# After a deploy, see the check it produced and the verdict +kubectl get hc -n production -l holmesgpt.dev/triggered-by=verify-checkout-deploys +kubectl describe thc verify-checkout-deploys -n production +``` + +The `{{ .new.image }}` / `{{ .old.image }}` tokens are substituted with the rollout's +before/after images, so the investigation knows exactly what changed. See +[Triggered Health Checks](triggered-health-checks.md) for the full field reference and +[how long the check waits](triggered-health-checks.md#how-long-to-wait) after a rollout. + +!!! tip "Catch slow-burn regressions too" + + Some problems (memory leaks, connection-pool exhaustion) only appear after the new + version has run for a while. Add a second trigger with a delay — e.g. + `delaySeconds: 86400` — to re-investigate the same rollout a day later, or use a + [ScheduledHealthCheck](scheduled-health-checks.md) for continuous coverage. + +## Gating CI/CD with an inline HealthCheck + +When a pipeline must **wait for the verdict** before promoting a release, include a +one-time `HealthCheck` in the same manifest as your deployment. It runs immediately after +`kubectl apply` and reports whether the new version started correctly. ```yaml # app-deployment.yaml @@ -45,17 +110,11 @@ spec: channel: "#deploy-alerts" ``` -Apply both together: - ```bash kubectl apply -f app-deployment.yaml ``` -If pods crash or fail readiness, the check fails and alerts your team. - -## Gating CI/CD on the Result - -After applying, poll for the result to gate your pipeline: +Then poll for the result to gate the pipeline: ```bash # Wait for the check to complete, then read the result @@ -75,13 +134,27 @@ echo "Timed out waiting for health check" exit 1 ``` -## Ongoing Monitoring with ScheduledHealthCheck +If pods crash or fail readiness, the check fails and the pipeline stops. + +## When to use which -One-time deploy checks catch immediate failures, but some problems only appear later — memory leaks, connection pool exhaustion, gradual performance degradation. [Scheduled Health Checks](scheduled-health-checks.md) run on a cron schedule to catch these regressions automatically. +| | TriggeredHealthCheck | Inline HealthCheck | +|---|---|---| +| Runs on | *Every* rollout, automatically | Only when you apply it | +| Setup | Declare once per service | Added to each deploy manifest/step | +| Covers out-of-band deploys (`kubectl set image`, GitOps, rollback) | Yes | No | +| Blocks a CI/CD pipeline | No (fire-and-forget) | Yes (poll the result to gate) | -## Tips for One-Time HealthChecks +## Tips -- **Version the check name** (e.g., `checkout-api-deploy-v2-4-1`) so each deploy creates a distinct resource and you keep an audit trail. This applies to one-time `HealthCheck` resources only — `ScheduledHealthCheck` resources use a fixed name and create child HealthChecks automatically. -- **Set a longer timeout** (60–120s) to give the rollout time to complete before Holmes evaluates. -- **Use labels** like `deploy-version` to query checks for a specific release: `kubectl get hc -l deploy-version=v2.4.1`. -- **Combine with ArgoCD**: If you use ArgoCD, the query can reference sync status — e.g., *"Is the ArgoCD application 'checkout-api' synced and healthy with no degraded resources?"* — since Holmes has access to the [ArgoCD toolset](../data-sources/builtin-toolsets/argocd.md). +- **Version the inline check name** (e.g., `checkout-api-deploy-v2-4-1`) so each deploy + creates a distinct resource and you keep an audit trail. This applies to one-time + `HealthCheck` resources only — `TriggeredHealthCheck` and `ScheduledHealthCheck` use a + fixed name and create child HealthChecks automatically. +- **Set a longer timeout** (60–120s) to give the investigation time to gather data. +- **Use labels** like `deploy-version` to query checks for a specific release: + `kubectl get hc -l deploy-version=v2.4.1`. +- **Combine with ArgoCD**: the query can reference sync status — e.g., *"Is the ArgoCD + application 'checkout-api' synced and healthy with no degraded resources?"* — since + Holmes has access to the [ArgoCD toolset](../data-sources/builtin-toolsets/argocd.md). + diff --git a/docs/operator/index.md b/docs/operator/index.md index 8b2405404d..e2e501ff5f 100644 --- a/docs/operator/index.md +++ b/docs/operator/index.md @@ -21,6 +21,7 @@ Under the hood, it uses Kubernetes CRDs to declaratively define one-time and sch - **[Deployment Verification](deployment-verification.md)**: Deploy a HealthCheck alongside your app to verify the new version is healthy — and gate CI/CD on the result - **[One-time Health Checks](health-checks.md)**: Create `HealthCheck` resources that run immediately and report results - **[Scheduled Health Checks](scheduled-health-checks.md)**: Create `ScheduledHealthCheck` resources that run on cron schedules for continuous monitoring +- **[Triggered Health Checks](triggered-health-checks.md)**: Create `TriggeredHealthCheck` resources that run automatically whenever a matching Deployment is rolled out - **Not just Kubernetes**: Health checks can query any connected data source — Prometheus, Datadog, AWS, databases, and [more](../data-sources/builtin-toolsets/index.md) - **Kubernetes-native**: Uses standard CRDs with kubectl support - **Status Tracking**: Full execution history and results stored in resource status @@ -154,6 +155,7 @@ kubectl describe hc example-check - **[Deployment Verification](deployment-verification.md)** - Verify new deploys are healthy and gate CI/CD pipelines on the result - **[Health Checks](health-checks.md)** - Learn how to create and manage one-time HealthCheck resources - **[Scheduled Health Checks](scheduled-health-checks.md)** - Set up recurring health checks with cron schedules +- **[Triggered Health Checks](triggered-health-checks.md)** - Run checks automatically on every Deployment rollout - **[Alert Destinations](destinations.md)** - Configure Slack and PagerDuty notifications - **[Configuration](configuration.md)** - Explore advanced configuration options - **[Development Guide](development.md)** - Build and test operator changes locally diff --git a/docs/operator/triggered-health-checks.md b/docs/operator/triggered-health-checks.md new file mode 100644 index 0000000000..466aa5521b --- /dev/null +++ b/docs/operator/triggered-health-checks.md @@ -0,0 +1,152 @@ +# Triggered Health Checks + +A `TriggeredHealthCheck` runs an investigation **automatically when a Deployment rolls +out a new version** — no per-deploy wiring, no CI polling. Declare it once, and every +rollout of a matching Deployment (from CI, Argo, `kubectl set image`, or a rollback) +fires a check. + +It is the event-driven sibling of the [ScheduledHealthCheck](scheduled-health-checks.md): +both are self-contained (they embed the check definition inline) and both spawn a +[HealthCheck](health-checks.md) per run, which becomes the execution record. + +!!! info "Alpha" + + `TriggeredHealthCheck` currently supports a single trigger type — `deploymentRollout`. + More event sources (pod crashloops, failed Jobs, alerts) are planned. + +## How it works + +1. The operator watches Deployments in namespaces where `TriggeredHealthCheck` + resources exist. +2. When a matching Deployment's **pod template changes** (a rollout), the operator waits + `delaySeconds` (default 5 minutes) and then runs the check. The wait gives the rollout + time to finish and gives any crashes or errors time to show up. +3. It creates a `HealthCheck` (owned by the trigger) with your query, having + substituted the rollout context into it. +4. Holmes investigates using every connected data source; in `alert` mode it notifies + your [destinations](destinations.md) on failure. + +## Example + +```yaml +apiVersion: holmesgpt.dev/v1alpha1 +kind: TriggeredHealthCheck +metadata: + name: verify-checkout-rollouts + namespace: production +spec: + deploymentRollout: + selector: + matchLabels: + app: checkout-api + delaySeconds: 300 # wait 5m after the rollout, then check (default) + cooldownSeconds: 600 # don't re-fire for the same Deployment within 10m + query: | + checkout-api was rolled out to {{ .new.image }} (was {{ .old.image }}). + Compare error rates, latency, restarts, and logs before vs after the rollout + and flag any regressions. + timeout: 120 + mode: alert + destinations: + - type: slack + config: + channel: "#deploy-alerts" +``` + +Apply it once: + +```bash +kubectl apply -f triggeredhealthcheck.yaml + +# List triggers (short name: thc) +kubectl get thc + +# See fire history and the HealthChecks each rollout produced +kubectl describe thc verify-checkout-rollouts +kubectl get hc -l holmesgpt.dev/triggered-by=verify-checkout-rollouts +``` + +## Query and rollout context + +The rollout facts are **always** prepended to the query the check runs, so even a terse +query like `"Is the new version healthy?"` gives the model what it needs: + +``` +This health check was triggered automatically by a Kubernetes Deployment rollout. Use this context when investigating: +- Deployment: checkout-api +- Namespace: production +- Previous image(s): myregistry/checkout-api:v2.4.0 +- New image(s): myregistry/checkout-api:v2.4.1 + + +``` + +You can also reference the same facts inline in your `query` with these tokens: + +| Token | Replaced with | +|-------|---------------| +| `{{ .deployment }}` | Name of the Deployment that rolled out | +| `{{ .namespace }}` | Its namespace | +| `{{ .old.image }}` | Container image(s) before the rollout | +| `{{ .new.image }}` | Container image(s) after the rollout | + +## Spec reference + +| Field | Default | Description | +|-------|---------|-------------| +| `enabled` | `true` | Whether the trigger is active | +| `deploymentRollout.selector.matchLabels` | `{}` | Deployment labels that must all match. **Empty matches every Deployment in the namespace.** | +| `delaySeconds` | `300` | How long to wait after a rollout before running the check. Default 5 minutes; `0` checks immediately; `86400` checks a day later (max 7 days). See [How long to wait](#how-long-to-wait). | +| `cooldownSeconds` | `0` | Suppress re-firing for the same Deployment within this window. `0` disables. | +| `query` | — | Natural-language investigation (supports the tokens above). Required. | +| `timeout` | `120` | Check execution timeout in seconds. | +| `mode` | `monitor` | `alert` notifies destinations on failure; `monitor` only records the result. | +| `model` | — | Override the default LLM model for this check. | +| `destinations` | `[]` | Alert destinations (used in `alert` mode). See [Destinations](destinations.md). | + +## How long to wait + +There is one knob: **`delaySeconds`** — how long to wait after a rollout before running +the check. That's it. + +The check does not run the instant a new version is deployed, because a brand-new +rollout hasn't finished and problems haven't surfaced yet. So Holmes waits `delaySeconds` +first, then looks. Pick a value that fits what you want to catch: + +- **`300` (default, 5 minutes)** — enough time for the rollout to finish and for crashes + or startup errors to appear. +- **`0`** — check right away (useful if you only care that the deploy was accepted). +- **`86400` (a day)** — catch slower problems like memory leaks or resource creep that + only show up after the version has been running a while. + +If you set the wait too short for your app's rollout, the check may run while pods are +still starting — Holmes will simply report that the rollout hasn't finished yet. + +The wait is saved on the resource, so it still completes if the operator restarts — even +a day-long wait. If the same Deployment is rolled out again while a check is still +waiting, the waiting check is replaced so only the newest version is checked. + +A common setup is one trigger with the default 5-minute wait, plus a second with +`delaySeconds: 86400` to re-check the same rollout a day later. + +## Notes & limitations + +- **Rollout = pod-template change.** Scaling and HPA changes (which only touch + `spec.replicas`) do **not** fire the trigger; only changes to the pod template do. +- **Restart behavior.** A check that was *already scheduled* before a restart still runs + (the pending fire is persisted in `status.pending`). *Detecting* new rollouts, however, + relies on an in-memory baseline of each Deployment's last-seen pod template (the operator + does not annotate your Deployments). After a restart the first observation of each + Deployment just re-establishes that baseline, so a rollout that happens *during* the + restart window is not detected. Use a + [ScheduledHealthCheck](scheduled-health-checks.md) for continuous coverage. +- **Cost.** Every fire is at least one LLM call. Use `cooldownSeconds` and a specific + `selector` to bound spend on busy namespaces. + +## Next Steps + +- **[Deployment Verification](deployment-verification.md)** — patterns for gating and + verifying deploys +- **[Health Checks](health-checks.md)** — the one-time checks this spawns +- **[Alert Destinations](destinations.md)** — Slack and PagerDuty configuration + diff --git a/helm/holmes/crds/triggeredhealthcheck.yaml b/helm/holmes/crds/triggeredhealthcheck.yaml new file mode 100644 index 0000000000..df36138309 --- /dev/null +++ b/helm/holmes/crds/triggeredhealthcheck.yaml @@ -0,0 +1,205 @@ +apiVersion: apiextensions.k8s.io/v1 +kind: CustomResourceDefinition +metadata: + name: triggeredhealthchecks.holmesgpt.dev +spec: + group: holmesgpt.dev + names: + kind: TriggeredHealthCheck + listKind: TriggeredHealthCheckList + plural: triggeredhealthchecks + singular: triggeredhealthcheck + shortNames: + - thc + scope: Namespaced + versions: + - name: v1alpha1 + served: true + storage: true + subresources: + status: {} + schema: + openAPIV3Schema: + type: object + properties: + spec: + type: object + required: + - deploymentRollout + - query + properties: + enabled: + type: boolean + description: "Whether the trigger is active" + default: true + deploymentRollout: + type: object + description: "Fire when a matching Deployment rolls out a new pod template" + properties: + selector: + type: object + description: "Select which Deployments to watch (empty matches all in the namespace)" + properties: + matchLabels: + type: object + description: "Deployment labels that must all match" + additionalProperties: + type: string + delaySeconds: + type: integer + description: "How long to wait after a rollout before running the check (seconds). Gives the rollout time to finish and problems time to surface. Default 300 (5 min); 0 = check immediately; up to 604800 (7 days). The wait is saved on the resource and survives operator restarts." + default: 300 + minimum: 0 + maximum: 604800 + cooldownSeconds: + type: integer + description: "Suppress re-firing for the same Deployment within this many seconds (0 = disabled)" + default: 0 + minimum: 0 + query: + type: string + description: "Natural language question. Supports tokens {{ .deployment }}, {{ .namespace }}, {{ .old.image }}, {{ .new.image }}" + minLength: 1 + maxLength: 5000 + timeout: + type: integer + description: "Execution timeout in seconds" + default: 120 + minimum: 1 + maximum: 300 + mode: + type: string + description: "Execution mode: 'alert' sends notifications on failure, 'monitor' logs only" + default: monitor + enum: + - alert + - monitor + model: + type: string + description: "Override default LLM model for this check" + destinations: + type: array + description: "Alert destinations for failed checks (only used in alert mode)" + items: + type: object + required: + - type + properties: + type: + type: string + description: "Destination type (e.g., 'slack', 'pagerduty')" + config: + type: object + description: "Destination-specific configuration" + x-kubernetes-preserve-unknown-fields: true + status: + type: object + properties: + lastTriggerTime: + type: string + format: date-time + description: "Last time the trigger fired" + lastTriggerDeployment: + type: string + description: "Deployment whose rollout last fired the trigger" + triggerCount: + type: integer + description: "Total number of times the trigger has fired" + cooldowns: + type: array + description: "Per-Deployment last-fire times used to enforce cooldownSeconds" + items: + type: object + required: + - deployment + - lastTriggerTime + properties: + deployment: + type: string + lastTriggerTime: + type: string + format: date-time + pending: + type: array + description: "Delayed checks scheduled to run later (delaySeconds), persisted to survive operator restarts" + items: + type: object + required: + - deployment + - fireAt + properties: + deployment: + type: string + fireAt: + type: string + format: date-time + scheduledAt: + type: string + format: date-time + oldImage: + type: string + newImage: + type: string + history: + type: array + description: "Recent triggered executions (last N)" + items: + type: object + required: + - triggerTime + - deployment + - checkName + properties: + triggerTime: + type: string + format: date-time + deployment: + type: string + checkName: + type: string + oldImage: + type: string + newImage: + type: string + conditions: + type: array + description: "Standard Kubernetes conditions" + items: + type: object + required: + - type + - status + properties: + type: + type: string + description: "Condition type (e.g., 'Ready', 'TriggerFailed')" + status: + type: string + description: "Condition status: True, False, or Unknown" + enum: + - "True" + - "False" + - Unknown + lastTransitionTime: + type: string + format: date-time + reason: + type: string + message: + type: string + additionalPrinterColumns: + - name: Enabled + type: boolean + description: Whether the trigger is enabled + jsonPath: .spec.enabled + - name: Triggers + type: integer + description: Number of times fired + jsonPath: .status.triggerCount + - name: Last Fired + type: date + description: Last time the trigger fired + jsonPath: .status.lastTriggerTime + - name: Age + type: date + jsonPath: .metadata.creationTimestamp diff --git a/helm/holmes/templates/operator-rbac.yaml b/helm/holmes/templates/operator-rbac.yaml index 4660114804..a39086fd18 100644 --- a/helm/holmes/templates/operator-rbac.yaml +++ b/helm/holmes/templates/operator-rbac.yaml @@ -46,6 +46,19 @@ rules: resources: ["scheduledhealthchecks/status"] verbs: ["get", "patch", "update"] +# TriggeredHealthCheck CRD permissions +- apiGroups: ["holmesgpt.dev"] + resources: ["triggeredhealthchecks"] + verbs: ["get", "list", "watch", "patch", "update"] +- apiGroups: ["holmesgpt.dev"] + resources: ["triggeredhealthchecks/status"] + verbs: ["get", "patch", "update"] + +# Watch Deployments to fire rollout triggers +- apiGroups: ["apps"] + resources: ["deployments"] + verbs: ["get", "list", "watch"] + # Events for audit trail - apiGroups: [""] resources: ["events"] diff --git a/holmes_operator/handlers/triggeredhealthcheck.py b/holmes_operator/handlers/triggeredhealthcheck.py new file mode 100644 index 0000000000..4f85bbc55f --- /dev/null +++ b/holmes_operator/handlers/triggeredhealthcheck.py @@ -0,0 +1,306 @@ +"""Kopf handlers for the TriggeredHealthCheck CRD (deployment-rollout trigger). + +A TriggeredHealthCheck is the event-driven sibling of ScheduledHealthCheck: instead +of a cron schedule, it watches Deployments and spawns a HealthCheck when a matching +Deployment rolls out a new pod template. +""" + +import asyncio +import logging +from typing import Any, Dict + +import kopf + +from holmes_operator import context, trigger_executor +from holmes_operator.models import ( + ConditionStatus, + HealthCheckCondition, + TriggeredHealthCheckConditionType, + TriggeredHealthCheckSpec, +) +from holmes_operator.utils import get_current_time_iso + +logger = logging.getLogger(__name__) + +GROUP = "holmesgpt.dev" +VERSION = "v1alpha1" + + +@kopf.on.create(GROUP, VERSION, "triggeredhealthchecks") # type: ignore[arg-type] +async def on_triggeredhealthcheck_create( + *, + spec: Dict[str, Any], + name: str, + namespace: str, + logger: kopf.Logger, + **kwargs: Any, +) -> None: + """Validate a TriggeredHealthCheck and mark it Ready (or not).""" + logger.info(f"Creating TriggeredHealthCheck: {namespace}/{name}") + await _validate_and_set_ready(spec, name, namespace, logger) + + +@kopf.on.update(GROUP, VERSION, "triggeredhealthchecks") # type: ignore[arg-type] +async def on_triggeredhealthcheck_update( + *, + new: Dict[str, Any], + name: str, + namespace: str, + logger: kopf.Logger, + **kwargs: Any, +) -> None: + """Re-validate on spec changes and refresh the Ready condition.""" + logger.info(f"Updating TriggeredHealthCheck: {namespace}/{name}") + await _validate_and_set_ready(new.get("spec", {}), name, namespace, logger) + + +async def _validate_and_set_ready( + spec: Dict[str, Any], name: str, namespace: str, logger: kopf.Logger +) -> None: + try: + parsed = TriggeredHealthCheckSpec(**spec) + except Exception as e: + logger.error(f"Invalid TriggeredHealthCheck {namespace}/{name}: {e}") + await set_triggeredhealthcheck_condition( + name=name, + namespace=namespace, + condition_type=TriggeredHealthCheckConditionType.TRIGGER_FAILED, + status=ConditionStatus.TRUE, + reason="InvalidSpec", + message=str(e), + ) + raise + + if parsed.enabled: + reason, message = "Watching", "Watching Deployments for matching rollouts" + status = ConditionStatus.TRUE + else: + reason, message = "Disabled", "Trigger is disabled" + status = ConditionStatus.FALSE + + await set_triggeredhealthcheck_condition( + name=name, + namespace=namespace, + condition_type=TriggeredHealthCheckConditionType.READY, + status=status, + reason=reason, + message=message, + ) + + +@kopf.on.event("apps", "v1", "deployments") # type: ignore[arg-type] +async def on_deployment_event( + *, + event: Dict[str, Any], + body: Dict[str, Any], + name: str, + namespace: str, + meta: Dict[str, Any], + logger: kopf.Logger, + **kwargs: Any, +) -> None: + """Watch Deployments and fan rollouts out to matching TriggeredHealthChecks. + + Uses a low-level event handler (no kopf diff annotations/finalizers are written + to user Deployments) with an in-memory baseline cache to detect pod-template + changes. + """ + key = f"{namespace}/{name}" + + if event.get("type") == "DELETED": + trigger_executor.forget_deployment(key) + return + + rollout = trigger_executor.detect_rollout(key, body) + if rollout is None: + return # baseline observation or non-rollout change (status/scale) + + old_image, new_image = rollout + deployment_labels = meta.get("labels", {}) or {} + + try: + triggers = await asyncio.to_thread( + context.k8s_api.list_namespaced_custom_object, + group=GROUP, + version=VERSION, + namespace=namespace, + plural="triggeredhealthchecks", + ) + except Exception as e: + logger.error(f"Failed to list TriggeredHealthChecks in {namespace}: {e}") + return + + for trigger in triggers.get("items", []): + trigger_meta = trigger.get("metadata", {}) + trigger_name = trigger_meta.get("name") + trigger_uid = trigger_meta.get("uid", "") + + try: + spec = TriggeredHealthCheckSpec(**trigger.get("spec", {})) + except Exception as e: + logger.warning( + f"Skipping invalid TriggeredHealthCheck {namespace}/{trigger_name}: {e}" + ) + continue + + if not spec.enabled: + continue + if not trigger_executor.selector_matches( + spec.deploymentRollout.selector.matchLabels, deployment_labels + ): + continue + if trigger_executor.is_in_cooldown( + trigger.get("status", {}), name, spec.cooldownSeconds + ): + logger.info( + f"TriggeredHealthCheck {namespace}/{trigger_name} skipped for " + f"{namespace}/{name}: within cooldown window" + ) + continue + + if spec.delaySeconds > 0: + fire_at = trigger_executor.compute_fire_at(spec.delaySeconds) + await trigger_executor.add_pending( + api=context.k8s_api, + trigger_name=trigger_name, + namespace=namespace, + deployment=name, + fire_at=fire_at, + old_image=old_image, + new_image=new_image, + ) + logger.info( + f"Rollout of {namespace}/{name} matched TriggeredHealthCheck " + f"{namespace}/{trigger_name}; check scheduled for {fire_at} " + f"(delay {spec.delaySeconds}s)" + ) + continue + + logger.info( + f"Rollout of {namespace}/{name} matched TriggeredHealthCheck " + f"{namespace}/{trigger_name}" + ) + task = asyncio.create_task( + trigger_executor.spawn_check( + trigger_name=trigger_name, + namespace=namespace, + trigger_uid=trigger_uid, + spec=spec, + deployment=name, + old_image=old_image, + new_image=new_image, + k8s_api=context.k8s_api, + ) + ) + trigger_executor.track_task(task) + + +@kopf.on.timer(GROUP, VERSION, "triggeredhealthchecks", interval=15.0) # type: ignore[arg-type] +async def on_triggeredhealthcheck_timer( + *, + spec: Dict[str, Any], + status: Dict[str, Any], + name: str, + namespace: str, + uid: str, + logger: kopf.Logger, + **kwargs: Any, +) -> None: + """Fire delayed checks whose scheduled time has arrived. + + Pending entries are claimed (removed from status) before spawning so a check is + never run twice, even across operator restarts. + """ + pending = (status or {}).get("pending", []) + if not pending: + return + + due = trigger_executor.due_pending(pending) + if not due: + return + + try: + parsed = TriggeredHealthCheckSpec(**spec) + except Exception as e: + logger.warning(f"Skipping timer for invalid {namespace}/{name}: {e}") + return + + # Claim the due entries first so a slow spawn can't be double-processed. + await trigger_executor.remove_pending(context.k8s_api, name, namespace, due) + + for entry in due: + logger.info( + f"Delayed check for {namespace}/{entry.get('deployment')} is due; " + f"running TriggeredHealthCheck {namespace}/{name}" + ) + task = asyncio.create_task( + trigger_executor.spawn_check( + trigger_name=name, + namespace=namespace, + trigger_uid=uid, + spec=parsed, + deployment=entry.get("deployment"), + old_image=entry.get("oldImage") or "", + new_image=entry.get("newImage") or "", + k8s_api=context.k8s_api, + ) + ) + trigger_executor.track_task(task) + + +async def set_triggeredhealthcheck_condition( + name: str, + namespace: str, + condition_type: TriggeredHealthCheckConditionType, + status: ConditionStatus, + reason: str, + message: str, +) -> None: + """Add or update a condition on a TriggeredHealthCheck resource.""" + condition = HealthCheckCondition( + type=condition_type, + status=status, + lastTransitionTime=get_current_time_iso(), + reason=reason, + message=message, + ) + + try: + resource = await asyncio.to_thread( + context.k8s_api.get_namespaced_custom_object, + group=GROUP, + version=VERSION, + namespace=namespace, + plural="triggeredhealthchecks", + name=name, + ) + + conditions = resource.get("status", {}).get("conditions", []) + condition_dict = { + "type": condition.type, + "status": condition.status.value, + "lastTransitionTime": condition.lastTransitionTime, + "reason": condition.reason, + "message": condition.message, + } + + existing_idx = next( + (i for i, c in enumerate(conditions) if c.get("type") == condition.type), + None, + ) + if existing_idx is not None: + conditions[existing_idx] = condition_dict + else: + conditions.append(condition_dict) + + await asyncio.to_thread( + context.k8s_api.patch_namespaced_custom_object_status, + group=GROUP, + version=VERSION, + namespace=namespace, + plural="triggeredhealthchecks", + name=name, + body={"status": {"conditions": conditions}}, + ) + except Exception as e: + logger.error(f"Failed to set condition on {namespace}/{name}: {e}") diff --git a/holmes_operator/models.py b/holmes_operator/models.py index f5ebd05d0b..07afce3ec8 100644 --- a/holmes_operator/models.py +++ b/holmes_operator/models.py @@ -55,6 +55,13 @@ class ScheduledHealthCheckConditionType(str, Enum): EXECUTION_FAILED = "ExecutionFailed" +class TriggeredHealthCheckConditionType(str, Enum): + """TriggeredHealthCheck condition types.""" + + READY = "Ready" + TRIGGER_FAILED = "TriggerFailed" + + class HealthCheckConditionType(str, Enum): """HealthCheck condition types.""" @@ -167,6 +174,83 @@ class ScheduledHealthCheckStatus(BaseModel): conditions: List[HealthCheckCondition] = Field(default_factory=list) +class TriggerSelector(BaseModel): + """Label selector for matching resources that fire a trigger.""" + + matchLabels: dict = Field(default_factory=dict) + + +class DeploymentRolloutTrigger(BaseModel): + """Fire when a Deployment matching the selector rolls out a new pod template.""" + + selector: TriggerSelector = Field(default_factory=TriggerSelector) + + +class TriggeredHealthCheckSpec(BaseModel): + """TriggeredHealthCheck CRD spec. + + Self-contained, mirroring ScheduledHealthCheck: it embeds the check definition + inline and spawns a HealthCheck child when the trigger fires (rather than + referencing a separate HealthCheck/template). + """ + + enabled: bool = Field(default=True) + deploymentRollout: DeploymentRolloutTrigger + # How long to wait after a new version is rolled out before running the check. + # Gives the rollout time to finish and gives crashes/errors time to show up. + # Default 5 minutes; 0 checks immediately; up to 7 days (e.g. 86400 = a day later). + # The wait is saved on the resource, so it still completes if the operator restarts. + delaySeconds: int = Field(default=300, ge=0, le=604800) + # Suppress re-firing for the same Deployment within this many seconds. 0 disables. + cooldownSeconds: int = Field(default=0, ge=0) + # Inline HealthCheck definition (same fields as HealthCheckSpec) + query: str = Field(..., min_length=1, max_length=5000) + timeout: int = Field(default=120, ge=1, le=300) + mode: CheckMode = Field(default=CheckMode.MONITOR) + model: Optional[str] = None + destinations: List[DestinationConfig] = Field(default_factory=list) + + +class TriggeredDeploymentCooldown(BaseModel): + """Last fire time for a Deployment, used to enforce cooldownSeconds.""" + + deployment: str + lastTriggerTime: str + + +class TriggeredCheckHistoryEntry(BaseModel): + """History entry for a triggered check execution.""" + + triggerTime: str + deployment: str + checkName: str + oldImage: Optional[str] = None + newImage: Optional[str] = None + + +class PendingCheck(BaseModel): + """A check scheduled to run later (delaySeconds), persisted in status so it + survives operator restarts.""" + + deployment: str + fireAt: str + scheduledAt: str + oldImage: Optional[str] = None + newImage: Optional[str] = None + + +class TriggeredHealthCheckStatus(BaseModel): + """TriggeredHealthCheck CRD status.""" + + lastTriggerTime: Optional[str] = None + lastTriggerDeployment: Optional[str] = None + triggerCount: int = 0 + cooldowns: List[TriggeredDeploymentCooldown] = Field(default_factory=list) + pending: List[PendingCheck] = Field(default_factory=list) + history: List[TriggeredCheckHistoryEntry] = Field(default_factory=list) + conditions: List[HealthCheckCondition] = Field(default_factory=list) + + class CheckResponse(BaseModel): status: CheckStatus message: str diff --git a/holmes_operator/operator.py b/holmes_operator/operator.py index 084f80fff8..57ebb42469 100644 --- a/holmes_operator/operator.py +++ b/holmes_operator/operator.py @@ -12,6 +12,7 @@ # Import handlers to register them with kopf from holmes_operator.handlers import healthcheck # noqa: F401 from holmes_operator.handlers import scheduledhealthcheck # noqa: F401 +from holmes_operator.handlers import triggeredhealthcheck # noqa: F401 # Configure logging logging.basicConfig( diff --git a/holmes_operator/trigger_executor.py b/holmes_operator/trigger_executor.py new file mode 100644 index 0000000000..24e32b2c36 --- /dev/null +++ b/holmes_operator/trigger_executor.py @@ -0,0 +1,473 @@ +"""Execution logic for TriggeredHealthCheck (deployment-rollout trigger). + +Mirrors scheduler/job_executor.py: when a trigger fires, this spawns a HealthCheck +child (owned by the TriggeredHealthCheck) which goes through the normal HealthCheck +execution path, and records the trigger in the TriggeredHealthCheck status. +""" + +import asyncio +import hashlib +import json +import logging +import re +from datetime import datetime, timedelta, timezone +from typing import Dict, List, Optional, Tuple +from uuid import uuid4 + +from kubernetes import client + +from holmes_operator import context +from holmes_operator.models import TriggeredHealthCheckSpec +from holmes_operator.utils import get_current_time_iso + +logger = logging.getLogger(__name__) + +GROUP = "holmesgpt.dev" +VERSION = "v1alpha1" + +_active_tasks: set[asyncio.Task] = set() + +# In-memory cache of the last seen pod-template per Deployment, keyed by +# "namespace/name" -> (template_hash, images). Used to detect rollouts from raw +# watch events without writing kopf diff annotations onto every Deployment in the +# cluster. Lost on operator restart by design: the first event for a Deployment +# after (re)start only establishes a baseline and is not treated as a rollout. +_last_template: Dict[str, Tuple[str, str]] = {} + + +def _log_task_exception(task: asyncio.Task) -> None: + """Log any exception from a background task and drop it from the registry.""" + _active_tasks.discard(task) + if task.cancelled(): + return + exc = task.exception() + if exc is not None: + logger.error(f"Trigger background task raised an exception: {exc}", exc_info=exc) + + +def track_task(task: asyncio.Task) -> None: + """Keep a strong reference to a background task until it completes.""" + _active_tasks.add(task) + task.add_done_callback(_log_task_exception) + + +def extract_images(deployment_body: dict) -> str: + """Return a stable, human-readable representation of a Deployment's images.""" + containers = ( + deployment_body.get("spec", {}) + .get("template", {}) + .get("spec", {}) + .get("containers", []) + ) + images = [c.get("image", "") for c in containers if c.get("image")] + return ", ".join(images) + + +def _template_hash(template: dict) -> str: + return hashlib.sha256( + json.dumps(template, sort_keys=True, default=str).encode("utf-8") + ).hexdigest() + + +def detect_rollout(key: str, deployment_body: dict) -> Optional[Tuple[str, str]]: + """Detect whether a Deployment event represents a rollout (pod-template change). + + Updates the in-memory baseline cache and returns ``(old_images, new_images)`` + when the pod template changed relative to the last seen value, or ``None`` when + this is a baseline observation or a non-template change (e.g. status/scale). + """ + template = (deployment_body.get("spec") or {}).get("template") + if template is None: + return None + + new_hash = _template_hash(template) + new_images = extract_images(deployment_body) + + previous = _last_template.get(key) + _last_template[key] = (new_hash, new_images) + + if previous is None: + return None # baseline only + + prev_hash, prev_images = previous + if prev_hash == new_hash: + return None # template unchanged + + return (prev_images, new_images) + + +def forget_deployment(key: str) -> None: + """Drop a Deployment from the baseline cache (e.g. on deletion).""" + _last_template.pop(key, None) + + +def clear_rollout_cache() -> None: + """Reset the baseline cache (used in tests).""" + _last_template.clear() + + +def selector_matches(match_labels: dict, labels: dict) -> bool: + """Return True if ``labels`` contains all of ``match_labels``. + + An empty selector matches every Deployment in the namespace. + """ + if not match_labels: + return True + return all(labels.get(k) == v for k, v in match_labels.items()) + + +def render_query( + template: str, + deployment: str, + namespace: str, + old_image: str, + new_image: str, +) -> str: + """Substitute trigger context tokens into the query template. + + Supported tokens (whitespace-insensitive): ``{{ .deployment }}``, + ``{{ .namespace }}``, ``{{ .old.image }}``, ``{{ .new.image }}``. + """ + replacements = { + r"\{\{\s*\.deployment\s*\}\}": deployment, + r"\{\{\s*\.namespace\s*\}\}": namespace, + r"\{\{\s*\.old\.image\s*\}\}": old_image or "unknown", + r"\{\{\s*\.new\.image\s*\}\}": new_image or "unknown", + } + rendered = template + for pattern, value in replacements.items(): + rendered = re.sub(pattern, lambda _m, v=value: v, rendered) + return rendered + + +def compose_query( + template: str, + deployment: str, + namespace: str, + old_image: str, + new_image: str, +) -> str: + """Build the spawned check's query. + + The rollout facts are always prepended as a structured context header — so the + model knows which Deployment/namespace and what changed even if the author's + query is terse and uses none of the tokens — followed by the author's query with + any tokens substituted. + """ + header = ( + "This health check was triggered automatically by a Kubernetes Deployment " + "rollout. Use this context when investigating:\n" + f"- Deployment: {deployment}\n" + f"- Namespace: {namespace}\n" + f"- Previous image(s): {old_image or 'unknown'}\n" + f"- New image(s): {new_image or 'unknown'}\n\n" + ) + return header + render_query(template, deployment, namespace, old_image, new_image) + + +def is_in_cooldown(status: dict, deployment: str, cooldown_seconds: int) -> bool: + """Return True if ``deployment`` fired within the cooldown window.""" + if cooldown_seconds <= 0: + return False + for entry in status.get("cooldowns", []): + if entry.get("deployment") != deployment: + continue + raw = entry.get("lastTriggerTime") + if not raw: + return False + try: + last = datetime.fromisoformat(raw) + except ValueError: + return False + if last.tzinfo is None: + last = last.replace(tzinfo=timezone.utc) + return (datetime.now(timezone.utc) - last).total_seconds() < cooldown_seconds + return False + + +def compute_fire_at(delay_seconds: int) -> str: + """ISO timestamp ``delay_seconds`` from now.""" + return (datetime.now(timezone.utc) + timedelta(seconds=delay_seconds)).isoformat() + + +def due_pending(pending: List[dict]) -> List[dict]: + """Return the pending entries whose fireAt is now or in the past.""" + now = datetime.now(timezone.utc) + due = [] + for entry in pending: + raw = entry.get("fireAt") + if not raw: + continue + try: + fire_at = datetime.fromisoformat(raw) + except ValueError: + continue + if fire_at.tzinfo is None: + fire_at = fire_at.replace(tzinfo=timezone.utc) + if fire_at <= now: + due.append(entry) + return due + + +async def add_pending( + api: client.CustomObjectsApi, + trigger_name: str, + namespace: str, + deployment: str, + fire_at: str, + old_image: str, + new_image: str, +) -> None: + """Schedule a delayed check, debounced per Deployment (a newer rollout replaces + any still-pending entry for the same Deployment).""" + + def modify_status(resource: dict) -> dict: + status = resource.get("status", {}) + pending = [ + p for p in status.get("pending", []) if p.get("deployment") != deployment + ] + pending.append( + { + "deployment": deployment, + "fireAt": fire_at, + "scheduledAt": get_current_time_iso(), + "oldImage": old_image, + "newImage": new_image, + } + ) + return {"pending": pending} + + await _patch_status_with_retry(api, trigger_name, namespace, modify_status) + + +async def remove_pending( + api: client.CustomObjectsApi, + trigger_name: str, + namespace: str, + entries: List[dict], +) -> None: + """Remove the given pending entries (matched by deployment + fireAt).""" + keys = {(e.get("deployment"), e.get("fireAt")) for e in entries} + + def modify_status(resource: dict) -> dict: + status = resource.get("status", {}) + pending = [ + p + for p in status.get("pending", []) + if (p.get("deployment"), p.get("fireAt")) not in keys + ] + return {"pending": pending} + + await _patch_status_with_retry(api, trigger_name, namespace, modify_status) + + +def generate_check_name(trigger_name: str) -> str: + timestamp = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + return f"{trigger_name}-{timestamp}-{uuid4().hex[:6]}" + + +def build_healthcheck_object( + check_name: str, + namespace: str, + trigger_name: str, + trigger_uid: str, + spec: TriggeredHealthCheckSpec, + deployment: str, + old_image: str, + new_image: str, +) -> dict: + healthcheck = { + "apiVersion": f"{GROUP}/{VERSION}", + "kind": "HealthCheck", + "metadata": { + "name": check_name, + "namespace": namespace, + "labels": { + "holmesgpt.dev/triggered-by": trigger_name, + "holmesgpt.dev/trigger-type": "deployment-rollout", + "holmesgpt.dev/deployment": deployment, + }, + "ownerReferences": [ + { + "apiVersion": f"{GROUP}/{VERSION}", + "kind": "TriggeredHealthCheck", + "name": trigger_name, + "uid": trigger_uid, + "controller": True, + "blockOwnerDeletion": True, + } + ], + }, + "spec": { + "query": compose_query( + spec.query, deployment, namespace, old_image, new_image + ), + "timeout": spec.timeout, + "mode": spec.mode.value, + }, + } + + if spec.model: + healthcheck["spec"]["model"] = spec.model + if spec.destinations: + healthcheck["spec"]["destinations"] = [ + d.model_dump() for d in spec.destinations + ] + return healthcheck + + +async def spawn_check( + trigger_name: str, + namespace: str, + trigger_uid: str, + spec: TriggeredHealthCheckSpec, + deployment: str, + old_image: str, + new_image: str, + k8s_api: client.CustomObjectsApi, +) -> None: + """Create a HealthCheck for this rollout and record it in the trigger status.""" + try: + check_name = generate_check_name(trigger_name) + logger.info( + f"TriggeredHealthCheck {namespace}/{trigger_name} fired by rollout of " + f"{namespace}/{deployment}; creating HealthCheck {check_name}", + extra={ + "trigger_name": trigger_name, + "namespace": namespace, + "deployment": deployment, + "check_name": check_name, + }, + ) + + healthcheck = build_healthcheck_object( + check_name, + namespace, + trigger_name, + trigger_uid, + spec, + deployment, + old_image, + new_image, + ) + + await asyncio.to_thread( + k8s_api.create_namespaced_custom_object, + group=GROUP, + version=VERSION, + namespace=namespace, + plural="healthchecks", + body=healthcheck, + ) + + await record_trigger( + api=k8s_api, + trigger_name=trigger_name, + namespace=namespace, + deployment=deployment, + check_name=check_name, + old_image=old_image, + new_image=new_image, + ) + + except Exception as e: + logger.error( + f"Failed to execute TriggeredHealthCheck {namespace}/{trigger_name}: {e}", + exc_info=True, + ) + + +async def record_trigger( + api: client.CustomObjectsApi, + trigger_name: str, + namespace: str, + deployment: str, + check_name: str, + old_image: str, + new_image: str, +) -> None: + """Update TriggeredHealthCheck status with this trigger (history + cooldown).""" + + def modify_status(resource: dict) -> dict: + status = resource.get("status", {}) + history = status.get("history", []) + cooldowns = status.get("cooldowns", []) + now_iso = get_current_time_iso() + + # Upsert the cooldown entry for this deployment + cooldowns = [c for c in cooldowns if c.get("deployment") != deployment] + cooldowns.append({"deployment": deployment, "lastTriggerTime": now_iso}) + + history.insert( + 0, + { + "triggerTime": now_iso, + "deployment": deployment, + "checkName": check_name, + "oldImage": old_image, + "newImage": new_image, + }, + ) + max_history = context.config.max_history_items if context.config else 10 + history = history[:max_history] + + return { + "lastTriggerTime": now_iso, + "lastTriggerDeployment": deployment, + "triggerCount": (status.get("triggerCount") or 0) + 1, + "cooldowns": cooldowns, + "history": history, + } + + await _patch_status_with_retry(api, trigger_name, namespace, modify_status) + + +async def _patch_status_with_retry( + api: client.CustomObjectsApi, + trigger_name: str, + namespace: str, + modify_fn, + max_retries: int = 5, +) -> None: + """Read-modify-write TriggeredHealthCheck status with conflict retry.""" + for attempt in range(max_retries): + try: + resource = await asyncio.to_thread( + api.get_namespaced_custom_object, + group=GROUP, + version=VERSION, + namespace=namespace, + plural="triggeredhealthchecks", + name=trigger_name, + ) + + status_updates = modify_fn(resource) + resource_version = resource.get("metadata", {}).get("resourceVersion") + + await asyncio.to_thread( + api.patch_namespaced_custom_object_status, + group=GROUP, + version=VERSION, + namespace=namespace, + plural="triggeredhealthchecks", + name=trigger_name, + body={ + "metadata": {"resourceVersion": resource_version}, + "status": status_updates, + }, + ) + return + + except client.exceptions.ApiException as e: + if e.status == 409: + logger.debug( + f"Conflict updating {namespace}/{trigger_name} status " + f"(attempt {attempt + 1}/{max_retries}), retrying..." + ) + if attempt == max_retries - 1: + raise Exception( + f"Max retries ({max_retries}) exceeded for status update" + ) from e + await asyncio.sleep(0.1 * (attempt + 1)) + else: + raise diff --git a/tests/holmes_operator/test_triggeredhealthcheck_component.py b/tests/holmes_operator/test_triggeredhealthcheck_component.py new file mode 100644 index 0000000000..6f6964f875 --- /dev/null +++ b/tests/holmes_operator/test_triggeredhealthcheck_component.py @@ -0,0 +1,286 @@ +"""Component tests for the TriggeredHealthCheck (deployment-rollout) trigger. + +Pure helpers are tested directly; the execution path is tested with the Kubernetes +API mocked and the rollout-settle wait patched out. +""" + +from unittest.mock import MagicMock + +import pytest + +from holmes_operator import context, trigger_executor +from holmes_operator.config import OperatorConfig +from holmes_operator.models import TriggeredHealthCheckSpec + + +@pytest.fixture +def mock_config(): + return OperatorConfig( + holmes_api_url="http://mock-holmes-api:80", + holmes_api_timeout=300, + log_level="INFO", + max_history_items=10, + cleanup_completed_checks=False, + completed_check_ttl_hours=24, + ) + + +@pytest.fixture +def mock_k8s_api(): + api = MagicMock() + api.create_namespaced_custom_object = MagicMock() + api.patch_namespaced_custom_object_status = MagicMock() + api.get_namespaced_custom_object = MagicMock( + return_value={ + "metadata": {"name": "verify-rollouts", "resourceVersion": "1"}, + "status": {}, + } + ) + return api + + +@pytest.fixture +def setup_context(mock_config, mock_k8s_api): + context.config = mock_config + context.k8s_api = mock_k8s_api + trigger_executor.clear_rollout_cache() + yield + context.config = None + context.k8s_api = None + trigger_executor.clear_rollout_cache() + + +def _deployment(images, labels=None): + return { + "metadata": {"labels": labels or {}}, + "spec": { + "template": { + "spec": {"containers": [{"name": "app", "image": i} for i in images]} + } + }, + } + + +class TestHelpers: + def test_selector_matches_subset(self): + assert trigger_executor.selector_matches( + {"app": "checkout"}, {"app": "checkout", "tier": "web"} + ) + + def test_selector_no_match(self): + assert not trigger_executor.selector_matches( + {"app": "checkout"}, {"app": "payments"} + ) + + def test_empty_selector_matches_all(self): + assert trigger_executor.selector_matches({}, {"app": "anything"}) + + def test_extract_images_joins_containers(self): + body = _deployment(["repo/app:v2", "repo/sidecar:v1"]) + assert trigger_executor.extract_images(body) == "repo/app:v2, repo/sidecar:v1" + + def test_render_query_substitutes_tokens(self): + rendered = trigger_executor.render_query( + "{{ .deployment }} in {{ .namespace }}: {{ .old.image }} -> {{.new.image}}", + deployment="checkout", + namespace="prod", + old_image="repo/app:v1", + new_image="repo/app:v2", + ) + assert rendered == "checkout in prod: repo/app:v1 -> repo/app:v2" + + def test_render_query_unknown_image(self): + rendered = trigger_executor.render_query( + "was {{ .old.image }}", "d", "n", "", "repo/app:v2" + ) + assert rendered == "was unknown" + + def test_compose_query_injects_context_for_terse_query(self): + # A query with no tokens still gets the rollout facts. + query = trigger_executor.compose_query( + "Is the new version healthy?", + deployment="checkout", + namespace="prod", + old_image="repo/app:v1", + new_image="repo/app:v2", + ) + assert "Is the new version healthy?" in query + assert "- Deployment: checkout" in query + assert "- Namespace: prod" in query + assert "- Previous image(s): repo/app:v1" in query + assert "- New image(s): repo/app:v2" in query + + +class TestDetectRollout: + def test_baseline_then_change(self): + key = "prod/checkout" + # First observation establishes a baseline -> not a rollout + assert trigger_executor.detect_rollout(key, _deployment(["app:v1"])) is None + # Same template -> not a rollout + assert trigger_executor.detect_rollout(key, _deployment(["app:v1"])) is None + # Template change -> rollout with old/new images + result = trigger_executor.detect_rollout(key, _deployment(["app:v2"])) + assert result == ("app:v1", "app:v2") + + def test_forget_resets_baseline(self): + key = "prod/checkout" + trigger_executor.detect_rollout(key, _deployment(["app:v1"])) + trigger_executor.forget_deployment(key) + # After forgetting, next observation is a baseline again + assert trigger_executor.detect_rollout(key, _deployment(["app:v2"])) is None + + +class TestCooldown: + def test_no_cooldown_when_disabled(self): + assert not trigger_executor.is_in_cooldown({}, "checkout", 0) + + def test_within_cooldown(self): + status = { + "cooldowns": [ + { + "deployment": "checkout", + "lastTriggerTime": trigger_executor.get_current_time_iso(), + } + ] + } + assert trigger_executor.is_in_cooldown(status, "checkout", 600) + + def test_outside_cooldown(self): + status = { + "cooldowns": [ + { + "deployment": "checkout", + "lastTriggerTime": "2000-01-01T00:00:00+00:00", + } + ] + } + assert not trigger_executor.is_in_cooldown(status, "checkout", 600) + + +class TestPendingQueue: + def test_due_pending_selects_past_entries(self): + pending = [ + {"deployment": "a", "fireAt": "2000-01-01T00:00:00+00:00"}, + {"deployment": "b", "fireAt": "2999-01-01T00:00:00+00:00"}, + ] + due = trigger_executor.due_pending(pending) + assert [e["deployment"] for e in due] == ["a"] + + def test_compute_fire_at_in_future(self): + from datetime import datetime, timezone + + fire_at = datetime.fromisoformat(trigger_executor.compute_fire_at(3600)) + assert fire_at > datetime.now(timezone.utc) + + async def test_add_pending_debounces_per_deployment( + self, setup_context, mock_k8s_api + ): + # Resource already has a pending entry for "checkout" + mock_k8s_api.get_namespaced_custom_object.return_value = { + "metadata": {"resourceVersion": "1"}, + "status": { + "pending": [ + {"deployment": "checkout", "fireAt": "2999-01-01T00:00:00+00:00"} + ] + }, + } + + await trigger_executor.add_pending( + mock_k8s_api, + trigger_name="verify-rollouts", + namespace="prod", + deployment="checkout", + fire_at="2999-06-01T00:00:00+00:00", + old_image="app:v1", + new_image="app:v2", + ) + + patched = mock_k8s_api.patch_namespaced_custom_object_status.call_args[1][ + "body" + ]["status"]["pending"] + # Old entry replaced, not duplicated + assert len(patched) == 1 + assert patched[0]["fireAt"] == "2999-06-01T00:00:00+00:00" + assert patched[0]["newImage"] == "app:v2" + + async def test_remove_pending_matches_deployment_and_fireat( + self, setup_context, mock_k8s_api + ): + mock_k8s_api.get_namespaced_custom_object.return_value = { + "metadata": {"resourceVersion": "1"}, + "status": { + "pending": [ + {"deployment": "checkout", "fireAt": "2999-01-01T00:00:00+00:00"}, + {"deployment": "payments", "fireAt": "2999-02-01T00:00:00+00:00"}, + ] + }, + } + + await trigger_executor.remove_pending( + mock_k8s_api, + trigger_name="verify-rollouts", + namespace="prod", + entries=[ + {"deployment": "checkout", "fireAt": "2999-01-01T00:00:00+00:00"} + ], + ) + + patched = mock_k8s_api.patch_namespaced_custom_object_status.call_args[1][ + "body" + ]["status"]["pending"] + assert [p["deployment"] for p in patched] == ["payments"] + + +class TestSpawnCheck: + async def test_spawns_healthcheck_and_records_status( + self, setup_context, mock_k8s_api + ): + spec = TriggeredHealthCheckSpec( + deploymentRollout={"selector": {"matchLabels": {"app": "checkout"}}}, + query="checkout rolled out to {{ .new.image }} (was {{ .old.image }})", + mode="alert", + destinations=[{"type": "slack", "config": {"channel": "#deploys"}}], + ) + + await trigger_executor.spawn_check( + trigger_name="verify-rollouts", + namespace="prod", + trigger_uid="thc-uid-1", + spec=spec, + deployment="checkout", + old_image="repo/app:v1", + new_image="repo/app:v2", + k8s_api=mock_k8s_api, + ) + + # A HealthCheck was created with the rendered query and owner reference + mock_k8s_api.create_namespaced_custom_object.assert_called_once() + create_kwargs = mock_k8s_api.create_namespaced_custom_object.call_args[1] + assert create_kwargs["plural"] == "healthchecks" + hc = create_kwargs["body"] + assert hc["kind"] == "HealthCheck" + query = hc["spec"]["query"] + # The author's query (with tokens substituted) is present... + assert "checkout rolled out to repo/app:v2 (was repo/app:v1)" in query + # ...and the rollout context is auto-injected so a terse query still works. + assert "- Deployment: checkout" in query + assert "- Namespace: prod" in query + assert "- Previous image(s): repo/app:v1" in query + assert "- New image(s): repo/app:v2" in query + assert hc["spec"]["mode"] == "alert" + assert hc["spec"]["destinations"][0]["type"] == "slack" + owner = hc["metadata"]["ownerReferences"][0] + assert owner["kind"] == "TriggeredHealthCheck" + assert owner["uid"] == "thc-uid-1" + assert hc["metadata"]["labels"]["holmesgpt.dev/triggered-by"] == "verify-rollouts" + + # Status recorded the trigger (history + cooldown + counters) + mock_k8s_api.patch_namespaced_custom_object_status.assert_called() + status = mock_k8s_api.patch_namespaced_custom_object_status.call_args[1]["body"][ + "status" + ] + assert status["lastTriggerDeployment"] == "checkout" + assert status["triggerCount"] == 1 + assert status["history"][0]["checkName"].startswith("verify-rollouts-") + assert status["history"][0]["newImage"] == "repo/app:v2" + assert status["cooldowns"][0]["deployment"] == "checkout"