Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 36 additions & 0 deletions migrations/1779800000026_sponsor-nonce-indexes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
import type { ColumnDefinitions, MigrationBuilder } from 'node-pg-migrate';

export const shorthands: ColumnDefinitions | undefined = undefined;

/**
* Re-keys the sponsor indexes on `txs` and `mempool_txs` from `(sponsor_address, nonce)` to
* `(sponsor_address, sponsor_nonce)`.
*
* A sponsored transaction consumes two account nonces: its origin's `nonce` and its sponsor's
* `sponsor_nonce`. Replace-by-fee detection and mempool pruning/restoring look up conflicting
* transactions by the sponsor's nonce slot, which the old index could only serve as a scan over
* every transaction the sponsor ever paid for (its `nonce` column is the origin's nonce, not the
* sponsor's). Every other query over these indexes filters by `sponsor_address` alone or reads
* `MAX(sponsor_nonce)`, which the new key serves at least as well, so the old indexes are replaced
* rather than kept alongside: index count and write-path cost are unchanged.
*
* The partial predicates are kept as they were. `mempool_txs` is small enough to index in full; the
* `txs` index only covers canonical sponsored transactions, a small share of the table.
*/
export function up(pgm: MigrationBuilder): void {
pgm.dropIndex('txs', ['sponsor_address', 'nonce']);
pgm.createIndex('txs', ['sponsor_address', 'sponsor_nonce'], {
where: 'sponsor_address IS NOT NULL AND canonical = true AND microblock_canonical = true',
});
pgm.dropIndex('mempool_txs', ['sponsor_address', 'nonce']);
pgm.createIndex('mempool_txs', ['sponsor_address', 'sponsor_nonce']);
}

export function down(pgm: MigrationBuilder): void {
pgm.dropIndex('mempool_txs', ['sponsor_address', 'sponsor_nonce']);
pgm.createIndex('mempool_txs', ['sponsor_address', 'nonce']);
pgm.dropIndex('txs', ['sponsor_address', 'sponsor_nonce']);
pgm.createIndex('txs', ['sponsor_address', 'nonce'], {
where: 'sponsor_address IS NOT NULL AND canonical = true AND microblock_canonical = true',
});
}
211 changes: 96 additions & 115 deletions src/datastore/pg-write-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -151,8 +151,22 @@ type TransactionHeader = {
sponsor_address?: string;
sponsored: boolean;
nonce: number;
sponsor_nonce?: number;
};

/**
* Account nonces a transaction consumes. Every transaction consumes its origin's `nonce`, and a
* sponsored transaction also consumes its sponsor's `sponsor_nonce`. Two transactions conflict when
* they share any of these slots, which mirrors the unique constraints the Stacks node's mempool
* enforces.
*/
function txNonceSlots(tx: TransactionHeader): [string, number][] {
const slots: [string, number][] = [[tx.sender_address, tx.nonce]];
if (tx.sponsored && tx.sponsor_address && tx.sponsor_nonce != null)
slots.push([tx.sponsor_address, tx.sponsor_nonce]);
return slots;
}

/**
* Extends `PgStore` to provide data insertion functions. These added features are usually called by
* the `EventServer` upon receiving blockchain events from a Stacks node. It also deals with chain data
Expand Down Expand Up @@ -325,6 +339,7 @@ export class PgWriteStore extends PgStore {
sponsor_address: d.tx.sponsor_address,
sponsored: d.tx.sponsored,
nonce: d.tx.nonce,
sponsor_nonce: d.tx.sponsor_nonce,
}));
await this.pruneMempoolTxs(sql, prunableTxs);
}
Expand Down Expand Up @@ -3667,54 +3682,60 @@ export class PgWriteStore extends PgStore {
): Promise<void> {
if (txIds.length === 0) return;

// If a transaction with equal nonce was confirmed in a block, mark all conflicting mempool txs
// as RBF. Otherwise, look for the one with the highest fee in the mempool and RBF all the
// others.
// Every transaction occupies one nonce slot for its origin and, if sponsored, a second one for
// its sponsor (see `txNonceSlots`). For each slot touched by these txs: if a transaction in it
// was confirmed in a block, mark all conflicting mempool txs as RBF. Otherwise, look for the
// one with the highest fee in the mempool and RBF all the others. A tx that loses in any of its
// slots is RBFd, since the node's mempool can't keep it alongside that slot's winner.
//
// Note that we're not filtering by `pruned` when we look at the mempool, because we want the
// RBF data to be retroactively applied to all conflicting txs we've ever seen.
// Note that we're not filtering by `pruned` when we look for losers, because we want the RBF
// data to be retroactively applied to all conflicting txs we've ever seen. Only unpruned txs
// can win in the mempool, though: a tx pruned for losing one of its slots must not keep
// blocking its other slot.
for (const batch of batchIterate(txIds, INSERT_BATCH_SIZE)) {
await sql`
WITH input_txids (tx_id) AS (
VALUES ${sql(batch.map(id => [id.replace('0x', '\\x')]))}
),
source_txs AS (
SELECT DISTINCT
tx_id,
(CASE sponsored WHEN true THEN sponsor_address ELSE sender_address END) AS address,
nonce
SELECT sender_address, nonce, sponsored, sponsor_address, sponsor_nonce
FROM ${mempool ? sql`mempool_txs` : sql`txs`}
WHERE tx_id IN (SELECT tx_id::bytea FROM input_txids)
),
affected_groups AS (
SELECT DISTINCT address, nonce
affected_slots AS (
SELECT sender_address AS address, nonce
FROM source_txs
UNION
SELECT sponsor_address AS address, sponsor_nonce AS nonce
FROM source_txs
WHERE sponsored = true AND sponsor_nonce IS NOT NULL
),
same_nonce_mempool_txs AS (
same_slot_mempool_txs AS (
SELECT m.tx_id, m.fee_rate, m.receipt_time, m.pruned, g.address, g.nonce
FROM mempool_txs m
INNER JOIN affected_groups g
INNER JOIN affected_slots g
ON m.sender_address = g.address AND m.nonce = g.nonce
UNION
SELECT m.tx_id, m.fee_rate, m.receipt_time, m.pruned, g.address, g.nonce
FROM mempool_txs m
INNER JOIN affected_groups g
ON m.sponsor_address = g.address AND m.nonce = g.nonce
INNER JOIN affected_slots g
ON m.sponsor_address = g.address AND m.sponsor_nonce = g.nonce
WHERE m.sponsored = true
),
mined_txs AS (
SELECT t.tx_id, g.address, g.nonce,
t.block_height, t.microblock_sequence, t.tx_index
FROM txs t
INNER JOIN affected_groups g
INNER JOIN affected_slots g
ON t.sender_address = g.address AND t.nonce = g.nonce
WHERE t.canonical = true AND t.microblock_canonical = true
UNION
SELECT t.tx_id, g.address, g.nonce,
t.block_height, t.microblock_sequence, t.tx_index
FROM txs t
INNER JOIN affected_groups g
ON t.sponsor_address = g.address AND t.nonce = g.nonce
WHERE t.canonical = true AND t.microblock_canonical = true
INNER JOIN affected_slots g
ON t.sponsor_address = g.address AND t.sponsor_nonce = g.nonce
WHERE t.sponsored = true AND t.canonical = true AND t.microblock_canonical = true
),
latest_mined_txs AS (
SELECT DISTINCT ON (address, nonce) tx_id, address, nonce
Expand All @@ -3723,37 +3744,37 @@ export class PgWriteStore extends PgStore {
),
highest_fee_mempool_txs AS (
SELECT DISTINCT ON (address, nonce) tx_id, address, nonce
FROM same_nonce_mempool_txs
FROM same_slot_mempool_txs
WHERE pruned = false
ORDER BY address, nonce, fee_rate DESC, receipt_time DESC
Comment thread
rafa-stacks marked this conversation as resolved.
Comment thread
rafa-stacks marked this conversation as resolved.
),
winning_txs AS (
SELECT DISTINCT
SELECT
g.address,
g.nonce,
COALESCE(l.tx_id, h.tx_id) AS tx_id
FROM affected_groups g
COALESCE(l.tx_id, h.tx_id) AS tx_id,
l.tx_id IS NOT NULL AS mined
FROM affected_slots g
LEFT JOIN latest_mined_txs l USING (address, nonce)
LEFT JOIN highest_fee_mempool_txs h USING (address, nonce)
),
txs_to_prune AS (
SELECT
-- A tx can lose in both of its slots, so pick one replacement deterministically, favoring
-- a confirmed one.
SELECT DISTINCT ON (s.tx_id)
s.tx_id,
s.pruned
FROM same_nonce_mempool_txs s
s.pruned,
w.tx_id AS replaced_by_tx_id
FROM same_slot_mempool_txs s
INNER JOIN winning_txs w USING (address, nonce)
WHERE s.tx_id <> w.tx_id
ORDER BY s.tx_id, w.mined DESC, w.tx_id
),
pruned AS (
UPDATE mempool_txs m
SET pruned = TRUE,
status = ${DbTxStatus.DroppedReplaceByFee},
replaced_by_tx_id = (
SELECT w.tx_id
FROM winning_txs w
INNER JOIN same_nonce_mempool_txs s ON w.address = s.address AND w.nonce = s.nonce
WHERE s.tx_id = m.tx_id
LIMIT 1
)
replaced_by_tx_id = p.replaced_by_tx_id
FROM txs_to_prune p
WHERE m.tx_id = p.tx_id
RETURNING m.tx_id
Expand Down Expand Up @@ -4209,6 +4230,37 @@ export class PgWriteStore extends PgStore {
block_height = EXCLUDED.block_height
`;
}

/**
* CTEs ending in `affected_mempool_tx_ids`: the given transactions plus every mempool transaction
* that shares a nonce slot with any of them (see {@link txNonceSlots}).
*/
private sameNonceSlotMempoolTxIdsCte(sql: PgSqlClient, transactions: TransactionHeader[]) {
const txIds = transactions.map(t => [t.txId.replace('0x', '\\x')]);
const slots = transactions.flatMap(t => txNonceSlots(t));
return sql`
input_tx_ids (tx_id) AS (
VALUES ${sql(txIds)}
),
input_slots (address, nonce) AS (
VALUES ${sql(slots)}
),
affected_mempool_tx_ids AS (
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN input_slots i ON m.sender_address = i.address AND m.nonce = i.nonce::int
UNION
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN input_slots i
ON m.sponsor_address = i.address AND m.sponsor_nonce = i.nonce::int
WHERE m.sponsored = true
UNION
SELECT tx_id::bytea FROM input_tx_ids
)
`;
}

/**
* Restore transactions in the mempool table. This should be called when mined transactions are
* marked from canonical to non-canonical.
Expand All @@ -4225,47 +4277,10 @@ export class PgWriteStore extends PgStore {
`Restoring mempool tx: ${tx.txId} sender: ${tx.sender_address} nonce: ${tx.nonce}`
);

// Restore new non-canonical txs into the mempool. Also restore transactions for the same
// senders/sponsors with the same `nonce`s. We will recalculate replace-by-fee ordering shortly
// afterwards.
const inputData = transactions.map(t => [
t.txId.replace('0x', '\\x'),
t.sender_address,
t.sponsor_address ?? 'null',
t.sponsored.toString(),
t.nonce,
]);
// Restore new non-canonical txs into the mempool. Also restore transactions that share a nonce
// slot with them. We will recalculate replace-by-fee ordering shortly afterwards.
const updatedRows = await sql<{ tx_id: string }[]>`
WITH input_data (tx_id, sender_address, sponsor_address, sponsored, nonce) AS (
VALUES ${sql(inputData)}
),
sponsored_inputs AS (SELECT * FROM input_data WHERE sponsored::boolean),
non_sponsored_inputs AS (SELECT * FROM input_data WHERE NOT sponsored::boolean),
affected_sponsored AS (
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN sponsored_inputs i ON m.sponsor_address = i.sponsor_address AND m.nonce = i.nonce::int
UNION
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN sponsored_inputs i ON m.sender_address = i.sponsor_address AND m.nonce = i.nonce::int
),
affected_non_sponsored AS (
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN non_sponsored_inputs i ON m.sponsor_address = i.sender_address AND m.nonce = i.nonce::int
UNION
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN non_sponsored_inputs i ON m.sender_address = i.sender_address AND m.nonce = i.nonce::int
),
affected_mempool_tx_ids AS (
SELECT tx_id FROM affected_sponsored
UNION
SELECT tx_id FROM affected_non_sponsored
UNION
SELECT tx_id::bytea FROM input_data
),
WITH ${this.sameNonceSlotMempoolTxIdsCte(sql, transactions)},
restored AS (
UPDATE mempool_txs
SET pruned = false, status = ${DbTxStatus.Pending}, replaced_by_tx_id = NULL
Expand Down Expand Up @@ -4326,47 +4341,10 @@ export class PgWriteStore extends PgStore {
`Pruning mempool tx: ${tx.txId} sender: ${tx.sender_address} nonce: ${tx.nonce}`
);

// Prune confirmed txs from the mempool. Also prune transactions for the same senders/sponsors
// with the same `nonce`s. We'll recalculate replaced-by-fee data later when new block data is
// written to the DB.
const inputData = transactions.map(t => [
t.txId.replace('0x', '\\x'),
t.sender_address,
t.sponsor_address ?? 'null',
t.sponsored.toString(),
t.nonce,
]);
// Prune confirmed txs from the mempool. Also prune transactions that share a nonce slot with
// them. We'll recalculate replaced-by-fee data later when new block data is written to the DB.
const updateResults = await sql<{ tx_id: string }[]>`
WITH input_data (tx_id, sender_address, sponsor_address, sponsored, nonce) AS (
VALUES ${sql(inputData)}
),
sponsored_inputs AS (SELECT * FROM input_data WHERE sponsored::boolean),
non_sponsored_inputs AS (SELECT * FROM input_data WHERE NOT sponsored::boolean),
affected_sponsored AS (
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN sponsored_inputs i ON m.sponsor_address = i.sponsor_address AND m.nonce = i.nonce::int
UNION
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN sponsored_inputs i ON m.sender_address = i.sponsor_address AND m.nonce = i.nonce::int
),
affected_non_sponsored AS (
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN non_sponsored_inputs i ON m.sponsor_address = i.sender_address AND m.nonce = i.nonce::int
UNION
SELECT m.tx_id
FROM mempool_txs m
INNER JOIN non_sponsored_inputs i ON m.sender_address = i.sender_address AND m.nonce = i.nonce::int
),
affected_mempool_tx_ids AS (
SELECT tx_id FROM affected_sponsored
UNION
SELECT tx_id FROM affected_non_sponsored
UNION
SELECT tx_id::bytea FROM input_data
),
WITH ${this.sameNonceSlotMempoolTxIdsCte(sql, transactions)},
pruned AS (
UPDATE mempool_txs
SET pruned = true, replaced_by_tx_id = NULL
Expand Down Expand Up @@ -4453,14 +4431,16 @@ export class PgWriteStore extends PgStore {
sponsor_address: string | null;
sponsored: boolean;
nonce: number;
sponsor_nonce: number | null;
update_balances_count: number;
}[]
>`
WITH updated_txs AS (
UPDATE txs
SET canonical = ${canonical}
WHERE index_block_hash = ${indexBlockHash} AND canonical != ${canonical}
RETURNING tx_id, sender_address, nonce, sponsor_address, fee_rate, sponsored, canonical
RETURNING tx_id, sender_address, nonce, sponsor_address, sponsor_nonce, fee_rate, sponsored,
canonical
),
affected_addresses AS (
SELECT
Expand Down Expand Up @@ -4495,7 +4475,7 @@ export class PgWriteStore extends PgStore {
SET balance = ft_balances.balance + EXCLUDED.balance
RETURNING ft_balances.address
)
SELECT tx_id, sender_address, sponsor_address, sponsored, nonce,
SELECT tx_id, sender_address, sponsor_address, sponsored, nonce, sponsor_nonce,
(SELECT COUNT(*)::int FROM update_ft_balances) AS update_balances_count
FROM updated_txs
`;
Expand All @@ -4505,6 +4485,7 @@ export class PgWriteStore extends PgStore {
sponsor_address: row.sponsor_address ?? undefined,
sponsored: row.sponsored,
nonce: row.nonce,
sponsor_nonce: row.sponsor_nonce ?? undefined,
}));
if (canonical) {
updatedEntities.markedCanonical.txs += txResult.count;
Expand Down
Loading
Loading