Repository navigation
Move each group of tables onto the connection that writes its parents - #712
Conversation
893685d to
4f8541b
Compare
4f8541b to
ecec219
Compare
ecec219 to
f8f2749
Compare
f8f2749 to
85e5637
Compare
There was a problem hiding this comment.
Pull request overview
Expands live ingestion into seven parallel write groups while preserving foreign-key ownership and cursor-last commit ordering.
Changes:
- Splits bulk tables and balance families across dedicated transactions.
- Fixes transaction keying, buffer capacity, and startup reconciliation.
- Optimizes balance upserts and updates tests/mocks.
Reviewed changes
Copilot reviewed 21 out of 21 changed files in this pull request and generated 3 comments.
Show a summary per file
| File | Description |
|---|---|
internal/services/token_ingestion.go |
Splits balance persistence APIs. |
internal/services/token_ingestion_test.go |
Updates balance service tests. |
internal/services/sep41/processor.go |
Clarifies history transaction behavior. |
internal/services/protocol_processor.go |
Documents live history commit semantics. |
internal/services/mocks.go |
Updates token ingestion mock methods. |
internal/services/ingest.go |
Separates parent and participant COPY calls. |
internal/services/ingest_test.go |
Updates persistence tests and adds batching coverage. |
internal/services/ingest_live.go |
Implements seven sibling write streams. |
internal/services/ingest_live_test.go |
Adds contract-data memo coverage. |
internal/services/blend/processor.go |
Clarifies history transaction behavior. |
internal/indexer/indexer_buffer.go |
Keys transactions by ToID. |
internal/indexer/indexer_buffer_test.go |
Covers duplicate hashes across ToIDs. |
internal/db/db.go |
Raises default pool capacity. |
internal/data/trustline_balances.go |
Sorts and guards trustline upserts. |
internal/data/transactions.go |
Splits transaction and account-link COPYs. |
internal/data/transactions_test.go |
Adapts transaction COPY tests. |
internal/data/operations.go |
Splits operation and account-link COPYs. |
internal/data/operations_test.go |
Adapts operation COPY tests. |
internal/data/native_balances.go |
Sorts and guards native-balance upserts. |
internal/data/ingest_store.go |
Bounds startup cleanup by close time. |
internal/data/ingest_store_test.go |
Covers bounded and fallback cleanup. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 21 out of 21 changed files in this pull request and generated 5 comments.
Suppressed comments (3)
internal/db/db.go:25
- The new persist path requires nine simultaneously held pool connections in production (seven siblings, one coordinator, and the advisory-lock session), but an explicitly configured
--db-max-connsbelow nine is still accepted. Because all sibling connections are acquired before any work starts, such a deployment blocks forever waiting for a connection that this same call cannot release. Raising only the default does not protect existing overrides; validate the live-ingest pool size at startup (and report the minimum) or restructure acquisition so undersized pools cannot deadlock.
// siblings + the balances and trustlines siblings + the coordinator),
// plus the advisory-lock session and up to 2 transient pool-side
// classification reads; 12 leaves headroom so no persist-path Acquire
// queues behind an unrelated consumer.
DefaultMaxConns int32 = 12
internal/services/ingest_live.go:282
- The PR description says the commit barrier is detached from cancellation, but the barrier directly below still calls every
Commitwith the pipelinectx. A SIGTERM after one sibling commits can therefore cancel a later commit and manufactureErrPartialPersist, splitting a batch during normal shutdown. Use acontext.WithoutCancel(ctx)commit context for all sibling commits and the final coordinator commit once staging has succeeded.
// Nothing has committed: the deferred rollbacks discard every
// transaction and the batch is cleanly retryable.
internal/services/token_ingestion.go:88
- This documentation is incorrect:
liquidity_pool_balances.pool_idhas a deferred foreign key toliquidity_pools.pool_id. The method is safe because it writes pools before pool-share balances in the same transaction, not because all target tables are FK-free.
// ProcessNativeAndPoolChanges applies native-balance and liquidity-pool changes; the
// target tables have no foreign keys, so any transaction may carry them.
85e5637 to
1f5031c
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 21 out of 21 changed files in this pull request and generated no new comments.
Suppressed comments (3)
internal/services/token_ingestion.go:88
- This is incorrect:
liquidity_pool_balances.pool_idhas a deferred foreign key toliquidity_pools.pool_id. The ordering and shared transaction are required for pool-share writes, rather than these tables being safe on any transaction.
// ProcessNativeAndPoolChanges applies native-balance and liquidity-pool changes; the
// target tables have no foreign keys, so any transaction may carry them.
internal/services/ingest_live.go:204
- Detach the commit barrier from the pipeline context. This PR expands the barrier to seven sequential sibling commits, but the commits below still use
ctx; if shutdown cancels it after an early sibling commits, the remaining commits fail locally and turn a healthy batch intoErrPartialPersist, forcing a fatal restart and reconciliation. Use a context derived fromcontext.WithoutCancel(ctx)for every sibling commit and the final coordinator commit.
{"balances", func(ctx context.Context, dbTx pgx.Tx, it *persistItem) error {
return m.tokenIngestionService.ProcessNativeAndPoolChanges(ctx, dbTx,
it.buffer.GetAccountChanges(),
it.buffer.GetLiquidityPoolShareChanges(),
it.buffer.GetLiquidityPoolChanges(),
)
internal/db/db.go:25
- Raising only the default leaves an explicit
--db-max-connsbelow 9 able to deadlock live ingestion. The advisory-lock session is held for the run, thenpersistLedgerDataholds the coordinator plus each sibling connection while acquiring the next; with seven siblings, a pool capped at 8 blocks forever acquiring the final sibling because none of the already-held connections can be released. Validate the live-mode pool size and fail fast (minimum 9, plus any desired headroom) rather than relying on the default.
// Live persist alone holds 8 connections at its commit barrier (5 COPY
// siblings + the balances and trustlines siblings + the coordinator),
// plus the advisory-lock session and up to 2 transient pool-side
// classification reads; 12 leaves headroom so no persist-path Acquire
// queues behind an unrelated consumer.
DefaultMaxConns int32 = 12
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 25 out of 25 changed files in this pull request and generated 3 comments.
Suppressed comments (1)
internal/services/token_ingestion.go:88
- This implementation comment is incorrect:
liquidity_pool_balances.pool_idhas a deferred foreign key toliquidity_pools.pool_id. The interface comment above correctly explains why both writes must share this transaction; keep the implementation documentation consistent so callers do not infer that these tables can be split safely.
// ProcessNativeAndPoolChanges applies native-balance and liquidity-pool changes; the
// target tables have no foreign keys, so any transaction may carry them.
29108cf to
858d2b1
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 25 out of 25 changed files in this pull request and generated 1 comment.
Suppressed comments (3)
internal/services/ingest_live.go:294
- The PR says the commit barrier runs on
context.WithoutCancel, but the barrier below still calls everyCommitwith the cancellable request context. A SIGTERM after the first sibling commit can therefore cancel the remaining commits and split the batch, exactly the shutdown failure the description says is prevented. Use a bounded context derived fromcontext.WithoutCancel(ctx)for all sibling and coordinator commits.
// Nothing has committed: the deferred rollbacks discard every
// transaction and the batch is cleanly retryable.
internal/ingest/ingest.go:186
- This rejects backfill runs configured with fewer than nine connections even though only live mode opens the seven-sibling barrier.
cfg.IngestionModeis already available here, and a low-concurrency backfill can safely use a smaller pool; this turns a previously valid explicit backfill configuration into a startup error. Apply the floor only for live ingestion.
if err := validateIngestPoolConfig(poolCfg); err != nil {
return nil, nil, err
internal/services/token_ingestion.go:88
- This contract is incorrect:
liquidity_pool_balances.pool_idhas a deferred foreign key toliquidity_pools, which is why this method must upsert pools before pool-share balances on the same transaction. Saying any transaction may carry these writes contradicts the interface documentation and could lead a future caller to split them again.
// ProcessNativeAndPoolChanges applies native-balance and liquidity-pool changes; the
// target tables have no foreign keys, so any transaction may carry them.
aristidesstaffieri
left a comment
There was a problem hiding this comment.
I flagged a few of the copilot comments that are still relevant but otherwise this lgtm.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Mutable balance siblings can still commit ahead of the cursor without rollback or startup reconciliation after a partial commit.
Review effort: Balanced
Findings: 2
Open (6)
Keying by ToID now deliberately persists multiple rows with the same hash, but… Pre-committing these mutable current-state tables is not crash-safe merely because their upserts… Warn before unbounded cleanup fallback Indexer buffer rotation overallocates one retained buffer The new ordering does not skip unchanged trustline rows: the conflict clause below still executes… Sorting reduces probe randomness, but the followingON CONFLICT DO UPDATEis still unconditional,…
198d1b1 to
616727a
Compare
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
Mutable sibling commits lack crash-safe recovery, and duplicate transaction hashes leave the singular hash lookup ambiguous.
Review effort: Balanced
Findings: 2
Open (6)
Keying by ToID now deliberately persists multiple rows with the same hash, but… Pre-committing these mutable current-state tables is not crash-safe merely because their upserts… Warn before unbounded cleanup fallback Indexer buffer rotation overallocates one retained buffer The new ordering does not skip unchanged trustline rows: the conflict clause below still executes… Sorting reduces probe randomness, but the followingON CONFLICT DO UPDATEis still unconditional,…
…lings The live persist path's sibling COPY streams grow from three to five: transactions_accounts and operations_accounts — the two largest index-maintenance payloads, each carrying an account_id-leading unique PK — stream concurrently with their parent tables instead of serially behind them. TransactionModel and OperationModel each split BatchCopy into a parent-table COPY and BatchCopyAccounts for the link table, with the link COPY's duration and batch size recorded under its own metric labels rather than folded into the parent's. Each sibling now writes exactly one table, so the disjointness at the commit barrier is per-table; the backfill path calls the five inserts in the old order and is behaviorally unchanged. Persist holds six connections at its commit barrier; the pool default rises to 12 so the persist path never queues on Acquire behind the advisory-lock session or pool-side classification reads.
… random order The native_balances and trustline_balances batch upserts arrive in Go map-iteration order and unconditionally rewrite every conflicting row. Both cost heap buffer touches far above what the row count warrants and bloat the table with dead tuples. Two changes: rows sort by the primary-key columns before the UNNEST arrays are built, so the btree descent and heap touches run in key order; and the DO UPDATE carries an IS DISTINCT FROM guard, so an identical row produces no new tuple version, no dead tuple, no WAL, and no index churn.
The buffer stored transactions keyed by hash while tracking their participants keyed by ToID. A merged ledger can carry the same envelope at several tx-set positions — distinct ToIDs, one hash — and the two maps then disagreed: duplicate transactions were silently dropped from the transactions COPY while their participant links survived, whose ledger_created_at lookup missed and COPYed the zero timestamp into a year-0001 chunk (observed on the rig: 6 orphaned transactions_accounts rows, 5 missing transaction rows across bootstrap ledgers). Real networks cannot repeat a hash within or across ledgers, so both maps now share the ToID key domain by construction, mirroring the operations pair. BatchCopyAccounts on both link tables now errors on a ToID/opID with no parent row instead of silently writing a zero timestamp.
…eachable A buffer stays checked out from the moment the process stage takes it until the batch carrying it commits, so the rotation needs 2*cap buffers to sustain a full batch: cap for the batch in flight, cap-1 queued in processed, one being filled. It held cap+1, so process starved on freeBuffers after filling (cap+1)-k and the batch settled at the fixed point k = (cap+1)-k. The batch size was therefore pinned below the configured cap at every load, never reaching it however deep the backlog. The new test pins the old behaviour: with the pool sized the old way, batches settle below a cap of 3.
…ling The native-balance and liquidity-pool upserts ran serially inside the coordinating transaction, lengthening its critical path while the COPY siblings streamed concurrently. Their tables carry no foreign keys, so a sixth sibling now carries them; the trustline and SAC balances stay on the coordinating transaction because their parent rows (trustline_assets, contract_tokens) are staged uncommitted there, and a sibling committing first would fail FK checks against committed state. Reconciliation (DeleteRowsAboveLedger) does not cover the balance tables: a crash can leave balance rows above the committed cursor, which is safe because the upserts are idempotent and reapply when those ledgers re-ingest.
…bling Each balance family now rides the transaction that stages its foreign-key parents: trustline_assets and trustline_balances move from the coordinating transaction onto a new sibling, alongside the balances sibling that carries pools and pool-share balances under the new pool foreign key. Every FK is checked at its own transaction's commit, and the coordinator's serial path shrinks to contracts, classification, protocol state, SAC balances, and the cursor. SAC balances stay coordinated: their parent (contract_tokens) is also written by the classification path there, and a same-key insert from two concurrent transactions could deadlock at the commit barrier.
…mo semantics BatchCopyAccounts on transactions and operations returned a link-row count every caller discarded — the count is already observed into the BatchSize metric inside the method — so both now return plain error, matching the balance models' BatchCopy shape. contractDataMemo's contract gets a unit test: get() always yields a rangeable non-nil map, and the extraction walk runs at most once no matter how many retry attempts share the memo. The startup-reconciliation comment now says a crash orphans at most the persist batch past the cursor; batching had outgrown "the single ledger".
…barrier Live persist now holds 8 connections at once, the coordinating transaction plus one for each of the 7 sibling transactions that commit before it, and the advisory-lock session holds a 9th for the process lifetime without ever returning it. The startup floor that rejects a db-max-conns too small for the barrier moves from 5 to 9 to match.
The IS DISTINCT FROM guard compared last_modified_ledger, which is fed from the ingested ledger sequence rather than the entry's own lastModifiedLedgerSeq (processors/accounts.go, processors/trustlines.go). That value always advances in forward ingestion, so the row was always DISTINCT and the update always fired: the promised no-op — no new tuple version, no dead tuple, no WAL, no index churn — could not occur. The one path that did suppress an update is re-ingesting an already-ingested ledger, which is bounded by the uncommitted cursor gap and already safe without the guard: every balance write is an absolute-value upsert or delete and is idempotent either way. What remained was a six-column comparison per conflicting row on the hot persist path, applied to native and trustline balances but not to pool-share or SAC balances. The PK sort in the same upserts is untouched — descending the btree in key order instead of scattering random heap probes is the part that pays.
The buffer rotation is sized 2*batchCap+1, so peak ledger-buffer memory grows at roughly twice the flag's value. That is invisible at the default of 1, but the deployments that raise the cap are the high-TPS ones already running close to GOMEMLIMIT — exactly where the doubling matters.
The five bulk-COPY tables were enumerated independently in startup reconciliation, the hypertable settings pass, the backfill recompressor, and a test's own copy — and a table missing from DeleteRowsAboveLedger leaves orphans that crash-loop re-ingest on PK collisions. data.BulkCopyTables (table + TOID column) is now the single definition all of them consume.
…bles liquidity_pool_balances.pool_id references liquidity_pools, so the comment claiming the targets have no foreign keys was false. Pools are written first on the same transaction, which is why any transaction may still carry them.
Backfill never opens the commit barrier, so a resource-constrained backfill run must not fail startup over a connection floor that only live persist needs.
Startup reconciliation can delete bulk-COPY rows a crash strands above the cursor, but a balance or trustline row can only be overwritten by replay, so those two groups must commit last, right before the cursor-carrying coordinating transaction, to keep the window in which such a row is visible ahead of the cursor as short as the barrier allows. The order was incidental slice order; it is now a named list with a test.
A participant keyed to a ToID or operation ID with no parent row would have landed a year-0001 partition; the Go-side guards are the only defense, since no schema constraint catches an orphan link row. Both guards now have a table case that fails without them.
616727a to
53bbd43
Compare
MinIngestMaxConns is the barrier's held connections plus the advisory-lock session, but only a comment linked it to the number of siblings. The order test now asserts the relation, so adding a sibling without raising the floor fails a test instead of wedging a pool at exactly the old floor. The batch cap flag's usage names the buffer count it scales.
…key parent persistLedgerData writes trustline assets with trustline balances and liquidity pools with pool-share balances on the same sibling transaction, because the foreign keys are deferred and checked at commit against committed state. Every persist test mocked the token ingestion service, so moving a parent onto another transaction would fail in production with SQLSTATE 23503 while the suite stayed green. This test persists one ledger through the real service and reads the four tables back.
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
The expanded multi-transaction commit barrier and accepted crash-consistency windows warrant final human validation.
Review effort: Balanced
Findings: None
Resolved since last review (6)
Keying by ToID now deliberately persists multiple rows with the same hash, but… Pre-committing these mutable current-state tables is not crash-safe merely because their upserts… Warn before unbounded cleanup fallback Indexer buffer rotation overallocates one retained buffer The new ordering does not skip unchanged trustline rows: the conflict clause below still executes… Sorting reduces probe randomness, but the followingON CONFLICT DO UPDATEis still unconditional,…


Start here: the table-group list in
persistLedgerData. This PR takes it from 3 groups to 7.The rule: each group is written by the transaction that also writes its foreign-key parents.
3 groups become 7
transactions_accounts,operations_accountsSAC balances stay on the coordinating transaction. Their parent
contract_tokensis written there too, and two transactions inserting the same key could deadlock.Protocol history moves to the
state_changesconnection. It was writingstate_changesfrom the coordinating transaction, so two of our transactions hit the same hypertable at once. At a chunk boundary TimescaleDB makes one wait — and the coordinating transaction only commits after the others finish, so that wait hangs ingestion forever. Postgres can't see it as a deadlock.Balances and trustlines commit before the cursor, and that is a trade-off. The two mutable groups (native and pool balances, trustline assets and balances) commit as siblings, last among them, right before the coordinator carries the cursor. A crash between those two commits leaves them one ledger ahead of the cursor until the restart replays that ledger. Every write in those groups is an absolute value (
EXCLUDED.*upserts, deletes,DO NOTHINGinserts), so the replay converges to the same rows; the only delta write, SEP-41 balances, stays on the coordinator behind the CAS cursor. Withlive-persist-max-batch-sizeabove 1 a replay can cut batches differently, so a row changed twice inside one batch can briefly show the earlier value again. Keeping these groups on the coordinator would avoid the window at the cost of serializing their upserts behind everything else; the window is accepted instead.The pool now needs 9 connections
8 are held at the commit barrier (7 siblings + coordinator), plus the advisory lock, which is held until the process exits. Default
db-max-connsgoes 10 → 12.Set it below 9 and a sibling waits forever: no error, no crash-loop, just a climbing lag gauge. Ingest now refuses to start.
Bugs this uncovered
2*cap, notcap+1. Processing ran out of free buffers, so batches never reached the configured size. Raisinglive-persist-max-batch-sizenow costs ~2× its value in peak buffer memory.Also: balance upserts are sorted by primary key, so each batch descends the btree in order instead of scattering random heap probes.
No performance numbers yet — please review design and correctness.
PR 5 of 6 replacing #684 · schema → fixes → pipeline → parallel writes → table groups → final fixes