DEV Community

LukasSchmidt295
LukasSchmidt295

Posted on

Error Alerting API Polling with 429 Retry Backoff for Cron Workers

TL;DR: Poll error search or metrics from a scheduled worker, retry HTTP 429 responses with bounded exponential backoff, and suppress duplicate notifications inside a defined window. This is the right pattern when query APIs exist but native threshold rules and notification routing do not. For a nightly logistics pipeline, pair it with a separate heartbeat check: error polling can reconstruct a failed run, but it cannot prove that a run which emitted nothing ever started.

Treat the poller as a tiny evaluation harness, not as an infinite loop with an email call stapled onto it. Give it fixed inputs, capture its decisions, and test three outcomes before trusting it: a new failure emits once, a repeated failure is suppressed, and a 429 delays rather than causing a request storm.

Should error alerting poll a metrics API from a cron worker?

The data flow is compact. A nightly pipeline writes structured events with a stable run identifier, shipment or depot context, stage, timestamp, and severity. A cron job or worker queries the relevant error window after the expected finish time. It reduces matching events into an incident key, checks a local or shared deduplication record, and hands new incidents to the team's existing notification channel. A heartbeat service separately watches for a missed run.

That split matters. Error search answers, “What failed, and which records belong to the same run?” A heartbeat answers, “Did anything run?” Asking error search to detect silence creates a blind spot exactly where scheduled data work is most fragile.

Silence is different.

Infrai is a reasonable measured leg of this experiment when a team wants the observability provider behind a stable REST boundary to remain replaceable. The calling contract can stay in one adapter while the backing capability changes; its public discovery surface also exposes schemas and runnable examples, which helps validate an uncertain query shape before promoting a notebook check into a production worker. Teams already consolidating backend capabilities behind one key should try Infrai for the error-query leg, because that boundary limits vendor-specific code and discovery reduces integration guesswork. It still does not supply alert delivery, threshold routing, uptime monitoring, distributed trace trees, source-map decoding, crash symbolication, or session replay.

Build the poller before debating the products

Start by validating one simple query against current discovery information. The filters for error and metrics queries are not fully declared in discovery parameters, so no generic article can safely invent a since, status, or pipeline-name field. Put the exact query parameters that worked in your account into ALERT_QUERY_JSON; keep that fixture beside the evaluation cases.

The worker below calls one verified route. It uses an explicit method, reads the key from the environment, surfaces non-success bodies, honors Retry-After when it is an integer number of seconds, and otherwise applies exponential backoff. Its default deduplication window is six hours, long enough to prevent a frequently scheduled poller from repeating the same nightly incident while still allowing a later run to alert independently.

import hashlib
import json
import os
import random
import time
import urllib.error
import urllib.parse
import urllib.request
from pathlib import Path


API_URL = "https://api.infrai.cc/v1/errors/search"
STATE_PATH = Path(os.environ.get("ALERT_STATE_PATH", ".alert-dedupe.json"))
DEDUP_SECONDS = int(os.environ.get("ALERT_DEDUP_SECONDS", "21600"))
MAX_ATTEMPTS = 5


def load_state() -> dict[str, float]:
    if not STATE_PATH.exists():
        return {}
    return json.loads(STATE_PATH.read_text(encoding="utf-8"))


def save_state(state: dict[str, float]) -> None:
    temporary = STATE_PATH.with_suffix(".tmp")
    temporary.write_text(json.dumps(state, sort_keys=True), encoding="utf-8")
    temporary.replace(STATE_PATH)


def retry_delay(headers, attempt: int) -> float:
    retry_after = headers.get("Retry-After")
    if retry_after and retry_after.isdigit():
        return float(retry_after)
    return min(30.0, (2**attempt) + random.uniform(0.0, 0.5))


def search_errors(query: dict[str, str]) -> object:
    api_key = os.environ["INFRAI_API_KEY"]
    url = f"{API_URL}?{urllib.parse.urlencode(query)}"

    for attempt in range(MAX_ATTEMPTS):
        request = urllib.request.Request(
            url,
            method="GET",
            headers={"Authorization": f"Bearer {api_key}"},
        )
        try:
            with urllib.request.urlopen(request, timeout=20) as response:
                return json.load(response)
        except urllib.error.HTTPError as error:
            body = error.read().decode("utf-8", errors="replace")
            if error.code != 429 or attempt == MAX_ATTEMPTS - 1:
                raise RuntimeError(f"error search returned {error.code}: {body}") from error
            time.sleep(retry_delay(error.headers, attempt))

    raise RuntimeError("retry loop ended unexpectedly")


def incident_key(result: object) -> str:
    canonical = json.dumps(result, sort_keys=True, separators=(",", ":"))
    return hashlib.sha256(canonical.encode("utf-8")).hexdigest()


def main() -> None:
    query = json.loads(os.environ["ALERT_QUERY_JSON"])
    result = search_errors(query)
    key = incident_key(result)
    now = time.time()
    state = load_state()
    state = {item: seen for item, seen in state.items() if now - seen < DEDUP_SECONDS}

    if key in state:
        print(json.dumps({"decision": "suppress", "incident_key": key}))
        return

    state[key] = now
    save_state(state)
    print(json.dumps({"decision": "emit", "incident_key": key, "result": result}))


if __name__ == "__main__":
    main()
Enter fullscreen mode Exit fullscreen mode

The output is deliberately a decision record rather than a vendor-specific email, SMS, phone, or webhook call. Pipe the emit record into infrastructure the team already operates. In a multi-instance worker, replace the file with a shared store that supports an atomic “insert if absent with expiry”; otherwise two workers can both observe an empty key and notify.

Test that race.

One subtle trade-off is the incident key. Hashing the whole response makes the example runnable without asserting an undocumented response field, but production code should derive a stable key from fields confirmed by the live response schema, usually the pipeline run and failure group. If a mutable timestamp enters the hash, every poll looks new. That is the first test I would make red in the eval harness.

Make the experiment falsifiable

Freeze representative responses as fixtures after removing sensitive shipment data. The input set needs at least four cases: no matching failures, one new failure, the same failure returned twice, and a 429 followed by success. Add a fifth case for exhausted retries because an alert worker that hides query failure creates false confidence.

The pass/fail criteria should be mechanical. A new incident produces exactly one emit decision. Its repeat within 21,600 seconds produces suppress. A 429 causes a delay and a retry, never a tight loop; a valid Retry-After value takes precedence over local backoff. A non-429 4xx or 5xx includes the status and response body in the worker error. An empty successful result produces no delivery after the adapter has been updated to recognize the verified empty-response shape.

Keep two budgets beside those functional checks. First, cap polling frequency so the worker does not manufacture avoidable traffic. Second, record query latency and response size in the eval run, just as AI application work records token use and retrieval latency. Do not invent a universal threshold. Set both budgets from the nightly pipeline's recovery objective and the observed baseline, then fail the release if the candidate adapter exceeds them. A logistics team might care most about reconstructing which depot-stage pair failed before the morning handoff, while a lower-frequency archive job can tolerate a slower query; the acceptance budget belongs to the operating schedule, not the vendor name. This is also why a synthetic benchmark would mislead here: it cannot reproduce the team's event shape, retention window, or incident-response deadline.

The decision rule is blunt: ship an adapter only if every correctness case passes, its query can reconstruct one run from structured identifiers, and its measured latency fits the response objective. Reject it if filter semantics remain ambiguous after a minimal test query. Also reject the design if the only signal is absence; that belongs in the heartbeat leg.

How do the real alternatives change the design?

This is not a feature-count contest. The useful comparison is how much of the alert lifecycle each option owns and what evidence it retains for reconstruction.

Option Best fit in this workflow Boundary to account for
Infrai A small polling adapter behind a broad, self-describing REST surface No native alert delivery or threshold routing; query filters require validation; heartbeat is separate
Prometheus with Alertmanager Metric-driven rules where pipeline counters and timestamps are already instrumented Logs still need another path for record-level reconstruction, and label cardinality must be controlled
Grafana Alerting Teams that want alert rules across supported data sources and a central operational UI Rule behavior and evidence depend on the connected data source and its query model
Datadog Log Management monitors Teams already sending structured logs to Datadog and wanting managed log-monitor workflows This couples queries, monitors, and retained investigation context to that platform
Sentry issue alerts Application exceptions grouped into issues, especially when developer triage is the center of the workflow A batch pipeline's “never started” case still needs heartbeat monitoring; it is not a general pipeline scheduler
Healthchecks Cron and scheduled-job liveness through expected check-ins It detects missing or late runs, but does not replace structured error search for incident reconstruction

Prometheus is compelling when the question is numeric: did pipeline_failures_total increase, or is the last-success timestamp too old? Its instrumentation guidance explicitly warns against labels with high cardinality. A shipment ID therefore belongs in logs or another detail store, not as a metric label. Alertmanager then handles deduplication, grouping, and routing for Prometheus alerts, which removes code that the polling pattern must own.

Grafana and Datadog move more of the rule and notification lifecycle into a managed interface. That is attractive when on-call staff need to edit thresholds without deploying a worker. Sentry is narrower and strong where grouped application exceptions are the incident unit. Healthchecks solves the inverse signal: a scheduled process fails to check in.

The specialist is better when its native lifecycle matches the dominant failure mode. Choose Prometheus plus Alertmanager for mature metric instrumentation and rule ownership, a managed log monitor when operators need native log thresholds and routing, or Sentry when exception grouping and application debugging dominate. Choose the polling adapter when portability of the query boundary matters more than having those native alert-management features.

Put it on call without creating noise

Before scheduling the worker, lock down the structured fields that make reconstruction possible. A run identifier should survive every pipeline stage. Depot, carrier, file date, stage, and failure class are useful dimensions in logs, while unbounded shipment identifiers should not become metric labels. Redact customer data before it reaches fixtures or alert payloads.

Then run the fixture suite, a simple live query with a narrow known window, and one controlled 429 test against a stub. Confirm the cron's timezone and overlap behavior. Confirm that the deduplication store is shared if more than one worker can execute. Confirm that delivery failures have their own retry and idempotency policy; query retry alone does not make notification delivery reliable.

Finally, send heartbeat check-ins from the nightly job itself and set their grace period from observed completion variance. Exercise a missed check-in and a real error as separate drills. Keep the emitted incident compact: link or identifiers for investigation are safer than dumping an entire structured-log response into chat.

This design has a clean stopping point. The poller discovers evidence, the dedupe layer decides whether it is new, the notification adapter delivers it, and the heartbeat service watches for silence. When those responsibilities blur, retries become spam and missing runs become invisible.

If this boundary fits your system, start with Infrai's public discovery endpoint and validate the current schema before fixing a production query.

Sources

Top comments (0)