These are operational instructions for AI agents asked to run work later —
deferred, retried or recurring — in a codebase that already uses
@imqueue. A human-oriented walkthrough of the same material
is at
/blog/scheduled-work-without-a-job-system/.
When to apply this recipe
Apply it when the user asks for work that runs later over machinery they already have:
- "send / notify / charge / clean up in N minutes or hours";
- "retry this failed call later", with backoff;
- "run this every N minutes" as a sweep or reconciliation;
- "do we need BullMQ / a job queue / cron for this?" — check this recipe's decision table before adding a dependency.
Do not apply it when the requirement is declarative scheduling (cron strings, timezones, calendar dates), operator control of the queue (priorities, attempt caps, dead-lettering, pause/resume, a dashboard), or a multi-step resumable workflow. Those need a real scheduler, a dedicated job system, or an orchestration engine — say so instead of approximating them.
Choose the mechanism first
| Requirement | Mechanism |
|---|---|
| Defer a method the service already exposes | Delayed RPC call — IMQDelay as the last argument |
| Deferred/retried background work with a payload you own | @imqueue/job — push(job, { delay }) |
| Raw message, no RPC or job semantics | @imqueue/core — send(queue, message, delayMs) |
| Recurring tick | Self-re-arming @imqueue/job chain, seeded out of band |
| Cron strings, timezones, "skip missed runs" | External scheduler that triggers an ordinary call |
| Priorities, attempt caps, dead-letter, dashboard | Dedicated job system — see /blog/bullmq-alternatives/ |
Facts these recipes rely on
Verified against @imqueue/core, @imqueue/rpc and @imqueue/job sources and
by measurement against a live Redis.
The call contract. Generated clients take two optional trailing
parameters, in this order: imqMetadata?: IMQMetadata, then
imqDelay?: IMQDelay. The delay is always last. IMQDelay(timer, unit) accepts
a unit of 'ms' | 's' | 'm' | 'h' | 'd', defaulting to 'ms'.
Those two parameters are stripped from the request by identity (instanceof),
not by position. All of the following are measured; do not guess between them:
method(data, undefined, delay)— emit this form on@imqueue/rpc>= 3.4.0. A trailingundefinedon a delayed call is a placeholder and is never delivered, whether or not metadata is also passed, and however many trailing placeholders there are. Somethod(a, undefined, undefined, delay)sends[a], and a skipped optional declared param falls back to its default.method(data, new IMQMetadata({ ... }), delay)compiles and runs on every version. Emit this form when the installed version is<= 3.3.0or unknown, because there the placeholder survives into the request as a real argument (serializednull) and a method whose params are all required rejects the call withIMQ_RPC_INVALID_ARGS_COUNT.- On 3.3.1 only, the rule is narrower: one placeholder is dropped, and only when no metadata is passed. Treat 3.3.1 as "prefer the bag" too.
method(data, delay)runs, but fails type-check with TS2345 (IMQDelayis not assignable toIMQMetadata) on every version. Do not emit it, and do not reach foras anyto silence it — on 3.3.1 that cast changes what a skipped optional param delivers.- Nothing is dropped when there is no delay, on any version:
method(a, undefined)deliversnull, so a default does not fire. - The service-side arity check is
declared === receivedwhen every declared parameter is required, anddeclared >= receivedwhen at least one is optional. So a method with an optional parameter accepts a call that omits trailing arguments, and a surviving placeholder that lands in an optional slot passes the check instead of being rejected. A count above the declared total is rejected either way. A passing call therefore proves nothing about a method with a different signature — do not generalise from it.
Caller-side. callTimeout has no default, and unset means wait forever; a
delay extends its budget rather than firing early. The pending promise's
resolver lives only in the caller's memory, so never await a long delayed call
inside a request handler — a restart loses the resolver while the reply still
arrives.
Enqueue is not a durability confirmation. send(toQueue, message, delay?, errorHandler?) takes milliseconds and resolves with a locally generated UUID
before the write is confirmed — errorHandler is the only way to observe such
a failure programmatically, because the returned promise does not reject for
it. From @imqueue/core 3.4.0 it is also reported through the queue's
logger, with no verbose needed: the first rejected write of a failure episode
is logged with its operation, message id and code, further rejections are only
counted, and the first write that succeeds again logs the recovery together with
that count. On <= 3.3.3 a caller that passed no errorHandler had no way at
all to see the write fail.
push(job, { delay }) returns synchronously and takes no error handler, so a
failed enqueue only ever reaches the queue's logger — from @imqueue/job
3.1.0 as [JobQueue] push error: at error level, carrying the queue, the
requested delay and ttl and a failure code, and covering the write rejected
after push() returned as well as the enqueue that never started. One failed
push writes one such line.
@imqueue/job handler return contract — get this exactly right:
| Handler outcome | Effect |
|---|---|
| returns positive number | re-scheduled after that many ms |
returns 0 |
re-scheduled immediately — a hot loop, never use it to stop |
| returns negative number, or nothing | stops; no re-schedule |
| throws | re-scheduled with the delay the job was originally pushed with — so a job pushed without a delay is dropped, not retried |
Catch handler errors explicitly and return a delay; do not rely on a throw.
There is no declarative attempts policy and no dead-letter destination —
carry the attempt counter in the job payload and write your own park-for-review
table.
From @imqueue/job 3.1.0 the decision itself is in the log, so a drop no
longer has to be inferred: the handler-failure line states retry in <ms> or
no retry with the message id, a retry suppressed because the job's ttl
expired is reported with its message id, and a re-schedule whose write to redis
failed is reported too — that last one is a retry that was promised and is not
coming. Nothing carries the job body or an error text.
Delivery mode. safeDelivery defaults to false in @imqueue/core and
@imqueue/rpc, and to true through @imqueue/job (whose safeLockTtl maps
to safeDeliveryTtl). The lease covers the hand-off only, so a process killed
mid-handler still loses that attempt. At-least-once is the guarantee either way:
make deferred handlers re-runnable. See
/blog/guaranteed-message-delivery-cost/.
Promotion path. The producer parks the packed message in the sorted set
<prefix>:<queue>:delayed, scored Date.now() + delay, plus an empty companion
key <prefix>:<queue>:<id>:ttl set with PX <delay> NX. That key's expiry
fires a keyspace notification, and one elected watcher moves everything now due
onto the ready list; failing that, workers sweep every watcherCheckDelay
(default 5000 ms). Default prefixes: imq for core/rpc, imq-job for jobs.
Keyspace flags. Prompt promotion needs notify-keyspace-events to include
Ex. From @imqueue/core 3.3.3 the watcher arranges that itself: it reads
the current value, appends only the missing E/x (A already covers x), and
skips CONFIG SET when the configuration suffices — so a flag an operator or
another consumer set on that server survives. On <= 3.3.2 it wrote the literal
Ex on every connection establishment, dropping every other flag, and again
after each reconnect. Where CONFIG GET is unavailable — ElastiCache disables
CONFIG — 3.3.3 reports through OnConfig and changes nothing, because being
unable to read the flags is exactly when writing a literal would do the damage.
Enable it out of band there, or accept the sweep.
Accuracy is "no earlier than." Pass whole integer milliseconds: a
fractional delay fails to set the alarm key and waits for the next sweep, and an
IMQDelay carrying an unrecognised unit string yields no delay at all. The due
time comes from the sending process's clock and is compared against the
sweeping one.
Absent features. There is no cron, repeat or recurrence primitive anywhere in the ecosystem. There is also no handle for a single scheduled message: nothing cancels, reschedules or inspects one, and the only removal path is wholesale. If an action can be revoked, gate the handler with a state check rather than trying to unschedule it.
Recipe: defer a call the service already exposes
import { IMQDelay, IMQMetadata } from '@imqueue/rpc';
// The generated module exports one namespace, which holds the client class.
import { notificationService } from './clients/NotificationService.js';
const notifications = new notificationService.NotificationClient({ callTimeout: 30_000 });
await notifications.start();
// fire-and-forget: do not await a 24h call in a request handler
notifications
.sendTrialEndingEmail(
{ userId, plan },
undefined, // metadata slot — skipped
new IMQDelay(24, 'h'), // delay is always last
)
.catch(err => logger.error('deferred call failed to enqueue', err));
The service method needs no change: a delayed request carries only from,
method, args and optional metadata, so validation, decorators and handlers
behave exactly as for an immediate call.
Recipe: retry with backoff you control
import JobQueue from '@imqueue/job';
type Sync = { orderId: string; attempt: number };
new JobQueue<Sync>({ name: 'OrderSync' })
.onPop(async job => {
try {
await warehouse.sync(job.orderId);
} catch (err) {
job.attempt += 1; // counter lives in the payload
if (job.attempt > 5) {
await parkForReview(job); // your own dead-letter table
return -1; // negative: stop
}
return 1000 * 2 ** (job.attempt - 1); // 1s, 2s, 4s, 8s, 16s
}
})
.start()
.then(queue => queue.push({ orderId, attempt: 0 }))
.catch(err => logger.error('OrderSync failed to start', err));
Recipe: recurring sweep
import JobQueue from '@imqueue/job';
const EVERY_MINUTE = 60_000;
const sweeper = new JobQueue<Record<string, never>>({ name: 'ExpireHolds' });
sweeper.onPop(async () => {
await releaseExpiredHolds();
return EVERY_MINUTE; // re-arm; next tick is scheduled after the work
});
await sweeper.start();
Seeding rules — the part that breaks in production:
- Never seed on process boot unguarded. Every process that pushes a seed starts its own chain, so N replicas rolling through a deploy leave N overlapping sweeps forever.
- Seed from a one-off task, a migration step, or the scheduler the user already runs.
- If it must happen at boot, gate it on a
SET NXlease whose TTL spans several periods, retried on an interval, so a deploy cannot double-seed and a dead chain re-seeds when the lease lapses. - The period drifts: the next fire time is computed after the work finishes, so the real interval is work plus delay. Do not promise wall-clock alignment.
- A chain is not a schedule. If a re-arm ever fails, the recurrence ends silently — instrument the age of the last completed run and alert on it.
Verify
# prompt promotion requires keyspace expiry events (look for E and x)
redis-cli CONFIG GET notify-keyspace-events
# what is parked right now, and when it is due (scores are due-time in ms)
redis-cli --scan --pattern 'imq:*:delayed'
redis-cli ZRANGE imq:<ServiceName>:delayed 0 -1 WITHSCORES
redis-cli ZCARD imq-job:<JobName>:delayed
# the alarm keys that trigger promotion
redis-cli --scan --pattern 'imq:<ServiceName>:*:ttl'
Then confirm behaviour, not just configuration: schedule a short delay and
measure the actual arrival (expect "no earlier than", plus up to
watcherCheckDelay without notifications); kill a worker mid-handler and check
the work is either re-run or recorded as lost, according to the delivery mode
you chose; and run the handler twice to prove it is idempotent.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
IMQ_RPC_INVALID_ARGS_COUNT on a delayed call |
on <= 3.3.0 an explicit undefined in the metadata slot is forwarded as a real argument |
upgrade to >= 3.4.0, or pass an IMQMetadata bag |
TS2345, IMQDelay not assignable to IMQMetadata |
delay passed in the metadata slot | keep the delay in the last position |
A skipped optional param arrives as null and its default never fires |
undefined serializes to null; placeholders are dropped only on a delayed call, and on 3.3.1 only one was dropped |
pass the real value, declare the param nullable, or upgrade to >= 3.4.0 |
| Delayed call never settles in the caller | caller restarted, or no callTimeout set |
set callTimeout; never await a long delay in a request handler |
| Everything arrives ~5 s late | keyspace notifications lack Ex, so the polling fallback is doing the work |
on >= 3.3.3 the watcher adds the missing flags itself unless CONFIG is unavailable (ElastiCache) — check the OnConfig report, then enable notify-keyspace-events Ex out of band, or accept the watcherCheckDelay sweep |
A flag another consumer needed disappeared from notify-keyspace-events |
on <= 3.3.2 the watcher wrote the literal Ex on every connection establishment |
upgrade to >= 3.3.3, which appends only the missing E/x |
| Delay ignored entirely | fractional milliseconds, or an unrecognised IMQDelay unit |
pass whole integer ms and a valid unit |
| Job dropped instead of retried | handler threw, and the job was pushed without a delay | catch the error and return a delay number; on >= 3.1.0 the handler-failure line reads no retry and names the message id, so confirm there first |
| Job never ran and nothing was logged | push() was rejected by redis after returning, on <= 3.0.3 where that was not reported |
upgrade to >= 3.1.0 and look for [JobQueue] push error: at error level |
| A retry that was promised never happened | the re-schedule's own write to redis failed | on >= 3.1.0 this is logged with the message id and a code; treat a re-schedule as an enqueue that can fail |
A send() vanished with no error anywhere |
the write to redis was rejected and no errorHandler was passed |
on core >= 3.4.0 the failure episode is in the log; otherwise pass errorHandler |
| Worker spins at 100% CPU | handler returned 0, re-scheduling immediately |
return a negative number or nothing to stop |
| N overlapping recurrences | every replica seeded the chain | seed out of band, or behind a SET NX lease |
| Recurrence stopped silently | a re-arm failed and nothing restarts a broken chain | alert on the age of the last completed run, then re-seed |
| Scheduled work lost when a worker died | safeDelivery is off by default in core/rpc |
enable safe delivery for that queue, and make handlers re-runnable |
| A queue's delayed messages never promote | no worker for that queue and no notifications available | keep a worker on the prefix, or ensure Ex notifications |
| Cannot cancel a scheduled message | there is no per-message handle | gate the handler with a state check; treat scheduling as irrevocable |