DEV Community

NorbertChristensen3183
NorbertChristensen3183

Posted on

Cohort Failure Alerting: How to Poll Node.js Logs and Errors for 5xx Cron Jobs

Poll exception groups and failed-job or HTTP 5xx logs on a schedule, evaluate thresholds per tenant cohort, and send notifications through a separate router; add a heartbeat monitor as the third signal, because a job that never starts cannot write the error or log that a poller expects to find. The right design is a small, auditable detector with three independent signals, not a claim of exactly-once alert delivery.

TL;DR: for a media experiment, keep control, trial, and enterprise cohorts separate; store a durable cursor and alert state; require repeated threshold breaches when noise is expensive; and attach evidence identifiers to every decision. The selected API can supply error and log evidence, but the worker owns alert rules and notification routing. Log trace_id and span_id fields support manual correlation, not a distributed trace or span tree.

ADR: What must remain true?

The decision is to run one scheduled detector outside the Node.js application fleet. It polls error groups and log search, normalizes the returned records into cohort counts, evaluates policy, then hands an idempotent notification to Slack, email, or another webhook. In a media experiment, this prevents a noisy trial cohort from paging the team responsible for a stable enterprise cohort, while still making a broad regression visible.

Four invariants govern the implementation. First, a polling window has a stable identity, composed from the cohort, signal, and window end, so a retry cannot create a second logical alert. Second, the cursor advances only after the evidence and decision are durably recorded. Third, notifications are effects, not evidence: losing Slack does not erase the detected breach. Fourth, every alert record retains the query fingerprint, threshold, observed count, window, and destination outcome needed for later reconciliation.

Exactly-once processing is an aspiration here, not a transport guarantee. A process can stop after Slack accepts a message but before local state is committed. The defensible contract is at-least-once execution plus an idempotency key at the notification boundary when that destination supports one, and a local outbox otherwise. Reconciliation closes the remaining gap.

The compliance boundary matters too. Logs have no per-user deletion route, and no bulk export or subscription route is available; retention and cold-storage configuration are not exposed through a configuration entry point. Do not place unnecessary personal data in alert evidence. If a tenant invokes a deletion right, an application-owned evidence store needs its own deletion and retention controls rather than assuming the logging system supplies them.

How should Node.js failure alerting poll logs and errors?

There are three signal classes, and only two come from the polled APIs. Error groups or events cover exceptions. Log search can cover failed background jobs and HTTP 5xx patterns after the query shape has been validated. A Healthchecks-style heartbeat covers absence: the nightly cohort aggregation that never ran produces neither an exception nor a completion log.

No event is still an event operationally.

The most dangerous boundary is query semantics. The logs.search filter parameters are not fully declared in discovery, so production code should not bake in a guessed status=500 or tenant_cohort=trial parameter. Validate the exact raw query string and the JSON array path against representative data, review it like a schema migration, and deploy its fingerprint with the threshold policy. The sample below therefore accepts both values as configuration and refuses to count an ambiguous response.

Source maps, crash symbolication, Electron minidumps, and Session Replay are outside this design. So is a distributed trace query. Those are not minor presentation features; they determine whether the evidence can answer the question an incident responder will ask.

Compare the operating models

The useful comparison is signal quality versus operational noise, not feature count or price. Each option draws a different ownership boundary.

Option Best fit for this decision Boundary to account for
Infrai Teams that value one key and one REST API across 295 routes in 20 modules No native alert rules or notification routing; it is not suitable when the vendor must own monitor evaluation and paging
Sentry Exception-centered triage where grouped errors and application context are the primary evidence A separate heartbeat is still needed for a task that never runs; verify the selected plan's alert workflow
Datadog Teams already operating log monitors and notification integrations in a broader monitoring platform Monitor configuration and tagging discipline become part of the control plane; cohort tags must be governed
Grafana Cloud Organizations that want Loki-style log queries connected to Grafana-managed alerting Query cardinality and label design affect the quality and maintainability of cohort rules
Healthchecks Cron and scheduled-task heartbeat monitoring with explicit late or missing signals It complements exception and 5xx evidence rather than replacing either source

The first row's advantage is breadth behind a simple surface: adding another backend capability remains another endpoint under the same contract instead of another SDK and credential integration. The self-describing discovery surface is public without a key, and documented capabilities have runnable examples in 10 languages. Its limitation remains decisive, however: those facts do not turn polling into native alerting, and they do not remove the need for an outbox. At write boundaries, 171 of 294 capabilities declare idempotency, with an Idempotency-Key convention and a 24-hour default deduplication window; that convention is valuable for effects, while the detector still has to reconcile delivery independently.

Sentry is a reasonable default when exception investigation is the center of gravity. Datadog or Grafana Cloud is a better fit when the organization wants the monitoring platform itself to own rule evaluation and routing. Healthchecks should remain beside any of them for scheduled-work liveness. A fair architecture may use two products; forcing all three signals through one tool can lower signal quality merely to simplify procurement. This is the central trade-off.

Implement the critical path

The following Go program is intentionally narrow. It polls the two verified read routes, passes only a previously tested query string, resolves a configured dot-separated path to an array, and evaluates one threshold. Run separate instances for control, trial, and enterprise, each with its own query and state file. This makes cohort policy visible rather than hiding it in a large conditional.

package main

import (
    "crypto/sha256"
    "encoding/hex"
    "encoding/json"
    "errors"
    "fmt"
    "io"
    "net/http"
    "net/url"
    "os"
    "strconv"
    "strings"
    "time"
)

type config struct {
    APIKey, Route, Query, ArrayPath string
    Cohort, Signal, WebhookURL      string
    Threshold                      int
}

type alert struct {
    Key         string `json:"key"`
    Cohort      string `json:"cohort"`
    Signal      string `json:"signal"`
    WindowEnd   string `json:"window_end"`
    QueryHash   string `json:"query_hash"`
    Threshold   int    `json:"threshold"`
    Observed    int    `json:"observed"`
}

func required(name string) string {
    v := os.Getenv(name)
    if v == "" {
        panic(name + " is required")
    }
    return v
}

func load() config {
    n, err := strconv.Atoi(required("ALERT_THRESHOLD"))
    if err != nil || n < 1 {
        panic("ALERT_THRESHOLD must be a positive integer")
    }
    route := required("INFRAI_ROUTE")
    if route != "/errors/groups" && route != "/logs/search" {
        panic("INFRAI_ROUTE must be /errors/groups or /logs/search")
    }
    return config{required("INFRAI_API_KEY"), route,
        required("INFRAI_QUERY"), required("INFRAI_ARRAY_PATH"),
        required("COHORT"), required("SIGNAL"),
        required("ALERT_WEBHOOK_URL"), n}
}

func getJSON(c config) (any, error) {
    endpoint := strings.TrimRight(required("INFRAI_BASE_URL"), "/") + c.Route
    if c.Query != "" {
        if _, err := url.ParseQuery(c.Query); err != nil {
            return nil, fmt.Errorf("invalid INFRAI_QUERY: %w", err)
        }
        endpoint += "?" + c.Query
    }
    req, err := http.NewRequest(http.MethodGet, endpoint, nil)
    if err != nil { return nil, err }
    req.Header.Set("Authorization", "Bearer "+c.APIKey)

    client := &http.Client{Timeout: 20 * time.Second}
    for attempt := 0; attempt < 5; attempt++ {
        resp, err := client.Do(req)
        if err != nil { return nil, err }
        body, readErr := io.ReadAll(io.LimitReader(resp.Body, 4<<20))
        resp.Body.Close()
        if readErr != nil { return nil, readErr }
        if resp.StatusCode == http.StatusTooManyRequests {
            delay := time.Duration(1<<attempt) * time.Second
            if seconds, e := strconv.Atoi(resp.Header.Get("Retry-After")); e == nil {
                delay = time.Duration(seconds) * time.Second
            }
            time.Sleep(delay)
            continue
        }
        if resp.StatusCode < 200 || resp.StatusCode >= 300 {
            return nil, fmt.Errorf("poll returned %s: %s", resp.Status, body)
        }
        var document any
        if err := json.Unmarshal(body, &document); err != nil { return nil, err }
        return document, nil
    }
    return nil, errors.New("poll remained rate limited after retries")
}

func arrayAt(document any, path string) ([]any, error) {
    current := document
    for _, part := range strings.Split(path, ".") {
        object, ok := current.(map[string]any)
        if !ok { return nil, fmt.Errorf("%q is not an object", part) }
        current, ok = object[part]
        if !ok { return nil, fmt.Errorf("path component %q is absent", part) }
    }
    items, ok := current.([]any)
    if !ok { return nil, errors.New("configured path does not resolve to an array") }
    return items, nil
}

func send(c config, a alert) error {
    body, err := json.Marshal(a)
    if err != nil { return err }
    req, err := http.NewRequest(http.MethodPost, c.WebhookURL, strings.NewReader(string(body)))
    if err != nil { return err }
    req.Header.Set("Content-Type", "application/json")
    req.Header.Set("Idempotency-Key", a.Key)
    resp, err := (&http.Client{Timeout: 20 * time.Second}).Do(req)
    if err != nil { return err }
    defer resp.Body.Close()
    if resp.StatusCode < 200 || resp.StatusCode >= 300 {
        body, _ := io.ReadAll(io.LimitReader(resp.Body, 64<<10))
        return fmt.Errorf("webhook returned %s: %s", resp.Status, body)
    }
    return nil
}

func main() {
    c := load()
    document, err := getJSON(c)
    if err != nil { panic(err) }
    items, err := arrayAt(document, c.ArrayPath)
    if err != nil { panic(err) }
    if len(items) < c.Threshold { return }

    windowEnd := time.Now().UTC().Truncate(5 * time.Minute)
    querySum := sha256.Sum256([]byte(c.Query))
    identity := fmt.Sprintf("%s|%s|%s", c.Cohort, c.Signal, windowEnd.Format(time.RFC3339))
    keySum := sha256.Sum256([]byte(identity))
    a := alert{hex.EncodeToString(keySum[:]), c.Cohort, c.Signal,
        windowEnd.Format(time.RFC3339), hex.EncodeToString(querySum[:]),
        c.Threshold, len(items)}
    if err := send(c, a); err != nil { panic(err) }
}
Enter fullscreen mode Exit fullscreen mode

The code honors Retry-After when it is an integer number of seconds and otherwise applies exponential backoff. It checks every response status and caps response reads. Configure INFRAI_BASE_URL with the documented versioned API base; keeping it outside an unlinked comparison prevents the article from becoming an acquisition link while preserving a runnable deployment contract. The program does not pretend that a webhook universally deduplicates Idempotency-Key; the receiver must document that behavior. For a destination without that contract, replace direct delivery with a durable outbox whose unique key is alert.Key, then record attempts and acknowledgements. Store the outbox row before attempting delivery, enforce a unique constraint on that key, record the response status and time, and let a separate reconciler retry unacknowledged rows. That longer path is warranted because the irreducible crash window sits between the remote acknowledgement and the local commit.

Test the boundary first.

Before enabling notifications, run the detector in record-only mode around known cohort traffic and inspect the raw response separately. Confirm that the configured array represents records rather than pagination metadata, that the query excludes synthetic traffic, and that the five-minute windows do not overlap. Then require two consecutive breaches for a noisy trial cohort while allowing a single severe exception group to alert for enterprise tenants. These are policy examples, not measured universal thresholds; use the experiment's baseline to choose the actual numbers.

Why reject an in-process Node.js timer?

An in-process setInterval looks smaller, but it couples detection to the service being observed. Deploys reset its memory, horizontal replicas duplicate work, and an unhealthy event loop can disable both the application and its detector. A separate scheduler plus worker gives the cursor and outbox an explicit owner.

The rejected option still has a valid use case: a single, non-critical development service where missed alerts are acceptable and the timer emits only diagnostic notifications. Once the message can page a human or affect an experiment decision, durable state and reconciliation are justified.

This ADR also rejects a logs-only design. Exceptions carry grouping value that string matching discards, while heartbeats express expected execution that neither logs nor errors can infer. Keep all three signals, preserve their distinct semantics, and correlate them in the alert record rather than flattening them into one noisy count.

References

Top comments (0)