High-throughput, low-latency block ingestion loop for Avalanche C-Chain. It combines:
- sliding-window backfill over an adjustable range [LUB, LIB]
- real-time new-head subscription
- unified concurrency control with reserved capacity for real-time
- idempotent marking and sliding of watermarks
This repo is a minimal, production-oriented skeleton you can adapt by plugging in your own persistence and processing logic.
Requirements:
- Go 1.24+
- A WebSocket endpoint to C-Chain (e.g.
wss://api.avax-test.network/ext/bc/C/wsor your own AvalancheGo node)
Run:
go run ./...Adjust parameters in main.go as needed:
wsEndpoint: C-Chain WebSocket endpointlub: initial Lowest Unprocessed BlockmaxConcurrency: total worker capacitybackfillPriority: number of threads reserved for backfill (the rest is implicitly reserved for real-time)newBlocksCapacity: buffer size for the real-time headers channel
Components:
- Manager (
internal/manager) orchestrates the pipeline: scans the [LUB, LIB] window, starts workers, reserves capacity, and handles real-time events. - Processor (
internal/worker) fetches a block by height and executes the processing logic (currently: fetch + log summary). - Subscriber (
internal/subscriber) subscribes to new heads via WebSocket and pushes headers to the Manager. - PipelineStateRepository (
internal/domain+internal/repository) stores/wraps the sliding window and processed state; an in-memory implementation is provided. - RPC client (
pkg/rpc) wrapscustomethclientreusing the Coreth types.
Data flow:
[Subscriber] --headers--> [Manager] --dispatch--> [Processor]
^ |
| mark + slide |
+-----[PipelineState]---+
-
Sliding window with watermarks:
- LUB (Lowest Unprocessed Block): the left edge of the window, advanced only when all contiguous heights from LUB are processed.
- LIB (Largest Ingested Block): the right edge (best known tip). Realtime updates can move LIB forward.
-
Unified dispatcher with backfill vs. real-time capacity:
maxConcurrencybounds total concurrent workers.backfillPrioritybounds the number of concurrent backfill workers.- This implicitly reserves
maxConcurrency - backfillPrioritycapacity for real-time events.
-
Drop-on-overload for realtime headers with eventual processing:
- If all worker slots are occupied when a real-time header arrives, the event is dropped (by design). Backfill will ingest it later because LIB has moved forward.
-
Idempotent processing and progress:
- Blocks are marked processed;
AdvanceLUBslides LUB as far as contiguous processed blocks allow. - Transient worker errors leave the height unprocessed and it will be retried by backfill.
- Blocks are marked processed;
-
Aggressive backfill
- Continuously scans [LUB, LIB] to find the next unprocessed, non-inflight height.
- Dispatches up to
backfillPriorityworkers, subject to the globalmaxConcurrency.
-
Realtime handling
- On new header
h, best-effort bumpLIB := max(LIB, h)so backfill can pick it up even if the immediate realtime attempt is skipped. - If a worker slot is available, process
himmediately for lowest latency. - If no slot is available, the header is dropped; the block will be picked up by backfill.
- On new header
-
Completion and signals
- On successful processing:
MarkProcessed(h), then tryAdvanceLUB(). - Notifies backfill loop (to refill capacity or re-scan after window movement).
- On successful processing:
Concurrency controls:
workerSemcaps total concurrent workers (both backfill and realtime).backfillSemcaps concurrent backfill workers.inflighttracks dispatched-but-not-finished heights to avoid duplicates.
internal/domain/pipeline_state.go defines the PipelineStateRepository interface. The in-memory implementation is in internal/repository/in_memory_pipeline_state.go.
Key operations and guarantees:
Window()returns [LUB, LIB].SetLIB(newLIB)must not set LIB below LUB.ResetLUB(newLUB)can move LUB backward/forward (forward must not exceed LIB). When LUB moves forward, processed marks below LUB are dropped.MarkProcessed(h)records a block as processed (heights below current LUB are implicitly processed).AdvanceLUB()slides LUB forward while contiguous heights are processed.
internal/worker/worker.go shows a placeholder processor:
- fetches block by number via
customethclient - waits 300ms (simulates work), then logs
height,hash, andtx count
Replace this with your real logic (e.g., mapping to internal/domain types and persisting to storage). Helpers already exist for mapping primary EVM and atomic cross-chain payloads in internal/domain/mappers.go and internal/domain/atomic_transaction.go.
internal/subscriber/subscriber.go subscribes to new heads and writes headers to the Manager’s sink channel. The Manager owns the channel and controls its buffer via newBlocksCapacity.
wsEndpoint: WebSocket endpoint for C-Chain.lub: initial LUB (starting height to backfill from).maxConcurrency: global maximum number of concurrent workers.backfillPriority: maximum number of concurrent backfill workers (realtime is the remainder).newBlocksCapacity: buffer for realtime headers; larger buffers can absorb short spikes but do not replace the drop-on-overload policy.
Example bootstrap (excerpt):
wsEndpoint := "wss://api.avax-test.network/ext/bc/C/ws"
lub := uint64(48662004)
maxConcurrency := 10
backfillPriority := 8 // reserve 8 for backfill, 2 for realtime
newBlocksCapacity := 64 // subscriber header buffer- Worker fetch failures: logged and left unprocessed; backfill retries later.
- Realtime overload: events may be dropped, but backfill still ingests because LIB has moved forward.
- Context cancellation / shutdown: the loop stops on
ctx.Done(); in-flight workers return and release capacity.
- Persistence: implement
PipelineStateRepositorybacked by your datastore (e.g., DynamoDB, PostgreSQL, Badger). Swap it in for the in-memory version. - Processing: replace the placeholder
Processor.Processwith your business logic, using the mapping helpers ininternal/domainas needed. - Observability: add structured logging/metrics around dispatch, success/failure counts, latency, and window sizes.
- Resilience: add backoff policies, classify retryable/non-retryable errors, and improve shutdown behavior with context awareness.
From the scratch notes in this repo:
- Send full blocks through the realtime channel (not just headers); adjust
Processaccordingly. - Add richer error handling for workers.
- Improve context awareness across components.
- LUB: Lowest Unprocessed Block — the earliest height not yet fully processed.
- LIB: Largest Ingested Block — the furthest known height to ingest up to.
main.go: wiring and application bootstrapinternal/manager/manager.go: core orchestration, concurrency controls, sliding windowinternal/worker/worker.go: placeholder processor (fetch + log)internal/subscriber/subscriber.go: new-head subscription (realtime)internal/domain/*.go: block/tx types, mappers, pipeline state interfaceinternal/repository/in_memory_pipeline_state.go: in-memory pipeline statepkg/rpc/rpc.go: light wrapper around Corethcustomethclient
- Backfill-first + reserved capacity keeps pipelines caught up while still responding quickly to new blocks.
- Drop-on-overload for realtime events avoids unbounded queues and simplifies backpressure. Eventual ingestion is guaranteed by backfill due to the watermark advance.
- Sliding window with idempotent marking eliminates fragile “one-shot” flows and supports restarts/resumes. Improve:
- Send blocks through the real time channel, not just height (adjust Process func accordingly)
- Add error handling for workers
- Add context awareness