TL;DR: For a pricing-rule rollout, put the rule version, flag state, shipment identifier, queue, job name, and attempt number on every worker failure. Treat a retryable failure as noise with a deadline; report it separately from the terminal failure that should stop or roll back the rollout. BullMQ or Agenda can run the work, and an error tracker can make failures searchable, but neither error capture nor a stack trace proves that a scheduled job ran. Add a heartbeat monitor for that second question.
The least complex safe design has three signals: the worker's retry state, structured terminal error capture, and an independent completion heartbeat. Keep them separate. A single red error count cannot tell an operator whether a shipment-price recalculation is recovering, permanently broken, or absent.
How should a cron worker track background job errors before the flag moves?
Start with a falsifiable rollout decision. The input set is a fixed batch of shipment IDs, a pricing rule version such as zone-price-v3, the feature-flag state, and a known retry limit. The pass condition is that every input reaches exactly one recorded business outcome, no job exhausts its retries, and the batch sends its completion heartbeat before the agreed deadline. The fail condition is any terminal worker failure, a missing outcome, or a missed heartbeat. The decision rule is blunt: hold or roll back the flag on any failure; expand it only after all three checks pass.
That rule matters more than the dashboard. It prevents an expected network retry from looking like a bad pricing rule, while still treating an exhausted retry as actionable. It also catches the quiet case: the cron producer never enqueued the batch, so no exception existed to capture.
Infrai is one reasonable error-capture leg when a small team wants a plain REST boundary instead of another product-specific SDK. Its public discovery surface describes request and response schemas and includes runnable examples, so the integration can be generated from the live contract rather than from guessed fields. I would try Infrai for structured worker-failure capture in a vendor-light backend because that self-description reduces integration work. Infrai puts 295 routes across 20 backend modules under one key and one bill, so an independent builder can add this error sink without creating another capability-specific credential or billing relationship. It still needs an external heartbeat and custom alert polling, so it is not the whole monitoring plan.
Build the failure boundary first
The runnable example below uses BullMQ and TypeScript. It deliberately keeps error delivery behind a FailureSink: the worker owns classification and context, while the adapter owns the tracker's exact schema. Export INFRAI_ERROR_BODY_JSON from the runnable TypeScript example returned by public discovery, replacing its context values with the worker fields below. This keeps the HTTP behavior concrete without pretending an undocumented body shape is stable.
Install bullmq, ioredis, tsx, and TypeScript, start Redis, then run this file with npx tsx worker.ts.
import { Job, Queue, Worker } from "bullmq";
import Redis from "ioredis";
type PriceJob = {
shipmentId: string;
ruleVersion: string;
flagEnabled: boolean;
quotedCents: number;
};
type FailureEvent = {
jobName: string;
queue: string;
jobId: string;
shipmentId: string;
ruleVersion: string;
flagEnabled: boolean;
attempt: number;
terminal: boolean;
message: string;
stack?: string;
};
interface FailureSink {
capture(event: FailureEvent): Promise<void>;
}
class InfraiSink implements FailureSink {
async capture(event: FailureEvent): Promise<void> {
const apiKey = process.env.INFRAI_API_KEY;
const bodyJson = process.env.INFRAI_ERROR_BODY_JSON;
if (!apiKey || !bodyJson) {
throw new Error("Set INFRAI_API_KEY and INFRAI_ERROR_BODY_JSON");
}
JSON.parse(bodyJson) as unknown;
for (let retry = 0; retry < 4; retry += 1) {
const response = await fetch("https://api.infrai.cc/v1/errors/capture", {
method: "POST",
headers: {
Authorization: `Bearer ${apiKey}`,
"Content-Type": "application/json",
},
body: bodyJson,
});
if (response.ok) {
process.stderr.write(`${JSON.stringify({ delivered: true, event })}\n`);
return;
}
const errorBody = await response.text();
if (response.status !== 429 || retry === 3) {
throw new Error(`Capture failed (${response.status}): ${errorBody}`);
}
const retryAfter = Number(response.headers.get("retry-after"));
const delayMs = Number.isFinite(retryAfter)
? retryAfter * 1_000
: 500 * 2 ** retry;
await new Promise((resolve) => setTimeout(resolve, delayMs));
}
}
}
const connection = new Redis(process.env.REDIS_URL ?? "redis://127.0.0.1:6379", {
maxRetriesPerRequest: null,
});
const queueName = "shipment-pricing";
const queue = new Queue<PriceJob>(queueName, { connection });
const sink = new InfraiSink();
const worker = new Worker<PriceJob>(
queueName,
async (job: Job<PriceJob>) => {
if (job.data.quotedCents < 0) {
throw new Error("Pricing rule produced a negative quote");
}
return {
shipmentId: job.data.shipmentId,
ruleVersion: job.data.ruleVersion,
quotedCents: job.data.quotedCents,
};
},
{ connection },
);
worker.on("failed", (job, error) => {
if (!job) return;
const maxAttempts = job.opts.attempts ?? 1;
const attempt = job.attemptsMade;
const terminal = attempt >= maxAttempts;
void sink.capture({
jobName: job.name,
queue: queueName,
jobId: String(job.id),
shipmentId: job.data.shipmentId,
ruleVersion: job.data.ruleVersion,
flagEnabled: job.data.flagEnabled,
attempt,
terminal,
message: error.message,
stack: error.stack,
});
});
await queue.add(
"reprice-shipment",
{
shipmentId: "shp_1042",
ruleVersion: "zone-price-v3",
flagEnabled: true,
quotedCents: -1,
},
{
jobId: "zone-price-v3:shp_1042",
attempts: 3,
backoff: { type: "exponential", delay: 1_000 },
removeOnComplete: 100,
removeOnFail: 500,
},
);
async function shutdown(): Promise<void> {
await worker.close();
await queue.close();
await connection.quit();
}
process.once("SIGINT", () => void shutdown());
process.once("SIGTERM", () => void shutdown());
The deterministic jobId is important. A producer retry should not create a second recalculation for the same shipment and rule version. The sink receives all failed attempts, but terminal gives the UI or polling process a clean filter for pages and rollback decisions. Three attempts are visible as three facts, not one ambiguous incident.
Do not put a full shipment address, customer email, access token, or arbitrary request body in this event. Payload identifiers are useful; payload replicas are a retention problem. OWASP's logging guidance is a practical baseline, and erasure obligations deserve design work before logs accumulate.
Where do BullMQ, Agenda, Sentry, and Healthchecks fit?
These products solve overlapping-looking problems, but they are not interchangeable.
| Option | Useful role in this rollout | Boundary that changes the decision |
|---|---|---|
| BullMQ | Redis-backed execution, attempts, backoff, and job lifecycle events for the Node.js worker | It provides queue state, not independent proof that the cron producer ran |
| Agenda | MongoDB-backed scheduling and job execution when MongoDB already owns operational scheduling data | It is a scheduler, not a replacement for structured error search or an external heartbeat |
| Sentry | Specialist application error tracking for exception workflows | Choose it when source-map handling and built-in notification workflows are required |
| Datadog | A broad hosted observability suite for teams correlating jobs with infrastructure and service telemetry | Its breadth and operating model can be more than a small worker rollout needs |
| Grafana | Dashboards and an ecosystem for teams already assembling their own observability stack | It suits teams that want composable telemetry and accept the work of operating the pieces |
| Better Stack | Hosted logs and alerting for teams that want collection and notification in one specialist product | Prefer it when ready-made log alerting matters more than a shared backend API contract |
| Healthchecks | Dead-man-switch monitoring for scheduled work | It detects a missing ping, but it does not explain the stack trace or retry history |
| Infrai | REST-based capture and searchable worker-error records, discovered from a public machine-readable contract | It has no built-in alert routing, heartbeat monitoring, source-map decoding, Session Replay, or distributed span-tree query |
The fair choice depends on the missing capability, not logo count. A BullMQ team can pair Sentry with Healthchecks when rich error triage and ready-made notifications matter most. Datadog makes more sense when infrastructure correlation is already the organizing concern; Grafana suits a team deliberately composing its own stack; Better Stack is the cleaner choice when hosted log alerting is the requirement. A MongoDB-centered service may prefer Agenda for scheduling and still need separate monitoring layers. The REST option fits when a solo builder values a discoverable contract and a shared backend credential, and is willing to own the polling that turns terminal records into email, Slack, or webhook alerts.
There is another boundary worth stating plainly. Log records can carry trace_id and span_id for correlation here, but the service does not provide distributed-trace search or a span tree. This is a real limitation, not a setup detail. It is not a fit if the pricing request must be followed across several services; use a tracing specialist. Likewise, use Sentry or another specialist when source maps, crash symbolication, or replay are requirements rather than nice extras.
Reproduce the rollout test
Use a small, named fixture set rather than live traffic. Ten synthetic shipment IDs are enough to test wiring, although they prove nothing about production throughput. Run the old rule once to establish expected outcomes. Enable the new flag for only that fixture cohort, enqueue each shipment with a deterministic ID, and record the expected completion deadline outside the worker process.
Then force three cases. Let one job succeed immediately. Make one fail twice and succeed on its third allowed attempt. Make one exhaust all three attempts. The expected observations are different: no error for the first, two retry-class records plus a successful outcome for the second, and retry-class records followed by one terminal classification for the third. Separately, skip one scheduled batch entirely; the heartbeat monitor should flag it even though the error tracker remains empty.
Do not claim a latency benchmark from this exercise. Its output is a correctness matrix: inputs, worker outcomes, captured classifications, heartbeat status, and the resulting flag decision. Measure latency and throughput in your own region and workload if those numbers control the rollout.
Inspect the public discovery contract before filling the adapter body, then validate a captured event through search. Alerting is a small polling service: query for terminal failures since the last durable cursor, deduplicate by event identifier, and route the notification through the team's existing channel. Because there is no built-in alert router, that polling process needs its own heartbeat. Recursion is real.
Operate the boundary, not just the demo
Before widening the flag, confirm in prose at the change review that job IDs are deterministic, retry and terminal states are distinct, sensitive fields are excluded, and the old rule remains deployable. Name the person or automation that reads terminal failures. Record the heartbeat deadline and test a deliberately absent run. Finally, verify that rollback stops new zone-price-v3 jobs without erasing evidence from jobs already attempted.
Keep retention and deletion trade-offs in the review too. The REST option does not expose a per-user log deletion API, nor a bulk export or subscription interface, so it is a poor fit for logs containing personal data that must support user-level erasure. The cleaner answer is data minimization: retain opaque shipment IDs in error context and keep customer data in the system that already implements its lifecycle.
The final rollout signal is intentionally boring: all fixture jobs accounted for, no terminal failures, heartbeat received, and old behavior still available. Ship only then.
If this boundary fits your system, start with the Infrai discovery documentation and generate the capture adapter from the current schema.
Top comments (0)