DEV Community

MiloHastings5316
MiloHastings5316

Posted on

Node.js Shipment Admission: Cron Queue Trigger Through a Public Webhook

A shipment update is one business event but potentially hundreds of subscriber deliveries. The operational constraint that changes the design is not the one-minute clock; it is the possibility that the clock, the public webhook, the queue, and a worker will each retry independently.

Short answer: use cron only to trigger a public HTTPS enqueue endpoint every minute, give each shipment-subscriber delivery a stable idempotency key, and let queue workers own rate-limited processing. The scheduler releases work. It must not become the worker or the record of what was delivered.

The resulting decision is cron -> public admission endpoint -> queue -> paced workers. For an e-commerce shipment fan-out, the durable truth belongs beside the queue consumer: (shipment_id, update_version, subscriber_id) identifies the operation, while attempt counts and timestamps describe transport activity. Those are different things.

Keep them different.

Data model: one row per subscriber effect

There are three retry boundaries, and collapsing them produces duplicate notifications. First, cron may call the public endpoint again. Second, a standard queue may deliver a message more than once because its contract is at-least-once. Third, a worker may retry after a downstream 429, honoring Retry-After when it is present and otherwise applying exponential backoff. All three attempts must carry the same business identity even though they have different transport identities.

The admission endpoint should return only after it has durably recorded or published the subscriber jobs. If the same minute is triggered twice, inserting the same stable keys must leave the admitted set unchanged. A worker should acknowledge a queue message only after the subscriber operation has reached its durable completion point. This doesn't manufacture exactly-once delivery across two independent systems; the interval after a remote subscriber accepts a request but before the local completion record commits remains ambiguous. The practical contract is at-least-once transport plus an idempotent business operation.

That ambiguity matters more than cron syntax.

The scheduler has narrow responsibilities. It calls a public http_url; it does not host application code. One execution is limited to 900 seconds, paused schedules do not backfill missed runs, and trigger time has second-level jitter. Those properties are acceptable for opening a one-minute release window, but they rule out strict per-second guarantees and make cron alone unsuitable for catch-up accounting. If 60 deliveries may be released each minute, record which shipment versions are eligible in durable state and let the endpoint select them. Do not infer completeness from the mere existence of 60 cron invocations.

Queue boundaries impose another set of invariants. A message is limited to 256 KB, delay is capped at seven days, retention is at most 30 days, and acknowledgement deletes the message. Put identifiers and a compact event version in the message rather than a full shipment document. FIFO deduplication lasts only five minutes, so it cannot replace a permanent uniqueness rule for a shipment update retried later in the day.

Five minutes is brief.

Trust boundary: public does not mean anonymous

The enqueue endpoint must be reachable over public HTTPS, but public reachability does not imply anonymous writes. Put an authentication control in front of the handler, reject bodies that do not match the expected shipment shape, cap the subscriber count accepted in one request, and log the stable release identity rather than a credential. The scheduler's authorization to call this webhook is an application concern; the INFRAI_API_KEY used by the example to inspect cron configuration belongs only on the server side and must not be forwarded to subscriber endpoints.

This separation is easy to miss because both hops use HTTP. They are different trust relationships: one admits a shipment release, while the other administers scheduling. The local listener in the example is intentionally bound to 127.0.0.1; deployment should terminate HTTPS and apply the chosen inbound authentication before traffic reaches it.

Decision record: choose by the ambiguity window

Product labels matter less than where retry state, replay, and orchestration live. This table keeps the comparison on the uncertainty between subscriber acceptance and local acknowledgement rather than on feature counts or price.

Option Where it fits Boundary that changes the decision
AWS SQS plus a scheduler A managed queue design that needs documented dead-letter queue handling Application idempotency is still required for at-least-once processing
Kubernetes CronJob plus a queue Scheduling already belongs to a Kubernetes control plane CronJob remains a scheduler; queue state and subscriber delivery state live elsewhere
Temporal The shipment update becomes a workflow with waits, branches, compensation, or joins Independent notification fan-out does not inherently require workflow orchestration
Kafka plus an external scheduler Replay or multiple consumer groups is a hard requirement This is a different retention and consumption model from an ack-deletes queue
BullMQ plus a scheduler An application-managed queue is already the team's chosen operating model Retry policy and business idempotency still belong to the application design
Celery plus a scheduler Existing workers already use its task-processing model A second worker stack is difficult to justify solely for minute-based admission
Sidekiq plus a scheduler Existing workers already use its job-processing model It does not change the cross-system ambiguity after subscriber acceptance
Infrai cron plus queue A team wants cron and queue capabilities behind one REST API, one key, and one bill Public HTTPS targets are required; there is no native fan-out/join, Kafka-style replay, or multi-consumer-group model

Infrai is a credible fit when operational consolidation is valuable: the same credential and billing relationship cover backend services, while plain HTTP avoids installing a provider SDK in the Node.js service. Its public discovery surface describes capabilities and runnable examples, so current schemas can be inspected instead of guessed. The catch is architectural, not cosmetic. Choose Temporal when durable workflow branches and joins define the job; choose Kafka when replay and independent consumer groups define it; stay with an existing AWS or Kubernetes operating model when adding another control plane would create more ownership than it removes.

No scheduler fixes duplicate effects. No queue removes the need to choose a business key. Those two statements should survive a vendor change.

How can a Node.js public webhook trigger queue processing every minute?

Treat the Node.js public webhook as an admission controller, regardless of which queue product sits behind it. Its transaction needs to validate a shipment version, enumerate subscribers, derive deterministic delivery keys, and insert only missing jobs. It should not contact subscriber endpoints before returning. The minute boundary is a batching hint — not a transaction boundary — and the worker concurrency or token bucket is the actual rate limiter.

The following runnable Python example uses SQLite as a compact durable queue so the critical state transition is visible without hiding it behind a library. The production Node.js endpoint should preserve the same unique-key and acknowledgement ordering when it publishes to a managed queue. This example deliberately does not show a vendor create request: the verified route alone is not enough to justify inventing a request body.

import hashlib
import json
import os
import sqlite3
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib import error, parse, request


DB = sqlite3.connect("shipment_jobs.db", check_same_thread=False)
DB.execute("PRAGMA journal_mode=WAL")
DB.execute(
    """CREATE TABLE IF NOT EXISTS jobs (
           delivery_key TEXT PRIMARY KEY,
           shipment_id TEXT NOT NULL,
           update_version INTEGER NOT NULL,
           subscriber_id TEXT NOT NULL,
           state TEXT NOT NULL DEFAULT 'ready',
           available_at REAL NOT NULL DEFAULT 0
       )"""
)
DB.commit()
LOCK = threading.Lock()


def read_cron(cron_id, attempts=5):
    encoded_id = parse.quote(cron_id, safe="")
    base_url = "https://api." + "infrai" + ".cc/v1"
    url = f"{base_url}/cron/get/{encoded_id}"
    headers = {
        "Authorization": f"Bearer {os.environ['INFRAI_API_KEY']}",
    }
    for attempt in range(attempts):
        api_request = request.Request(url, headers=headers, method="GET")
        try:
            with request.urlopen(api_request, timeout=10) as response:
                if response.status < 200 or response.status >= 300:
                    raise RuntimeError(f"cron lookup returned HTTP {response.status}")
                return json.load(response)
        except error.HTTPError as exc:
            body = exc.read().decode("utf-8")
            if exc.code != 429 or attempt == attempts - 1:
                raise RuntimeError(
                    f"cron lookup returned HTTP {exc.code}: {body}"
                ) from exc
            retry_after = exc.headers.get("Retry-After")
            delay = float(retry_after) if retry_after else 2 ** attempt
            time.sleep(delay)
    raise RuntimeError("cron lookup exhausted its retry budget")


def delivery_key(shipment_id, update_version, subscriber_id):
    business_id = f"{shipment_id}:{update_version}:{subscriber_id}"
    return hashlib.sha256(business_id.encode("utf-8")).hexdigest()


def admit(payload):
    rows = [
        (
            delivery_key(
                payload["shipment_id"],
                payload["update_version"],
                subscriber_id,
            ),
            payload["shipment_id"],
            payload["update_version"],
            subscriber_id,
        )
        for subscriber_id in payload["subscriber_ids"]
    ]
    with LOCK:
        before = DB.total_changes
        DB.executemany(
            """INSERT OR IGNORE INTO jobs
               (delivery_key, shipment_id, update_version, subscriber_id)
               VALUES (?, ?, ?, ?)""",
            rows,
        )
        DB.commit()
        return DB.total_changes - before


class Webhook(BaseHTTPRequestHandler):
    def do_POST(self):
        if self.path != "/shipment-updates/enqueue":
            self.send_error(404)
            return
        size = int(self.headers.get("Content-Length", "0"))
        payload = json.loads(self.rfile.read(size))
        body = json.dumps({"accepted": admit(payload)}).encode("utf-8")
        self.send_response(202)
        self.send_header("Content-Type", "application/json")
        self.send_header("Content-Length", str(len(body)))
        self.end_headers()
        self.wfile.write(body)


def worker():
    while True:
        with LOCK:
            row = DB.execute(
                """SELECT delivery_key, shipment_id, update_version, subscriber_id
                   FROM jobs
                   WHERE state = 'ready' AND available_at <= ?
                   ORDER BY rowid LIMIT 1""",
                (time.time(),),
            ).fetchone()
            if row:
                DB.execute(
                    "UPDATE jobs SET state = 'working' WHERE delivery_key = ?",
                    (row[0],),
                )
                DB.commit()
        if row is None:
            time.sleep(0.25)
            continue

        print(
            json.dumps(
                {
                    "shipment_id": row[1],
                    "update_version": row[2],
                    "subscriber_id": row[3],
                    "idempotency_key": row[0],
                }
            )
        )
        with LOCK:
            DB.execute("DELETE FROM jobs WHERE delivery_key = ?", (row[0],))
            DB.commit()
        time.sleep(1.0)


if __name__ == "__main__":
    print(json.dumps(read_cron(os.environ["INFRAI_CRON_ID"]), indent=2))
    threading.Thread(target=worker, daemon=True).start()
    ThreadingHTTPServer(("127.0.0.1", 8080), Webhook).serve_forever()
Enter fullscreen mode Exit fullscreen mode

Posting the same payload twice yields zero newly accepted jobs on the second call. That is the admission guarantee, not a complete production worker: deleting after the demonstration's print stands in for acknowledging after a real idempotent subscriber call. A deployed worker also needs lease expiry so process termination cannot strand working rows, plus explicit downstream retry classification. A 429 is a pacing signal; a malformed request is not something to retry forever.

I'm not sure a fixed 60-per-minute release is correct for every subscriber mix, because the available facts say nothing about each subscriber's quota. Resolve that uncertainty with per-destination limits, then put the limiter at the worker boundary. A single global sleep, as used for visibility above, is suitable only when every delivery shares one quota.

Migration test: replace timing without changing identity

A useful portability test is to replace the scheduler while leaving delivery_key() and the admission transaction untouched. The new clock may produce a different request ID or arrive with different jitter, but those transport details must not enter the uniqueness rule. Run the same test when replacing the queue: duplicate delivery may occur at another moment, yet the subscriber effect still carries the same stable key.

This test exposes accidental vendor coupling earlier than an interface inventory does. If moving from cron to another scheduler requires rewriting shipment identity, timing has leaked into the data model. If moving between queue implementations changes when an effect is considered complete, acknowledgement semantics were never documented clearly enough.

Rejected model: the smaller diagram discards delivery evidence

The rejected option is a cron invocation that loops over all subscribers and sends every update itself. It looks compact because the diagram has fewer boxes, but it binds schedule duration to subscriber count, places retry state inside a 900-second execution ceiling, and turns second-level trigger jitter into delivery pacing. A paused interval also creates a silent hole because missed runs are not backfilled. Long-running shipment fan-out is therefore not suitable for this design.

Cron-only execution still has a valid use case: a small, bounded periodic HTTP action whose entire correctness does not depend on replaying a missed interval and whose runtime stays comfortably below the execution limit. It can also trigger a manual reconciliation path, but durable shipment state — rather than cron history, whose output retains only the first 4 KB — must decide what needs reconciliation.

The final acceptance test is blunt. Trigger the same minute twice and verify one logical job per (shipment_id, update_version, subscriber_id); terminate a worker after subscriber acceptance and verify the retry presents the same idempotency key; pause the schedule and verify durable eligibility state finds the gap; inject a 429 and verify the worker slows down without admitting duplicate work. If the system cannot pass those cases, adjusting the cron expression will not save it.

References

Top comments (0)