{"_id":"@aprismlab/thinrsmq","name":"@aprismlab/thinrsmq","dist-tags":{"latest":"0.1.2"},"versions":{"0.1.2":{"name":"@aprismlab/thinrsmq","version":"0.1.2","description":"Thin Redis Streams message queue layer with retry, DLQ, and claim monitoring","type":"module","main":"./dist/index.cjs","module":"./dist/index.js","types":"./dist/index.d.ts","exports":{".":{"import":{"types":"./dist/index.d.ts","default":"./dist/index.js"},"require":{"types":"./dist/index.d.cts","default":"./dist/index.cjs"}}},"sideEffects":false,"engines":{"node":">=18"},"scripts":{"build":"tsup src/index.ts --format esm,cjs --dts","test":"vitest run","test:unit":"vitest run --dir test/unit","test:integration":"vitest run --dir test/integration","prepublishOnly":"npm run build","typecheck":"tsc --noEmit"},"peerDependencies":{"ioredis":">=5.0.0"},"peerDependenciesMeta":{"ioredis":{"optional":false}},"devDependencies":{"ioredis":"^5.4.1","commander":"^12.1.0","typescript":"^5.5.4","vitest":"^2.0.5","tsup":"^8.2.4","@types/node":"^20.14.15"},"keywords":["redis","streams","message-queue","mq","dlq","dead-letter-queue","retry","consumer-group","claim-monitor"],"license":"MIT","repository":{"type":"git","url":"git+https://github.com/hunetmoducoding/thinrsmq-node.git","directory":"thinrsmq-node"},"homepage":"https://github.com/hunetmoducoding/thinrsmq-node/tree/main/thinrsmq-node#readme","bugs":{"url":"https://github.com/hunetmoducoding/thinrsmq-node/issues"},"publishConfig":{"access":"public","registry":"https://registry.npmjs.org"},"gitHead":"a9bca5de012171904da6edc5c367e955dc2952ad","_id":"@aprismlab/thinrsmq@0.1.2","_nodeVersion":"26.0.0","_npmVersion":"11.12.1","dist":{"integrity":"sha512-CGAArgyC2/o5kM3qaSd4GovF6fOhepO50xj5jqp6tLZCpFZpM4ZCWmCT8VyeEJ3vicFz/84xBU3U8aq/riT95Q==","shasum":"cb1a1ad0a9fc630e2006937e7bcfb6ee630404cc","tarball":"https://registry.npmjs.org/@aprismlab/thinrsmq/-/thinrsmq-0.1.2.tgz","fileCount":7,"unpackedSize":111545,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEQCIEZHnRjcPUEzIH2ebrvVO1XKfh6sAUhg/Biu26KXdbEoAiBuBBFNKU2YQkOnF8m7+F3wTzngXy32zJbG+Ssk0r3aYw=="}]},"_npmUser":{"name":"netscout","email":"netscout82@naver.com"},"directories":{},"maintainers":[{"name":"netscout","email":"netscout82@naver.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/thinrsmq_0.1.2_1780021998100_0.1400795787633773"},"_hasShrinkwrap":false}},"time":{"created":"2026-05-29T02:33:17.926Z","0.1.2":"2026-05-29T02:33:18.246Z","modified":"2026-05-29T02:33:18.423Z"},"maintainers":[{"name":"netscout","email":"netscout82@naver.com"}],"description":"Thin Redis Streams message queue layer with retry, DLQ, and claim monitoring","homepage":"https://github.com/hunetmoducoding/thinrsmq-node/tree/main/thinrsmq-node#readme","keywords":["redis","streams","message-queue","mq","dlq","dead-letter-queue","retry","consumer-group","claim-monitor"],"repository":{"type":"git","url":"git+https://github.com/hunetmoducoding/thinrsmq-node.git","directory":"thinrsmq-node"},"bugs":{"url":"https://github.com/hunetmoducoding/thinrsmq-node/issues"},"license":"MIT","readme":"# thinrsmq (Node.js)\n\nThin Redis Streams message queue layer with retry, DLQ, and claim monitoring. Wire-compatible with the Go implementation.\n\n## Installation\n\n```bash\nnpm install thinrsmq ioredis\n```\n\nNote: `ioredis` is a peer dependency. You must install it separately.\n\n## Quick Start\n\n### Producer\n\n```typescript\nimport Redis from 'ioredis';\nimport { Producer, withDefaults } from 'thinrsmq';\n\nconst redis = new Redis();\nconst config = withDefaults({ namespace: 'myapp' });\nconst producer = new Producer(redis, config);\n\n// Publish a single message\nconst id = await producer.publish('orders', {\n  type: 'order.created',\n  payload: JSON.stringify({ orderId: '123', amount: 99.99 }),\n  version: '1',\n  producedAt: '',\n  traceId: 'trace-abc-123', // optional\n  producer: 'order-service', // optional\n});\n\nconsole.log('Published message:', id);\n\n// Publish a batch\nconst ids = await producer.publishBatch('orders', [\n  { type: 'order.created', payload: '{\"orderId\":\"124\"}', version: '1', producedAt: '' },\n  { type: 'order.created', payload: '{\"orderId\":\"125\"}', version: '1', producedAt: '' },\n]);\n\nconsole.log('Published batch:', ids);\n```\n\n### Consumer\n\n```typescript\nimport Redis from 'ioredis';\nimport { Consumer, withDefaults } from 'thinrsmq';\n\nconst redis = new Redis();\nconst config = withDefaults({ namespace: 'myapp' });\nconst consumer = new Consumer(redis, config);\n\n// Subscribe to a topic with a handler\nawait consumer.subscribe('orders', 'order-processor', async (msg) => {\n  console.log('Processing order:', msg.payload);\n\n  // Throw error to trigger retry\n  // if (shouldRetry) throw new Error('Temporary failure');\n\n  // Success - message will be acknowledged\n});\n\n// Optional: handle skipped messages (nil/trimmed or unknown version)\nconsumer.onSkip((id, reason) => {\n  console.warn('Skipped message:', id, reason);\n});\n\n// Graceful shutdown\nprocess.on('SIGTERM', async () => {\n  await consumer.stop();\n  await redis.quit();\n});\n```\n\n### Claim Monitor\n\nThe claim monitor detects idle messages in the pending entries list (PEL) and routes them to retry or DLQ based on attempt count.\n\n```typescript\nimport Redis from 'ioredis';\nimport { ClaimMonitor, withDefaults } from 'thinrsmq';\n\nconst redis = new Redis();\nconst config = withDefaults({\n  namespace: 'myapp',\n  monitor: {\n    enabled: true,\n    scanIntervalMs: 5000,\n    minIdleTimeMs: 10000,\n  },\n});\n\nconst monitor = new ClaimMonitor(redis, config);\nawait monitor.start();\n\n// Graceful shutdown\nprocess.on('SIGTERM', async () => {\n  await monitor.stop();\n  await redis.quit();\n});\n```\n\n### Dead Letter Queue (DLQ)\n\n```typescript\nimport Redis from 'ioredis';\nimport { DLQ, withDefaults } from 'thinrsmq';\n\nconst redis = new Redis();\nconst config = withDefaults({ namespace: 'myapp' });\nconst dlq = new DLQ(redis, config);\n\n// Peek at DLQ messages\nconst messages = await dlq.peek('orders', 10);\nconsole.log('DLQ messages:', messages);\n\n// Replay a message (with guard)\nawait dlq.replay('orders', messageId);\n\n// Purge all DLQ messages for a topic\nawait dlq.purge('orders');\n\n// Get DLQ size\nconst size = await dlq.size('orders');\nconsole.log('DLQ size:', size);\n```\n\n### Admin\n\n```typescript\nimport Redis from 'ioredis';\nimport { Admin, withDefaults } from 'thinrsmq';\n\nconst redis = new Redis();\nconst config = withDefaults({ namespace: 'myapp' });\nconst admin = new Admin(redis, config);\n\n// Get pending stats for a stream\nconst pendingInfo = await admin.pendingStats('orders', 'order-processor');\nconsole.log('Pending messages:', pendingInfo.count);\n\n// Get consumer details\nconst consumerInfo = await admin.consumerInfo('orders', 'order-processor');\nconsole.log('Consumers:', consumerInfo);\n\n// Get stream details\nconst streamInfo = await admin.streamInfo('orders');\nconsole.log('Stream length:', streamInfo.length);\n```\n\n## Configuration\n\n### Configuration Reference\n\n| Field | Type | Default | Description |\n|-------|------|---------|-------------|\n| `namespace` | string | (required) | Namespace for all Redis keys |\n| `redis.address` | string | `\"localhost:6379\"` | Redis server address |\n| `redis.password` | string | `\"\"` | Redis password |\n| `redis.db` | number | `0` | Redis database number |\n| `redis.poolSize` | number | `10` | Connection pool size |\n| `redis.readTimeoutMs` | number | `3000` | Read timeout in milliseconds |\n| `redis.writeTimeoutMs` | number | `3000` | Write timeout in milliseconds |\n| `redis.useTls` | boolean | `false` | Enable TLS for Redis connection |\n| `streams.defaultMaxLen` | number | `10000` | Default MAXLEN for streams |\n| `consumer.batchSize` | number | `10` | Number of messages to read per batch |\n| `consumer.blockMs` | number | `5000` | Block duration for XREADGROUP |\n| `consumer.consumerName` | string | `\"thinrsmq\"` | Prefix for consumer name generation |\n| `consumer.shutdownTimeoutMs` | number | `30000` | Graceful shutdown timeout |\n| `retry.maxAttempts` | number | `3` | Maximum retry attempts before DLQ |\n| `retry.baseDelayMs` | number | `1000` | Base delay for exponential backoff |\n| `retry.maxDelayMs` | number | `60000` | Maximum delay cap |\n| `retry.jitter` | boolean | `true` | Add random jitter to delays |\n| `monitor.enabled` | boolean | `false` | Enable claim monitor |\n| `monitor.scanIntervalMs` | number | `10000` | Interval between PEL scans |\n| `monitor.minIdleTimeMs` | number | `60000` | Minimum idle time to claim |\n| `monitor.claimBatchSize` | number | `100` | Max messages to claim per scan |\n| `dlq.maxReplays` | number | `3` | Max replays before freezing |\n\n### Using Environment Variables\n\n```typescript\nimport Redis from 'ioredis';\nimport { configFromEnv, withDefaults, Producer } from 'thinrsmq';\n\n// Read Redis connection from environment variables\nconst envConfig = configFromEnv();\nconst config = withDefaults({ ...envConfig, namespace: 'myapp' });\n\nconst [host, port] = config.redis.address.split(':');\nconst redis = new Redis({\n  host,\n  port: parseInt(port, 10),\n  password: config.redis.password || undefined,\n  tls: config.redis.useTls ? { servername: host } : undefined,\n});\n\nconst producer = new Producer(redis, config);\n```\n\n### Environment Variables\n\n| Variable | Description | Default |\n|----------|-------------|---------|\n| `REDIS_HOST` | Redis hostname | `\"localhost\"` |\n| `REDIS_PORT` | Redis port | `\"6379\"` |\n| `REDIS_PASSWORD` | Redis password | `\"\"` |\n| `REDIS_USE_TLS` | Enable TLS (`\"true\"` or `\"1\"`) | `false` |\n\n## API Reference\n\n### Producer\n\n#### `new Producer(redis, config)`\n\nCreates a new producer instance.\n\n- `redis`: ioredis client instance\n- `config`: Configuration object\n\n#### `publish(topic: string, msg: Omit<Message, 'id'>): Promise<string>`\n\nPublishes a single message. Returns the stream entry ID.\n\n- `topic`: Topic name (stream will be created as `{namespace}:{topic}`)\n- `msg`: Message object without the `id` field\n  - `type`: Message type identifier\n  - `payload`: Message payload (string)\n  - `version`: Envelope version (use `\"1\"`)\n  - `producedAt`: ISO 8601 timestamp (auto-set if empty)\n  - `traceId`: Optional trace ID\n  - `producer`: Optional producer identifier\n\n#### `publishBatch(topic: string, msgs: Omit<Message, 'id'>[]): Promise<string[]>`\n\nPublishes multiple messages in a pipeline. Returns array of stream entry IDs.\n\n### Consumer\n\n#### `new Consumer(redis, config)`\n\nCreates a new consumer instance with a unique name.\n\n#### `subscribe(topic: string, group: string, handler: Handler): Promise<void>`\n\nSubscribes to a topic with a consumer group. Starts processing messages.\n\n- `topic`: Topic name\n- `group`: Consumer group name\n- `handler`: `async (msg: Message) => void` - Message handler function\n\nThrows error to trigger retry. Success (no throw) acknowledges the message.\n\n#### `stop(): Promise<void>`\n\nGracefully stops the consumer. Waits for in-flight messages to complete (up to `shutdownTimeoutMs`).\n\n#### `onSkip(handler: SkipHandler): void`\n\nRegisters a callback for skipped messages (nil/trimmed or unknown version).\n\n- `handler`: `(id: string, reason: string) => void`\n\n#### `getConsumerName(): string`\n\nReturns the unique consumer name (format: `{prefix}-{hostname}-{pid}-{uuid}`).\n\n### ClaimMonitor\n\n#### `new ClaimMonitor(redis, config)`\n\nCreates a new claim monitor instance.\n\n#### `start(): Promise<void>`\n\nStarts the claim monitor. Periodically scans for idle messages and routes them to retry or DLQ.\n\n#### `stop(): Promise<void>`\n\nStops the claim monitor.\n\n#### `scanOnce(): Promise<void>`\n\nPerforms a single scan cycle. Useful for testing.\n\n### DLQ\n\n#### `new DLQ(redis, config)`\n\nCreates a new DLQ instance.\n\n#### `moveToDLQ(topic: string, msg: Message, failureMetadata): Promise<void>`\n\nMoves a message to the DLQ.\n\n- `failureMetadata`:\n  - `originalStream`: Original stream key\n  - `failedAt`: Failure timestamp\n  - `totalAttempts`: Total retry attempts\n  - `lastError`: Error message\n  - `consumerGroup`: Consumer group name\n\n#### `peek(topic: string, count: number): Promise<DLQMessage[]>`\n\nReturns up to `count` messages from the DLQ (oldest first).\n\n#### `replay(topic: string, messageId: string): Promise<void>`\n\nReplays a DLQ message back to the original stream. Increments `replay_count` and freezes if `>= maxReplays`.\n\n#### `purge(topic: string): Promise<void>`\n\nDeletes all messages from the DLQ for a topic.\n\n#### `size(topic: string): Promise<number>`\n\nReturns the number of messages in the DLQ.\n\n### Admin\n\n#### `new Admin(redis, config)`\n\nCreates a new admin instance.\n\n#### `pendingStats(topic: string, group: string): Promise<PendingInfo>`\n\nReturns pending message statistics.\n\n- Returns: `{ count, minId, maxId, consumers }`\n\n#### `consumerInfo(topic: string, group: string): Promise<ConsumerDetail[]>`\n\nReturns consumer details for a group.\n\n- Returns: Array of `{ name, pending, idle }`\n\n#### `streamInfo(topic: string): Promise<StreamDetail>`\n\nReturns stream metadata.\n\n- Returns: `{ length, firstEntryId, lastEntryId, groups }`\n\n#### DLQ Operations\n\nAdmin also provides DLQ operations:\n\n- `dlqSize(topic)`: Alias for `DLQ.size()`\n- `dlqPeek(topic, count)`: Alias for `DLQ.peek()`\n- `dlqReplay(topic, messageId)`: Alias for `DLQ.replay()`\n- `dlqPurge(topic)`: Alias for `DLQ.purge()`\n\n### RetryStore\n\n#### `new RetryStore(redis, config)`\n\nCreates a new retry store instance.\n\n#### `initIfNotExists(topic: string, messageId: string): Promise<void>`\n\nInitializes retry hash if it doesn't exist (HSETNX for single-writer rule).\n\n#### `get(topic: string, messageId: string): Promise<RetryInfo | null>`\n\nGets retry information for a message.\n\n- Returns: `{ attempt, lastAttemptAt }` or `null`\n\n#### `set(topic: string, messageId: string, info: RetryInfo): Promise<void>`\n\nSets retry information.\n\n#### `delete(topic: string, messageId: string): Promise<void>`\n\nDeletes retry hash.\n\n#### `incrementAttempt(topic: string, messageId: string): Promise<number>`\n\nIncrements attempt count and returns new value.\n\n### Helper Functions\n\n#### `configFromEnv(): { redis: Partial<RedisConfig> }`\n\nReads Redis connection settings from environment variables.\n\n#### `withDefaults(config: Partial<Config>): Config`\n\nMerges provided config with defaults.\n\n#### `validateConfig(config: Config): void`\n\nValidates configuration. Throws `ConfigError` on validation failure.\n\n#### `computeDelay(attempt: number, cfg: BackoffConfig): number`\n\nComputes exponential backoff delay in milliseconds.\n\n- Formula: `min(baseDelay * 2^(attempt-1) + jitter, maxDelay)`\n\n### Key Functions\n\n#### `streamKey(namespace: string, topic: string): string`\n\nReturns stream key: `\"{namespace}:{topic}\"`\n\n#### `dlqKey(namespace: string, topic: string): string`\n\nReturns DLQ key: `\"{namespace}:{topic}:dlq\"`\n\n#### `retryKey(namespace: string, topic: string, messageId: string): string`\n\nReturns retry hash key: `\"{namespace}:{topic}:retries:{messageId}\"`\n\n### Types\n\n#### `Message`\n\n```typescript\ntype Message = {\n  id: string;\n  version: string;\n  type: string;\n  payload: string;\n  traceId?: string;\n  producedAt: string;\n  producer?: string;\n};\n```\n\n#### `DLQMessage`\n\n```typescript\ntype DLQMessage = Message & {\n  originalId: string;\n  originalStream: string;\n  failedAt: string;\n  totalAttempts: number;\n  lastError: string;\n  consumerGroup: string;\n  replayCount: number;\n};\n```\n\n### Error Types\n\n- `ThinrsmqError`: Base error class\n- `ConfigError`: Configuration validation error\n- `UnknownVersionError`: Thrown when envelope version is not \"1\"\n- `MessageTrimmedError`: Thrown when message was trimmed from stream\n- `MissingFieldError`: Thrown when required field is missing\n\n## CLI Tools\n\nCLI tools are provided for development and testing.\n\n### Producer CLI\n\n```bash\nnpx tsx src/cli/producer.ts --topic orders --type order.created --payload '{\"orderId\":\"123\"}'\n```\n\n### Consumer CLI\n\n```bash\nnpx tsx src/cli/consumer.ts --topic orders --group order-processor\n```\n\n### Admin CLI\n\n```bash\n# Pending stats\nnpx tsx src/cli/admin.ts --topic orders --group order-processor --command pending\n\n# Consumer info\nnpx tsx src/cli/admin.ts --topic orders --group order-processor --command consumers\n\n# Stream info\nnpx tsx src/cli/admin.ts --topic orders --command stream\n```\n\n## Wire Format\n\nMessages use envelope version \"1\" with the following fields:\n\n| Field | Required | Description |\n|-------|----------|-------------|\n| `v` | Yes | Envelope version (always \"1\") |\n| `type` | Yes | Message type identifier |\n| `payload` | Yes | Message payload (string) |\n| `produced_at` | Yes | ISO 8601 timestamp |\n| `trace_id` | No | Trace ID for distributed tracing |\n| `producer` | No | Producer identifier |\n\nDLQ messages include additional enrichment fields:\n\n- `original_id`: Original stream entry ID\n- `original_stream`: Original stream key\n- `failed_at`: Failure timestamp\n- `total_attempts`: Total retry attempts\n- `last_error`: Error message (truncated to 1000 chars)\n- `consumer_group`: Consumer group name\n- `replay_count`: Number of times replayed from DLQ\n\n## Cross-Language Interoperability\n\nThis Node.js implementation is wire-compatible with the Go implementation. Messages produced by one can be consumed by the other seamlessly.\n\nSee [Go README](../thinrsmq-go/README.md) for the Go implementation.\n\n## License\n\nMIT License. See [LICENSE](./LICENSE) for details.\n","readmeFilename":"README.md","_rev":"1-5e4dbe3d0beb4688253598a78e7b5955"}