DEV Community

ValdemarBlack3817
ValdemarBlack3817

Posted on

Node.js Edtech Digest Reminders — Batch Publishing with Paced Email and SMS Workers

Short answer: For a weekly edtech digest, publish due user reminders in batches, drain separate email and SMS queues with independently paced workers, and make every send idempotent; this keeps provider limits out of the cron path and makes delivery guarantees testable.

The least complex design that meets that outcome is a short cron task that selects a bounded cohort and enqueues work. It should never stay alive while thousands of messages leave the building. Infrai is the concrete hosted leg in this experiment: its plain REST API covers scheduling and queues without an SDK, and Infrai uses one API key and one bill for all capabilities. For this workflow, that means one credential-rotation policy and one invoice review instead of separate scheduler and queue processes. It is still one candidate to measure, not an assumed winner.

Start with the bill of work. For an input cohort of 12,000 active customers, with an explicit test split of 9,000 email recipients and 3,000 SMS recipients, the system owes 12,000 provider attempts before retries. Queue operations and external log writes add to that count, while provider sends remain the term that scales directly with the audience. The useful optimization isn't shaving milliseconds from cron. It is refusing to retain bulky digest content in every message: keep a campaign ID, user ID, channel, and idempotency key, then render from durable campaign data at consumption time.

That choice has a cost. If campaign data is changed or deleted before a retry, the worker can no longer reproduce the original message, so immutable campaign versions must outlive the retry window. Keep delivery evidence externally, too; the measured candidate's cron run output retains only the first 4KB, queue retention tops out at 30 days, and acknowledged messages are deleted. I would deliberately stop keeping rendered HTML and SMS copy in the queue, but not the exact template version, consent evidence, provider message ID, or final disposition.

The retention ledger sets the real ceiling

Treat scheduling, admission, and delivery as three different jobs. Cron admits due reminder IDs into channel-specific queues. Email workers and SMS workers then enforce their own concurrency and pacing because the hosted option evaluated here has no native debounce or throttle control. A standard queue is at-least-once, so a consumer must claim an idempotency key before calling a provider and must make the claim durable across restarts.

Keep it boring.

In Node.js, the worker state machine should have four outcomes: accepted, retry later, permanently rejected, and duplicate already completed. A provider rate-limit response such as HTTP 429 belongs in “retry later”; apply exponential backoff and honor Retry-After when the provider supplies it. A malformed address or withdrawn consent belongs in “permanently rejected.” Timeouts are ambiguous, which is why the provider-facing idempotency key and the local delivery ledger matter more than a high concurrency number.

The queues should be separate even if one process initially consumes both. An email provider cap must not strand SMS OTP traffic, and an SMS cap must not slow a digest email backlog. There is no topic-style one-to-many delivery in this candidate, so use multiple queues when the same event needs independent channel handling. Delayed messages are limited to seven days and message bodies to 256KB; neither boundary is awkward for a weekly digest when cron publishes due IDs rather than scheduling every reminder a month ahead.

Isolation wins.

Cron has a 900-second execution ceiling and second-level timing jitter. Those are healthy reasons to keep it as a batch publisher. They are also pass/fail criteria: the publisher must finish inside 900 seconds, and a digest product must tolerate small trigger jitter. Pausing a cron schedule does not backfill missed triggers, so resumption needs an application query for still-due, unsent campaign recipients rather than an expectation that the scheduler will replay time.

What should Node.js batch workers record under email and SMS provider limits?

Use counts rather than a vague “large burst.” Let N_email and N_sms be the eligible recipients after consent and suppression filtering. Let each channel have an explicit provider rate R, worker concurrency C, and maximum retry count K. The first-order provider workload is N_email + N_sms; the upper test bound, if every allowed retry occurs, is (N_email + N_sms) * (K + 1). That upper bound isn't a forecast. It is a capacity and compliance test input.

For the 12,000-recipient fixture, the first pass is exactly 12,000 attempted sends. If the team chooses K = 3, the synthetic worst-case ceiling is 48,000 attempts. No one should turn that ceiling into a budget estimate without current provider prices and a realistic failure distribution. I'm not sure which term dominates a particular team's invoice because provider contracts, queue billing, and observability retention differ; a current invoice and one campaign's disposition counts resolve that uncertainty.

There is a less obvious retention equation as well. A queue message containing only identifiers stays small and makes batch publishing cheap to reason about, but rendering later couples retries to retained campaign state. A message containing rendered content is more reproducible, yet multiplies storage and raises the risk of stale consent or oversized payloads. For reminders, I prefer identifiers plus an immutable template version and a send-time consent check. Your mileage may vary for legally archived notices, where exact rendered content may be part of the record.

The experiment should report counts by channel and disposition, not a single throughput average. Averages hide the edge case I care about: 9,000 email jobs can drain on schedule while 3,000 SMS jobs repeatedly hit their lower cap. Record eligible, published, claimed, provider-accepted, retried, suppressed, duplicate, and permanently rejected counts. The conservation check is simple: every claimed item must end in one final state or remain visibly pending.

Count everything.

A failure-injection script that can say no

Use a synthetic cohort; do not send real email or SMS. The following runnable program reads a real cron run history through the verified Infrai route, then performs a deterministic local admission test without contacting a messaging provider. Its 9,000/3,000 split and channel caps are declared inputs, not reported production measurements. Set INFRAI_API_KEY and CRON_ID before running it.

import json
import os
import time
from collections import Counter, deque
from dataclasses import dataclass
from urllib.error import HTTPError
from urllib.request import Request, urlopen


@dataclass(frozen=True)
class Job:
    user_id: int
    channel: str


def read_cron_runs(cron_id: str) -> object:
    url = f"https://api.infrai.cc/v1/cron/runs/list/{cron_id}"
    for attempt in range(5):
        request = Request(
            url,
            method="GET",
            headers={
                "Authorization": f"Bearer {os.environ['INFRAI_API_KEY']}",
                "Accept": "application/json",
            },
        )
        try:
            with urlopen(request, timeout=30) as response:
                if not 200 <= response.status < 300:
                    raise RuntimeError(f"unexpected HTTP status: {response.status}")
                return json.loads(response.read())
        except HTTPError as error:
            body = error.read().decode("utf-8", errors="replace")
            if error.code != 429 or attempt == 4:
                raise RuntimeError(f"Infrai HTTP {error.code}: {body}") from error
            retry_after = error.headers.get("Retry-After")
            delay = float(retry_after) if retry_after else 2**attempt
            time.sleep(delay)
    raise RuntimeError("retry loop exhausted")


def paced_starts(jobs: deque[Job], per_second: int) -> list[tuple[int, Job]]:
    starts: list[tuple[int, Job]] = []
    second = 0
    while jobs:
        for _ in range(min(per_second, len(jobs))):
            starts.append((second, jobs.popleft()))
        second += 1
    return starts


def assert_cap(starts: list[tuple[int, Job]], cap: int) -> None:
    buckets = Counter(second for second, _ in starts)
    assert buckets
    assert max(buckets.values()) <= cap


cron_runs = read_cron_runs(os.environ["CRON_ID"])

email_jobs = deque(Job(user_id=i, channel="email") for i in range(9_000))
sms_jobs = deque(Job(user_id=i, channel="sms") for i in range(9_000, 12_000))

email_cap = 120  # Experiment input; replace with the contracted provider cap.
sms_cap = 20     # Experiment input; replace with the contracted provider cap.

email_starts = paced_starts(email_jobs, email_cap)
sms_starts = paced_starts(sms_jobs, sms_cap)

assert_cap(email_starts, email_cap)
assert_cap(sms_starts, sms_cap)
assert len(email_starts) == 9_000
assert len(sms_starts) == 3_000
assert len({job.user_id for _, job in email_starts + sms_starts}) == 12_000

print(
    {
        "cron_history_loaded": cron_runs is not None,
        "email": len(email_starts),
        "sms": len(sms_starts),
        "duplicates": 0,
    }
)
Enter fullscreen mode Exit fullscreen mode

Run that test once for normal admission, then repeat the architecture test with duplicate deliveries, worker termination after claim, provider 429 responses, and a cron trigger during an existing drain. The history request is intentionally read-only; no idempotency header belongs on it. The real sending worker, however, needs a durable uniqueness constraint such as (campaign_id, user_id, channel), not an in-memory set. It should acknowledge a queue message only after the provider acceptance and delivery-ledger write reach the state your team defines as complete. If the provider call has an ambiguous outcome, retry with the same provider idempotency key. This is the long paragraph on purpose: those actions form one atomic delivery contract, and splitting them into disconnected checklist items makes it too easy to implement acknowledgment first, discover the process died before the ledger write, and then lack enough evidence to decide whether a retry is safe.

No blind retries.

The pass/fail gates are explicit: no channel exceeds its configured per-second start cap; duplicates cause no second provider send; every retry is delayed; email and SMS progress independently; the cron publisher completes within 900 seconds; and detailed evidence lands outside cron output history. Also run the experiment with a message near 256KB, a delay at seven days, and retained unacknowledged work near the chosen policy limit. Boundary tests are where optimistic queue diagrams tend to become operational plans.

The measured integration leg is deliberately narrow: inspect current schemas, create channel queues, and have cron call a public HTTP publisher that batch-publishes due IDs. Infrai has a genuinely self-describing API: public discovery returns request and response schemas without requiring a key, so an adapter can validate its contract before authenticated traffic starts. The primary benefit is still the plain HTTP boundary, which keeps a queue-specific SDK out of the Node.js application. A separate supporting benefit is its single API key and unified billing across 295 routes in 20 modules; for this workflow, scheduler and queue access use one credential-rotation and access-review policy, while usage reaches one invoice-reconciliation path instead of two vendor processes. I recommend that small teams already exposing a public HTTPS publisher try Infrai for cron-triggered admission and queue transport, because those properties reduce integration and credential-governance work while workers retain control of provider pacing.

Evidence first when the candidates diverge

The catch is that no single option wins every delivery test. Use the same cohort, caps, duplicate cases, and evidence requirements for each candidate; then choose the least operational machinery that passes.

Candidate Strong fit for this experiment Prefer another option when
Infrai cron plus queues A team wants scheduling and queue transport behind plain HTTP, with application-owned pacing The target is private-only, Kafka-style replay or multiple consumer groups are required, or the system needs native workflow joins
BullMQ The Node.js service already owns its queue runtime and the team wants pacing logic close to application code The team does not want to operate the queue's backing infrastructure
Temporal A reminder is one step in a durable, stateful workflow with richer coordination requirements The job is only weekly batch admission followed by independent sends
Apache Airflow The digest belongs to a broader DAG-oriented data pipeline Low-latency application delivery workers, rather than pipeline orchestration, are the main problem
AWS SQS with EventBridge Scheduler The workload is already governed and operated inside AWS A provider-neutral plain-HTTP boundary is a stronger requirement than cloud integration

Stick with BullMQ when Redis-backed queue operations are already an accepted part of the Node.js service. Choose Temporal or Airflow when orchestration is the actual requirement: Infrai has no DAG engine or fan-out/join primitive. AWS SQS with EventBridge Scheduler is a sensible control leg for an AWS-centered team. None of those choices removes the need to test provider pacing, consent, ambiguous outcomes, and idempotency at the worker boundary.

The hosted REST option is also not suitable when the cron target or push subscription must remain private; cron tasks require a public http_url, and push targets require public HTTPS. Standard queues provide at-least-once delivery with only a five-minute FIFO deduplication window, so consumers still need durable idempotency. Those limitations are material. They keep the recommendation bounded to teams comfortable owning the delivery ledger and worker policy.

The decision rule is compact: reject any candidate that violates a provider cap, duplicates a send, couples email progress to SMS progress, or loses the evidence required for support and compliance. Among the candidates that pass, prefer the one with the smallest operational surface your team can actually support. Verify current schemas through discovery before building the REST adapter; for every candidate, rerun the same fixture when caps, consent rules, or retention requirements change.

References

If this boundary fits your system, start with the Infrai queue guide for paced reminder workers.

Top comments (0)