{"_id":"@aimf/event-mesh","_rev":"2-a0ea9739cb1721676fefb242bb9bbefd","name":"@aimf/event-mesh","dist-tags":{"latest":"0.1.1"},"versions":{"0.1.0":{"name":"@aimf/event-mesh","version":"0.1.0","keywords":["aimf","event-mesh","pubsub","events","messaging","reactive"],"license":"MIT","_id":"@aimf/event-mesh@0.1.0","maintainers":[{"name":"vetriselvans","email":"vetriselvans029@gmail.com"}],"dist":{"shasum":"b4d1c021cb71aa92fb481993160c3d92ceab079c","tarball":"https://registry.npmjs.org/@aimf/event-mesh/-/event-mesh-0.1.0.tgz","fileCount":30,"integrity":"sha512-uj7KIW4EFJy1oNNg+SB9LUJqP1BdHCQXBUOAz8hyWM8CqG/53az4KGs1GeVjOEg+LpdaQrhHzo//O4de/5GJdg==","signatures":[{"sig":"MEUCIQDT51hvuYP/vwNF4SR1OILIROEB1qi7bkMXTEb+z2LWQgIgMb4jg7Ali6hzoFVV9r5kOT/IA5ugLtGjN2raDVcdW7o=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":113124},"main":"dist/index.js","type":"module","_from":"file:aimf-event-mesh-0.1.0.tgz","types":"dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js"}},"scripts":{"lint":"eslint src/","test":"vitest run","build":"tsc","clean":"rimraf dist","test:watch":"vitest"},"_npmUser":{"name":"vetriselvans","email":"vetriselvans029@gmail.com"},"_resolved":"C:\\Users\\I9159~1.LAP\\AppData\\Local\\Temp\\59d443a9378a3c17a0a82ab59cbb10e5\\aimf-event-mesh-0.1.0.tgz","_integrity":"sha512-uj7KIW4EFJy1oNNg+SB9LUJqP1BdHCQXBUOAz8hyWM8CqG/53az4KGs1GeVjOEg+LpdaQrhHzo//O4de/5GJdg==","_npmVersion":"10.9.0","description":"Real-time Event Mesh for Pub/Sub Messaging in AIMF","directories":{},"_nodeVersion":"22.12.0","dependencies":{"@aimf/shared":"0.1.0"},"_hasShrinkwrap":false,"devDependencies":{"rimraf":"^5.0.0","vitest":"^2.1.9","typescript":"^5.9.3"},"_npmOperationalInternal":{"tmp":"tmp/event-mesh_0.1.0_1765128989808_0.06223435485458162","host":"s3://npm-registry-packages-npm-production"}},"0.1.1":{"name":"@aimf/event-mesh","version":"0.1.1","description":"Real-time Event Mesh for Pub/Sub Messaging in AIMF","type":"module","main":"dist/index.js","types":"dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js"}},"dependencies":{"@aimf/shared":"0.1.1"},"devDependencies":{"typescript":"^5.9.3","vitest":"^2.1.9","rimraf":"^5.0.0"},"keywords":["aimf","event-mesh","pubsub","events","messaging","reactive"],"license":"MIT","scripts":{"build":"tsc","test":"vitest run","test:watch":"vitest","lint":"eslint src/","clean":"rimraf dist"},"_id":"@aimf/event-mesh@0.1.1","_integrity":"sha512-169z/FnN3HTGhF8yn6HR/utFa6IWoJzOA3PVomAkh9kAOqusG5OuWTHyAXMuGUaZpG86JrdHnpVQFkd74xoilg==","_resolved":"C:\\Users\\I9159~1.LAP\\AppData\\Local\\Temp\\055421e82d480637da17f373074a773e\\aimf-event-mesh-0.1.1.tgz","_from":"file:aimf-event-mesh-0.1.1.tgz","_nodeVersion":"22.12.0","_npmVersion":"10.9.0","dist":{"integrity":"sha512-169z/FnN3HTGhF8yn6HR/utFa6IWoJzOA3PVomAkh9kAOqusG5OuWTHyAXMuGUaZpG86JrdHnpVQFkd74xoilg==","shasum":"ef873e7822fc4a855986b77ae676bbbcfcb395ed","tarball":"https://registry.npmjs.org/@aimf/event-mesh/-/event-mesh-0.1.1.tgz","fileCount":30,"unpackedSize":113124,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIQDt1BxjNR1+FxJStCMGta0hcjbS4zFWwigHBe2BO/z02wIgb5SdLL4+FuwVQVEz/juHPjJFv30grGahfi1WykBEh00="}]},"_npmUser":{"name":"vetriselvans","email":"vetriselvans029@gmail.com"},"directories":{},"maintainers":[{"name":"vetriselvans","email":"vetriselvans029@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/event-mesh_0.1.1_1765130855851_0.3830236966049805"},"_hasShrinkwrap":false}},"time":{"created":"2025-12-07T17:36:29.717Z","modified":"2025-12-07T18:07:36.306Z","0.1.0":"2025-12-07T17:36:29.944Z","0.1.1":"2025-12-07T18:07:36.048Z"},"license":"MIT","keywords":["aimf","event-mesh","pubsub","events","messaging","reactive"],"description":"Real-time Event Mesh for Pub/Sub Messaging in AIMF","maintainers":[{"name":"vetriselvans","email":"vetriselvans029@gmail.com"}],"readme":"# @aimf/event-mesh\r\n\r\nReal-time Event Mesh for Pub/Sub Messaging in the AI-MCP Framework.\r\n\r\n## Overview\r\n\r\nThe Event Mesh package provides a high-performance, topic-based publish/subscribe messaging system with advanced features like pattern matching, typed channels, request/reply, and event aggregation.\r\n\r\n## Features\r\n\r\n- **Topic-based Pub/Sub**: Flexible topic patterns with wildcards (`*` and `#`)\r\n- **Typed Channels**: Type-safe event channels with TypeScript generics\r\n- **Request/Reply**: RPC-style communication over the event bus\r\n- **Event Aggregation**: Batch and window events by count or time\r\n- **Dead Letter Queue**: Capture and handle failed events\r\n- **Priority Handling**: Execute handlers in priority order\r\n- **Event Filtering**: Filter events before they reach handlers\r\n\r\n## Installation\r\n\r\n```bash\r\npnpm add @aimf/event-mesh\r\n```\r\n\r\n## Quick Start\r\n\r\n```typescript\r\nimport { createEventBus, createChannel } from \"@aimf/event-mesh\";\r\n\r\n// Create an event bus\r\nconst bus = createEventBus();\r\n\r\n// Simple pub/sub\r\nbus.subscribe(\"users.created\", (data) => {\r\n  console.log(\"User created:\", data);\r\n});\r\n\r\nawait bus.publish(\"users.created\", { id: 1, name: \"Alice\" });\r\n\r\n// Pattern matching with wildcards\r\nbus.subscribe(\"users.*\", (data) => {\r\n  console.log(\"User event:\", data);\r\n});\r\n\r\nbus.subscribe(\"users.#\", (data) => {\r\n  console.log(\"Any user event:\", data);\r\n});\r\n\r\n// Clean up\r\nawait bus.destroy();\r\n```\r\n\r\n## Topic Patterns\r\n\r\nThe event mesh supports two types of wildcards:\r\n\r\n- `*` - Matches exactly one segment\r\n- `#` - Matches zero or more segments\r\n\r\n```typescript\r\n// Match any direct child topic\r\nbus.subscribe(\"users.*\", handler);\r\n// Matches: users.created, users.updated, users.deleted\r\n// Does NOT match: users.profile.updated\r\n\r\n// Match any descendant topic\r\nbus.subscribe(\"users.#\", handler);\r\n// Matches: users, users.created, users.profile.updated\r\n\r\n// Complex patterns\r\nbus.subscribe(\"*.events.#\", handler);\r\n// Matches: users.events.login, orders.events.created.success\r\n```\r\n\r\n## Typed Channels\r\n\r\nCreate type-safe channels for specific event types:\r\n\r\n```typescript\r\nimport { createChannel, createChannelRegistry } from \"@aimf/event-mesh\";\r\n\r\ninterface UserCreatedEvent {\r\n  id: number;\r\n  name: string;\r\n  email: string;\r\n}\r\n\r\n// Create a typed channel\r\nconst userChannel = createChannel<UserCreatedEvent>(bus, \"users.created\");\r\n\r\n// Subscribe with type safety\r\nuserChannel.subscribe((user) => {\r\n  console.log(user.name); // TypeScript knows this is a string\r\n});\r\n\r\n// Publish with type checking\r\nawait userChannel.publish({\r\n  id: 1,\r\n  name: \"Alice\",\r\n  email: \"alice@example.com\"\r\n});\r\n\r\n// Channel Registry for managing multiple channels\r\nconst registry = createChannelRegistry(bus);\r\n\r\nregistry.register<UserCreatedEvent>(\"users.created\");\r\nregistry.register<{ orderId: string }>(\"orders.created\");\r\n\r\nconst channel = registry.get(\"users.created\");\r\n```\r\n\r\n## Request/Reply Pattern\r\n\r\nImplement RPC-style communication:\r\n\r\n```typescript\r\nimport { createRequestClient, createRequestResponder } from \"@aimf/event-mesh\";\r\n\r\n// Server side - create responder\r\nconst responder = createRequestResponder(bus, \"math.add\");\r\nresponder.handle<{ a: number; b: number }, number>(({ a, b }) => a + b);\r\n\r\n// Client side - send requests\r\nconst client = createRequestClient(bus);\r\n\r\nconst result = await client.request<{ a: number; b: number }, number>(\r\n  \"math.add\",\r\n  { a: 5, b: 3 }\r\n);\r\nconsole.log(result); // 8\r\n\r\n// Clean up\r\nresponder.stop();\r\n```\r\n\r\n## Event Aggregation\r\n\r\nBatch events by count or time window:\r\n\r\n```typescript\r\nimport { createEventAggregator, aggregators } from \"@aimf/event-mesh\";\r\n\r\nconst aggregator = createEventAggregator(bus);\r\n\r\n// Count-based aggregation\r\naggregator.addRule({\r\n  id: \"batch-orders\",\r\n  pattern: \"orders.#\",\r\n  windowType: \"count\",\r\n  windowSize: 10, // Emit after 10 events\r\n  outputTopic: \"orders.batch\",\r\n  aggregator: aggregators.collect,\r\n});\r\n\r\n// Time-based aggregation\r\naggregator.addRule({\r\n  id: \"metrics-window\",\r\n  pattern: \"metrics.#\",\r\n  windowType: \"time\",\r\n  windowSize: 5000, // Emit every 5 seconds\r\n  outputTopic: \"metrics.batch\",\r\n  aggregator: aggregators.sum,\r\n});\r\n\r\n// Subscribe to batched events\r\nbus.subscribe(\"orders.batch\", (batch) => {\r\n  console.log(`Processing ${batch.length} orders`);\r\n});\r\n\r\n// Built-in aggregators\r\n// - aggregators.collect - Collect into array\r\n// - aggregators.sum - Sum numeric values\r\n// - aggregators.average - Average numeric values\r\n// - aggregators.count - Count events\r\n// - aggregators.first - Get first value\r\n// - aggregators.last - Get last value\r\n// - aggregators.min - Get minimum value\r\n// - aggregators.max - Get maximum value\r\n// - aggregators.groupBy(key) - Group by object key\r\n\r\n// Clean up\r\naggregator.destroy();\r\n```\r\n\r\n## Event Filtering\r\n\r\nFilter events before they reach handlers:\r\n\r\n```typescript\r\nbus.subscribe(\r\n  \"users.#\",\r\n  (user) => {\r\n    console.log(\"Admin user:\", user);\r\n  },\r\n  {\r\n    filter: (data) => data.role === \"admin\",\r\n  }\r\n);\r\n```\r\n\r\n## Priority Handling\r\n\r\nControl handler execution order:\r\n\r\n```typescript\r\n// Higher priority handlers execute first\r\nbus.subscribe(\"events.important\", handler1, { priority: 1 });\r\nbus.subscribe(\"events.important\", handler2, { priority: 10 }); // Runs first\r\nbus.subscribe(\"events.important\", handler3, { priority: 5 });\r\n```\r\n\r\n## Dead Letter Queue\r\n\r\nCapture and handle failed events:\r\n\r\n```typescript\r\n// Listen for dead letters\r\nbus.onDeadLetter((deadLetter) => {\r\n  console.error(`Event failed: ${deadLetter.topic}`, deadLetter.error);\r\n  // Retry or log the failure\r\n});\r\n\r\n// Get dead letter queue\r\nconst deadLetters = bus.getDeadLetterQueue();\r\n\r\n// Clear the queue\r\nbus.clearDeadLetterQueue();\r\n```\r\n\r\n## Waiting for Events\r\n\r\nWait for specific events with timeout:\r\n\r\n```typescript\r\n// Wait for an event (with timeout)\r\nconst envelope = await bus.waitFor(\"users.verified\", 5000);\r\nconsole.log(\"User verified:\", envelope.data);\r\n\r\n// One-time subscription\r\nbus.once(\"startup.complete\", () => {\r\n  console.log(\"Application started\");\r\n});\r\n```\r\n\r\n## Event Envelope\r\n\r\nAccess full event metadata:\r\n\r\n```typescript\r\nbus.subscribe(\"events.test\", (data, envelope) => {\r\n  console.log(\"Event ID:\", envelope.metadata.id);\r\n  console.log(\"Timestamp:\", envelope.metadata.timestamp);\r\n  console.log(\"Correlation ID:\", envelope.metadata.correlationId);\r\n  console.log(\"Source:\", envelope.metadata.source);\r\n});\r\n\r\n// Publish with metadata\r\nawait bus.publish(\"events.test\", data, {\r\n  correlationId: \"request-123\",\r\n  source: \"user-service\",\r\n});\r\n```\r\n\r\n## Configuration\r\n\r\n```typescript\r\nconst bus = createEventBus({\r\n  maxDeadLetterSize: 1000, // Maximum dead letters to keep\r\n  generateId: () => crypto.randomUUID(), // Custom ID generator\r\n});\r\n```\r\n\r\n## Statistics\r\n\r\nMonitor event bus activity:\r\n\r\n```typescript\r\nconst stats = bus.getStatistics();\r\nconsole.log(\"Published:\", stats.publishedCount);\r\nconsole.log(\"Subscriptions:\", stats.subscriptionCount);\r\nconsole.log(\"Dead Letters:\", stats.deadLetterCount);\r\n\r\n// Aggregator statistics\r\nconst aggStats = aggregator.getStatistics(\"rule-id\");\r\nconsole.log(\"Events Processed:\", aggStats.eventsProcessed);\r\nconsole.log(\"Batches Emitted:\", aggStats.batchesEmitted);\r\n```\r\n\r\n## Namespaced Channels\r\n\r\nCreate isolated channel namespaces:\r\n\r\n```typescript\r\nimport { createNamespacedChannels } from \"@aimf/event-mesh\";\r\n\r\nconst appChannels = createNamespacedChannels(bus, \"myapp\");\r\n\r\nconst usersChannel = appChannels.channel<UserEvent>(\"users\");\r\n// Topic: myapp.users\r\n\r\nconst ordersChannel = appChannels.channel<OrderEvent>(\"orders\");\r\n// Topic: myapp.orders\r\n```\r\n\r\n## API Reference\r\n\r\n### EventBus\r\n\r\n| Method | Description |\r\n|--------|-------------|\r\n| `subscribe(topic, handler, options?)` | Subscribe to events |\r\n| `unsubscribe(subscriptionId)` | Unsubscribe by ID |\r\n| `publish(topic, data, metadata?)` | Publish an event |\r\n| `once(topic, handler)` | Subscribe for one event |\r\n| `waitFor(topic, timeout?)` | Wait for an event |\r\n| `getSubscriptions()` | Get all subscriptions |\r\n| `getTopics()` | Get unique topics |\r\n| `hasSubscribers(topic)` | Check for subscribers |\r\n| `getDeadLetterQueue()` | Get failed events |\r\n| `clearDeadLetterQueue()` | Clear dead letters |\r\n| `onDeadLetter(handler)` | Handle dead letters |\r\n| `getStatistics()` | Get bus statistics |\r\n| `destroy()` | Clean up resources |\r\n\r\n### Topic Matching\r\n\r\n| Function | Description |\r\n|----------|-------------|\r\n| `parseTopicPattern(pattern)` | Parse a topic pattern |\r\n| `matchTopic(pattern, topic)` | Match topic against pattern |\r\n| `validateTopic(topic, allowPatterns?)` | Validate topic format |\r\n| `isPattern(topic)` | Check if topic contains wildcards |\r\n\r\n### Channel Functions\r\n\r\n| Function | Description |\r\n|----------|-------------|\r\n| `createChannel(bus, topic)` | Create a typed channel |\r\n| `createChannelRegistry(bus)` | Create a channel registry |\r\n| `createTypedChannel(bus, topic, options)` | Create with options |\r\n| `createNamespacedChannels(bus, namespace)` | Create namespaced factory |\r\n\r\n### Request/Reply\r\n\r\n| Function | Description |\r\n|----------|-------------|\r\n| `createRequestClient(bus, options?)` | Create request client |\r\n| `createRequestResponder(bus, topic)` | Create responder |\r\n| `createRequestReply(bus, topic)` | Create paired client/responder |\r\n\r\n### Aggregator\r\n\r\n| Function | Description |\r\n|----------|-------------|\r\n| `createEventAggregator(bus)` | Create aggregator |\r\n| `aggregators.collect` | Collect events into array |\r\n| `aggregators.sum` | Sum numeric events |\r\n| `aggregators.average` | Average numeric events |\r\n| `aggregators.count` | Count events |\r\n| `aggregators.first` | Get first event |\r\n| `aggregators.last` | Get last event |\r\n| `aggregators.min` | Get minimum |\r\n| `aggregators.max` | Get maximum |\r\n| `aggregators.groupBy(key)` | Group by object key |\r\n\r\n## License\r\n\r\nMIT\r\n","readmeFilename":"README.md"}