{"_id":"@cloudamqp/a2a-amqp","_rev":"2-fbe43718e30849ea27f13193cef9e077","name":"@cloudamqp/a2a-amqp","dist-tags":{"latest":"0.1.1"},"versions":{"0.1.0":{"name":"@cloudamqp/a2a-amqp","version":"0.1.0","keywords":["a2a","agent-to-agent","amqp","lavinmq","eventbus","ai-agents","work-queue","event-sourcing"],"author":{"name":"cloudamqp"},"license":"MIT","_id":"@cloudamqp/a2a-amqp@0.1.0","maintainers":[{"name":"carlhoerberg","email":"carl.hoerberg@gmail.com"},{"name":"antondalgren","email":"anton@84codes.com"},{"name":"patrik84c","email":"patrik@84codes.com"},{"name":"abaelter","email":"abaelter@gmail.com"}],"homepage":"https://github.com/cloudamqp/a2a-amqp#readme","bugs":{"url":"https://github.com/cloudamqp/a2a-amqp/issues"},"dist":{"shasum":"e8c184d5efd2852be6e235d9f3dfc74b2f3bc461","tarball":"https://registry.npmjs.org/@cloudamqp/a2a-amqp/-/a2a-amqp-0.1.0.tgz","fileCount":54,"integrity":"sha512-Z7h0DATEoTNAuxHR55nXNRSczchMC2vRoH5ATatCledQUMiy9ZWr26Cv6QN3uI+Wj4eQPxakypWPh3swS7n/wQ==","signatures":[{"sig":"MEUCIQDzczKxMwHLYaOaA5AiFDABrbMIMdw2ds9J8V5d5Yt4CwIgQUqzdhfPAbTukoKAuJRegsykAFn66vKT9sAvZlxsUgE=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":209831},"main":"./dist/index.js","type":"module","types":"./dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js"}},"gitHead":"75b3c367fefbdc4b7b3e8ca6358a674e1dca5565","scripts":{"test":"bun run test:unit && bun run test:integration","build":"tsc","server":"bun run src/examples/http-server.ts","worker":"bun run src/examples/worker.ts","prepare":"bun run build","test:ui":"vitest --ui","test:unit":"vitest run --config vitest.config.ts","test:watch":"vitest watch --config vitest.config.ts","type-check":"tsc --project tsconfig.typecheck.json","release:major":"bun run scripts/release.ts major","release:minor":"bun run scripts/release.ts minor","release:patch":"bun run scripts/release.ts patch","test:coverage":"vitest run --config vitest.config.ts --coverage","test:integration":"vitest run --config vitest.integration.config.ts"},"_npmUser":{"name":"antondalgren","email":"anton@84codes.com"},"repository":{"url":"git+https://github.com/cloudamqp/a2a-amqp.git","type":"git"},"_npmVersion":"10.9.2","description":"AMQP EventBus adapter for A2A protocol with LavinMQ streams support","directories":{"test":"tests"},"_nodeVersion":"22.14.0","dependencies":{"zod":"^4.1.12","uuid":"^13.0.0","@cloudamqp/amqp-client":"^3.4.0"},"_hasShrinkwrap":false,"devDependencies":{"vitest":"^4.0.9","@types/bun":"latest","@vitest/ui":"^4.0.9","typescript":"^5.9.3","@types/express":"^5.0.5","@vitest/coverage-v8":"^4.0.9"},"peerDependencies":{"express":"*","@a2a-js/sdk":"*"},"_npmOperationalInternal":{"tmp":"tmp/a2a-amqp_0.1.0_1763458858781_0.652843005137177","host":"s3://npm-registry-packages-npm-production"}},"0.1.1":{"name":"@cloudamqp/a2a-amqp","version":"0.1.1","description":"AMQP EventBus adapter for A2A protocol with LavinMQ streams support","type":"module","main":"./dist/index.js","types":"./dist/index.d.ts","exports":{".":{"import":"./dist/index.js","types":"./dist/index.d.ts"}},"scripts":{"build":"tsc","type-check":"tsc --project tsconfig.typecheck.json","prepare":"bun run build","server":"bun run src/examples/http-server.ts","worker":"bun run src/examples/worker.ts","test":"bun run test:unit && bun run test:integration","test:unit":"vitest run --config vitest.config.ts","test:integration":"vitest run --config vitest.integration.config.ts","test:watch":"vitest watch --config vitest.config.ts","test:coverage":"vitest run --config vitest.config.ts --coverage","test:ui":"vitest --ui","release:patch":"bun run scripts/release.ts patch","release:minor":"bun run scripts/release.ts minor","release:major":"bun run scripts/release.ts major"},"repository":{"type":"git","url":"git+https://github.com/cloudamqp/a2a-amqp.git"},"keywords":["a2a","agent-to-agent","amqp","lavinmq","eventbus","ai-agents","work-queue","event-sourcing"],"author":{"name":"cloudamqp"},"license":"MIT","publishConfig":{"access":"public"},"dependencies":{"@cloudamqp/amqp-client":"^3.4.0","uuid":"^13.0.0","zod":"^4.1.12"},"devDependencies":{"@types/bun":"latest","@types/express":"^5.0.5","@vitest/coverage-v8":"^4.0.9","@vitest/ui":"^4.0.9","typescript":"^5.9.3","vitest":"^4.0.9"},"peerDependencies":{"@a2a-js/sdk":"*","express":"*"},"directories":{"test":"tests"},"bugs":{"url":"https://github.com/cloudamqp/a2a-amqp/issues"},"homepage":"https://github.com/cloudamqp/a2a-amqp#readme","gitHead":"60f9eb5854d77d1c85072ee32dc1f1d5213503f1","_id":"@cloudamqp/a2a-amqp@0.1.1","_nodeVersion":"20.19.5","_npmVersion":"11.6.2","dist":{"integrity":"sha512-OqSZmHEyu/+ULPyTBu/1jbTxfYH8m6FEpv+nlHkne9Zpe8sMEd5UoHt1ilNIchCwmwLAVq9GRdNkEO6oOApTWg==","shasum":"de9476cf850acb9e3b84dc14fea334f5f19185fb","tarball":"https://registry.npmjs.org/@cloudamqp/a2a-amqp/-/a2a-amqp-0.1.1.tgz","fileCount":42,"unpackedSize":130676,"attestations":{"url":"https://registry.npmjs.org/-/npm/v1/attestations/@cloudamqp%2fa2a-amqp@0.1.1","provenance":{"predicateType":"https://slsa.dev/provenance/v1"}},"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIQDJVYAF32UKvLICnfzU3f4cltd+ByFMnpqnpmzVURBxGAIgLWebhR86SNdsLU/fdd33TNDdxq/GgFpnc6q2iIMCS9Q="}]},"_npmUser":{"name":"GitHub Actions","email":"npm-oidc-no-reply@github.com","trustedPublisher":{"id":"github","oidcConfigId":"oidc:ce4ae718-d1ea-41c6-8577-3c8e62de79c3"}},"maintainers":[{"name":"carlhoerberg","email":"carl.hoerberg@gmail.com"},{"name":"antondalgren","email":"anton@84codes.com"},{"name":"patrik84c","email":"patrik@84codes.com"},{"name":"abaelter","email":"abaelter@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/a2a-amqp_0.1.1_1763460419729_0.13477567212031727"},"_hasShrinkwrap":false}},"time":{"created":"2025-11-18T09:40:58.636Z","modified":"2025-11-18T10:07:00.444Z","0.1.0":"2025-11-18T09:40:58.958Z","0.1.1":"2025-11-18T10:06:59.945Z"},"bugs":{"url":"https://github.com/cloudamqp/a2a-amqp/issues"},"author":{"name":"cloudamqp"},"license":"MIT","homepage":"https://github.com/cloudamqp/a2a-amqp#readme","keywords":["a2a","agent-to-agent","amqp","lavinmq","eventbus","ai-agents","work-queue","event-sourcing"],"repository":{"type":"git","url":"git+https://github.com/cloudamqp/a2a-amqp.git"},"description":"AMQP EventBus adapter for A2A protocol with LavinMQ streams support","maintainers":[{"name":"carlhoerberg","email":"carl.hoerberg@gmail.com"},{"name":"antondalgren","email":"anton@84codes.com"},{"name":"patrik84c","email":"patrik@84codes.com"},{"name":"abaelter","email":"abaelter@gmail.com"}],"readme":"# a2a-amqp\n\nAMQP-backed EventBus and WorkQueue for scaling A2A agents with long-running tasks.\n\n## Why?\n\nA2A agents often need to handle long-running tasks (LLM calls, complex processing, etc.). Running these tasks inline in HTTP handlers causes:\n\n- **Timeout issues**: HTTP connections timeout on long tasks\n- **Scaling problems**: Single server bottlenecks\n- **Resource waste**: Servers blocked waiting for tasks to complete\n\nThis library solves these problems by:\n\n1. **Queuing tasks** via AMQP instead of processing inline\n2. **Distributing work** across multiple worker processes\n3. **Event sourcing** all task events for replay and recovery\n4. **Streaming results** back via SSE while workers process in the background\n\n## Architecture\n\n```\nHTTP Request → Server (enqueues task) → Returns immediately\n                    ↓\n               AMQP Queue\n                    ↓\n          Worker Pool (scales horizontally)\n                    ↓\n          Process task & publish events\n                    ↓\n          AMQP Stream (event sourcing)\n                    ↓\n          Client streams results via SSE\n```\n\n## Installation\n\n```bash\n# Installing using bun\nbun add @cloudamqp/a2a-amqp @a2a-js/sdk @cloudamqp/amqp-client\n\n# Or with npm\nnpm install @cloudamqp/a2a-amqp @a2a-js/sdk @cloudamqp/amqp-client\n```\n\n**Requires**: LavinMQ or RabbitMQ with stream support\n\n```bash\ndocker run -d -p 5672:5672 -p 15672:15672 cloudamqp/lavinmq:latest\n```\n\n## Quick Start\n\n### 1. HTTP Server (enqueues tasks)\n\n```typescript\nimport { AMQPAgentBackend, QueuingRequestHandler } from \"@84codes/a2a-amqp\";\nimport { A2AExpressApp } from \"@a2a-js/sdk/server/express\";\nimport express from \"express\";\n\n// Create AMQP backend\nconst backend = await AMQPAgentBackend.create({\n  url: \"amqp://localhost:5672\",\n  agentName: \"my-agent\",\n});\n\n// Create request handler (handles task queuing + event projection)\nconst requestHandler = new QueuingRequestHandler(agentCard, backend);\nawait requestHandler.initialize();\n\n// Setup Express with A2A routes\nconst app = express();\nnew A2AExpressApp(requestHandler).setupRoutes(app, \"/\");\napp.listen(3000);\n```\n\n### 2. Worker Process (processes tasks)\n\n```typescript\nimport { AMQPAgentBackend, WorkerEventBus } from \"@84codes/a2a-amqp\";\nimport { AgentExecutor, RequestContext } from \"@a2a-js/sdk/server\";\n\n// Create backend with same agent name as server\nconst backend = await AMQPAgentBackend.create({\n  url: \"amqp://localhost:5672\",\n  agentName: \"my-agent\",\n});\n\n// Initialize work queue\nawait backend.workQueue.initialize();\n\nclass MyExecutor implements AgentExecutor {\n  async execute(context: RequestContext, eventBus: ExecutionEventBus) {\n    // Your long-running task logic here\n    eventBus.publish({\n      kind: \"status-update\",\n      taskId: context.taskId,\n      contextId: context.contextId,\n      status: { state: \"working\", timestamp: new Date().toISOString() },\n      final: false,\n    });\n\n    // ... do work ...\n\n    eventBus.publish({\n      kind: \"status-update\",\n      taskId: context.taskId,\n      contextId: context.contextId,\n      status: { state: \"completed\", timestamp: new Date().toISOString() },\n      final: true,\n    });\n    eventBus.finished();\n  }\n}\n\nconst executor = new MyExecutor();\n\n// Start consuming with async generator pattern\nconst messages = backend.workQueue.start();\n\nfor await (const taskMessage of messages) {\n  const { taskId, contextId, requestContext } = taskMessage;\n\n  // Create request context\n  const context = new RequestContext(\n    requestContext.userMessage,\n    taskId,\n    contextId,\n    requestContext.task,\n    requestContext.referenceTasks\n  );\n\n  // Create event bus for publishing task events\n  const eventBus = new WorkerEventBus(backend.amqpConnection, taskId, contextId);\n\n  // Execute task\n  await executor.execute(context, eventBus);\n}\n```\n\n### 3. Scale horizontally\n\nRun multiple workers to process tasks in parallel:\n\n```bash\n# Terminal 1: HTTP Server\nbun run server\n\n# Terminal 2-N: Workers (scale as needed)\nbun run worker\nbun run worker  # Add more workers for higher throughput\n```\n\n## Features\n\n- **Work Queue**: Distribute tasks across multiple worker processes\n- **Event Sourcing**: All task events stored in AMQP streams for replay\n- **In-Memory Projection**: Fast task lookups with automatic recovery from streams\n- **SSE Streaming**: Automatic streaming of task events back to clients\n- **Horizontal Scaling**: Add more workers to increase throughput\n- **Graceful Shutdown**: Clean consumer and connection handling\n- **Type-Safe**: Full TypeScript support with Zod validation\n\n## Configuration\n\n```typescript\ninterface AMQPAgentBackendConfig {\n  url: string;                    // AMQP broker URL\n  agentName: string;              // Agent identifier\n  streamRetention?: string;       // Event retention (default: \"7d\")\n  streamMaxBytes?: number;        // Max stream size (default: 1GB)\n  workQueueName?: string;         // Custom work queue name\n  exchangeName?: string;          // Custom exchange name\n  logger?: Logger;                // Custom logger\n  connection?: {\n    heartbeat?: number;           // Heartbeat interval in seconds\n    reconnectDelay?: number;      // Reconnection delay in ms\n    maxReconnectAttempts?: number;// Max reconnection attempts\n  };\n  publishing?: {\n    persistent?: boolean;         // Persistent messages (default: true)\n    confirmMode?: boolean;        // Publisher confirms (default: true)\n    messageTtl?: number;          // Message TTL in ms (0 = no expiration)\n  };\n}\n```\n\n## Examples\n\nSee complete working examples:\n\n- `src/examples/http-server.ts` - HTTP server with queuing\n- `src/examples/worker.ts` - Worker process\n\n```bash\n# Run the example\nbun run server  # Terminal 1\nbun run worker  # Terminal 2\n\n# Send a request\ncurl -X POST http://localhost:3000/ \\\n  -H \"Content-Type: application/json\" \\\n  -d '{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"messages/send\",\"params\":{\"message\":{\"kind\":\"message\",\"role\":\"user\",\"messageId\":\"1\",\"contextId\":\"ctx-1\",\"parts\":[{\"kind\":\"text\",\"text\":\"Hello\"}]}}}'\n```\n\n## Testing\n\n```bash\nbun run test             # Run all tests (unit + integration)\nbun run test:unit        # Unit tests only\nbun run test:integration # Integration tests only\nbun run test:watch       # Watch mode\nbun run test:coverage    # With coverage\n```\n\n## License\n\nMIT\n","readmeFilename":"README.md"}