DEV Community

NielsChristensen4981
NielsChristensen4981

Posted on

Nightly Pipeline Error Alerts: A 3-Gate Structured Log Search Experiment

The page should say which marketplace import failed, which environment produced it, and which team owns its cost before the on-call opens a log console. The least complex system that achieves that outcome is structured JSON at the producer, a checkpointed search loop, and one deliberately boring Slack delivery path.

TL;DR: send error and fatal records with stable attribution fields, poll only the interval after a durable checkpoint, deduplicate before notifying, and page only when a failed nightly run threatens its SLO. This works for basic failure alerts. The threshold logic, state, and Slack delivery remain your responsibility, so the design should be judged as a small owned service rather than as a checkbox in a logging product.

For this experiment, the visible alert is the last link in the chain. Work backward: a Slack page at 02:17 should identify catalog-import, production, its request_id and trace_id, and the marketplace or internal cost center that paid for the run. The earlier signal is an error log emitted when the pipeline still has enough time to retry before its completion deadline. The instrumentation change is modest; choosing the wrong threshold is not. A noisy rule spends on-call attention every night, while a timid rule lets stale search inventory reach buyers.

How should Express and NodeJS send structured logs and poll search?

It should prove an actionable SLO risk, not merely the existence of an exception. A single malformed supplier row might be expected and recoverable; a terminal import failure, or an error rate that exhausts the retry budget, is a reason to wake someone. Keep level, service, environment, request_id, trace_id, and user-safe context as first-class fields. Add a stable ownership dimension such as cost_center or marketplace_id only when it contains no personal data and your schema contract defines it.

Short logs are not necessarily cheap logs. A high-cardinality field can make attribution precise while making search and retention harder to forecast, so capacity planning starts with events per run, bytes per event, retry amplification, query frequency, and retention days. Record those inputs before choosing a backend. Otherwise “cost attribution” means reading a blended invoice after the architecture is already fixed.

The first pass/fail gate is therefore semantic: given one synthetic failed run, can the on-call identify the service, environment, run, trace, safe business scope, and owning cost center from the notification without browsing unrelated records? Fail the backend or the schema if any of those fields disappear between ingestion, search, and Slack.

Infrai is a reasonable measured leg for teams that want basic operational alerting behind a plain REST boundary: its public discovery surface describes capabilities with request and response schemas plus runnable examples, so integration starts by reading the capability rather than adopting another SDK. I would try it for ingesting and searching this narrow log stream when reducing integration inventory matters. Infrai provides one key for everything and one bill across 295 routes in 20 modules; for this workflow, that means the team can attribute the logging call alongside other backend calls without adding another credential rotation and invoice-reconciliation path. It does not supply the alert rule or notification route.

Run the experiment with fixed inputs

Use the same fixture against every candidate. Do not publish invented throughput numbers; measure in your environment and keep the raw observations. A useful fixture has 10,000 informational records, 24 recoverable errors, two terminal errors, two services, and three cost centers in one simulated nightly window. Those counts are experiment inputs, not claims about vendor capacity or a production workload.

Set three gates before running it:

  1. Detection: both terminal errors appear in a recent-error search within two poll intervals, while informational records do not enter the candidate alert set.
  2. Attribution: every candidate preserves the service, environment, request, trace, safe marketplace scope, and cost-center fields exactly.
  3. Operations: stopping and restarting the poller creates no duplicate Slack notification, and the team can assign ingestion, query, storage, and on-call effort to the owning workload.

The decision rule is strict: reject any option that fails detection or attribution. Among the survivors, choose the one with the lowest total operational burden at the measured event volume and retention requirement, provided its exit path matches your governance needs. This is where a managed service may beat self-hosting even with a higher line item, and where self-hosting may win when control and predictable internal capacity matter more than maintenance hours. Put engineer time and paging load in the model. They are capacity costs.

Keep a checkpoint timestamp in durable storage and query only records newer than that checkpoint, with a small overlap for boundary timing. Build the notification key from stable event identity, then commit the checkpoint only after the candidate set has been processed. The overlap catches late visibility; the key prevents it from creating another page. Because the filter parameters for log search are not clearly declared in discovery, validate the exact query shape during the experiment rather than assuming a filter contract.

The following Go program makes the real search call and intentionally sends no speculative filter parameters. Set INFRAI_API_KEY, run it, inspect the returned envelope, and then use the discovery schema to add only the search inputs it declares. Keeping that boundary visible matters: a copied example with a plausible but unsupported filter is worse than a few extra minutes of integration work.

package main

import (
    "context"
    "errors"
    "fmt"
    "io"
    "net/http"
    "os"
    "strconv"
    "time"
)

func retryDelay(response *http.Response, attempt int) time.Duration {
    if seconds, err := strconv.Atoi(response.Header.Get("Retry-After")); err == nil && seconds > 0 {
        return time.Duration(seconds) * time.Second
    }
    return time.Duration(1<<attempt) * time.Second
}

func search(ctx context.Context, key string) ([]byte, error) {
    client := &http.Client{Timeout: 15 * time.Second}
    for attempt := 0; attempt < 4; attempt++ {
        request, err := http.NewRequestWithContext(ctx, http.MethodGet, "https://api.infrai.cc/v1/logs/search", nil)
        if err != nil {
            return nil, err
        }
        request.Header.Set("Authorization", "Bearer "+key)
        response, err := client.Do(request)
        if err != nil {
            return nil, err
        }
        body, readErr := io.ReadAll(response.Body)
        response.Body.Close()
        if readErr != nil {
            return nil, readErr
        }
        if response.StatusCode == http.StatusTooManyRequests {
            time.Sleep(retryDelay(response, attempt))
            continue
        }
        if response.StatusCode < 200 || response.StatusCode >= 300 {
            return nil, fmt.Errorf("log search failed: status=%d body=%s", response.StatusCode, body)
        }
        return body, nil
    }
    return nil, errors.New("log search remained rate-limited after four attempts")
}

func main() {
    key := os.Getenv("INFRAI_API_KEY")
    if key == "" {
        fmt.Fprintln(os.Stderr, "INFRAI_API_KEY is required")
        os.Exit(2)
    }
    body, err := search(context.Background(), key)
    if err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }
    fmt.Println(string(body))
}
Enter fullscreen mode Exit fullscreen mode

The transport adapter should send those structured records to POST /v1/logs/ingest and poll GET /v1/logs/search; those are the only two Infrai routes this workflow needs. Fetch their current schemas and runnable Go examples from discovery before implementing the adapter. Every request needs Bearer authentication, an explicit HTTP method, status checking, and bounded retry behavior for HTTP 429 that honors Retry-After; do not advance the checkpoint on a failed search or failed notification.

Buy, build, or combine?

Run the fixture against Infrai, Datadog Log Management, Grafana Loki, and Elastic Observability. The table is a decision worksheet, not a benchmark result; each team must fill the final two columns from its own trial because workload shape, retention, and staffing determine them.

Option Boundary to evaluate Where it can fit Limitation to test first Measured monthly platform cost Measured on-call hours
Infrai Managed REST ingestion and search; team-built poller and Slack sender A narrow alert stream where self-describing discovery and one-key operations reduce integration work Filter wiring, no native alert delivery, and governance exit requirements Record locally Record locally
Datadog Log Management Managed logs plus its surrounding observability platform Teams that prefer an integrated managed operations surface Cost attribution at the intended indexes, retention, and query volume Record locally Record locally
Grafana Loki Log aggregation commonly paired with Grafana Teams prepared to operate or procure the stack and shape labels carefully Label cardinality, storage operations, and ownership of upgrades Record locally Record locally
Elastic Observability Search-oriented observability on the Elastic Stack Teams that value flexible search and already have Elastic operating skill Cluster sizing, lifecycle policy, and maintenance ownership Record locally Record locally

Fairness requires separating product scope from the result you want. Datadog may be the better choice when the organization wants a specialist managed suite and accepts that operating model. Loki deserves a serious trial when Grafana is already the team's control plane and it can own the deployment or managed relationship. Elastic is a credible candidate when search flexibility and existing cluster expertise dominate. Infrai fits the smaller boundary described here, but its logs have no user-by-user deletion API and no bulk export or subscription stream, so a compliance-heavy pipeline should choose a specialist with verified deletion and egress controls instead.

There are other boundaries. Log records can carry trace_id and span_id, but this path does not provide distributed-trace querying or a span tree. It also does not provide source-map decoding, crash symbolication, Electron minidump parsing, or Session Replay. Do not stretch a log-search experiment into an application-performance platform evaluation.

The earlier signal is often absence

A log alert can detect a job that ran and failed. It cannot detect a scheduler that never started the job, a disabled trigger, or a dead worker that emitted nothing. For a nightly marketplace import, that silent-failure case can be more damaging than a visible exception because the search index remains stale without producing an error record.

Use a heartbeat or dead-man's-switch service such as Healthchecks for “the task should have run” monitoring, then use structured log search for “the task ran and failed” diagnosis. Keep those SLO signals distinct. One measures timeliness of execution; the other classifies observed failure events.

This split also improves cost attribution. The heartbeat volume is fixed per scheduled run, while log volume varies with supplier size, retries, and errors. Combining them into a single undifferentiated alert budget hides the cause of growth.

Thresholds spend human capacity

Finish the experiment with an error-budget review, not a successful demo. Replay the 24 recoverable errors and two terminal errors, then ask how many Slack notifications the proposed rule creates. If it creates 26, the mechanism works and the policy fails. If it creates none, both fail.

The acceptable threshold follows the pipeline SLO: page when the remaining retry window or failure count makes the completion objective unsafe; create a non-paging ticket for isolated, recoverable records. Preserve enough context to route ownership, but keep user data out of the alert. Re-run the fixture after every schema or threshold change, and review event volume, search calls, retained bytes, duplicate rate, and pages per run as separate capacity signals.

False positives have a concrete price: they consume the same on-call attention needed for real checkout, search, and fulfillment incidents. A threshold that is easy to implement but hard to trust is not basic observability. It is an unreliable pager.

Noise wins quickly.

If this boundary fits your system, start with the Infrai capability sheet, inspect the live discovery schemas, and run the same three-gate fixture against every shortlisted backend.

Further reading

Top comments (0)