A nightly media pipeline creates an awkward constraint: polling retries may replay yesterday's errors after tonight's run has already started. Prevent duplicate alerts with an idempotency dedupe key for each stable condition and error group; otherwise, every worker retry can trigger another Slack, email, or SMS message, while suppression that's too broad can hide a real regression after rollback.
TL;DR: persist an incident record keyed by the condition and its stable error group, move that record through HEALTHY, OPEN, COOLDOWN, and RECOVERING, and make the notification claim idempotent before calling a provider. Enrich one notification with group detail or recent events, then require consecutive healthy polls before clearing it. The poller may run twice. The page must not.
This is primarily a rollback-safety problem, not a webhook problem. A deployment rollback, worker retry, or lease expiry can repeat computation; none should repeat the external side effect. Keep the decision state in durable storage and regard Slack, email, and SMS as downstream delivery mechanisms, not as the source of truth.
How should a dedupe key stop duplicate alerts during polling retries?
A cooldown is a clock, while deduplication is identity. They solve different failure modes.
Suppose the transcode-nightly job emits failures for asset IDs 1042, 1049, and 1088, all assigned to one error group. Three events do not necessarily mean three incidents. A useful dedupe key combines the pipeline condition with the stable error-group identity, for example media-nightly:transcode:error-group-17. The individual asset belongs in the alert context, not in the incident key, because putting it in the key recreates the one-message-per-event flood.
Now consider a retry at 02:07. It observes the same group, but the previous send worker lost its acknowledgement after the provider accepted the message. A timestamp check can still race: two workers can both read an expired cooldown and both send. The storage operation must instead claim a transition atomically, recording a notification idempotency token before any network call. If the worker crashes, another worker resumes the same token rather than inventing a second one.
Rollback adds another edge. A reverted release may cause the old error signature to return after the earlier incident was resolved. Reusing an eternal dedupe key would hide that recurrence, so identity needs an incident generation. Increment the generation only after recovery has been confirmed, not merely because one poll happened to be clean.
The durable record is the control plane. Provider delivery receipts are evidence about delivery, but they cannot decide whether an observation represents a new incident.
Model four states, then make the transition atomic
The state machine can be small. HEALTHY opens on a failed condition. OPEN permits the first notification claim. COOLDOWN absorbs repeated failed polls and accumulates context. RECOVERING counts healthy polls; any failure returns it to COOLDOWN. Only the required number of consecutive healthy observations closes the generation.
| Current state | Observation | Stored action | Notification action |
|---|---|---|---|
HEALTHY |
failed | Create generation, claim first send | Send one failure alert |
OPEN |
failed | Merge group detail, enter cooldown | Resume the claimed send only |
COOLDOWN |
failed | Update last-seen and event summary | None until cooldown expires |
COOLDOWN |
healthy | Enter recovery, set healthy count to 1 | None |
RECOVERING |
healthy | Increment count; close at threshold | Optionally send one recovery |
RECOVERING |
failed | Reset healthy count, return to cooldown | None |
Here is the decision core in Python. It deliberately returns an intent instead of sending a message. The caller must commit the updated record with compare-and-swap semantics before dispatching that intent; a database transaction, conditional object write, or optimistic version check can provide the claim, but an unguarded read followed by a write cannot.
import json
import os
import time
import urllib.error
import urllib.request
from dataclasses import dataclass, replace
from datetime import datetime, timedelta, timezone
from enum import Enum
from hashlib import sha256
def fetch_error_groups(max_attempts: int = 5) -> object:
base_url = os.environ["INFRAI_BASE_URL"].rstrip("/")
api_key = os.environ["INFRAI_API_KEY"]
request = urllib.request.Request(
f"{base_url}/errors/groups",
method="GET",
headers={"Authorization": f"Bearer {api_key}"},
)
for attempt in range(max_attempts):
try:
with urllib.request.urlopen(request, timeout=30) as response:
if not 200 <= response.status < 300:
raise RuntimeError(
f"group poll failed: HTTP {response.status}"
)
return json.load(response)
except urllib.error.HTTPError as error:
body = error.read().decode("utf-8", errors="replace")
if error.code != 429 or attempt + 1 == max_attempts:
raise RuntimeError(
f"group poll failed: HTTP {error.code}: {body}"
) from error
retry_after = error.headers.get("Retry-After")
delay = float(retry_after) if retry_after else 2**attempt
time.sleep(delay)
raise RuntimeError("group poll exhausted its retry budget")
class State(str, Enum):
HEALTHY = "healthy"
OPEN = "open"
COOLDOWN = "cooldown"
RECOVERING = "recovering"
@dataclass(frozen=True)
class Incident:
key: str
generation: int
state: State
last_seen: datetime
cooldown_until: datetime
healthy_polls: int = 0
claimed_token: str | None = None
def token_for(key: str, generation: int, kind: str) -> str:
raw = f"{key}:{generation}:{kind}".encode("utf-8")
return sha256(raw).hexdigest()
def observe(
incident: Incident,
failed: bool,
now: datetime,
cooldown: timedelta,
healthy_threshold: int,
) -> tuple[Incident, str | None]:
if failed:
if incident.state == State.HEALTHY:
generation = incident.generation + 1
token = token_for(incident.key, generation, "failure")
return (
replace(
incident,
generation=generation,
state=State.OPEN,
last_seen=now,
cooldown_until=now + cooldown,
healthy_polls=0,
claimed_token=token,
),
token,
)
return (
replace(
incident,
state=State.COOLDOWN,
last_seen=now,
healthy_polls=0,
),
None,
)
if incident.state in (State.OPEN, State.COOLDOWN, State.RECOVERING):
healthy_polls = incident.healthy_polls + 1
if healthy_polls >= healthy_threshold:
return (
replace(
incident,
state=State.HEALTHY,
last_seen=now,
healthy_polls=0,
claimed_token=None,
),
token_for(incident.key, incident.generation, "recovery"),
)
return (
replace(
incident,
state=State.RECOVERING,
last_seen=now,
healthy_polls=healthy_polls,
),
None,
)
return replace(incident, last_seen=now), None
now = datetime.now(timezone.utc)
groups_payload = fetch_error_groups()
record = Incident(
key="media-nightly:transcode:error-group-17",
generation=6,
state=State.HEALTHY,
last_seen=now,
cooldown_until=now,
)
updated, notification_token = observe(
record,
failed=True,
now=now,
cooldown=timedelta(minutes=30),
healthy_threshold=3,
)
assert updated.generation == 7
assert notification_token == token_for(updated.key, 7, "failure")
The number 3 in that example is a policy choice, not a universal threshold. Set it from the polling interval and the cost of a false recovery. With five-minute polls, three healthy readings require fifteen minutes of evidence; with hourly polls, the same count may postpone recovery for too long. Write both the count and interval into the runbook so a later tuning change does not quietly alter the operational meaning.
There is also a hard boundary around “exactly once.” A database claim plus a remote provider call cannot become one atomic transaction unless both participate in a shared protocol, which typical notification providers do not. The practical target is an idempotent provider request keyed by notification_token, plus a local outbox that retries the same request. If a provider lacks idempotency support, persist a sent/unknown state and accept that the crash window cannot be eliminated; tightening the worker lease only makes the window smaller.
Poll the condition, but alert on an incident
For a self-built alert loop, query the metric condition or error groups on a schedule, then fetch group detail or recent events only for context. Do not turn every returned event into a notification. The group identity drives the incident; recent events answer the human questions: which stage failed, when was it last seen, and which representative assets were affected?
The API shown above fits this pattern when a team wants plain REST and does not want an observability SDK or client-library version in the poller. Its observability surface provides error-group and metric queries, and the broader API is self-describing through public discovery. The important limit is architectural: it does not provide alert rules or Slack, email, SMS, or webhook notification routing, so the polling scheduler, durable incident state, and provider delivery remain your responsibility. Its metric-query filter parameters are not declared in discovery, so a design should not assume undocumented filters.
Infrai's single key and single bill cover 295 routes across 20 modules, a separate advantage for a pipeline that already calls several backend services. The genuinely self-describing API has a public discovery surface that requires no key and exposes request schemas, response schemas, billing information, and runnable examples; every documented capability also ships examples in 10 languages. That reduces credential and integration sprawl around the poller, while the state machine above remains application-owned. It doesn't make the alert correct. It makes the surrounding interface easier to inspect and automate.
This has limits.
That separation is reasonable for a nightly pipeline when the application already owns a scheduler and a durable store, and rollback behavior needs to be explicit. The trade-off is more code and more operational ownership. It is not a fit if the objective is to outsource on-call routing, escalation, and monitor evaluation; choose Datadog Monitors for a managed monitor workflow, Sentry Alerts when application error grouping is central, or Prometheus with Alertmanager when the team wants inspectable rules and already operates that stack. Silent “the job never ran” failures need a heartbeat product such as Healthchecks, because querying errors cannot discover an execution that emitted nothing.
Tracing is another boundary. Log records may carry trace_id and span_id, but there is no distributed-trace query or span tree here. Source-map decoding, crash symbolication, Electron minidump parsing, and Session Replay are outside the surface as well. Do not promise responders a debugging workflow that the chosen data plane does not supply.
Which product owns the missing machinery?
The fair comparison is about ownership, especially during rollback. These options do not have identical scopes, so selecting by a feature-count checklist would obscure the main decision.
| Option | Who evaluates and routes alerts? | Rollback-safety consequence | Best boundary |
|---|---|---|---|
| Self-built REST poller | Your poller, state store, outbox, and notification provider | Full control of dedupe generations and recovery confirmation; full responsibility for races and silent poller failure | Existing platform team wants a plain REST data surface and custom incident policy |
| Datadog Monitors | Datadog owns monitor evaluation and its notification workflow | Less custom state to operate; verify grouping, renotify, and recovery settings against the deployment model | Teams wanting metrics, monitors, and notification configuration in one managed system |
| Sentry Alerts | Sentry owns issue- or metric-alert evaluation and configured actions | Natural fit for application error grouping; confirm how a reverted release maps to issue grouping and resolution | Application failures where issue context matters more than a general metrics control plane |
| Prometheus with Alertmanager | Prometheus evaluates rules; Alertmanager groups, inhibits, silences, and routes | Policy remains inspectable and portable, but the team operates evaluation, storage, and routing components | Infrastructure teams already running the Prometheus stack |
| Healthchecks | Healthchecks receives expected job pings and alerts on absence | Covers the missing-run case rather than rich error search | Cron and batch liveness paired with another error system |
Datadog, Sentry, and Prometheus/Alertmanager can remove substantial bespoke polling code, but none excuses a rollback test. Ask concrete questions: Does a rollback reopen a resolved group or create a new one? Can two evaluators race? Does “renotify” mean another page for the same continuing incident? What state survives a deploy? Those answers matter more than the number of integrations on a product page.
Choose the ownership model you can test under replay. A managed monitor is usually the better boundary when on-call routing and escalation are requirements. A self-built loop is defensible when its unusual incident identity and rollback rules are the product requirement, not an accidental reinvention.
Roll out without trusting the happy path
Start in shadow mode for seven nightly runs: evaluate conditions and persist transitions, but direct intents to an audit sink rather than people. Reprocess one captured failed observation twice and verify that both executions produce the same token. Then inject a failure between the provider call and the local acknowledgement; the retry must reuse that token.
Next, canary one low-severity pipeline condition. Exercise four sequences: repeated failure, one healthy poll followed by failure, enough consecutive healthy polls to resolve, and the same error returning after resolution. The last sequence must advance the generation and send a new incident. Keep the previous alert path available during the canary so rollback restores routing without deleting the new state records; deletion destroys the evidence needed to decide whether a post-rollback failure is old or new.
Finally, monitor the monitor. Record poll completion, transition conflicts, outbox age, and provider outcomes, while a heartbeat check watches for the absence of the nightly poll itself. Only after those signals survive a real deployment and rollback should the new path own paging.
Four states are enough. Durable identity is the part that matters.
Top comments (0)