Skip to content

fix(indexer): share one per-batch pipeline between the single and parallel ingest paths - #219

Merged
Miracle656 merged 1 commit into
Miracle656:mainfrom
DevTobis:fix/parallel-ingest-parity
Sep 30, 2026
Merged

Miracle656 merged 1 commit into
Miracle656:mainfrom
DevTobis:fix/parallel-ingest-parity

Conversation

@DevTobis

Copy link
Copy Markdown
Contributor

closes #203

Summary

INGEST_WORKERS > 1 switched to pollParallel, which only fetched SAC contracts and only stored fungible transfers. It silently skipped NFT transfers and metadata, LP shares, host-fn logs, SAC tagging, token metadata and account summaries. Raising the worker count changed what was indexed, not just how fast.

Changes

  • New src/indexer/batch.ts: processEventBatch(events, { network, knownLpPools }), the body of pollOnce moved verbatim (fungible + SAC tagging + token metadata + account summaries + WS emit, host-fn logs, LP shares, NFT + metadata). It returns per-type insert counts; batchTotal sums them.
  • pollOnce now fetches and calls processEventBatch. pollParallel workers do the same: its last parameter is a required { fetchEvents, processBatch }, so a sharded run cannot exist without a batch handler.
  • New exported ingestWindow(loop, from, to, workers = INGEST_WORKERS) in src/indexer.ts picks the path; the poll loop calls it. The sharded branch now shards allContractIds (SAC + NFT, previously SAC only), and fetches through the loop's source switcher (previously raw fetchEventsSafe, bypassing Horizon failover). createLoopState / LoopState are exported for the test.
  • .env.example documents INGEST_WORKERS.

Parallel-unsafe work

None found. Workers process disjoint contracts, and every table the pipeline writes is keyed by contract (eventId unique; AccountSummary, NftMetadata, token metadata include contractId in their key). The only shared mutable state is the knownLpPools set, and a contract only ever adds itself to it. So nothing needed to be excluded from parallel mode.

Tests

src/__tests__/ingestParity.test.ts (in-memory db, fake source that honours the contract filter):

  • the same events through ingestWindow(..., 1) and ingestWindow(..., 4) leave identical TokenTransfer, NftTransfer, NftMetadata, account summary, host-fn, LP-share, token-metadata and cursor state, and the same totalIndexed;
  • both paths route every event through processEventBatch;
  • a sensitivity test wires a sharded run whose handler skips steps (the old bug) and asserts the snapshots differ, i.e. a step added to one path only is caught.

Fail-before: with the parallel branch temporarily restored to the old behaviour (SAC ids only, fungible-only handler), 2 of the 5 tests fail (parity and same-pipeline). Restored, all 5 pass. All fixtures are built with StrKey so they are valid by construction.

Verification

  • npx jest src/__tests__/ingestParity.test.ts src/__tests__/multiNetworkIndexer.test.ts src/__tests__/indexerSources.test.ts src/__tests__/lpShares.test.ts src/__tests__/tombstoneWiring.test.ts src/__tests__/nft.test.ts: 6 suites, 78 passed.
  • npm run typecheck: clean.
  • Full suite and vitest integration tests not run (no Docker; shared, heavily loaded machine).

Caveats

@drips-wave

drips-wave Bot commented Sep 30, 2026

Copy link
Copy Markdown

@DevTobis Great news! 🎉 Based on an automated assessment of this PR, the linked Wave issue(s) no longer count against your application limits.

You can now already apply to more issues while waiting for a review of this PR. Keep up the great work! 🚀

Learn more about application limits

@Miracle656 Miracle656 left a comment

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Verified locally on a worktree merged with current main.

  • npx tsc --noEmit -p tsconfig.test.json: clean
  • Full jest suite: 503 passed / 42 suites, 0 failed

This is the right shape for #203. Moving the pipeline body into src/indexer/batch.ts verbatim and making ParallelIo.processBatch a required parameter means the old failure mode is now unrepresentable rather than merely tested against — a sharded run cannot exist without a batch handler. Two other real bugs got fixed along the way that the issue did not name: the sharded branch was sharding sacContractIds rather than allContractIds (so NFT contracts were never even requested), and it called fetchEventsSafe directly, bypassing the source switcher's Horizon failover.

ingestParity.test.ts does what was asked: the same fixture through ingestWindow(..., 1) and ingestWindow(..., 4) produces identical TokenTransfer, NftTransfer, AccountSummary, NftMetadata, host-fn, LP-share, token-metadata and cursor state, and the sensitivity case proves the snapshot comparison actually catches a path that skips steps. Building every fixture with StrKey.encodeContract / encodeEd25519PublicKey from raw bytes, plus the explicit isValidContract assertion, is the right habit — hand-typed strkeys in this repo have been wrong before.

Prometheus labels stay {network, type}, so nothing unbounded moved into batch.ts.

One small thing for a follow-up, not worth holding this: newLoop() sets SAC_CONTRACT_IDS_TESTNET / NFT_CONTRACT_IDS_TESTNET on process.env and never restores them. Jest reuses a worker process across test files, so that leaks into whatever runs next in the same worker. An afterAll restoring the previous values would close it.

https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

@Miracle656
Miracle656 merged commit d04392d into Miracle656:main Sep 30, 2026
1 check passed
Miracle656 added a commit that referenced this pull request Sep 30, 2026
* Use token decimals for display amounts

* docs(wave): W094-W098 batch draft (published as #203-207)

Five issues, 800 points. The headline is W094: `src/indexer.ts:491` switches
between pollOnce and pollParallel on INGEST_WORKERS, and `indexer/parallel.ts`
imports only fetchEventsSafe, parseEvents, upsertTransfers, setLastIndexedLedger
and emitTransfer — no NFT parsing, no metadata, no account summaries. Raising
the worker count for throughput silently stops indexing whole categories, with
no error to notice.

Also: /offramp/orders/:orderId serves an order with no authorization, NFT
metadata is fetched one serial round trip at a time, and every paged query pays
for a full COUNT.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

* Docs: Mainnet deployment guide (#166) (#211)

* Docs: Mainnet deployment guide (#166)

* docs(mainnet): correct the native SAC ids, fix the backup link, add DIRECT_DATABASE_URL

Review fixes applied on top of #211:

- The mainnet block used CDLZFC3SY…, which is the *testnet* native XLM SAC
  (Asset.native().contractId(Networks.TESTNET)). Mainnet is CAS3J7GY…
  (Networks.PUBLIC). The testnet block used CDMLFMKMM…, which is not the
  native SAC on either network. Both corrected, with the derivation inlined
  so the next reader can check rather than trust.
- Added a warning that SAC_CONTRACT_IDS must be set explicitly on mainnet,
  because the built-in fallback in src/indexer.ts is the wrong address.
- ../W076 did not resolve to anything; pointed at ./backup-restore.md.
- Added DIRECT_DATABASE_URL to both env blocks — prisma/schema.prisma
  declares directUrl, and boot-time schema sync fails without it.
- Linked ./DUAL_NETWORK.md for the NETWORKS=testnet,mainnet option.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

---------

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

* fix: narrow XDR error classification (#209)

* fix(offramp): require a bearer token to read an order (#216)

* fix(offramp): require a bearer token to read an order

GET /offramp/orders/:orderId served any order to whoever held its id: the
bank payout amount, the deposit address and the rate.

Orders now get a random public id (ofr_ + 128 bits) in place of the
provider's, and the creator is handed a bearer token (oft_ + 256 bits) once,
in the create response. Only its SHA-256 is stored. Lookups need
Authorization: Bearer; a missing header is 401, and an unknown id, a
malformed id and a wrong token are all the same 404.

Failed lookups are rate-limited separately from the app-wide limiter (10 per
15 minutes per IP by default); successful polling does not count.

Replaying a creation with the same idempotencyKey now also needs the same
walletAddress (409 otherwise) and re-issues the token, since only its hash is
kept. The lookup response no longer names the provider (source: live) or
returns its id, and a database failure returns a generic 500 instead of
reaching the global handler.

* test(offramp): use a valid strkey for the second wallet fixture

GBBB...SAM was 56 characters but failed StrKey.isValidEd25519PublicKey.
The test passes either way today, because the route does not validate the
address, but a fixture that is not a real strkey breaks the moment it is
handed to anything that decodes one.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

---------

Co-authored-by: blockchain-maxis <267648998+blockchain-maxis@users.noreply.github.com>
Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

* fix: skip malformed events in parseEvents instead of wedging the indexer (#197)

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

* fix(indexer): share one per-batch pipeline between the single and parallel ingest paths (#219)

Co-authored-by: DevTobis <232918735+DevTobis@users.noreply.github.com>

* Bring docker-compose.yml in line with env contract (#194) (#220)

* Bring docker-compose.yml in line with env contract (#194)

- Remove obsolete version: "3.9" key
- Add all missing env keys from .env.example:
  - DIRECT_DATABASE_URL
  - STELLAR_NETWORK
  - SOROBAN_RPC_URL (with STELLAR_RPC_URL as backward-compat alias)
  - NETWORKS
  - SAC_CONTRACT_IDS (and per-network variants)
  - NFT_CONTRACT_IDS (and per-network variants)
  - RETENTION_DAYS
  - CACHE_ENABLED and all Redis cache config
  - TOMBSTONE_CHECK_EVERY_CYCLES
  - LP_POOL_CONTRACT_IDS variants
  - SKIP_INDEXER
- Keep CONTRACT_IDS as documented backward-compat alias
- Add optional redis service with cache profile for CACHE_ENABLED support
- Update README to use SAC_CONTRACT_IDS and add cache profile instructions
- Add test to validate compose file meets requirements

* test(compose): derive the env-contract check from .env.example

The drift test built envKeys from .env.example and then never used it,
asserting against a hand-maintained list instead — so a key added to
.env.example and forgotten in docker-compose.yml still passed, which is the
one failure #194 is about. The assertions were also substring matches on the
whole file: toContain("SAC_CONTRACT_IDS") is satisfied by
SAC_CONTRACT_IDS_TESTNET, and toContain("CONTRACT_IDS") by either, so they
held on a file declaring none of them.

Now it parses the wraith service's own environment block (anchored on the
service, since db has an environment block too) and compares declared keys
against every uncommented key in .env.example, with an explicit, empty
exclusion list. Verified it fails on an added key and passes without one.

Also: the docker compose config case ran unconditionally and called fail(),
which is not defined under jest-circus, so on a machine or CI runner without
Docker it failed with a ReferenceError. Unit tests here do not require
Docker (the integration suite is vitest + Docker), so it now skips instead.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

---------

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

* Fix/cache network key (#202)

* fix(cache): include resolved network in Redis cache key (#182)

defaultKeyFn now prepends req.network (set by networkMiddleware) to the cache key so mainnet and testnet requests never collide, even when the network is supplied via the X-Network header rather than ?network=.

The ?network= query param is filtered from the query segment since it is already captured by req.network, ensuring ?network=mainnet and X-Network: mainnet produce identical keys.

Existing keys change shape and will expire naturally over their TTL.

* fix(cache): include resolved network in Redis cache key (#182)

* chore: drop the unrelated package-lock.json change

The branch re-resolved fsevents and dropped its "dev": true marker. Nothing
in this PR touches dependencies, so restore the lockfile to main's.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

---------

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

* Add npm run db:seed with deterministic fixtures (#221)

- Create shared fixture module in src/fixtures.ts with deterministic test data
  across two addresses (ALICE, BOB, CAROL) and two contracts (CONTRACT_A, CONTRACT_B)
- Add seed script at scripts/seed.ts with --network flag support (testnet/mainnet)
- Add db:seed script to package.json
- Update integration tests to use shared fixture module instead of local copy
- Update README Quick Start with seed step and real sample output
- Add fixture validation test to verify data structure and coverage

The seed script uses skipDuplicates on unique constraints, making it safe to
run multiple times without duplicating rows. Fixtures include TokenTransfer,
NftTransfer, and AccountSummary rows with deterministic eventId values.

Resolves #193

* fix(seed): accept --network mainnet as well as --network=mainnet

Follow-up to #221: this fix was pushed to the PR branch but did not make it
into the squash. The docstring advertised the space-separated form while
parseArgs matched only --network=, so 'npm run db:seed -- --network mainnet'
silently seeded testnet. Both spellings now parse; an unknown value still
exits 1.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

* Fix token decimals review issues

* perf(api): drop the default COUNT from list queries, add hasMore and opt-in includeTotal (#218)

Co-authored-by: royalTreasure <295874283+royalTreasure@users.noreply.github.com>

* feat: instrument HTTP surface in Prometheus (#189) (#212)

* feat: instrument HTTP surface in Prometheus

* fix(metrics): instrument before networkMiddleware and the rate limiter

Mounted after them, the HTTP middleware never saw the requests those two
reject, so 429s and invalid-?network= 400s were absent from
http_requests_total. Moved to the top of the chain, right after cors().

Also documents the two new metrics in the README table and records why the
route label must stay req.route.path: req.baseUrl is the matched mount path,
and src/api/accounts.ts mounts a router at "/:address/transfers", so baseUrl
carries the real address and would make the label unbounded.

Claude-Session: https://claude.ai/code/session_01USgemLt4Rnz4SGB1Srf3GB

---------

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

* fix(indexer): bound NFT metadata lookups with a worker pool and per-cycle budget (#217)

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>

---------

Co-authored-by: Miracle656 <iupacnumen2020@gmail.com>
Co-authored-by: Gloria <glorious217@gmail.com>
Co-authored-by: Jemimah <ekongjemimah@gmail.com>
Co-authored-by: blockchain-maxis <blockchainmaxis@gmail.com>
Co-authored-by: blockchain-maxis <267648998+blockchain-maxis@users.noreply.github.com>
Co-authored-by: Apulupie <167634780+Frun1na@users.noreply.github.com>
Co-authored-by: Tobiz <deborahayoola2000@gmail.com>
Co-authored-by: DevTobis <232918735+DevTobis@users.noreply.github.com>
Co-authored-by: BOA <97275013+boalambo@users.noreply.github.com>
Co-authored-by: Bathoul Mohammed <funds0033@gmail.com>
Co-authored-by: royaldev <chiditreasure15@gmail.com>
Co-authored-by: royalTreasure <295874283+royalTreasure@users.noreply.github.com>
Co-authored-by: Collins Ezedike-egwom <62267326+collinsezedike@users.noreply.github.com>
Co-authored-by: That guy <120946193+ezedike-evan@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Parallel ingest silently stops indexing NFTs and account summaries

2 participants