{"_id":"better-effect-mq","_rev":"4-8755d0149681cf8ec3949b4836f871a1","name":"better-effect-mq","dist-tags":{"latest":"0.1.3"},"versions":{"0.1.0":{"name":"better-effect-mq","version":"0.1.0","keywords":["better-effect","better-result","message-queue","queue","typescript"],"license":"MIT","_id":"better-effect-mq@0.1.0","maintainers":[{"name":"nitoba","email":"nito.ba.dev@gmail.com"}],"homepage":"https://github.com/nitoba/better-effect#readme","bugs":{"url":"https://github.com/nitoba/better-effect/issues"},"dist":{"shasum":"331336f9b0c966548a1878d169cd226902cded04","tarball":"https://registry.npmjs.org/better-effect-mq/-/better-effect-mq-0.1.0.tgz","fileCount":27,"integrity":"sha512-JKdrtECAu9YA09hevVzeVGkWEeHqYEnwEDpq6vHzFfb1n5Yr0oTYcic/WuqZDEdRP4fBRIeDVLpaUwZYQGawAg==","signatures":[{"sig":"MEQCIF3qQIAnAzBtEfGGnS1h+ab2PteO0rpVec2l6nQn5TEcAiA2UrNNqoyEMrAsiaZVe2InbHNoYM7Mz6XMcU2K/n+E6A==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":2917831},"type":"module","_from":"file:/tmp/better-effect-mq-publish.XMRjWj/better-effect-mq-0.1.0.tgz","exports":{".":"./dist/index.mjs","./testing":"./dist/testing.mjs","./package.json":"./package.json"},"scripts":{"dev":"tsdown --watch","lint":"oxlint --type-aware .","test":"bun test","build":"tsdown","check":"bun run typecheck && bun run test:types && bun run test:types:performance:clean-dist && bun test && bun run format:check && bun run build && bun run typecheck:examples && bun run test:examples && bun run test:package-boundaries && bun run test:package-consumer && bun run publint && bun run release:dry && bun run lint","format":"oxfmt --write .","publint":"publint","lint:fix":"oxlint --type-aware --fix .","typecheck":"tsc --noEmit","test:types":"tsc -p tests/types/tsconfig.json --pretty false","release:dry":"bun ../../scripts/release-artifact.ts --package better-effect-mq","format:check":"oxfmt --check .","test:examples":"bun examples/producer-only/main.ts && bun examples/worker/main.ts && bun examples/testing/main.ts && bun examples/composition/main.ts","prepublishOnly":"bun run check","typecheck:examples":"tsc -p examples/producer-only/tsconfig.json --pretty false && tsc -p examples/worker/tsconfig.json --pretty false && tsc -p examples/testing/tsconfig.json --pretty false && tsc -p examples/composition/tsconfig.json --pretty false","test:package-consumer":"bun tests/package/external-consumer.ts","test:types:performance":"cd ../.. && bun scripts/type-system-performance.ts --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget","test:package-boundaries":"bun tests/package/boundaries.ts","test:types:performance:clean-dist":"cd ../.. && bun scripts/type-system-performance.ts --clean-dist --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget"},"_npmUser":{"name":"nitoba","email":"nito.ba.dev@gmail.com"},"_resolved":"/tmp/better-effect-mq-publish.XMRjWj/better-effect-mq-0.1.0.tgz","_integrity":"sha512-JKdrtECAu9YA09hevVzeVGkWEeHqYEnwEDpq6vHzFfb1n5Yr0oTYcic/WuqZDEdRP4fBRIeDVLpaUwZYQGawAg==","repository":{"url":"git+https://github.com/nitoba/better-effect.git","type":"git"},"_npmVersion":"11.19.0","description":"Experimental message-queue foundations for better-effect.","directories":{},"sideEffects":false,"_nodeVersion":"26.8.1","publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"oxfmt":"^0.63.0","oxlint":"^1.78.0","tsdown":"^0.22.14","publint":"^0.3.23","lefthook":"^2.1.10","@types/bun":"1.4.1","typescript":"^7.0.2","better-effect":"0.13.0","better-result":"^3.0.0","oxlint-tsgolint":"^7.0.2001"},"peerDependencies":{"typescript":">=6.0.0","better-effect":">=0.13.0 <0.14.0","better-result":"^3.0.0"},"_npmOperationalInternal":{"tmp":"tmp/better-effect-mq_0.1.0_1788918097034_0.09215831459754842","host":"s3://npm-registry-packages-npm-production"}},"0.1.1":{"name":"better-effect-mq","version":"0.1.1","keywords":["better-effect","better-result","message-queue","queue","typescript"],"license":"MIT","_id":"better-effect-mq@0.1.1","maintainers":[{"name":"nitoba","email":"nito.ba.dev@gmail.com"}],"homepage":"https://github.com/nitoba/better-effect#readme","bugs":{"url":"https://github.com/nitoba/better-effect/issues"},"dist":{"shasum":"d9cad3c6c9451e657ce440c53cbeea4452d5831b","tarball":"https://registry.npmjs.org/better-effect-mq/-/better-effect-mq-0.1.1.tgz","fileCount":27,"integrity":"sha512-5NPXu8PxPtzeBAGsOQTiEUP8t27Rzx76+QORnPQ1ujxobBn/sLvUa5LNOnV0wsyXFdfY5KXPR+m6xULewYOlkw==","signatures":[{"sig":"MEYCIQCiZhbgNCAIeJt5csuaZVanX1rtF80L97at7ArST16Q6wIhALJiyQS2HtcIi8gR5+5fUST9zK/gXdQEZdljEs3NXx1X","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"attestations":{"url":"https://registry.npmjs.org/-/npm/v1/attestations/better-effect-mq@0.1.1","provenance":{"predicateType":"https://slsa.dev/provenance/v1"}},"unpackedSize":2879872},"type":"module","exports":{".":"./dist/index.mjs","./testing":"./dist/testing.mjs","./package.json":"./package.json"},"gitHead":"3d65f9d212aa71c4b5b3fef26a20dbf14a054a9e","scripts":{"dev":"tsdown --watch","lint":"oxlint --type-aware .","test":"bun test","build":"tsdown","check":"bun run typecheck && bun run test:types && bun run test:types:performance:clean-dist && bun test && bun run format:check && bun run build && bun run typecheck:examples && bun run test:examples && bun run test:package-boundaries && bun run test:package-consumer && bun run publint && bun run release:dry && bun run lint","format":"oxfmt --write .","publint":"publint","lint:fix":"oxlint --type-aware --fix .","typecheck":"tsc --noEmit","test:types":"tsc -p tests/types/tsconfig.json --pretty false","release:dry":"bun ../../scripts/release-artifact.ts --package better-effect-mq","format:check":"oxfmt --check .","test:examples":"bun examples/producer-only/main.ts && bun examples/worker/main.ts && bun examples/testing/main.ts && bun examples/composition/main.ts && bun examples/flow/main.ts","prepublishOnly":"bun run check","typecheck:examples":"tsc -p examples/producer-only/tsconfig.json --pretty false && tsc -p examples/worker/tsconfig.json --pretty false && tsc -p examples/testing/tsconfig.json --pretty false && tsc -p examples/composition/tsconfig.json --pretty false && tsc -p examples/flow/tsconfig.json --pretty false","test:package-consumer":"bun tests/package/external-consumer.ts","test:types:performance":"cd ../.. && bun scripts/type-system-performance.ts --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget","test:package-boundaries":"bun tests/package/boundaries.ts","test:types:performance:clean-dist":"cd ../.. && bun scripts/type-system-performance.ts --clean-dist --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget"},"_npmUser":{"name":"GitHub Actions","email":"npm-oidc-no-reply@github.com","trustedPublisher":{"id":"github","oidcConfigId":"oidc:99968650-74f7-4f0c-a075-92aed7cb45da"}},"repository":{"url":"git+https://github.com/nitoba/better-effect.git","type":"git"},"_npmVersion":"11.6.2","description":"Experimental message-queue foundations for better-effect.","directories":{},"sideEffects":false,"_nodeVersion":"24.20.0","publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"oxfmt":"^0.63.0","oxlint":"^1.78.0","tsdown":"^0.22.14","publint":"^0.3.23","lefthook":"^2.1.10","@types/bun":"1.4.1","typescript":"^7.0.2","better-effect":"0.13.0","better-result":"^3.0.0","oxlint-tsgolint":"^7.0.2001"},"peerDependencies":{"typescript":">=6.0.0","better-effect":">=0.13.0 <0.14.0","better-result":"^3.0.0"},"_npmOperationalInternal":{"tmp":"tmp/better-effect-mq_0.1.1_1788987367799_0.17637732160552821","host":"s3://npm-registry-packages-npm-production"}},"0.1.2":{"name":"better-effect-mq","version":"0.1.2","keywords":["better-effect","better-result","message-queue","queue","typescript"],"license":"MIT","_id":"better-effect-mq@0.1.2","maintainers":[{"name":"nitoba","email":"nito.ba.dev@gmail.com"}],"homepage":"https://github.com/nitoba/better-effect#readme","bugs":{"url":"https://github.com/nitoba/better-effect/issues"},"dist":{"shasum":"f98280b5d9e5765a9a9a7d9f76a41a61ea433521","tarball":"https://registry.npmjs.org/better-effect-mq/-/better-effect-mq-0.1.2.tgz","fileCount":27,"integrity":"sha512-pzLpokGd3xdb+pIBeAXMc6dwpRFsZmL++8lL40tG+cEsi9VYy7MlN9jUXHXlXtYVnRuVfJnluE2MKu6RBjw88Q==","signatures":[{"sig":"MEYCIQDB4aucxxNNh/7w920aKTZgWzue86+NidVzu4LQc1BClwIhAL1a5x22jx+lmlRe3xJjdwx4rh5Ntj4rxfdF3//Y9uQU","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"attestations":{"url":"https://registry.npmjs.org/-/npm/v1/attestations/better-effect-mq@0.1.2","provenance":{"predicateType":"https://slsa.dev/provenance/v1"}},"unpackedSize":2879972},"type":"module","exports":{".":"./dist/index.mjs","./testing":"./dist/testing.mjs","./package.json":"./package.json"},"gitHead":"dd05239e5b35b351037f153803fc5367af9b2250","scripts":{"dev":"tsdown --watch","lint":"oxlint --type-aware .","test":"bun test","build":"tsdown","check":"bun run typecheck && bun run test:types && bun run test:types:performance:clean-dist && bun test && bun run format:check && bun run build && bun run typecheck:examples && bun run test:examples && bun run test:package-boundaries && bun run test:package-consumer && bun run publint && bun run release:dry && bun run lint","format":"oxfmt --write .","publint":"publint","lint:fix":"oxlint --type-aware --fix .","typecheck":"tsc --noEmit","test:types":"tsc -p tests/types/tsconfig.json --pretty false","release:dry":"bun ../../scripts/release-artifact.ts --package better-effect-mq","format:check":"oxfmt --check .","test:examples":"bun examples/producer-only/main.ts && bun examples/worker/main.ts && bun examples/testing/main.ts && bun examples/composition/main.ts && bun examples/flow/main.ts","prepublishOnly":"bun run check","typecheck:examples":"tsc -p examples/producer-only/tsconfig.json --pretty false && tsc -p examples/worker/tsconfig.json --pretty false && tsc -p examples/testing/tsconfig.json --pretty false && tsc -p examples/composition/tsconfig.json --pretty false && tsc -p examples/flow/tsconfig.json --pretty false","test:package-consumer":"bun tests/package/external-consumer.ts","test:types:performance":"cd ../.. && bun scripts/type-system-performance.ts --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget","test:package-boundaries":"bun tests/package/boundaries.ts","test:types:performance:clean-dist":"cd ../.. && bun scripts/type-system-performance.ts --clean-dist --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget"},"_npmUser":{"name":"GitHub Actions","email":"npm-oidc-no-reply@github.com","trustedPublisher":{"id":"github","oidcConfigId":"oidc:99968650-74f7-4f0c-a075-92aed7cb45da"}},"repository":{"url":"git+https://github.com/nitoba/better-effect.git","type":"git"},"_npmVersion":"11.6.2","description":"Experimental message-queue foundations for better-effect.","directories":{},"sideEffects":false,"_nodeVersion":"24.20.0","publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"oxfmt":"^0.63.0","oxlint":"^1.78.0","tsdown":"^0.22.14","publint":"^0.3.23","lefthook":"^2.1.10","@types/bun":"1.4.1","typescript":"^7.0.2","better-effect":"0.14.0","better-result":"^3.0.0","oxlint-tsgolint":"^7.0.2001"},"peerDependencies":{"typescript":">=6.0.0","better-effect":">=0.14.0 <0.15.0","better-result":"^3.0.0"},"_npmOperationalInternal":{"tmp":"tmp/better-effect-mq_0.1.2_1789064764943_0.1953273152837045","host":"s3://npm-registry-packages-npm-production"}},"0.1.3":{"_id":"better-effect-mq@0.1.3","bugs":{"url":"https://github.com/nitoba/better-effect/issues"},"dist":{"shasum":"0a5757eed2ed4d42ab641c54d450afb381360571","tarball":"https://registry.npmjs.org/better-effect-mq/-/better-effect-mq-0.1.3.tgz","fileCount":28,"integrity":"sha512-EybsNQ9JJf+2B5n1E+rysmmDhaKxux1DgX6nC23Q9+7pBb7NIbXRAU6oBlSS7RC/t2moWjPV5QSXnSvcHdiPNA==","signatures":[{"sig":"MEUCIQCJVt3Yu44DvGwFuRWqCVb3xPQ4D3AIJn+TdHLb2t6JEQIgeGsQ6EUo1viliA908J3kX9dxPVYn0ee9ToHdIDH0tnE=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"},{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIQDDExESUF5nhpICaPm6JG74Weu+eUD0P/Hve62pQBeFFAIgKbqTed1SP3cMBS8jX/i6hbSAzpbnpqqR/eBmXaWGM/0="}],"attestations":{"url":"https://registry.npmjs.org/-/npm/v1/attestations/better-effect-mq@0.1.3","provenance":{"predicateType":"https://slsa.dev/provenance/v1"}},"unpackedSize":2896062},"name":"better-effect-mq","type":"module","exports":{".":"./dist/index.mjs","./testing":"./dist/testing.mjs","./package.json":"./package.json"},"gitHead":"fb0f3b8e3deaea03724a6fa4e1be095e343ba7c5","license":"MIT","scripts":{"dev":"tsdown --watch","lint":"oxlint --type-aware .","test":"bun test","build":"tsdown","check":"bun run typecheck && bun run test:types && bun run test:types:performance:clean-dist && bun test && bun run format:check && bun run build && bun run typecheck:examples && bun run test:examples && bun run test:package-boundaries && bun run test:package-consumer && bun run publint && bun run release:dry && bun run lint","format":"oxfmt --write .","publint":"publint","lint:fix":"oxlint --type-aware --fix .","typecheck":"tsc --noEmit","test:types":"tsc -p tests/types/tsconfig.json --pretty false","release:dry":"bun ../../scripts/release-artifact.ts --package better-effect-mq","format:check":"oxfmt --check .","test:examples":"bun examples/producer-only/main.ts && bun examples/worker/main.ts && bun examples/testing/main.ts && bun examples/composition/main.ts && bun examples/flow/main.ts","prepublishOnly":"bun run check","typecheck:examples":"tsc -p examples/producer-only/tsconfig.json --pretty false && tsc -p examples/worker/tsconfig.json --pretty false && tsc -p examples/testing/tsconfig.json --pretty false && tsc -p examples/composition/tsconfig.json --pretty false && tsc -p examples/flow/tsconfig.json --pretty false","test:package-consumer":"bun tests/package/external-consumer.ts","test:types:performance":"cd ../.. && bun scripts/type-system-performance.ts --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget","test:package-boundaries":"bun tests/package/boundaries.ts","test:types:performance:clean-dist":"cd ../.. && bun scripts/type-system-performance.ts --clean-dist --scenarios=job-registry,job-store,job-producer,worker-handlers --job-sizes=10,50,100,250 --producer-sizes=10,50,100 --worker-sizes=10,50,100 --check-budget"},"version":"0.1.3","_npmUser":{"name":"GitHub Actions","email":"npm-oidc-no-reply@github.com","trustedPublisher":{"id":"github","oidcConfigId":"oidc:99968650-74f7-4f0c-a075-92aed7cb45da"}},"homepage":"https://github.com/nitoba/better-effect#readme","keywords":["better-effect","better-result","message-queue","queue","typescript"],"repository":{"url":"git+https://github.com/nitoba/better-effect.git","type":"git"},"_npmVersion":"11.6.2","description":"Experimental message-queue foundations for better-effect.","directories":{},"maintainers":[{"name":"nitoba","email":"nito.ba.dev@gmail.com"}],"sideEffects":false,"_nodeVersion":"24.20.0","publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"oxfmt":"^0.63.0","oxlint":"^1.78.0","tsdown":"^0.22.14","publint":"^0.3.23","lefthook":"^2.1.10","@types/bun":"1.4.1","typescript":"^7.0.2","better-effect":"0.14.0","better-result":"^3.0.0","oxlint-tsgolint":"^7.0.2001"},"peerDependencies":{"typescript":">=6.0.0","better-effect":">=0.14.0 <0.15.0","better-result":"^3.0.0"},"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/better-effect-mq_0.1.3_1789157215765_0.7869354799152919"}}},"time":{"created":"2026-09-09T01:41:36.763Z","modified":"2026-09-11T20:06:56.230Z","0.1.0":"2026-09-09T01:41:37.197Z","0.1.1":"2026-09-09T20:56:07.953Z","0.1.2":"2026-09-10T18:26:05.090Z","0.1.3":"2026-09-11T20:06:55.883Z"},"bugs":{"url":"https://github.com/nitoba/better-effect/issues"},"license":"MIT","homepage":"https://github.com/nitoba/better-effect#readme","keywords":["better-effect","better-result","message-queue","queue","typescript"],"repository":{"url":"git+https://github.com/nitoba/better-effect.git","type":"git"},"description":"Experimental message-queue foundations for better-effect.","maintainers":[{"name":"nitoba","email":"nito.ba.dev@gmail.com"}],"readme":"# better-effect-mq\n\nTyped, storage-neutral building blocks for durable background work with\n[`better-effect`](https://github.com/nitoba/better-effect) and\n[`better-result`](https://github.com/nitoba/better-result).\n\nThe core package gives applications one vocabulary for defining jobs,\nenqueueing work, running workers, coordinating parent/child executions, and\nreading results. It does not choose a database, queue server, or dependency\ninjection container. A storage adapter implements the `JobStore` contract and\nis provided through a `better-effect` `Layer`.\n\n```text\nQueue.define → Job → enqueue / awaitResult\n                         ↓\n                    JobStore\n                         ↓\n                Worker.service → Worker.handle\n```\n\n## Install\n\n```bash\nbun add better-effect-mq better-effect better-result better-effect-schema zod\n```\n\nThe package expects `better-effect >=0.14 <0.15`, `better-result ^3`, and TypeScript\n6 or newer. `better-effect-schema` and `zod` are optional integration packages;\ninstall them when a job crosses an untrusted boundary and needs runtime\nvalidation. Use the package's `npm` or `pnpm` equivalent if that is how your\napplication manages dependencies.\n\n## Quick start: a schema-backed in-memory queue\n\nThis complete program uses Zod through `better-effect-schema` to validate the\nwire payload, construct a schema-backed class for the handler, and encode that\nclass back to JSON for the Job boundary. It composes a Layer-owned worker,\nsubmits one item, waits for its result, and disposes the Runtime. The same\napplication code works with a durable adapter after replacing the store Layer.\n\n```ts\nimport * as z from 'zod'\nimport { Effect, Layer, Runtime } from 'better-effect'\nimport { ClockLive } from 'better-effect/standard-services'\nimport { Result } from 'better-result'\nimport { Codec, JobEncodeFailure, JobStore, MemoryJobStore, Queue, Worker } from 'better-effect-mq'\nimport { Schema as CoreSchema } from 'better-effect-schema'\nimport { Schema } from 'better-effect-schema/zod'\n\nconst DateFromISOString = z.codec(z.iso.datetime(), z.date(), {\n  decode: (value) => new Date(value),\n  encode: (value) => value.toISOString()\n})\n\nclass UserEvent extends Schema.Class<UserEvent>('app/UserEvent')({\n  eventId: z.uuid(),\n  occurredAt: DateFromISOString,\n  kind: z.string().min(1),\n  payload: z.record(z.string(), z.string())\n}) {}\n\nconst DeliveryReceipt = z.object({\n  accepted: z.literal(true),\n  eventId: z.uuid()\n})\n\nconst DeliveryFailure = z.object({\n  code: z.string().min(1),\n  retryable: z.boolean()\n})\n\nconst UserEventCodec = Codec.standardSchema({\n  schema: UserEvent,\n  encode: (value) =>\n    CoreSchema.encode(UserEvent, value).mapError(\n      (error) => new JobEncodeFailure({ message: error.message, code: 'schema-encode' })\n    )\n})\n\nconst ReceiptCodec = Codec.standardSchema({ schema: DeliveryReceipt })\nconst FailureCodec = Codec.standardSchema({ schema: DeliveryFailure })\n\nconst Events = Queue.define('events')\nconst IngestEvent = Events.job('ingest-event', {\n  version: 1,\n  payload: UserEventCodec,\n  result: ReceiptCodec,\n  failure: FailureCodec,\n  idempotencyKey: ({ eventId }) => eventId,\n  retryable: ({ retryable }) => retryable\n})\n\nconst store = MemoryJobStore.make()\nconst handler = Worker.handle(IngestEvent, (event) =>\n  Effect.fn(async function* () {\n    console.log(`handling ${event.kind} at ${event.occurredAt.toISOString()}`)\n    return Result.ok({ accepted: true as const, eventId: event.eventId })\n  })\n)\n\nconst EventsWorker = Worker.service('@app/EventsWorker')\nconst EventsWorkerLive = EventsWorker.layer(() => ({\n  handlers: [handler] as const,\n  concurrency: 1,\n  pollIntervalMs: 10\n}))\n\nconst AppLive = Layer.complete(\n  Layer.merge(Layer.succeed(JobStore, JobStore.of(store)), Layer.merge(ClockLive, EventsWorkerLive))\n)\nconst runtime = await Runtime.make(AppLive)\n\ntry {\n  const started = await runtime.run(() =>\n    Effect.gen(async function* () {\n      return Result.ok(yield* EventsWorker)\n    })\n  )\n  if (Result.isError(started)) throw started.error\n\n  const completed = await runtime.run(() =>\n    Effect.gen(async function* () {\n      const jobId = yield* IngestEvent.enqueue({\n        eventId: '550e8400-e29b-41d4-a716-446655440000',\n        occurredAt: '2026-09-02T10:00:00.000Z',\n        kind: 'user.created',\n        payload: { source: 'example' }\n      })\n      const receipt = yield* IngestEvent.awaitResult(jobId)\n      return Result.ok({ jobId, receipt })\n    })\n  )\n  if (Result.isError(completed)) throw completed.error\n\n  console.log(completed.value)\n  await started.value.awaitIdle()\n} finally {\n  await runtime.dispose()\n}\n```\n\n`MemoryJobStore` is process-local and intentionally disposable. It is useful\nfor examples, tests, and local development; use a durable adapter when work\nmust survive a restart or be shared by multiple processes.\n\nThe schema-backed codec decodes persisted JSON into the class used by the\nhandler; the explicit `encode` callback delegates the wire projection to\n`CoreSchema.encode`. If a provider schema's output is already JSON-safe, omit\n`encode` and the codec uses that value for both sides.\n\n### Plain JSON escape hatch\n\nFor internal or otherwise simple data that is already JSON-safe, `Codec.json<T>()`\nis the smallest option:\n\n```ts\nconst InternalPing = Queue.define('internal').job('ping', {\n  version: 1,\n  payload: Codec.json<{ readonly requestId: string }>(),\n  result: Codec.string\n})\n```\n\nPrefer a schema-backed codec for HTTP, database, queue, or other untrusted\nboundaries where runtime validation, normalized failures, or a decoded domain\nclass matters.\n\n## Advanced: raw Standard Schema interoperability\n\n`Codec.standardSchema` can bridge a schema that implements the Standard Schema\ncontract when no provider adapter is available. Prefer a native provider such\nas Zod 4 through `better-effect-schema`; hand-writing a `StandardSchemaV1`\nobject is an adapter/interoperability escape hatch, not the recommended Job\ndefinition path.\n\n## Define jobs with `Queue` and `Job`\n\n`Queue.define` creates a namespace. A Job descriptor gives one work type a\nstable identity and declares the codecs used at the storage boundary:\n\n```ts\nimport * as z from 'zod'\nimport { Codec, JobEncodeFailure, Queue, Retry } from 'better-effect-mq'\nimport { Schema as CoreSchema } from 'better-effect-schema'\nimport { Schema } from 'better-effect-schema/zod'\n\nconst DateFromISOString = z.codec(z.iso.datetime(), z.date(), {\n  decode: (value) => new Date(value),\n  encode: (value) => value.toISOString()\n})\n\nconst Billing = Queue.define('billing')\nclass ChargeCardPayload extends Schema.Class<ChargeCardPayload>('app/ChargeCardPayload')({\n  paymentId: z.string().min(1),\n  amountCents: z.int().positive(),\n  requestedAt: DateFromISOString\n}) {}\nconst ChargeCardResult = z.object({ receiptId: z.string().min(1) })\nconst ChargeCardFailure = z.object({ code: z.string().min(1) })\nconst ChargeCardPayloadCodec = Codec.standardSchema({\n  schema: ChargeCardPayload,\n  encode: (value) =>\n    CoreSchema.encode(ChargeCardPayload, value).mapError(\n      (error) => new JobEncodeFailure({ message: error.message, code: 'schema-encode' })\n    )\n})\n\nconst ChargeCard = Billing.job('charge-card', {\n  version: 1,\n  payload: ChargeCardPayloadCodec,\n  result: Codec.standardSchema({ schema: ChargeCardResult }),\n  failure: Codec.standardSchema({ schema: ChargeCardFailure }),\n  defaults: {\n    attempts: 3,\n    backoff: Retry.exponential({\n      initialDelayMs: 1_000,\n      factor: 2,\n      maxDelayMs: 60_000,\n      maxAttempts: 3\n    })\n  },\n  idempotencyKey: ({ paymentId }) => paymentId,\n  retryable: ({ code }) => code !== 'card-declined'\n})\n```\n\nThe payload is the decoded value passed to the handler. Result and failure\ncodecs control what producers and operators can read back. Defaults provide\nretry and timeout policy; enqueue options can override the per-item schedule.\nThe descriptor is inert: defining a Job does not resolve a Service, create a\nworker, or register anything globally. For a payload whose in-memory value is\ndifferent from its wire value—such as a `Date` or a schema class—provide the\nexplicit encoder shown in the quick start.\n\nThe `version` identifies the persisted Job contract. Increment it when a\npayload, result, or typed failure changes incompatibly.\n\n## `JobStore`: the persistence seam\n\n`JobStore` is a yieldable Service, not a database client. The adapter contract\ncovers enqueueing, claiming and leasing, settlement, heartbeats, stalled-job\nrecovery, inspection, administration, and optional queue wake-ups. Provide the\nadapter through a Layer:\n\n```ts\nconst jobs = MemoryJobStore.make()\nconst AppStorage = Layer.succeed(JobStore, JobStore.of(jobs))\n```\n\nA durable adapter provides the same `JobStore` token through its own Layer.\nThe `Job`, `Worker`, `enqueue`, and `awaitResult` calls do not change when the\nstorage provider changes. Use `JobStore.named('billing')` when one Runtime\ncontains independent stores; bind each Job to the store it belongs to.\n\n## Workers and producer operations\n\n`Worker.handle(job, handler)` connects a Job to a typed `better-effect`\nprogram. A handler receives the decoded payload and can yield application\nServices. `JobContext` provides the job ID, attempt number, delivery count,\nmetadata, and worker identity for the current attempt.\n\n`Worker.service(tag)` returns a Layer-first Service. Acquiring its Layer starts\nthe supervisor; Runtime shutdown stops it. Workers support bounded concurrency,\nqueue and handler limits, retries, timeouts, lease heartbeats, stalled\nrecovery, and graceful stop.\n\nEvery Job exposes yieldable producer operations:\n\n| Operation                    | Use it for                                      |\n| ---------------------------- | ----------------------------------------------- |\n| `enqueue`                    | Submit one codec input and receive its `JobId`. |\n| `enqueueMany`                | Submit a batch while retaining input order.     |\n| `poll`                       | Read one job snapshot without waiting.          |\n| `awaitResult`                | Wait for a terminal result or typed failure.    |\n| `execute`                    | Enqueue and wait in one operation.              |\n| `attempts`                   | Read the delivery ledger for one job.           |\n| `cancel`, `retry`, `promote` | Apply explicit job administration.              |\n\n## Flow: coordinate a parent execution\n\nFlow solves the fan-out/fan-in problem: one durable parent Job can create\ntyped child Jobs, wait until every child reaches a terminal state, and publish\none aggregate parent result. For example, order fulfillment can reserve each\nline and charge the order in parallel while keeping one order-level lifecycle.\nThat gives operators one parent to inspect or retry, instead of a collection of\nunrelated jobs whose relationship only exists in application logs.\n\nThe example below is complete and compilable. It uses provider-backed Zod 4\nschemas at every persisted payload boundary, fans out to two different child\nJob types, implements both child handlers, collects useful fulfillment output,\nand awaits the terminal parent result. See\n[`examples/flow/main.ts`](./examples/flow/main.ts) for a small runnable,\nplain-JSON variant.\n\n```ts\nimport * as z from 'zod'\nimport { Effect, Layer, Runtime } from 'better-effect'\nimport { ClockLive } from 'better-effect/standard-services'\nimport { Result } from 'better-result'\nimport {\n  Codec,\n  Flow,\n  FlowStore,\n  JobEncodeFailure,\n  JobStore,\n  MemoryFlowStore,\n  MemoryJobStore,\n  Queue,\n  Worker\n} from 'better-effect-mq'\nimport { Schema as CoreSchema } from 'better-effect-schema'\nimport { Schema } from 'better-effect-schema/zod'\n\nconst OrderItem = z.object({ sku: z.string().min(1), quantity: z.int().positive() })\nconst OrderFailure = z.object({ code: z.string().min(1), message: z.string().min(1) })\n\nclass FulfillOrderPayload extends Schema.Class<FulfillOrderPayload>('app/FulfillOrderPayload')({\n  orderId: z.string().min(1),\n  currency: z.string().length(3),\n  totalCents: z.int().positive(),\n  items: z.array(OrderItem).min(1)\n}) {}\nconst FulfillOrderPayloadCodec = Codec.standardSchema({\n  schema: FulfillOrderPayload,\n  encode: (value) =>\n    CoreSchema.encode(FulfillOrderPayload, value).mapError(\n      (error) => new JobEncodeFailure({ message: error.message, code: 'schema-encode' })\n    )\n})\n\nclass ReserveInventoryPayload extends Schema.Class<ReserveInventoryPayload>(\n  'app/ReserveInventoryPayload'\n)({\n  orderId: z.string().min(1),\n  sku: z.string().min(1),\n  quantity: z.int().positive()\n}) {}\nconst ReserveInventoryPayloadCodec = Codec.standardSchema({\n  schema: ReserveInventoryPayload,\n  encode: (value) =>\n    CoreSchema.encode(ReserveInventoryPayload, value).mapError(\n      (error) => new JobEncodeFailure({ message: error.message, code: 'schema-encode' })\n    )\n})\nconst ReserveInventoryResult = z.object({\n  kind: z.literal('inventory'),\n  sku: z.string().min(1),\n  reservedQuantity: z.int().nonnegative()\n})\n\nclass ChargeOrderPayload extends Schema.Class<ChargeOrderPayload>('app/ChargeOrderPayload')({\n  orderId: z.string().min(1),\n  currency: z.string().length(3),\n  amountCents: z.int().positive()\n}) {}\nconst ChargeOrderPayloadCodec = Codec.standardSchema({\n  schema: ChargeOrderPayload,\n  encode: (value) =>\n    CoreSchema.encode(ChargeOrderPayload, value).mapError(\n      (error) => new JobEncodeFailure({ message: error.message, code: 'schema-encode' })\n    )\n})\nconst ChargeOrderResult = z.object({\n  kind: z.literal('payment'),\n  chargeId: z.string().min(1),\n  amountCents: z.int().positive()\n})\n\nconst FulfillOrderResult = z.object({\n  orderId: z.string().min(1),\n  requestedItems: z.int().nonnegative(),\n  reservedItems: z.int().nonnegative(),\n  payment: z.enum(['charged', 'not-charged']),\n  failedChildren: z.array(z.string())\n})\n\nconst Orders = Queue.define('orders')\nconst FulfillOrder = Orders.job('fulfill-order', {\n  version: 1,\n  payload: FulfillOrderPayloadCodec,\n  result: Codec.standardSchema({ schema: FulfillOrderResult }),\n  failure: Codec.standardSchema({ schema: OrderFailure })\n})\nconst ReserveInventory = Orders.job('reserve-inventory', {\n  version: 1,\n  payload: ReserveInventoryPayloadCodec,\n  result: Codec.standardSchema({ schema: ReserveInventoryResult }),\n  failure: Codec.standardSchema({ schema: OrderFailure })\n})\nconst ChargeOrder = Orders.job('charge-order', {\n  version: 1,\n  payload: ChargeOrderPayloadCodec,\n  result: Codec.standardSchema({ schema: ChargeOrderResult }),\n  failure: Codec.standardSchema({ schema: OrderFailure })\n})\n\nconst Fulfillment = Flow.define('order-fulfillment', {\n  parent: FulfillOrder,\n  children: [ReserveInventory, ChargeOrder] as const,\n  onChildFailure: 'continue'\n})\n\nconst FulfillmentHandler = Flow.handle(Fulfillment, {\n  fanOut: (payload) =>\n    Effect.fn(async function* () {\n      return Result.ok([\n        Flow.children(\n          ReserveInventory,\n          payload.items.map((item) => ({\n            key: `reserve:${item.sku}`,\n            payload: { orderId: payload.orderId, sku: item.sku, quantity: item.quantity }\n          }))\n        ),\n        Flow.children(ChargeOrder, [\n          {\n            key: 'charge',\n            payload: {\n              orderId: payload.orderId,\n              currency: payload.currency,\n              amountCents: payload.totalCents\n            }\n          }\n        ])\n      ] as const)\n    }),\n  collect: (payload, results) =>\n    Effect.fn(async function* () {\n      const children = yield* Result.await(results.all()())\n      const reservedItems = children.filter(\n        (child) => child.outcome === 'completed' && child.result?.kind === 'inventory'\n      ).length\n      const charged = children.some(\n        (child) => child.outcome === 'completed' && child.result?.kind === 'payment'\n      )\n      const payment = charged ? ('charged' as const) : ('not-charged' as const)\n      return Result.ok({\n        orderId: payload.orderId,\n        requestedItems: payload.items.length,\n        reservedItems,\n        payment,\n        failedChildren: children\n          .filter((child) => child.outcome !== 'completed')\n          .map((child) => child.childKey)\n      })\n    })\n})\n\nconst handlers = [\n  Worker.handle(ReserveInventory, (payload) =>\n    Effect.fn(async function* () {\n      // Reserve payload.sku in inventory; make that write idempotent by order ID and SKU.\n      return Result.ok({\n        kind: 'inventory' as const,\n        sku: payload.sku,\n        reservedQuantity: payload.quantity\n      })\n    })\n  ),\n  Worker.handle(ChargeOrder, (payload) =>\n    Effect.fn(async function* () {\n      // Pass payload.orderId as the payment provider's idempotency key.\n      return Result.ok({\n        kind: 'payment' as const,\n        chargeId: `charge:${payload.orderId}`,\n        amountCents: payload.amountCents\n      })\n    })\n  )\n] as const\n\nconst FulfillmentWorker = Worker.service('@app/FulfillmentWorker')\nconst FulfillmentWorkerLive = FulfillmentWorker.layer(() => ({\n  handlers,\n  flows: [FulfillmentHandler] as const,\n  concurrency: 4,\n  pollIntervalMs: 10,\n  flowSweepIntervalMs: 50,\n  flowBatchSize: 32\n}))\n\n// Replace these process-local providers with the matching durable adapter Layers in production.\nconst AppLive = Layer.complete(\n  Layer.merge(\n    Layer.succeed(JobStore, JobStore.of(MemoryJobStore.make())),\n    Layer.merge(\n      Layer.succeed(FlowStore, FlowStore.of(MemoryFlowStore.make())),\n      Layer.merge(ClockLive, FulfillmentWorkerLive)\n    )\n  )\n)\nconst runtime = await Runtime.make(AppLive)\n\ntry {\n  const execution = await runtime.run(() =>\n    Effect.gen(async function* () {\n      const parentId = yield* FulfillOrder.enqueue({\n        orderId: 'order-123',\n        currency: 'USD',\n        totalCents: 12_500,\n        items: [\n          { sku: 'coffee-beans', quantity: 2 },\n          { sku: 'pour-over-kit', quantity: 1 }\n        ]\n      })\n      const result = yield* FulfillOrder.awaitResult(parentId)\n      return Result.ok({ parentId, result })\n    })\n  )\n  if (Result.isError(execution)) throw execution.error\n  console.log(execution.value)\n} finally {\n  await runtime.dispose()\n}\n```\n\n`Flow.define` is an immutable descriptor. `Flow.handle` supplies `fanOut` and\n`collect` programs; `Flow.children` keeps each child payload tied to its Job\ndefinition, so adding a payment child cannot accidentally receive an inventory\npayload. The Flow store records the parent, child manifest, terminal reports,\nand relay work while the JobStore handles ordinary child delivery.\n\nWith `onChildFailure: 'continue'`, `collect` runs after all children settle and\ncan return a partial result such as two reservations and `payment:\n'not-charged'`. Use `onChildFailure: 'fail'` when the parent must fail as soon\nas a child fails and remaining work should be cancelled. Both policies keep\nthe parent-child relationship durable and replayable; choose based on whether\npartial fulfillment is useful to the caller.\n\n## Events and observing results\n\nPolling is the default and needs only a `JobStore`. To use a durable event log\nas a wake-up hint, provide the matching `JobEventStore` in the same Runtime:\n\n```ts\nimport { Effect, Layer } from 'better-effect'\nimport { Result } from 'better-result'\nimport { JobEventStore, JobStore, MemoryJobEventStore, MemoryJobStore } from 'better-effect-mq'\n\nconst events = MemoryJobEventStore.make({\n  retention: { count: 1_000, ageMs: 24 * 60 * 60 * 1_000 }\n})\nconst jobs = MemoryJobStore.make({ eventStore: events })\nconst AppStorage = Layer.merge(\n  Layer.succeed(JobStore, JobStore.of(jobs)),\n  Layer.succeed(JobEventStore, JobEventStore.of(events))\n)\n\nconst result = Effect.gen(async function* () {\n  return Result.ok(\n    yield* SendEmail.awaitResult(jobId, {\n      strategy: 'events',\n      eventStore: JobEventStore,\n      pollFallbackMs: 5_000\n    })\n  )\n})\n```\n\nAn event is a wake-up hint, not the source of truth: `awaitResult` rereads the\nJob record before decoding the terminal result or failure. Event retention is\nbounded, and a failed event read can fall back to polling. The\n`JobEventStore`, detailed `job.attempts(jobId)` ledger, and process-local\n`JobObserver` are complementary surfaces; none replaces the others.\n\n## Outbox: make a transaction handoff durable\n\nThe optional [`better-effect-mq-outbox`](../better-effect-mq-outbox/README.md)\nextension solves the database dual-write problem. It is not a second queue API:\nyou still define a Job with `Queue.define`, route a prepared request to the\ncore `JobStore`, and process it with the same `Worker.handle` handler. The\nextension adds a transaction-bound outbox record and a Runtime-owned publisher:\n\n```text\ndatabase transaction\n  ├─ write the order\n  └─ append a prepared Job request\n        │ commit\n        ▼\nOutboxPublisher → core JobStore → Worker.handle(SendConfirmation)\n```\n\nThe transaction-to-publisher shape is:\n\n```ts\nimport * as z from 'zod'\nimport { Effect, Layer, Runtime } from 'better-effect'\nimport { ClockLive } from 'better-effect/standard-services'\nimport { Codec, JobEncodeFailure, JobStore, Queue, Worker } from 'better-effect-mq'\nimport { Result } from 'better-result'\nimport { Schema as CoreSchema } from 'better-effect-schema'\nimport { Schema } from 'better-effect-schema/zod'\nimport { OutboxId, OutboxPublisher, OutboxRoutes, makeOutboxRecord } from 'better-effect-mq-outbox'\nimport { PostgresJobStore, PostgresOutbox, type Pool } from 'better-effect-mq-postgres'\n\ndeclare const pool: Pool\nconst DateFromISOString = z.codec(z.iso.datetime(), z.date(), {\n  decode: (value) => new Date(value),\n  encode: (value) => value.toISOString()\n})\n\nclass ConfirmationPayload extends Schema.Class<ConfirmationPayload>('app/ConfirmationPayload')({\n  orderId: z.string().min(1),\n  email: z.email(),\n  queuedAt: DateFromISOString\n}) {}\n\nconst ConfirmationPayloadCodec = Codec.standardSchema({\n  schema: ConfirmationPayload,\n  encode: (value) =>\n    CoreSchema.encode(ConfirmationPayload, value).mapError(\n      (error) => new JobEncodeFailure({ message: error.message, code: 'schema-encode' })\n    )\n})\n\nconst Orders = Queue.define('orders')\nconst SendConfirmation = Orders.job('send-confirmation', {\n  version: 1,\n  payload: ConfirmationPayloadCodec,\n  result: Codec.string,\n  defaults: { attempts: 5 },\n  idempotencyKey: ({ orderId }) => `order-confirmation:${orderId}`\n})\n\nconst confirmationHandler = Worker.handle(SendConfirmation, (payload) =>\n  Effect.fn(async function* () {\n    return Result.ok(`sent:${payload.email}`)\n  })\n)\nconst QueueWorker = Worker.service('@app/OrderWorker')\nconst QueueWorkerLive = QueueWorker.layer(() => ({\n  handlers: [confirmationHandler] as const,\n  concurrency: 2,\n  pollIntervalMs: 100\n}))\n\nconst Routes = OutboxRoutes.make({ jobs: JobStore })\nconst Publisher = OutboxPublisher.service('OrderOutboxPublisher')\nconst PublisherLive = Publisher.layer(() => ({\n  outboxes: [PostgresOutbox] as const,\n  routes: Routes,\n  concurrency: 4,\n  pollIntervalMs: 1_000\n}))\nconst AppLive = Layer.complete(\n  Layer.merge(\n    ClockLive,\n    PostgresJobStore.layer({ pool, namespace: 'orders' }),\n    PostgresOutbox.layer({ pool, namespace: 'orders' }),\n    PublisherLive,\n    QueueWorkerLive\n  )\n)\nconst runtime = await Runtime.make(AppLive)\nawait runtime.warmup()\n\nconst prepared = await runtime.run(() =>\n  Effect.gen(async function* () {\n    return Result.ok(\n      yield* SendConfirmation.prepare(\n        {\n          orderId: 'order-123',\n          email: 'ada@example.test',\n          queuedAt: '2026-09-09T12:00:00.000Z'\n        },\n        { jobId: 'order-confirmation:order-123' }\n      )\n    )\n  })\n)\nif (Result.isError(prepared)) throw prepared.error\nconst record = makeOutboxRecord({\n  id: OutboxId.make('order-confirmation:order-123').unwrap(),\n  target: 'jobs',\n  request: prepared.value,\n  attemptsMax: 5\n})\nif (Result.isError(record)) throw record.error\n\nconst persisted = await PostgresOutbox.transaction(\n  pool,\n  record.value,\n  async (transaction) => {\n    await transaction.query('INSERT INTO orders (id, email) VALUES ($1, $2)', [\n      'order-123',\n      'ada@example.test'\n    ])\n    return Result.ok(undefined)\n  },\n  { namespace: 'orders' }\n)\nif (Result.isError(persisted)) throw persisted.error\n\nconst completed = await runtime.run(() =>\n  Effect.gen(async function* () {\n    const result = yield* SendConfirmation.awaitResult('order-confirmation:order-123')\n    return Result.ok(result)\n  })\n)\nif (Result.isError(completed)) throw completed.error\nconsole.log(completed.value)\nawait runtime.dispose()\n```\n\nPrepare the request and record before calling the adapter helper. The adapter\nappends the record after the domain callback succeeds and owns the connection,\ncommit, rollback, and cleanup. After commit, the Runtime-owned publisher\nenqueues the prepared request into the routed `JobStore`; the Runtime-owned\nWorker then runs `SendConfirmation`, and `awaitResult` observes the terminal\nresult. The complete\nPostgreSQL setup, connection ownership rules, and adapter equivalents are in\nthe outbox extension's [transaction-to-publisher example](../better-effect-mq-outbox/README.md#end-to-end-example-with-postgresql).\n\n### Advanced: caller-owned transactions\n\nAdapters may expose `appendIn` as an escape hatch when application code already\nowns a native transaction. It is not part of the normal outbox path; use the\nadapter's record-first `transaction` helper for application writes.\n\n## Reliability\n\n- Delivery is at least once, not exactly once. Make handler side effects\n  idempotent with the Job ID or an application key.\n- Leases and fencing prevent an old worker from settling a newer delivery, but\n  they cannot undo an external side effect that already happened.\n- Typed failures can be retried according to `Retry` and the Job's `retryable`\n  policy. Timeouts, cancellations, codec failures, and defects remain\n  distinguishable for operators.\n- Worker shutdown is cooperative. Runtime disposal stops admitting new work,\n  lets active attempts settle according to its policy, and closes owned\n  resources.\n- Event retention is bounded. An EventLog is not an infinite archive or a\n  replacement for payload/result storage.\n\n## Testing and further reading\n\nThe package's runnable examples can be checked with:\n\n```bash\nbun run typecheck:examples\nbun run test:examples\nbunx tsc -p examples/flow/tsconfig.json --pretty false\nbun examples/flow/main.ts\n```\n\nThe [examples README](./examples/README.md) describes each runnable example.\nFor application-facing adapter composition, see the [composition guide](./docs/composition.md).\nFor adapter authors and maintainers, see the [driver author guide](./docs/writing-a-driver.md).\n","readmeFilename":"README.md"}