{"_id":"@0x0642/c-indexer-consumer","_rev":"3-7bd2807ce8ac0a5eeb127603bf2378b9","name":"@0x0642/c-indexer-consumer","dist-tags":{"latest":"1.0.3"},"versions":{"1.0.0":{"name":"@0x0642/c-indexer-consumer","version":"1.0.0","keywords":["blockchain","ethereum","kafka","events","web3","indexer","etl","consumer","smart-contract"],"author":{"name":"0x0642"},"license":"MIT","_id":"@0x0642/c-indexer-consumer@1.0.0","maintainers":[{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"}],"dist":{"shasum":"691ba3742798e4603e75a5e3f152747a2be9fd38","tarball":"https://registry.npmjs.org/@0x0642/c-indexer-consumer/-/c-indexer-consumer-1.0.0.tgz","fileCount":8,"integrity":"sha512-Q83RtMlaRD+pelv4u7zpEgxQlqxQEzy3oDwdYm3VO0CjZ0eWYEVvi7ob1PQ9Nb20olNrC/rIXB8ZES65j4oIDg==","signatures":[{"sig":"MEYCIQDVu6LYaLMiB12K6iGkNf4bWeuYVgnfy7szmvbzLMcogwIhAP20A8bLUiABBJVOrqj3u1ymda0omrkEIHzGfUPv0jvC","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":19452},"main":"src/index.ts","type":"module","types":"src/index.ts","module":"src/index.ts","engines":{"node":">=18.0.0"},"exports":{".":{"types":"./src/index.ts","import":"./src/index.ts"}},"_npmUser":{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"},"_npmVersion":"9.6.4","description":"A TypeScript library for consuming blockchain events through Kafka","directories":{},"_nodeVersion":"20.1.0","dependencies":{"ethers":"^6.15.0","kafkajs":"^2.2.4"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"latest"},"peerDependencies":{"typescript":"^5"},"_npmOperationalInternal":{"tmp":"tmp/c-indexer-consumer_1.0.0_1760542688887_0.018681401451922097","host":"s3://npm-registry-packages-npm-production"}},"1.0.2":{"name":"@0x0642/c-indexer-consumer","version":"1.0.2","keywords":["blockchain","ethereum","kafka","events","web3","indexer","etl","consumer","smart-contract"],"author":{"name":"0x0642"},"license":"MIT","_id":"@0x0642/c-indexer-consumer@1.0.2","maintainers":[{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"}],"dist":{"shasum":"2e33c27f9303141e28903fc8e66a4cc89e54d257","tarball":"https://registry.npmjs.org/@0x0642/c-indexer-consumer/-/c-indexer-consumer-1.0.2.tgz","fileCount":9,"integrity":"sha512-+UWvkhXxIq7Inw6z1i0rdtzYfsgiWYXbjBcgbgxGzHG1rNyqAuO5Pmz/7BlvEbhXAedcXAdNBJGvEy9sc/h+PA==","signatures":[{"sig":"MEYCIQDATuTPSpv2K24Z0LU8LI4gMm7/dHlIXwvZJ4czX+03+wIhANLItWKtQNNDAOHxCrjHEQzGJNRPSoTLy7L/6ww9ClOZ","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":25295},"main":"src/index.ts","type":"module","types":"src/index.ts","module":"src/index.ts","engines":{"node":">=18.0.0"},"exports":{".":{"types":"./src/index.ts","import":"./src/index.ts"}},"_npmUser":{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"},"_npmVersion":"9.6.4","description":"A TypeScript library for consuming blockchain events through Kafka","directories":{},"_nodeVersion":"20.1.0","dependencies":{"ethers":"^6.15.0","kafkajs":"^2.2.4"},"_hasShrinkwrap":false,"devDependencies":{"@types/bun":"latest"},"peerDependencies":{"typescript":"^5"},"_npmOperationalInternal":{"tmp":"tmp/c-indexer-consumer_1.0.2_1760544150004_0.15638128730032452","host":"s3://npm-registry-packages-npm-production"}},"1.0.3":{"name":"@0x0642/c-indexer-consumer","version":"1.0.3","description":"A TypeScript library for consuming blockchain events through Kafka","main":"src/index.ts","module":"src/index.ts","types":"src/index.ts","type":"module","exports":{".":{"import":"./src/index.ts","types":"./src/index.ts"}},"keywords":["blockchain","ethereum","kafka","events","web3","indexer","etl","consumer","smart-contract"],"author":{"name":"0x0642"},"license":"MIT","devDependencies":{"@types/bun":"latest"},"peerDependencies":{"typescript":"^5"},"dependencies":{"ethers":"^6.15.0","kafkajs":"^2.2.4"},"engines":{"node":">=18.0.0"},"_id":"@0x0642/c-indexer-consumer@1.0.3","gitHead":"1383d767d0bc3b211138508d493f4c8208d0e8dc","_nodeVersion":"22.19.0","_npmVersion":"10.9.3","dist":{"integrity":"sha512-SP6eZY750PjfnwcSNBrn/flYblAPq34vB+WoBaX+Qb670t5Sf/Tt9LsalPCaV1/R6tdrUi2Pv/TpxWCyZ35S0w==","shasum":"1e3bf36e78910fbf8df240fcbe0e950e7690cbfc","tarball":"https://registry.npmjs.org/@0x0642/c-indexer-consumer/-/c-indexer-consumer-1.0.3.tgz","fileCount":9,"unpackedSize":25272,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEYCIQCTrJZjKuy2Ko5pXQzyFm4swCP/I3XEz2+QIU8k2aMQDwIhAMUWy9xsKb+CaetvCR32/nYBNJFI58MOdzIXANq0g2EV"}]},"_npmUser":{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"},"directories":{},"maintainers":[{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/c-indexer-consumer_1.0.3_1760586487912_0.37717061443151856"},"_hasShrinkwrap":false}},"time":{"created":"2025-10-15T15:38:08.775Z","modified":"2025-10-16T03:48:08.296Z","1.0.0":"2025-10-15T15:38:09.362Z","1.0.2":"2025-10-15T16:02:30.176Z","1.0.3":"2025-10-16T03:48:08.101Z"},"author":{"name":"0x0642"},"license":"MIT","keywords":["blockchain","ethereum","kafka","events","web3","indexer","etl","consumer","smart-contract"],"description":"A TypeScript library for consuming blockchain events through Kafka","maintainers":[{"name":"0x0642.xyz","email":"ductrungnguyen98@gmail.com"}],"readme":"# @0x0642/c-indexer-consumer\n\nA TypeScript library for consuming blockchain events through Kafka, designed to work with blockchain ETL services. This package simplifies the process of registering blockchain event listeners and consuming them via Kafka.\n\n## Features\n\n- 🔗 **Easy Event Registration**: Register blockchain contract events with minimal configuration\n- 📨 **Kafka Integration**: Built-in Kafka consumer management with kafkajs\n- 🔧 **Type-Safe**: Full TypeScript support with comprehensive type definitions\n- 🎯 **Event Parsing**: Automatic event signature parsing using ethers.js v6\n- 🔄 **Transform Scripts**: Support for custom data transformation scripts\n- ⚡ **Efficient**: Optimized for handling high-throughput blockchain events\n\n## Installation\n\n```bash\nnpm install @0x0642/c-indexer-consumer\n```\n\nor with bun:\n\n```bash\nbun add @0x0642/c-indexer-consumer\n```\n\n## Prerequisites\n\n- A blockchain ETL service endpoint\n- Kafka broker(s) configured and running\n- Contract ABI for the events you want to monitor\n\n## Usage\n\n### Basic Example\n\n```typescript\nimport { BlockWatcherRegister, withTransformScript } from \"@0x0642/c-indexer-consumer\";\nimport { Kafka } from \"kafkajs\";\n\n// Initialize Kafka\nconst kafka = new Kafka({\n  clientId: \"my-app\",\n  brokers: [\"localhost:9092\"],\n});\n\n// Create a BlockWatcherRegister instance\nconst watcher = new BlockWatcherRegister({\n  etlBaseUrl: \"https://your-etl-service.com\",\n  defaultGroupId: \"my-consumer-group\", // optional, default: \"blockwatcher-group\"\n});\n\n// Define your contract ABI\nconst contractABI = `[\n  {\n    \"anonymous\": false,\n    \"inputs\": [\n      {\"indexed\": true, \"name\": \"from\", \"type\": \"address\"},\n      {\"indexed\": true, \"name\": \"to\", \"type\": \"address\"},\n      {\"indexed\": false, \"name\": \"value\", \"type\": \"uint256\"}\n    ],\n    \"name\": \"Transfer\",\n    \"type\": \"event\"\n  }\n]`;\n\n// Register an event handler\nawait watcher.register(\n  \"0x1234567890123456789012345678901234567890\", // contract address\n  contractABI,\n  \"Transfer\", // event name\n  1, // chain ID (1 = Ethereum mainnet)\n  async (payload) => {\n    // payload is already parsed - contains the event data\n    console.log(\"Transfer event:\", payload);\n  },\n  \"my-group-id\" // optional groupId, uses defaultGroupId if not provided\n);\n\n// Start consuming Kafka messages\nawait watcher.startKafka(kafka, {\n  groupId: \"my-consumer-group\", // optional, overrides defaultGroupId\n  fromBeginning: true, // optional, default: true\n});\n```\n\n### Advanced Usage with Custom Kafka Topic\n\n```typescript\nimport { \n  BlockWatcherRegister, \n  withKafkaTopic \n} from \"@0x0642/c-indexer-consumer\";\n\nconst watcher = new BlockWatcherRegister({\n  etlBaseUrl: \"https://your-etl-service.com\",\n});\n\n// Register with custom Kafka topic\nawait watcher.register(\n  \"0x1234567890123456789012345678901234567890\",\n  contractABI,\n  \"Transfer\",\n  1,\n  async (payload, meta) => {\n    console.log(\"Transfer event:\", payload);\n    console.log(\"Block number:\", meta.block_number);\n  },\n  undefined, // groupId (optional)\n  withKafkaTopic(\"my-custom-topic-name\") // custom Kafka topic\n);\n\nawait watcher.startKafka(kafka);\n```\n\n### Advanced Usage with Transform Scripts\n\n```typescript\nimport { \n  BlockWatcherRegister, \n  withTransformScript, \n  withTestData,\n  withKafkaTopic\n} from \"@0x0642/c-indexer-consumer\";\n\nconst watcher = new BlockWatcherRegister({\n  etlBaseUrl: \"https://your-etl-service.com\",\n});\n\n// Custom transformation script\nconst transformScript = `\nexport function transform(data, meta) {\n  return {\n    ...data,\n    timestamp: Date.now(),\n    chainId: meta.chain_id,\n    blockNumber: meta.block_number\n  };\n}\n`;\n\nawait watcher.register(\n  \"0x1234567890123456789012345678901234567890\",\n  contractABI,\n  \"Transfer\",\n  1,\n  async (payload, meta) => {\n    // Payload is already transformed by the script\n    console.log(\"Transformed event:\", payload);\n    console.log(\"With metadata:\", meta);\n  },\n  undefined, // groupId (optional)\n  withTransformScript(transformScript),\n  withTestData({ mockField: \"test\" }),\n  withKafkaTopic(\"custom-transfer-topic\") // custom Kafka topic (optional)\n);\n\nawait watcher.startKafka(kafka, {\n  groupId: \"my-consumer-group\",\n  fromBeginning: true,\n});\n```\n\n### Multiple Event Registration\n\n```typescript\nconst watcher = new BlockWatcherRegister({\n  etlBaseUrl: \"https://your-etl-service.com\",\n});\n\n// Handler functions\nconst handleTransfer = async (payload: any, meta: EtlMetaData) => {\n  console.log(`Transfer on chain ${meta.chain_id}:`, payload);\n};\n\nconst handleApproval = async (payload: any) => {\n  console.log(\"Approval event:\", payload);\n};\n\nconst handleSwap = async (payload: any, meta: EtlMetaData) => {\n  console.log(`Swap at block ${meta.block_number}:`, payload);\n};\n\n// Register multiple events\nawait watcher.register(\n  \"0xContractA...\",\n  contractABI_A,\n  \"Transfer\",\n  1,\n  handleTransfer\n);\n\nawait watcher.register(\n  \"0xContractB...\",\n  contractABI_B,\n  \"Approval\",\n  1,\n  handleApproval\n);\n\nawait watcher.register(\n  \"0xContractC...\",\n  contractABI_C,\n  \"Swap\",\n  1,\n  handleSwap\n);\n\n// Start consuming all registered events\nawait watcher.startKafka(kafka);\n```\n\n## API Reference\n\n### `BlockWatcherRegister`\n\nMain class for managing blockchain event registration and consumption.\n\n#### Constructor\n\n```typescript\nnew BlockWatcherRegister(args: BlockWatcherConfig)\n```\n\n**BlockWatcherConfig:**\n```typescript\ntype BlockWatcherConfig = {\n  etlBaseUrl: string;\n  defaultGroupId?: string; // optional, default: \"blockwatcher-group\"\n}\n```\n\n- `etlBaseUrl`: The base URL of your blockchain ETL service\n- `defaultGroupId`: Default Kafka consumer group ID for all registered events\n\n#### Methods\n\n##### `register()`\n\nRegister a blockchain event to monitor.\n\n```typescript\nasync register(\n  contractAddress: string,\n  contractABI: string,\n  eventName: string,\n  chainId: number,\n  handler: KafkaConsumerMessageHandlerFunc,\n  groupId?: string,\n  ...options: RegisterOption[]\n): Promise<void>\n```\n\n**Parameters:**\n- `contractAddress`: The smart contract address\n- `contractABI`: Contract ABI JSON string\n- `eventName`: Name of the event to monitor (e.g., \"Transfer\", \"Approval\")\n- `chainId`: Blockchain chain ID (1 for Ethereum mainnet, 137 for Polygon, etc.)\n- `handler`: Callback function to handle incoming events\n- `groupId`: Optional Kafka consumer group ID for this specific event (overrides defaultGroupId)\n- `options`: Optional configuration (transform scripts, test data)\n\n##### `startKafka()`\n\nInitialize Kafka consumer and start processing events.\n\n```typescript\nasync startKafka(\n  kafka: Kafka,\n  options?: {\n    groupId?: string;\n    fromBeginning?: boolean;\n  }\n): Promise<void>\n```\n\n**Parameters:**\n- `kafka`: KafkaJS Kafka instance\n- `options`: Optional configuration object\n  - `groupId`: Consumer group ID (overrides defaultGroupId)\n  - `fromBeginning`: Whether to read from the beginning of the topic (default: true)\n\n### Types\n\n#### `KafkaConsumerMessageHandlerFunc`\n\nThe handler function can have multiple signatures:\n\n```typescript\n// Option 1: Just payload\ntype Handler = (payload: any) => Promise<void> | void;\n\n// Option 2: Payload + metadata\ntype Handler = (payload: any, meta: EtlMetaData) => Promise<void> | void;\n\n// Option 3: Payload + raw Kafka message\ntype Handler = (payload: any, raw: KafkaMessage) => Promise<void> | void;\n```\n\n**Examples:**\n\n```typescript\n// Simple handler - just the payload\nasync (payload) => {\n  console.log(\"Event data:\", payload);\n}\n\n// Handler with metadata\nasync (payload, meta) => {\n  console.log(\"Event data:\", payload);\n  console.log(\"Chain ID:\", meta.chain_id);\n  console.log(\"Block number:\", meta.block_number);\n  console.log(\"Transaction hash:\", meta.tx_hash);\n}\n\n// Handler with raw message (for advanced use cases)\nasync (payload, rawMessage) => {\n  console.log(\"Event data:\", payload);\n  console.log(\"Kafka headers:\", rawMessage.headers);\n  console.log(\"Kafka offset:\", rawMessage.offset);\n}\n```\n\n#### `EtlMetaData`\n\nMetadata extracted from blockchain events:\n\n```typescript\ntype EtlMetaData = {\n  chain_id?: number;\n  _etl_event_log_entity_id?: number;\n  address?: string;\n  block_number?: number;\n  block_hash?: string;\n  tx_hash?: string;\n  tx_index?: number;\n  index?: number;\n  topics?: string[];\n  [k: string]: any;\n};\n```\n\n#### `RegisterOption`\n\n```typescript\ntype RegisterOption = {\n  transformScript?: string;\n  testData?: Record<string, any>;\n  kafkaTopic?: string;\n};\n```\n\n### Helper Functions\n\n#### `withTransformScript(script: string)`\n\nCreate a RegisterOption with a custom transform script.\n\n```typescript\nconst option = withTransformScript(`\n  export function transform(data, meta) {\n    return { ...data, processed: true };\n  }\n`);\n```\n\n#### `withTestData(testData: Record<string, any>)`\n\nCreate a RegisterOption with test data.\n\n```typescript\nconst option = withTestData({ mockField: \"value\" });\n```\n\n#### `withKafkaTopic(kafkaTopic: string)`\n\nCreate a RegisterOption with a custom Kafka topic name.\n\n```typescript\nconst option = withKafkaTopic(\"my-custom-topic\");\n```\n\n**Note:** If not specified, the default topic naming convention is:\n```\nblockchain.{chainId}.{eventName}.{contractAddress}.{topicHash4Bytes}\n```\n\n## How It Works\n\n1. **Event Registration**: When you register an event, the library:\n   - Parses the contract ABI using ethers.js\n   - Extracts the event signature and generates the topic hash\n   - Creates a unique Kafka topic name\n   - Sends a registration request to the ETL service\n\n2. **Kafka Topic Naming**: By default, topics are automatically generated in the format:\n   ```\n   blockchain.{chainId}.{eventName}.{contractAddress}.{topicHash4Bytes}\n   ```\n   You can override this by using `withKafkaTopic()` option.\n\n3. **Message Processing Flow**:\n   ```\n   Kafka Message\n       ↓\n   KafkaConsumerRegister (receives raw message)\n       ↓\n   KafkaConsumerHandlerSingleFunc (wrapper layer)\n       ↓\n   - Parses message.value JSON\n   - Extracts payload\n   - Extracts metadata (if needed)\n   - Converts string numbers to actual numbers\n       ↓\n   Your Handler Function (receives parsed data)\n   ```\n\n4. **Event Consumption**: The Kafka consumer subscribes to all registered topics and routes messages to the appropriate handlers.\n\n5. **Automatic Parsing**: The library automatically:\n   - Parses the Kafka message JSON\n   - Extracts the `payload` field\n   - Extracts and converts `meta` data (chain_id, block_number, etc.)\n   - Calls your handler with the appropriate parameters based on function signature\n\n## Configuration\n\n### Environment Variables\n\nYou can configure Kafka brokers and other settings through environment variables or directly in code:\n\n```typescript\nconst kafka = new Kafka({\n  clientId: process.env.KAFKA_CLIENT_ID || \"blockchain-consumer\",\n  brokers: process.env.KAFKA_BROKERS?.split(\",\") || [\"localhost:9092\"],\n});\n```\n\n## Error Handling\n\nThe library throws errors for:\n- Invalid contract ABIs\n- Non-existent event names\n- ETL service connection failures\n- Kafka connection issues\n\nAlways wrap registration and startup calls in try-catch blocks:\n\n```typescript\ntry {\n  await watcher.register(...);\n  await watcher.startKafka(kafka);\n} catch (error) {\n  console.error(\"Error:\", error.message);\n}\n```\n\n## Dependencies\n\n- **ethers** (^6.15.0): For ABI parsing and event signature generation\n- **kafkajs** (^2.2.4): For Kafka consumer functionality\n\n## License\n\nMIT\n\n## Contributing\n\nContributions are welcome! Please feel free to submit a Pull Request.\n\n## Support\n\nFor issues and questions, please open an issue on the GitHub repository.\n","readmeFilename":"README.md"}