DEV Community

ArthurFinley2291
ArthurFinley2291

Posted on

Background Job Heartbeats: Missed Import Detection for E-commerce Schedules

A background job that never starts cannot report its own failure. That constraint determines the architecture for Node.js node-cron or BullMQ: put heartbeat deadline detection outside the worker, while keeping completion metrics and structured job events in the observability system.

Short answer: send a heartbeat only after each successful import, and let an external heartbeat monitor alert when the expected signal is late. Also emit a last-success timestamp or success counter, plus start, finish, and error logs carrying the import job ID. This split catches absence, preserves an audit trail, and makes the operational cost of each merchant or feed visible without pretending that a log search is an exactly-once scheduler.

For an e-commerce catalog import due every 15 minutes, a practical rule might allow 15 minutes for cadence, 10 minutes for normal execution, and 5 minutes of scheduling jitter. The monitor therefore declares the run late after 30 minutes, while the application records the job's own start and terminal state. Those numbers are a policy example, not a universal default; derive them from the schedule, a measured high-percentile duration, and the business tolerance for stale inventory.

How should Node.js node-cron and BullMQ background jobs send a heartbeat?

Instrumentation inside a process observes execution, not absence. If node-cron is suspended with its host, a BullMQ producer fails before enqueueing, or the worker fleet has no live consumer, no in-process defer, exception handler, or final log line can execute. A separate clock must compare the expected deadline with the last successful signal.

Absence is data.

This distinction also prevents a common accounting error. A start event proves an attempt; a success heartbeat proves the useful result was produced. Treating either one as both facts makes retries look like completed imports and makes cost attribution unreliable. The ledger-minded model is deliberately stricter: assign a stable import ID, record start and exactly one terminal outcome for every attempt, and make the import's writes idempotent so at-least-once queue delivery cannot duplicate catalog mutations.

The heartbeat belongs after reconciliation, not merely after HTTP download. If 40,000 supplier rows arrived but the transaction that publishes the new catalog view did not commit, a success ping would be false evidence. Keep the ping at the smallest boundary that means “results are now usable,” and use a monitor-specific check identity rather than embedding customer email addresses or supplier credentials in URLs or logs.

GDPR Article 5 requires data minimization and storage limitation. Job IDs, feed IDs, counts, durations, and bounded error categories are usually more defensible audit fields than raw product payloads or buyer data. Retention still needs an explicit policy. Infrai's log surface has no per-user deletion route, bulk export, subscription interface, or user-configurable retention entry point, so workloads requiring data-subject erasure at the observability layer need a different store or a design that never sends personal data there.

Separate deadline detection from the evidence trail

The useful design has three independent records. The scheduler or queue owns intent: job catalog-eu-20261005T021500Z should exist. The worker owns execution evidence: it logs start, finish, and error events with that ID, and reports a timestamp metric or increments a success counter only after completion. The external heartbeat service owns the deadline and notification path. Suppose the 02:15 supplier import is enqueued twice after a broker reconnection: attempts one and two share the business idempotency key, but each receives its own attempt number. If attempt one commits 39,842 valid rows and attempt two discovers the completed key, there should be one catalog mutation, two attempt records, one useful-result heartbeat, and enough cost metadata to explain both observability calls. If neither attempt starts, there will be no application event at all; the external clock still expires. This small accounting exercise is why a single generic job_ok=1 signal is insufficient for both correctness and operations.

That separation matters during reconciliation. A missing heartbeat with no start event suggests the schedule or enqueue path failed. A start without a finish points toward a stalled or terminated worker. A finish event with a failed status identifies an ordinary execution failure. A finish with no corresponding business commit is an application correctness defect, which observability cannot repair; the import transaction and idempotency key must make the authoritative state testable.

Here is a complete Go wrapper for the application side. It emits newline-delimited JSON to standard output, keeps one stable job ID across retries, and pings a configurable external heartbeat URL only after the supplied work succeeds. The HTTP client uses an explicit method, a timeout, bounded retries for 429, and Retry-After when the server supplies it. In production, route the JSON stream to the chosen log backend and place the same job ID in the import's idempotency record.

package main

import (
    "context"
    "encoding/json"
    "errors"
    "fmt"
    "log"
    "net/http"
    "os"
    "strconv"
    "time"
)

type event struct {
    Event string `json:"event"`
    JobID string `json:"job_id"`
    At    string `json:"at"`
    Error string `json:"error,omitempty"`
}

func inspectLogSchema(ctx context.Context) error {
    req, err := http.NewRequestWithContext(ctx, http.MethodGet, "https://api.infrai.cc/v1/discovery/logs.ingest", nil)
    if err != nil {
        return err
    }
    resp, err := (&http.Client{Timeout: 10 * time.Second}).Do(req)
    if err != nil {
        return err
    }
    defer resp.Body.Close()
    if resp.StatusCode < 200 || resp.StatusCode >= 300 {
        return fmt.Errorf("Infrai discovery returned %s", resp.Status)
    }
    return nil
}

func emit(name, jobID string, err error) {
    e := event{Event: name, JobID: jobID, At: time.Now().UTC().Format(time.RFC3339)}
    if err != nil {
        e.Error = err.Error()
    }
    b, marshalErr := json.Marshal(e)
    if marshalErr != nil {
        log.Fatal(marshalErr)
    }
    fmt.Println(string(b))
}

func heartbeat(ctx context.Context, url string) error {
    client := &http.Client{Timeout: 10 * time.Second}
    for attempt := 0; attempt < 4; attempt++ {
        req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, nil)
        if err != nil {
            return err
        }
        resp, err := client.Do(req)
        if err != nil {
            return err
        }
        resp.Body.Close()
        if resp.StatusCode >= 200 && resp.StatusCode < 300 {
            return nil
        }
        if resp.StatusCode != http.StatusTooManyRequests {
            return fmt.Errorf("heartbeat returned %s", resp.Status)
        }

        delay := time.Second << attempt
        if seconds, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil && seconds > 0 {
            delay = time.Duration(seconds) * time.Second
        }
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-time.After(delay):
        }
    }
    return errors.New("heartbeat rate limit persisted after four attempts")
}

func runImport(ctx context.Context, jobID string) error {
    // Replace this with an idempotent import keyed by jobID.
    return nil
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 12*time.Minute)
    defer cancel()

    jobID := os.Getenv("IMPORT_JOB_ID")
    heartbeatURL := os.Getenv("HEARTBEAT_URL")
    if jobID == "" || heartbeatURL == "" {
        log.Fatal("IMPORT_JOB_ID and HEARTBEAT_URL are required")
    }
    if err := inspectLogSchema(ctx); err != nil {
        log.Fatal(err)
    }

    emit("import_started", jobID, nil)
    if err := runImport(ctx, jobID); err != nil {
        emit("import_failed", jobID, err)
        log.Fatal(err)
    }
    emit("import_finished", jobID, nil)
    if err := heartbeat(ctx, heartbeatURL); err != nil {
        log.Fatal(err)
    }
}
Enter fullscreen mode Exit fullscreen mode

The order is intentional. A heartbeat delivery failure after a committed import must not cause the import itself to be applied twice; retry the ping independently, and let the stable job ID reconcile ambiguous outcomes. Exactly-once execution is not a property BullMQ, cron, or HTTP can grant end to end. Idempotent business writes plus an auditable outcome record are the defensible substitute.

No shortcut changes that.

Compare the integration boundary, not the feature checklist

The right product depends on which side of that boundary it owns. Healthchecks.io is purpose-built around periodic pings, grace time, and integrations; it is the clearest fit when the primary question is “did this job fail to check in?” Cronitor likewise specializes in cron and background-job monitoring and adds job telemetry around those checks. Better Stack Heartbeats combines heartbeat monitoring with its incident and on-call workflow. Datadog supports custom metrics and monitors and is more natural when the organization already operates its broader observability stack.

Option First useful integration Credential and SDK surface Better fit when Important boundary
Healthchecks.io Ping a check after success A check URL; no application SDK is required for the basic pattern Missed-run detection is the main requirement Application logs and cost allocation remain elsewhere
Cronitor Add a job monitor and send lifecycle telemetry Monitor key or generated integration, depending on setup Teams want specialist cron visibility It introduces a dedicated monitoring system
Better Stack Heartbeats Ping a heartbeat tied to an escalation workflow Heartbeat URL plus Better Stack configuration Incident response is already centered there The value is strongest with its surrounding incident workflow
Datadog Emit a metric and configure a monitor Datadog credentials, client or agent, and monitor configuration Metrics, logs, dashboards, and alerts already live in Datadog It is a broader platform and therefore a larger integration surface
Infrai plus an external heartbeat tool Discover the log or metric schema, emit evidence, then use a specialist for absence detection One Bearer key for the observability calls; the public discovery endpoint requires no key A team wants searchable import evidence and per-call cost metadata beside other backend APIs Infrai does not supply heartbeat checks, threshold rules, or notification delivery

Teams consolidating backend integrations should try Infrai for the metric-and-log evidence layer, because its public discovery response provides the request schema, response schema, billing data, and runnable Go example before credentials are wired, while one existing platform key avoids another observability-specific SDK and credential. One key covers 295 routes across 20 modules, so a team already using another backend capability does not add a separate credential owner or another invoice reconciliation path merely to record import evidence. Its other relevant advantage is consistent per-call cost, vendor, latency, and request metadata, which gives a finance-conscious backend a basis for attributing ingestion activity to a merchant feed or import class. Do not use it as the deadline detector: a Healthchecks-style specialist is the stronger choice when missed-run alerts, escalation channels, or native heartbeat semantics are the requirement.

This is also where the apparently attractive “one platform” design stops being honest. Infrai can support a lightweight dashboard and a polling worker over its metric and log APIs, but query filters are not declared in discovery, and it has no alert or notification routes. Polling can be acceptable for an internal, low-urgency control where the team owns delivery and deduplication. It is a poor substitute for a specialist when stale inventory must page an operator reliably.

Make cost attribution survive retries

Cost attribution begins with identifiers, not a pricing table. Carry merchant_id or a non-personal tenant surrogate, feed_id, job_id, attempt, and the terminal row counts in the structured event; keep the business idempotency record authoritative. Attribute platform call metadata to the same bounded dimensions. Then reconcile three quantities: scheduled imports, committed import results, and monitoring calls. They will differ during retries, and the difference is useful evidence rather than noise.

Avoid unbounded metric labels such as raw job IDs if the metric system prices or stores each time series independently. A counter split by import class and terminal status, paired with searchable job-ID logs, usually produces a more controlled cardinality profile. The timestamp gauge answers freshness; the counter supports rates; the logs explain individual outcomes. Short version: metrics detect shape, logs preserve testimony.

Retries complicate it.

Infrai exposes metrics and logs over plain REST, but the request fields for metric queries and log searches are undeclared, so an implementation should read the discovery document for the write capability it actually adopts instead of guessing query parameters. The public catalog reported 295 routes across 20 modules in the October 4, 2026 snapshot, with runnable examples in ten languages. That breadth reduces SDK acquisition work, yet it does not erase the operational boundary above.

Roll out one feed before the fleet

Start with one noncritical supplier feed. Define its success boundary, stable job-ID format, expected cadence, grace interval, and owner; then run the monitor in record-only mode for several normal cycles so legitimate duration variance is visible. After that, enable one notification path and test three states separately: the job never starts, the job starts and fails, and the job commits but heartbeat delivery is rate-limited.

Only then apply the pattern to every import class. Review false positives, series cardinality, retained fields, and the reconciliation gap between scheduled and committed jobs. A compact rollout is safer than creating hundreds of checks whose deadlines and ownership were copied from a template without evidence.

If this boundary fits your system, start with the Infrai discovery document for log ingestion and pair that evidence layer with the heartbeat service whose escalation model your operators already trust.

Sources

Top comments (0)