{"_id":"@alkdev/pubsub","name":"@alkdev/pubsub","dist-tags":{"latest":"0.1.0"},"versions":{"0.1.0":{"name":"@alkdev/pubsub","version":"0.1.0","description":"Type-safe publish/subscribe with pluggable event target adapters (in-process, Redis, WebSocket, Worker)","type":"module","main":"./dist/index.cjs","module":"./dist/index.js","types":"./dist/index.d.ts","author":{"name":"alk.dev"},"repository":{"type":"git","url":"git+https://git.alk.dev/alkdev/pubsub.git"},"homepage":"https://git.alk.dev/alkdev/pubsub","bugs":{"url":"https://git.alk.dev/alkdev/pubsub/issues"},"exports":{".":{"import":{"types":"./dist/index.d.ts","default":"./dist/index.js"},"require":{"types":"./dist/index.d.cts","default":"./dist/index.cjs"}},"./event-target-redis":{"import":{"types":"./dist/event-target-redis.d.ts","default":"./dist/event-target-redis.js"},"require":{"types":"./dist/event-target-redis.d.cts","default":"./dist/event-target-redis.cjs"}},"./event-target-websocket-client":{"import":{"types":"./dist/event-target-websocket-client.d.ts","default":"./dist/event-target-websocket-client.js"},"require":{"types":"./dist/event-target-websocket-client.d.cts","default":"./dist/event-target-websocket-client.cjs"}},"./event-target-websocket-server":{"import":{"types":"./dist/event-target-websocket-server.d.ts","default":"./dist/event-target-websocket-server.js"},"require":{"types":"./dist/event-target-websocket-server.d.cts","default":"./dist/event-target-websocket-server.cjs"}},"./event-target-worker":{"import":{"types":"./dist/event-target-worker.d.ts","default":"./dist/event-target-worker.js"},"require":{"types":"./dist/event-target-worker.d.cts","default":"./dist/event-target-worker.cjs"}}},"publishConfig":{"access":"public"},"sideEffects":false,"scripts":{"build":"tsup","build:tsc":"tsc --noEmit","test":"vitest run","test:watch":"vitest","test:coverage":"vitest run --coverage","lint":"tsc --noEmit","prepublishOnly":"npm run build"},"keywords":["pubsub","typed-event-target","redis","websocket","iroh","quic"],"license":"MIT OR Apache-2.0","dependencies":{},"peerDependencies":{"ioredis":"^5.0.0"},"peerDependenciesMeta":{"ioredis":{"optional":true}},"devDependencies":{"@types/node":"^22.0.0","@vitest/coverage-v8":"^3.2.4","ioredis":"^5.10.1","tsup":"^8.5.1","typescript":"^5.7.0","vitest":"^3.1.0"},"engines":{"node":">=18.0.0"},"gitHead":"b3f598dffd027e2d8d4336aa30297fb6e7a1a379","_id":"@alkdev/pubsub@0.1.0","_nodeVersion":"25.8.1","_npmVersion":"11.11.0","dist":{"integrity":"sha512-CdCyBMJEWwrSDAalaZCsBHDscljijgx4PdfYjjpV33PB6ad3KNNX1hoIiYGlO0YgARyd466lXqmth2Kw8J5JNw==","shasum":"9609ca5201e0d943b7482eefca54e5de29bd72e7","tarball":"https://registry.npmjs.org/@alkdev/pubsub/-/pubsub-0.1.0.tgz","fileCount":52,"unpackedSize":235075,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEQCIG4cZu6gpQ/kx28mI01x4Oq3MO4wFffMm9vSztpS4dEhAiAo2qnJmx/yu/UXnyDfs53zzS5sMbfw9Y/ze5DbRqKEKg=="}]},"_npmUser":{"name":"alkdev","email":"admin@alk.dev"},"directories":{},"maintainers":[{"name":"alkdev","email":"admin@alk.dev"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/pubsub_0.1.0_1778260020960_0.7029324654380287"},"_hasShrinkwrap":false}},"time":{"created":"2026-05-08T17:07:00.780Z","0.1.0":"2026-05-08T17:07:01.095Z","modified":"2026-05-08T17:07:01.371Z"},"maintainers":[{"name":"alkdev","email":"admin@alk.dev"}],"description":"Type-safe publish/subscribe with pluggable event target adapters (in-process, Redis, WebSocket, Worker)","homepage":"https://git.alk.dev/alkdev/pubsub","keywords":["pubsub","typed-event-target","redis","websocket","iroh","quic"],"repository":{"type":"git","url":"git+https://git.alk.dev/alkdev/pubsub.git"},"author":{"name":"alk.dev"},"bugs":{"url":"https://git.alk.dev/alkdev/pubsub/issues"},"license":"MIT OR Apache-2.0","readme":"# @alkdev/pubsub\n\nType-safe publish/subscribe with pluggable event target adapters. Transport layer only — no call protocol or coordination semantics.\n\nEvery event is an `EventEnvelope<TType, TPayload>` with `{ type, id, payload }`. Adapters implement the `TypedEventTarget` interface so you can swap transports without changing your subscribe logic.\n\n## Install\n\n```bash\nnpm install @alkdev/pubsub\n```\n\nFor Redis transport:\n\n```bash\nnpm install ioredis\n```\n\nWebSocket and Worker adapters use built-in APIs — no additional dependencies.\n\n## Quick Start\n\n### In-Process (default)\n\n```ts\nimport { createPubSub } from \"@alkdev/pubsub\";\n\ntype EventMap = {\n  \"user.created\": { name: string };\n  \"order.placed\": { orderId: string };\n};\n\nconst pubsub = createPubSub<EventMap>();\n\npubsub.subscribe(\"user.created\", (_, payload) => {\n  console.log(`New user: ${payload.name}`);\n});\n\npubsub.publish(\"user.created\", \"id-1\", { name: \"Alice\" });\n```\n\n### Redis\n\n```ts\nimport { createPubSub } from \"@alkdev/pubsub\";\nimport { createRedisEventTarget } from \"@alkdev/pubsub/event-target-redis\";\nimport Redis from \"ioredis\";\n\nconst publishClient = new Redis();\nconst subscribeClient = new Redis();\n\nconst eventTarget = createRedisEventTarget({\n  publishClient,\n  subscribeClient,\n});\n\nconst pubsub = createPubSub({ eventTarget });\n```\n\n### WebSocket Client (browser/Node)\n\n```ts\nimport { createPubSub } from \"@alkdev/pubsub\";\nimport { createWebSocketClientEventTarget } from \"@alkdev/pubsub/event-target-websocket-client\";\n\nconst ws = new WebSocket(\"ws://localhost:8080\");\nconst eventTarget = createWebSocketClientEventTarget(ws);\n\nconst pubsub = createPubSub({ eventTarget });\n```\n\n### WebSocket Server (Node)\n\n```ts\nimport { createWebSocketServerEventTarget } from \"@alkdev/pubsub/event-target-websocket-server\";\n\nconst server = createWebSocketServerEventTarget({\n  onConnection(spoke, ws) { /* new client connected */ },\n  onDisconnection(spoke, ws) { /* client disconnected */ },\n  maxBufferedAmount: 1_048_576,\n  onBackpressure(ws, bufferedAmount) { /* optional backpressure signal */ },\n});\n\n// When a new WebSocket connects:\nserver.addConnection(ws);\n\n// When it disconnects:\nserver.removeConnection(ws);\n\n// Subscribe local handlers:\nserver.addEventListener(\"user.created:id-1\", (event) => {\n  // event.detail is the EventEnvelope\n});\n\n// Publish to subscribed connections:\nserver.dispatchEvent(new CustomEvent(\"user.created:id-1\", { detail: envelope }));\n```\n\n### Worker (Host ↔ Thread)\n\n```ts\n// Host (main thread)\nimport { createWorkerHostEventTarget } from \"@alkdev/pubsub/event-target-worker\";\n\nconst worker = new Worker(\"./worker.js\");\nconst eventTarget = createWorkerHostEventTarget(worker);\n```\n\n```ts\n// Worker thread\nimport { createWorkerThreadEventTarget } from \"@alkdev/pubsub/event-target-worker\";\n\nconst eventTarget = createWorkerThreadEventTarget();\n// Must be called inside a Worker context — throws if globalThis.postMessage is unavailable\n```\n\n## Lifecycle\n\nAll transport adapters provide a `close()` method for graceful teardown:\n\n```ts\nconst eventTarget = createRedisEventTarget({ publishClient, subscribeClient });\n// ... subscribe and publish ...\n\neventTarget.close(); // unsubscribes all channels, removes listener, clears state\n```\n\nAfter `close()`:\n- `addEventListener`, `removeEventListener`, and `dispatchEvent` are no-ops\n- Intercepted handlers (`onmessage`, `onclose`) are restored to their originals\n- Subscriptions are cleaned up (Redis channels unsubscribed, WebSocket `__unsubscribe` sent)\n- The underlying transport (Redis connection, WebSocket, Worker) is **not** destroyed — the caller owns it\n\n`close()` is idempotent. Calling it multiple times is safe.\n\n## Operators\n\nOperators transform `AsyncIterable` streams from `subscribe()`:\n\n```ts\nimport { pipe, filter, map, take, batch } from \"@alkdev/pubsub\";\n\nconst pubsub = createPubSub<EventMap>();\n\nconst stream = pubsub.subscribe(\"user.created\");\n\nfor await (const event of pipe(\n  stream,\n  filter((e) => e.payload.name.startsWith(\"A\")),\n  map((e) => e.payload.name),\n  take(5),\n)) {\n  console.log(event);\n}\n```\n\nAvailable operators: `filter`, `map`, `pipe`, `take`, `reduce`, `toArray`, `batch`, `dedupe`, `window`, `flat`, `groupBy`, `chain`, `join`.\n\n## EventEnvelope\n\nAll events are serialized as `EventEnvelope`:\n\n```ts\ninterface EventEnvelope<TType = string, TPayload = unknown> {\n  type: TType;\n  id: string;\n  payload: TPayload;\n}\n```\n\nThis is the cross-platform wire format. Adapters serialize/deserialize this automatically (JSON for Redis and WebSocket, structured clone for Worker).\n\n## Subscription Control Protocol\n\nEvent types starting with `__` are reserved for internal use. Adapters use `__subscribe` and `__unsubscribe` control events to manage topic subscriptions across connections. User code must not define event types with the `__` prefix.\n\n## TypeScript\n\nFull type inference through `EventMap`:\n\n```ts\ntype EventMap = {\n  \"user.created\": { name: string; role: string };\n  \"order.placed\": { orderId: string; total: number };\n};\n\nconst pubsub = createPubSub<EventMap>();\n\npubsub.publish(\"user.created\", \"id-1\", { name: \"Alice\", role: \"admin\" });\n//                                               ^ full type checking on payload\n```\n\n## Exports\n\n| Import | Description |\n|--------|-------------|\n| `@alkdev/pubsub` | Core: `createPubSub`, `EventEnvelope`, `Repeater`, operators |\n| `@alkdev/pubsub/event-target-redis` | Redis adapter (peer dep: `ioredis`) |\n| `@alkdev/pubsub/event-target-websocket-client` | WebSocket client adapter |\n| `@alkdev/pubsub/event-target-websocket-server` | WebSocket server adapter |\n| `@alkdev/pubsub/event-target-worker` | Worker host + thread adapters |\n\n## Upstream Attribution\n\nCore `createPubSub`, `TypedEventTarget`, and operators are adapted from [graphql-yoga](https://github.com/graphql-hive/graphql-yoga) (MIT). The `Repeater` class is inlined from [@repeaterjs/repeater](https://github.com/repeaterjs/repeater) (MIT).\n\n## License\n\nDual-licensed under [MIT](LICENSE-MIT) or [Apache-2.0](LICENSE-APACHE). Portions adapted from upstream projects retain their MIT attribution.","readmeFilename":"README.md","_rev":"1-b9458677660ec80f7fa5229ac9792534"}