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
9 changes: 7 additions & 2 deletions crates/common/src/local_db/pipeline/adapters/apply.rs
Original file line number Diff line number Diff line change
Expand Up @@ -375,14 +375,19 @@ mod tests {

// Expect derived state refreshes plus the watermark and ANALYZE when there is no work.
let texts: Vec<_> = batch.statements().iter().map(|s| s.sql().trim()).collect();
assert_eq!(texts.len(), 8);
assert_eq!(texts.len(), 9);
assert!(texts
.iter()
.any(|s| s.contains("DELETE FROM derived_vault_deltas")));
assert!(texts
.iter()
.any(|s| s.contains("INSERT OR REPLACE INTO derived_vault_deltas")));
assert!(texts.iter().any(|s| s.contains("vault_balance_changes")));
assert!(texts
.iter()
.any(|s| s.contains("DELETE FROM vault_balance_changes")));
assert!(texts
.iter()
.any(|s| s.contains("INSERT OR IGNORE INTO vault_balance_changes")));
assert!(texts
.iter()
.any(|s| s.contains("INSERT OR REPLACE INTO running_vault_balances")));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,17 +22,23 @@ filtered AS (
PARTITION BY vd.chain_id, vd.raindex_address, vd.owner, vd.token, vd.vault_id
ORDER BY vd.block_number, vd.log_index
) AS rn,
COALESCE(rvb.balance, FLOAT_ZERO_HEX()) AS prefix_balance
COALESCE(
(
SELECT FLOAT_SUM(prev.delta ORDER BY prev.block_number, prev.log_index)
FROM derived_vault_deltas prev
WHERE prev.chain_id = vd.chain_id
AND prev.raindex_address = vd.raindex_address
AND prev.owner = vd.owner
AND prev.token = vd.token
AND prev.vault_id = vd.vault_id
AND prev.block_number < p.start_block
),
FLOAT_ZERO_HEX()
) AS prefix_balance
FROM derived_vault_deltas vd
JOIN params p
ON p.chain_id = vd.chain_id
AND p.raindex_address = vd.raindex_address
LEFT JOIN running_vault_balances rvb
ON rvb.chain_id = vd.chain_id
AND rvb.raindex_address = vd.raindex_address
AND rvb.owner = vd.owner
AND rvb.token = vd.token
AND rvb.vault_id = vd.vault_id
WHERE vd.block_number BETWEEN p.start_block AND p.end_block
),
ordered AS (
Expand Down
Original file line number Diff line number Diff line change
@@ -1,18 +1,34 @@
WITH latest_blocks AS (
SELECT
WITH affected_keys AS (
SELECT DISTINCT
chain_id,
raindex_address,
owner,
token,
vault_id,
MAX(block_number) AS last_block
vault_id
FROM derived_vault_deltas vd
WHERE vd.chain_id = ?1
AND vd.raindex_address = ?2
AND vd.block_number BETWEEN ?3 AND ?4
GROUP BY chain_id, raindex_address, owner, token, vault_id
),
delta_batches AS (
latest_blocks AS (
SELECT
vd.chain_id,
vd.raindex_address,
vd.owner,
vd.token,
vd.vault_id,
MAX(vd.block_number) AS last_block
FROM derived_vault_deltas vd
JOIN affected_keys ak
ON ak.chain_id = vd.chain_id
AND ak.raindex_address = vd.raindex_address
AND ak.owner = vd.owner
AND ak.token = vd.token
AND ak.vault_id = vd.vault_id
WHERE vd.block_number <= ?4
GROUP BY vd.chain_id, vd.raindex_address, vd.owner, vd.token, vd.vault_id
),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
aggregated AS (
SELECT
vd.chain_id,
vd.raindex_address,
Expand All @@ -22,7 +38,7 @@ delta_batches AS (
COALESCE(
FLOAT_SUM(vd.delta ORDER BY vd.block_number, vd.log_index),
FLOAT_ZERO_HEX()
) AS balance_delta,
) AS balance,
lb.last_block,
(
SELECT MAX(vd2.log_index)
Expand All @@ -41,63 +57,8 @@ delta_batches AS (
AND lb.owner = vd.owner
AND lb.token = vd.token
AND lb.vault_id = vd.vault_id
WHERE vd.chain_id = ?1
AND vd.raindex_address = ?2
AND vd.block_number BETWEEN ?3 AND ?4
WHERE vd.block_number <= ?4
GROUP BY vd.chain_id, vd.raindex_address, vd.owner, vd.token, vd.vault_id, lb.last_block
),
existing_matching AS (
SELECT
mvb.chain_id,
mvb.raindex_address,
mvb.owner,
mvb.token,
mvb.vault_id,
mvb.balance AS balance_value,
mvb.last_block,
mvb.last_log_index
FROM running_vault_balances mvb
JOIN delta_batches db
ON db.chain_id = mvb.chain_id
AND db.raindex_address = mvb.raindex_address
AND db.owner = mvb.owner
AND db.token = mvb.token
AND db.vault_id = mvb.vault_id
),
combined AS (
SELECT
chain_id,
raindex_address,
owner,
token,
vault_id,
balance_delta AS contribution,
last_block,
last_log_index
FROM delta_batches
UNION ALL
SELECT
chain_id,
raindex_address,
owner,
token,
vault_id,
balance_value AS contribution,
last_block,
last_log_index
FROM existing_matching
),
aggregated AS (
SELECT
chain_id,
raindex_address,
owner,
token,
vault_id,
COALESCE(FLOAT_SUM(contribution), FLOAT_ZERO_HEX()) AS balance,
MAX(last_block) AS last_block
FROM combined
GROUP BY chain_id, raindex_address, owner, token, vault_id
)
INSERT OR REPLACE INTO running_vault_balances (
chain_id,
Expand All @@ -118,17 +79,6 @@ SELECT
a.vault_id,
a.balance,
a.last_block,
(
SELECT c.last_log_index
FROM combined c
WHERE c.chain_id = a.chain_id
AND c.raindex_address = a.raindex_address
AND c.owner = a.owner
AND c.token = a.token
AND c.vault_id = a.vault_id
AND c.last_block = a.last_block
ORDER BY c.last_log_index DESC
LIMIT 1
) AS last_log_index,
a.last_log_index,
(CAST(strftime('%s', 'now') AS INTEGER) * 1000) AS updated_at
FROM aggregated a;
Loading
Loading