{"_id":"@abernatskiy/hybrid-pipes-streams","_rev":"2-ad7f5f61c0f292d24df86048f1eaaaac","name":"@abernatskiy/hybrid-pipes-streams","dist-tags":{"latest":"0.0.1"},"versions":{"0.0.0":{"name":"@abernatskiy/hybrid-pipes-streams","version":"0.0.0","license":"MIT","_id":"@abernatskiy/hybrid-pipes-streams@0.0.0","maintainers":[{"name":"abernatskiy","email":"abernatskiy@toha.li"}],"dist":{"shasum":"7621cea9ba17d63dac8453c5b633ae6dc67cdc25","tarball":"https://registry.npmjs.org/@abernatskiy/hybrid-pipes-streams/-/hybrid-pipes-streams-0.0.0.tgz","fileCount":269,"integrity":"sha512-iV2yOUiVtW/zCrpGtXu4cO7IU20rzoZRQpPL66JV3kTjA4ubVRJsorFcOlVYwdnGkXEb5nJn9J4jCGo4GgZHRA==","signatures":[{"sig":"MEUCIQCWQEIL26RSJ3gr2uMFg+muNM3OcA4Qm6IHFBX528zK2QIgP4tXS7LlAdIqbztdNPtxZjnSM3ulBpcTTokXb0SuplQ=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":1964526},"main":"dist/index.js","types":"./dist/index.d.ts","engines":{"node":"<=23.9.0"},"gitHead":"ce5a30f388f75a6f3c3c3ff8a1571288ec6f9c33","scripts":{"dev":"ts-node src/index.ts","lint":"biome check .","build":"rm -rf dist/ && tsc","start":"node dist/index.js","watch":"tsc -w","lint:fix":"biome check . --write"},"_npmUser":{"name":"abernatskiy","email":"abernatskiy@toha.li"},"_npmVersion":"10.9.0","description":"Core package of the SQD Pipes ecosystem that provides specialized streams for blockchain data consumption.","directories":{},"_nodeVersion":"22.12.0","dependencies":{"pino":"^9.6.0","lodash":"^4.17.21","@solana/web3.js":"^1.98.4","@subsquid/borsh":"^0.3.0","@subsquid/evm-abi":"^0.3.1","@subsquid/evm-codec":"^0.3.0","@subsquid/http-client":"^1.6.0","@subsquid/portal-client":"0.2.0","@subsquid/solana-stream":"1.0.0-portal-api.a3f844","@subsquid/util-internal":"^3.2.0","@subsquid/solana-objects":"0.0.5-portal-api.a3f844","@abernatskiy/hybrid-pipes-core":"^0.0.1"},"_hasShrinkwrap":false,"devDependencies":{"ts-node":"^10.9.2","@types/bun":"^1.2.22","typescript":"^5.9.2","@types/node":"^22.13.10","@types/lodash":"^4.17.16"},"peerDependencies":{"@clickhouse/client":"^1.12.1"},"peerDependenciesMeta":{"@clickhouse/client":{"optional":true}},"_npmOperationalInternal":{"tmp":"tmp/hybrid-pipes-streams_0.0.0_1758580316777_0.10727039753113798","host":"s3://npm-registry-packages-npm-production"}},"0.0.1":{"name":"@abernatskiy/hybrid-pipes-streams","version":"0.0.1","main":"dist/index.js","license":"MIT","scripts":{"build":"rm -rf dist/ && tsc","start":"node dist/index.js","dev":"ts-node src/index.ts","watch":"tsc -w","lint":"biome check .","lint:fix":"biome check . --write"},"devDependencies":{"@types/bun":"^1.2.22","@types/lodash":"^4.17.16","@types/node":"^22.13.10","ts-node":"^10.9.2","typescript":"^5.9.2"},"dependencies":{"@abernatskiy/hybrid-pipes-core":"^0.0.2","@solana/web3.js":"^1.98.4","@subsquid/borsh":"^0.3.0","@subsquid/evm-abi":"^0.3.1","@subsquid/evm-codec":"^0.3.0","@subsquid/http-client":"^1.6.0","@subsquid/portal-client":"0.2.0","@subsquid/solana-objects":"0.0.5-portal-api.a3f844","@subsquid/solana-stream":"1.0.0-portal-api.a3f844","@subsquid/util-internal":"^3.2.0","lodash":"^4.17.21","pino":"^9.6.0"},"peerDependencies":{"@clickhouse/client":"^1.12.1"},"peerDependenciesMeta":{"@clickhouse/client":{"optional":true}},"engines":{"node":"<=23.9.0"},"_id":"@abernatskiy/hybrid-pipes-streams@0.0.1","gitHead":"00ac9c47c30b10b3ddfcab8fb512c2104f12b019","types":"./dist/index.d.ts","description":"Core package of the SQD Pipes ecosystem that provides specialized streams for blockchain data consumption.","_nodeVersion":"22.12.0","_npmVersion":"10.9.0","dist":{"integrity":"sha512-ywfSLelrDKaF6Lre/3QD3wZoz6lh1nam5iKlAYddorR72xBRWp4dkxyac0tlrp4THo1DgxHH5c/7yEy4kbWYOg==","shasum":"c7844e6522483a5f00a1bf6ffa98cf7871b67433","tarball":"https://registry.npmjs.org/@abernatskiy/hybrid-pipes-streams/-/hybrid-pipes-streams-0.0.1.tgz","fileCount":270,"unpackedSize":1880203,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEYCIQC5naJFjKo9wbwaSdJLVYLHieLHWR/Ak5FhEikWojeuQAIhAO50N5HvkNOkfrEJLb1FpjNGwiKA5tJX+pr3CiLXnutf"}]},"_npmUser":{"name":"abernatskiy","email":"abernatskiy@toha.li"},"directories":{},"maintainers":[{"name":"abernatskiy","email":"abernatskiy@toha.li"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/hybrid-pipes-streams_0.0.1_1758661135393_0.18187894389339987"},"_hasShrinkwrap":false}},"time":{"created":"2025-09-22T22:31:56.674Z","modified":"2025-09-23T20:58:55.782Z","0.0.0":"2025-09-22T22:31:56.998Z","0.0.1":"2025-09-23T20:58:55.595Z"},"license":"MIT","description":"Core package of the SQD Pipes ecosystem that provides specialized streams for blockchain data consumption.","maintainers":[{"name":"abernatskiy","email":"abernatskiy@toha.li"}],"readme":"# @sqd-pipes/streams\n\nCore package of the SQD Pipes ecosystem that provides specialized streams for blockchain data consumption.\n\n## Overview\n\nThis package offers a collection of stream implementations for different blockchain platforms:\n\n- **Stream Abstractions**: Provides a common interface for consuming blockchain data via the Subsquid Portal\n- **Chain-Specific Implementations**: Includes specialized streams for Solana and EVM blockchains\n- **State Management**: Utilities for tracking stream progress and handling checkpoints\n\nUnlike a general-purpose stream processing library, this package is specifically designed for consuming and processing blockchain-specific data types like swaps, liquidity events, and token metadata.\n\n## Key Components\n\n### Core Functionality\n- `PortalAbstractStream`: Base abstract class for all blockchain data streams\n- State management interfaces for tracking progress\n- Error handling and retry mechanisms\n\n### Solana Streams\n- `SolanaSwapsStream`: For collecting DEX swap events (Orca, Raydium, Meteora)\n- `SolanaLiquidityStream`: For tracking liquidity events\n- `MetaplexTokenStream` and `PumpfunTokenStream`: For token metadata\n\n### EVM Streams\n- `EVMSwapsStream`: For collecting Ethereum and EVM DEX swap events\n\n## Installation\n\n```bash\nyarn add @sqd-pipes/streams\n```\n\n## Usage\n\n### Solana Swaps Example\n\n```typescript\nimport { SolanaSwapsStream, ClickhouseState } from \"@sqd-pipes/streams\";\n\n// Create a state manager for checkpointing\nconst clickhouseState = new ClickhouseState(client, {\n  table: \"sync_status\",\n  id: \"my_indexer\",\n});\n\n// Initialize the swaps stream for specific DEXes\nconst stream = new SolanaSwapsStream({\n  portal: \"https://v2.archive.subsquid.io/datasets/solana-mainnet\",\n  blockRange: {\n    from: 12345678,\n  },\n  args: {\n    type: [\"orca_whirlpool\", \"raydium_amm\"],\n  },\n  state: clickhouseState,\n  logger,\n});\n\n// Process the stream data\nfor await (const swaps of await stream.stream()) {\n  // Each swap contains details like tokens, amounts, block info\n  console.log(`Processing ${swaps.length} swaps`);\n  \n  // After processing, update the checkpoint\n  await stream.ack();\n}\n```\n\n## State Management\n\nThe package provides state management interfaces:\n\n- `ClickhouseState`: Stores checkpoint data in ClickHouse\n- In-memory state tracking for development and testing\n\nState management handles critical functionality like:\n- Resuming from the last processed block\n- Handling blockchain reorganizations (forks)\n- Tracking indexer progress\n\n## Error Handling\n\nThe streams include built-in error handling for common blockchain data issues:\n- Network connectivity problems\n- Chain reorganizations\n- Rate limiting\n\n## License\n\nMIT\n","readmeFilename":"README.md"}