npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

@effectmq/core

v0.3.0

Published

Effect-based message queue

Downloads

168

Readme

effectmq

A typed, Redis-backed task queue for Effect 4. Define work with schemas, run handlers as Effects, and keep payloads, results, and failures typed from producer to worker.

Website · Documentation · Getting started · API reference · npm

effectmq is for background jobs that need durable Redis state without becoming a workflow engine: send an email, resize an image, refresh a cache, or materialize a scheduled report. It provides:

  • schema-checked payloads, successes, and failures;
  • at-least-once delivery with fenced attempts and stalled-worker recovery;
  • Effect Schedule retries, delayed offers, deduplication, and durable cron;
  • bounded local worker concurrency, graceful draining, and maintenance;
  • typed lifecycle streams plus wait and execute for durable results;
  • durable application events with named subscriptions, independent acknowledgements, optional deadlines, and deletion or archival.

Install

These docs target 0.3.0-rc.1 with Effect 4.0.0-rc.115. Check the current release record for publication status before installing. 0.3.0-rc.0 uses the older Effect beta dependency set.

pnpm add @effectmq/[email protected] [email protected] @effect/[email protected]

[!IMPORTANT] effectmq currently targets the Effect 4 release candidate and is not compatible with the stable Effect 3 release. Pin the versions shown above. Node.js 22.19 or newer is required; CI verifies Node.js 22 and 24.

The package includes its pooled NodeRedisPool implementation. @effect/platform-node is only needed by these examples for NodeRuntime.


Quick start

Define a task, enqueue work, process it. The whole loop:

import { Effect, Schema } from "effect";
import { NodeRuntime } from "@effect/platform-node";
import { Task, TaskEngine, TaskQueue } from "@effectmq/core";

const SendEmail = Task.make({
  name: "send-email",
  payload: { to: Schema.String, subject: Schema.String },
  success: Schema.String,
  error: Schema.Never,
});

const emails = TaskQueue.make("emails", SendEmail);

const program = Effect.gen(function* () {
  yield* TaskQueue.offer(emails, { to: "[email protected]", subject: "Welcome" });

  yield* TaskQueue.complete(emails, (task) =>
    Effect.succeed(`provider:${task.payload.to}`),
  );
});

// The engine + its Redis layer: the only wiring you need to run the above.
const AppLayer = TaskEngine.layer({
  redis: { url: "redis://localhost:6379" },
});

program.pipe(Effect.provide(AppLayer), NodeRuntime.runMain);

That is the complete producer-to-worker loop. For a clean-room walkthrough with Redis startup and expected output, follow the getting-started tutorial.


Runtime setup

TaskEngine.layer() is the complete Node live graph: it provides the engine, cryptographic identity generation, and the retained Redis pool, role, and health services:

import { TaskEngine } from "@effectmq/core";

const AppLayer = TaskEngine.layer({
  redis: { url: "redis://localhost:6379" },
});

Use TaskEngine.layerNoDeps() when composing a custom RedisPool implementation. NodeRedisPool.layer() remains available independently and accepts node-redis client options. It establishes separate producer, worker, and maintenance pools when the Layer starts. It supports standalone Redis and Sentinel; Redis Cluster fails startup because queue transitions use multi-key atomic scripts. See the operations runbook for TLS, ACL, bounded-pool, persistence, failover, health, and shutdown guidance.

The tested platform matrix is in the support policy, and reproducible throughput/tail-latency results are published as performance evidence.

TaskEngine is the machinery underneath: atomic Lua scripts, leases, and the lists tasks move between. Provide its layer once; application code normally lives in TaskQueue, Worker, and Scheduler.


Define a task

A task is a schema, not a function. You declare what goes in (payload), what a success looks like, and what a failure looks like. The idempotencyKey decides what "the same task" means: offer the same key twice and you get one task, not two.

A tagged error makes failures pattern-matchable downstream, so reach for Schema.TaggedError rather than a bare struct.

import { Effect, Schedule, Schema, Stream } from "effect";
import { Task, TaskQueue, type TaskHandler, Worker } from "@effectmq/core";

class EmailRejected extends Schema.TaggedError<EmailRejected>()(
  "EmailRejected",
  { reason: Schema.String },
) {}

const SendEmail = Task.make({
  name: "send-email",
  payload: { to: Schema.String, subject: Schema.String },
  success: Schema.String,
  error: EmailRejected,
  idempotencyKey: (p) => `email:${p.to}:${p.subject}`,
  retry: Schedule.exponential("1 second"),
});

const emails = TaskQueue.make("emails", SendEmail);

const sendViaProvider = (payload: { readonly to: string }) =>
  Effect.succeed(`provider:${payload.to}`);

Offer work, then do it

offer enqueues a payload. complete takes the next task, runs your handler, reports the outcome back to the engine, and returns the task's id (a failing handler is routed per the queue's failure policy). The task your handler receives is fully decoded: task.payload is the real object, not a JSON string.

const offerAndCompleteProgram = Effect.gen(function* () {
  yield* TaskQueue.offer(emails, {
    to: "[email protected]",
    subject: "Welcome",
  });

  // Take one task, run it, report the outcome. Resolves with the task id.
  const taskId = yield* TaskQueue.complete(emails, (task) =>
    Effect.gen(function* () {
      const id = yield* sendViaProvider(task.payload); // your code
      return id; // matches successSchema
      // ...or `yield* new EmailRejected({ reason })` to fail with the typed error
    }),
  );
  // Success resolves per onSuccessPolicy; a failing handler is routed per onFailurePolicy.
});

One complete processes one task. To process many, set your workers up accordingly.

Predefine the handler

Handlers are functions and workers are Effects, so both are values you can name once and reuse. Type a handler with TaskHandler to declare it next to the task definition before any queue exists. Bind it with complete for one task, or use the managed Worker shown below for a long-running process:

// Declared against the task definition — no queue in sight yet.
const handleSendEmail: TaskHandler<
  typeof SendEmail.payloadSchema,
  typeof SendEmail.successSchema,
  typeof SendEmail.errorSchema
> = (task) => sendViaProvider(task.payload);

// Bound to a queue: an Effect that takes one task and runs it to completion.
const sendEmailWorker = TaskQueue.complete(emails, handleSendEmail);

const repeatedWorkerProgram = Effect.gen(function* () {
  yield* sendEmailWorker; // process one task...
  yield* sendEmailWorker.pipe(Effect.repeat(Schedule.forever)); // ...or loop forever
});

Streaming & events

The engine publishes a lifecycle event to a per-queue Redis Stream every time a task changes state. TaskQueue.stream hands you those events as an Effect Stream, decoded against your queue's schemas: task.created and task.updated carry fully-typed tasks, task.failed carries your typed error or a built-in stalled/canceled error, task.completed carries your typed success value, and task.moved reports the list transition.

const watch = TaskQueue.stream(emails).pipe(
  Stream.runForEach((event) => Effect.log(event._tag, event.taskId)),
);

You can also wait on a specific generation. wait reads durable state and consumes lifecycle events until that generation settles, returning its success value or a TaskFailed error whose failure contains the typed handler or built-in failure. Retriable failure events do not end the wait. execute is the offer-and-wait shortcut.

// Offer + await the result in one call.
const executeMessage = TaskQueue.execute(emails, {
  to: "[email protected]",
  subject: "Welcome",
}); // resolves with the success value, or TaskFailed with an EmailRejected failure

// Or await a task you already offered.
const offerAndWait = Effect.gen(function* () {
  const task = yield* TaskQueue.offer(emails, {
    to: "[email protected]",
    subject: "Welcome",
  });
  return yield* TaskQueue.wait(emails, task.handle);
});

wait reads durable state, subscribes from the handle's authoritative Redis cursor, and rechecks state after subscription, so completion before or during subscription is observed. Streams use a blocking Redis read (default two-second block); persist a retained cursor when building a resumable event consumer.


Run a worker

complete processes exactly one task. For a long-running process, Worker provides bounded local concurrency, lease supervision, maintenance, and graceful draining:

const worker = Worker.make(emails, handleSendEmail, { concurrency: 5 });
const program = Worker.run(worker);

concurrency is local to one worker process (valid values: 1–1000). Run more processes to fan out. Distributed/global concurrency and rate limits require external coordination; effectmq does not pretend a process-local semaphore can enforce them. The queue fences each attempt, but handlers remain at-least-once, so make external side effects idempotent.


"But Effect already has Workflow"

It does, and it's excellent, for a different problem. Effect Workflow is durable execution: long-running, multi-step sagas that survive process death, resume exactly where they left off, and persist every intermediate step so the whole history can be replayed. It leans on clustering and sharding; nodes have to be live and coordinated; the durability is total because the use case demands it.

That power has a price that sometimes isn't worth paying. Sometimes you don't have a saga. You have a job. "Send this email." "Resize that image." "Spawn 5 AI agents to complete these tasks." There's no multi-step history worth replaying; there's a payload, a handler, and an outcome. Reaching for durable execution there is like renting a shipping container to mail a letter.

effectmq works whether you have a single worker running in a separate fiber or a hundred distributed across multiple processes.

So:

If Workflow is Temporal (durable, replayable, cluster-coordinated orchestration) then this is BullMQ: a queue. You put work in, workers take attempts, and Redis keeps unfinished work recoverable across retries and worker loss. No workflow replay log or shard map—just a queue with Effect's types and primitives.

Pick durable execution when the process is the thing you can't afford to lose. Pick a queue when the work is.


Scheduling

For recurring work, Scheduler.make durably materializes an ordinary queue task for each selected cron tick. The task id is derived from the schedule name and nominal tick time, so competing schedulers and crash recovery can safely re-offer it. A managed worker executes the task with the queue's normal leases, retries, and at-least-once delivery semantics.

import { Cron, Effect, Schema } from "effect";
import { Scheduler, Task, TaskQueue, Worker } from "@effectmq/core";

const reportTask = Task.make({
  name: "nightly-report-task",
  payload: { scheduledAt: Schema.String },
  success: Schema.Void,
  error: Schema.String,
});
const reportQueue = TaskQueue.make("nightly-reports", reportTask);
const schedule = Scheduler.make({
  name: "nightly-report",
  cron: Cron.parseUnsafe("0 2 * * *", "UTC"),
  queue: reportQueue,
  payload: (tick) => ({ scheduledAt: tick.scheduledAt.toISOString() }),
  missed: { _tag: "coalesce" },
});
const worker = Worker.make(reportQueue, ({ payload }) =>
  Effect.log(`Building report for ${payload.scheduledAt}`),
);

Notes

  • Completion policies. offer accepts onSuccessPolicy and onFailurePolicy, each one of delete | keep | mark-as-success | mark-as-failure. They decide where a finished task lands: gone, quietly retained, or parked on the success/failed list for inspection. Defaults are delete.
  • Retries. Declare retry on the task definition (Task.make) as an Effect Schedule — or a { while, until, times, schedule } options object. On failure the next run time is computed from the schedule and the task lands on the scheduled list until then; when the schedule is exhausted, the failure policy applies. The task-level maxRetries cap defaults to 5; pass null there for an intentionally unbounded cap. An offer may override the cap with a finite non-negative number. Built-in canceled or stalled failures are not retried by the handler schedule.
  • Idempotency. The idempotencyKey is the task id. By default, offering the same key returns the existing generation unchanged; replacement requires explicit new-generation mode.
  • Delays. offer(..., { delay }) schedules the task for the future; it sits on the scheduled list until its time comes.
  • The engine. TaskEngine is the low-level, Lua-backed layer all of this sits on. You provide its layer; you rarely call it directly.

Go deeper

| If you need to… | Read… | | --- | --- | | learn the library from a running example | Getting started | | process, schedule, retry, or await tasks | How-to guides | | look up exact API behavior | API reference | | understand delivery and task identity | Delivery guarantees · Idempotent offers | | operate Redis and plan upgrades | Operations · Upgrade and rollback | | inspect architecture and storage contracts | Architecture · Runtime boundaries · Storage protocol v1 | | evaluate support and performance | Support policy · Performance evidence · Soak evidence | | release the package | Release process · Current release record |


License

MIT.

Task progress history

Declare a progress schema on Task.make to enable task-owned progress and lifecycle history. Managed handlers receive context.progress(value) as their second argument; TaskQueue.readEvents(queue, handle, { after, limit }) reads typed pages while work runs and after completion when the task record is kept. History is unlimited by default, with optional oldest-first count trimming, and disappears when its task record is removed. See the runnable progress guide.