pump_swap correction.

This commit is contained in:
2026-06-20 08:04:06 +02:00
parent 58f9b36969
commit b26785a456
20 changed files with 2611 additions and 1019 deletions

View File

@@ -1484,7 +1484,16 @@ fn parse_pump_fees_instruction_data(
let decoded_fields_result = decode_borsh_fields(&decoded[8..], spec.fields);
let decoded_fields = match decoded_fields_result {
Ok(decoded_fields) => decoded_fields,
Err(error) => return Err(error),
Err(error) => {
tracing::debug!(
decoder = "pump_fees",
discriminator_hex = bytes_to_hex(discriminator.as_slice()),
payload_size = decoded.len().saturating_sub(8),
error = %error,
"ignoring malformed pump_fees instruction payload"
);
return Ok(None);
},
};
return Ok(Some(PumpFeesInstructionData {
spec,
@@ -1544,7 +1553,16 @@ fn decode_pump_fees_anchor_event_from_base64(
let fields_result = decode_borsh_fields(&payload[8..], spec.fields);
let fields_json = match fields_result {
Ok(fields_json) => fields_json,
Err(error) => return Err(error),
Err(error) => {
tracing::debug!(
decoder = "pump_fees",
discriminator_hex = bytes_to_hex(discriminator.as_slice()),
payload_size = payload.len().saturating_sub(8),
error = %error,
"ignoring malformed pump_fees anchor event payload"
);
return Ok(None);
},
};
return Ok(Some(PumpFeesAnchorEventData {
spec,
@@ -2422,6 +2440,37 @@ mod tests {
);
}
#[test]
fn ignores_truncated_known_instruction_payload() {
let mut data = std::vec::Vec::new();
data.extend_from_slice(&super::PUMP_FEES_GET_FEES_DISCRIMINATOR);
let encoded = bs58::encode(data).into_string();
let data_json = serde_json::json!(encoded).to_string();
let decoded_result = super::parse_pump_fees_instruction_data(Some(&data_json));
match decoded_result {
Ok(None) => {},
Ok(Some(decoded)) => panic!("unexpected decoded truncated instruction: {}", decoded.spec.name),
Err(error) => panic!("truncated instruction should be ignored: {}", error),
}
}
#[test]
fn ignores_truncated_known_anchor_event_payload() {
let data = pump_fees_anchor_event_prefix(
&super::PUMP_FEES_SOCIAL_FEE_PDA_CLAIMED_DISCRIMINATOR,
);
let encoded = {
use base64::Engine as _;
base64::engine::general_purpose::STANDARD.encode(data)
};
let decoded_result = super::decode_pump_fees_anchor_event_from_base64(encoded.as_str());
match decoded_result {
Ok(None) => {},
Ok(Some(decoded)) => panic!("unexpected decoded truncated anchor event: {}", decoded.spec.name),
Err(error) => panic!("truncated anchor event should be ignored: {}", error),
}
}
fn pump_fees_anchor_event_prefix(event_discriminator: &[u8; 8]) -> std::vec::Vec<u8> {
let mut data = std::vec::Vec::new();
data.extend_from_slice(&super::PUMP_FEES_ANCHOR_SELF_CPI_LOG_DISCRIMINATOR);

View File

@@ -1865,10 +1865,31 @@ fn build_anchor_event_payload(
}
fn is_materializable_pump_swap_anchor_event(event_name: &str) -> bool {
if event_name == "claim_token_incentives_event" {
return true;
}
return false;
return matches!(
event_name,
"admin_set_coin_creator_event"
| "admin_update_token_incentives_event"
| "buy_event"
| "claim_cashback_event"
| "claim_token_incentives_event"
| "close_user_volume_accumulator_event"
| "collect_coin_creator_fee_event"
| "create_config_event"
| "create_pool_event"
| "deposit_event"
| "disable_event"
| "extend_account_event"
| "init_user_volume_accumulator_event"
| "migrate_pool_coin_creator_event"
| "reserved_fee_recipients_event"
| "sell_event"
| "set_bonding_curve_coin_creator_event"
| "set_metaplex_coin_creator_event"
| "sync_user_volume_accumulator_event"
| "update_admin_event"
| "update_fee_config_event"
| "withdraw_event"
);
}
fn add_anchor_event_aliases(
@@ -3622,7 +3643,7 @@ mod tests {
}
#[test]
fn pump_swap_sync_user_volume_accumulator_anchor_event_is_decoded_audit_only_when_present() {
fn pump_swap_sync_user_volume_accumulator_anchor_event_is_materializable_reward_when_present() {
let decoder = crate::PumpSwapDecoder::new();
let encoded = make_pump_swap_anchor_event_base64_from_fields(
super::PUMP_SWAP_SYNC_USER_VOLUME_ACCUMULATOR_EVENT_DISCRIMINATOR,
@@ -3646,14 +3667,10 @@ mod tests {
found_anchor_event = true;
assert_eq!(
event.payload_json.get("anchorEventAuditOnly"),
Some(&serde_json::Value::Bool(true))
);
assert_eq!(
event.payload_json.get("skipCatalogReason"),
Some(&serde_json::Value::String(
"pump_swap_anchor_event_audit_only".to_string()
))
Some(&serde_json::Value::Bool(false))
);
assert_eq!(event.payload_json.get("skipCatalogReason"), None);
assert_eq!(event.payload_json.get("skipRewardReason"), None);
}
}
}
@@ -3746,7 +3763,7 @@ mod tests {
found_buy_anchor_event = true;
assert_eq!(
event.payload_json.get("anchorEventAuditOnly"),
Some(&serde_json::Value::Bool(true))
Some(&serde_json::Value::Bool(false))
);
}
},

View File

@@ -440,6 +440,48 @@ fn is_meteora_dlmm_pool_lifecycle_event_kind(event_kind: &str) -> bool {
);
}
fn is_pump_swap_anchor_swap_log_event_kind(event_kind: &str) -> bool {
return matches!(event_kind, "pump_swap.buy_event" | "pump_swap.sell_event");
}
fn is_pump_swap_admin_event_kind(event_kind: &str) -> bool {
if !event_kind.starts_with("pump_swap.") {
return false;
}
if is_pump_swap_anchor_swap_log_event_kind(event_kind) {
return false;
}
return matches!(
event_kind,
"pump_swap.admin_set_coin_creator"
| "pump_swap.admin_set_coin_creator_event"
| "pump_swap.admin_update_token_incentives"
| "pump_swap.admin_update_token_incentives_event"
| "pump_swap.create_config"
| "pump_swap.create_config_event"
| "pump_swap.disable"
| "pump_swap.disable_event"
| "pump_swap.extend_account"
| "pump_swap.extend_account_event"
| "pump_swap.migrate_pool_coin_creator"
| "pump_swap.migrate_pool_coin_creator_event"
| "pump_swap.reserved_fee_recipients_event"
| "pump_swap.set_bonding_curve_coin_creator_event"
| "pump_swap.set_coin_creator"
| "pump_swap.set_metaplex_coin_creator_event"
| "pump_swap.set_reserved_fee_recipient"
| "pump_swap.set_reserved_fee_recipients"
| "pump_swap.toggle_cashback_enabled"
| "pump_swap.toggle_mayhem_mode"
| "pump_swap.update_admin"
| "pump_swap.update_admin_event"
| "pump_swap.update_buyback_config"
| "pump_swap.update_fee_config"
| "pump_swap.update_fee_config_event"
);
}
fn is_meteora_dlmm_admin_event_kind(event_kind: &str) -> bool {
if !event_kind.starts_with("meteora_dlmm.") {
return false;
@@ -807,6 +849,9 @@ pub fn is_dex_orderbook_event_kind(event_kind: &str) -> bool {
/// Returns true for pool, pair, launch, mint, burn or migration lifecycle events.
pub fn is_dex_pool_lifecycle_event_kind(event_kind: &str) -> bool {
if is_pump_swap_anchor_swap_log_event_kind(event_kind) {
return true;
}
if is_meteora_dlmm_pool_lifecycle_event_kind(event_kind) {
return true;
}
@@ -1053,6 +1098,9 @@ pub fn is_dex_admin_event_kind(event_kind: &str) -> bool {
{
return true;
}
if is_pump_swap_admin_event_kind(event_kind) {
return true;
}
if event_kind.starts_with("pump_fees.")
&& (event_kind.contains("authority")
|| event_kind.contains("admin")
@@ -1839,4 +1887,23 @@ mod tests {
crate::UPSTREAM_REGISTRY_INSTRUCTION_MATCH_EVENT_KIND
));
}
#[test]
fn classifies_pump_swap_anchor_swap_logs_as_lifecycle_not_trade() {
assert!(super::is_dex_pool_lifecycle_event_kind("pump_swap.buy_event"));
assert!(super::is_dex_pool_lifecycle_event_kind("pump_swap.sell_event"));
assert!(!super::is_dex_trade_event_kind("pump_swap.buy_event"));
assert!(!super::is_dex_trade_event_kind("pump_swap.sell_event"));
}
#[test]
fn classifies_pump_swap_admin_instructions_and_events_explicitly() {
assert!(super::is_dex_admin_event_kind("pump_swap.update_fee_config"));
assert!(super::is_dex_admin_event_kind("pump_swap.update_fee_config_event"));
assert!(super::is_dex_admin_event_kind("pump_swap.set_coin_creator"));
assert!(super::is_dex_admin_event_kind("pump_swap.extend_account_event"));
assert!(!super::is_dex_admin_event_kind("pump_swap.buy_event"));
}
}

View File

@@ -1920,10 +1920,10 @@ fn infer_pump_swap_expected_db_target(
if entry_name == "buy" || entry_name == "sell" || entry_name == "buy_exact_quote_in" {
return Some(crate::DexEventCoverageEntryDto::DB_TARGET_TRADE_EVENTS.to_string());
}
if entry_name.starts_with("observed_unknown_") {
return Some(crate::DexEventCoverageEntryDto::DB_TARGET_DECODED_EVENTS_ONLY.to_string());
if entry_name == "buy_event" || entry_name == "sell_event" {
return Some(crate::DexEventCoverageEntryDto::DB_TARGET_POOL_LIFECYCLE_EVENTS.to_string());
}
if entry_name.ends_with("_event") && entry_name != "claim_token_incentives_event" {
if entry_name.starts_with("observed_unknown_") {
return Some(crate::DexEventCoverageEntryDto::DB_TARGET_DECODED_EVENTS_ONLY.to_string());
}
if entry_name == "deposit"
@@ -1998,7 +1998,7 @@ fn infer_pump_swap_event_family(
return Some("swap".to_string());
}
if entry_name == "buy_event" || entry_name == "sell_event" {
return Some("swap_event_audit".to_string());
return Some("swap_log".to_string());
}
if entry_name == "deposit" || entry_name == "deposit_event" {
return Some("liquidity_add".to_string());
@@ -3962,4 +3962,41 @@ mod tests {
};
assert!(!rows.is_empty());
}
#[test]
fn pump_swap_anchor_events_have_materialized_targets() {
assert_eq!(
super::infer_pump_swap_expected_db_target("buy_event", crate::ENTRY_KIND_EVENT),
Some(crate::DexEventCoverageEntryDto::DB_TARGET_POOL_LIFECYCLE_EVENTS.to_string())
);
assert_eq!(
super::infer_pump_swap_expected_db_target("sell_event", crate::ENTRY_KIND_EVENT),
Some(crate::DexEventCoverageEntryDto::DB_TARGET_POOL_LIFECYCLE_EVENTS.to_string())
);
assert_eq!(
super::infer_pump_swap_expected_db_target(
"collect_coin_creator_fee_event",
crate::ENTRY_KIND_EVENT,
),
Some(crate::DexEventCoverageEntryDto::DB_TARGET_FEE_EVENTS.to_string())
);
assert_eq!(
super::infer_pump_swap_expected_db_target("update_fee_config_event", crate::ENTRY_KIND_EVENT),
Some(crate::DexEventCoverageEntryDto::DB_TARGET_POOL_ADMIN_EVENTS.to_string())
);
}
#[test]
fn pump_swap_anchor_swap_events_are_swap_logs() {
assert_eq!(
super::infer_pump_swap_event_family("buy_event", crate::ENTRY_KIND_EVENT),
Some("swap_log".to_string())
);
assert_eq!(
super::infer_pump_swap_event_family("sell_event", crate::ENTRY_KIND_EVENT),
Some("swap_log".to_string())
);
}
}

View File

@@ -99,6 +99,71 @@ fn should_attempt_meteora_dbc_explicit_skip_materialization(
return false;
}
fn should_attempt_pump_swap_explicit_skip_materialization(
decoded_event: &crate::DexDecodedEventDto,
_payload: &serde_json::Value,
) -> bool {
if decoded_event.protocol_name != "pump_swap" {
return false;
}
return is_pump_swap_known_non_trade_materialization_candidate(
decoded_event.event_kind.as_str(),
);
}
fn is_pump_swap_known_non_trade_materialization_candidate(event_kind: &str) -> bool {
return matches!(
event_kind,
"pump_swap.admin_set_coin_creator"
| "pump_swap.admin_set_coin_creator_event"
| "pump_swap.admin_update_token_incentives"
| "pump_swap.admin_update_token_incentives_event"
| "pump_swap.buy_event"
| "pump_swap.claim_cashback"
| "pump_swap.claim_cashback_event"
| "pump_swap.claim_token_incentives"
| "pump_swap.claim_token_incentives_event"
| "pump_swap.close_user_volume_accumulator"
| "pump_swap.close_user_volume_accumulator_event"
| "pump_swap.collect_coin_creator_fee"
| "pump_swap.collect_coin_creator_fee_event"
| "pump_swap.create_config"
| "pump_swap.create_config_event"
| "pump_swap.create_pool"
| "pump_swap.create_pool_event"
| "pump_swap.deposit"
| "pump_swap.deposit_event"
| "pump_swap.disable"
| "pump_swap.disable_event"
| "pump_swap.extend_account"
| "pump_swap.extend_account_event"
| "pump_swap.init_user_volume_accumulator"
| "pump_swap.init_user_volume_accumulator_event"
| "pump_swap.migrate_pool_coin_creator"
| "pump_swap.migrate_pool_coin_creator_event"
| "pump_swap.reserved_fee_recipients_event"
| "pump_swap.sell_event"
| "pump_swap.set_bonding_curve_coin_creator_event"
| "pump_swap.set_coin_creator"
| "pump_swap.set_metaplex_coin_creator_event"
| "pump_swap.set_reserved_fee_recipient"
| "pump_swap.set_reserved_fee_recipients"
| "pump_swap.sync_user_volume_accumulator"
| "pump_swap.sync_user_volume_accumulator_event"
| "pump_swap.toggle_cashback_enabled"
| "pump_swap.toggle_mayhem_mode"
| "pump_swap.transfer_creator_fees_to_pump"
| "pump_swap.transfer_creator_fees_to_pump_v2"
| "pump_swap.update_admin"
| "pump_swap.update_admin_event"
| "pump_swap.update_buyback_config"
| "pump_swap.update_fee_config"
| "pump_swap.update_fee_config_event"
| "pump_swap.withdraw"
| "pump_swap.withdraw_event"
);
}
fn is_meteora_dbc_instruction_fee_materialization_candidate(event_kind: &str) -> bool {
return matches!(
event_kind,
@@ -131,6 +196,37 @@ enum FeeAmountRecoveryPolicy {
InnerSplTransfer,
}
fn is_pump_swap_forced_pool_admin_materialization_candidate(event_kind: &str) -> bool {
return matches!(
event_kind,
"pump_swap.admin_set_coin_creator"
| "pump_swap.admin_set_coin_creator_event"
| "pump_swap.admin_update_token_incentives"
| "pump_swap.admin_update_token_incentives_event"
| "pump_swap.create_config"
| "pump_swap.create_config_event"
| "pump_swap.disable"
| "pump_swap.disable_event"
| "pump_swap.extend_account"
| "pump_swap.extend_account_event"
| "pump_swap.migrate_pool_coin_creator"
| "pump_swap.migrate_pool_coin_creator_event"
| "pump_swap.reserved_fee_recipients_event"
| "pump_swap.set_bonding_curve_coin_creator_event"
| "pump_swap.set_coin_creator"
| "pump_swap.set_metaplex_coin_creator_event"
| "pump_swap.set_reserved_fee_recipient"
| "pump_swap.set_reserved_fee_recipients"
| "pump_swap.toggle_cashback_enabled"
| "pump_swap.toggle_mayhem_mode"
| "pump_swap.update_admin"
| "pump_swap.update_admin_event"
| "pump_swap.update_buyback_config"
| "pump_swap.update_fee_config"
| "pump_swap.update_fee_config_event"
);
}
fn fee_amount_recovery_policy_for_event_kind(event_kind: &str) -> FeeAmountRecoveryPolicy {
return match event_kind {
"pump_fees.crank_donation_fee_pda"
@@ -271,7 +367,30 @@ impl NonTradeEventMaterializationService {
continue;
},
};
if is_anchor_event_audit_only(&payload) {
if decoded_event.protocol_name == "pump_swap"
&& is_pump_swap_forced_pool_admin_materialization_candidate(
decoded_event.event_kind.as_str(),
)
{
let materialized = self
.materialize_pool_admin_event(
&transaction,
transaction_id,
decoded_event,
&payload,
)
.await;
match materialized {
Ok(was_materialized) => {
if was_materialized {
result.pool_admin_event_count += 1;
}
},
Err(error) => return Err(error),
}
continue;
}
if is_anchor_event_audit_only(decoded_event, &payload) {
continue;
}
if should_skip_non_trade_event_due_to_explicit_reason(decoded_event, &payload)
@@ -279,6 +398,7 @@ impl NonTradeEventMaterializationService {
decoded_event,
&payload,
)
&& !should_attempt_pump_swap_explicit_skip_materialization(decoded_event, &payload)
{
tracing::debug!(
event_kind = %decoded_event.event_kind,
@@ -297,6 +417,29 @@ impl NonTradeEventMaterializationService {
);
continue;
}
if decoded_event.protocol_name == "pump_swap"
&& is_pump_swap_forced_pool_admin_materialization_candidate(
decoded_event.event_kind.as_str(),
)
{
let materialized = self
.materialize_pool_admin_event(
&transaction,
transaction_id,
decoded_event,
&payload,
)
.await;
match materialized {
Ok(was_materialized) => {
if was_materialized {
result.pool_admin_event_count += 1;
}
},
Err(error) => return Err(error),
}
continue;
}
if crate::is_dex_pool_lifecycle_event_kind(decoded_event.event_kind.as_str()) {
let cleanup_result =
self.delete_stale_pool_admin_event_for_lifecycle(decoded_event).await;
@@ -3481,10 +3624,21 @@ fn is_pump_fun_payload(payload: &serde_json::Value) -> bool {
return false;
}
fn is_anchor_event_audit_only(payload: &serde_json::Value) -> bool {
fn is_anchor_event_audit_only(
decoded_event: &crate::DexDecodedEventDto,
payload: &serde_json::Value,
) -> bool {
if is_pump_fun_payload(payload) {
return false;
}
if decoded_event.protocol_name == "pump_swap"
&& is_pump_swap_known_non_trade_materialization_candidate(decoded_event.event_kind.as_str())
{
return false;
}
if is_pump_swap_materializable_anchor_payload(payload) {
return false;
}
if let Some(object) = payload.as_object() {
let flag = object.get("anchorEventAuditOnly");
if let Some(flag) = flag {
@@ -3502,6 +3656,26 @@ fn is_anchor_event_audit_only(payload: &serde_json::Value) -> bool {
return false;
}
fn is_pump_swap_materializable_anchor_payload(payload: &serde_json::Value) -> bool {
if !is_pump_swap_payload(payload) {
return false;
}
let event_kind = extract_first_string(payload, &["eventKind", "event_kind"]);
let event_kind = match event_kind {
Some(event_kind) => event_kind,
None => return false,
};
return is_pump_swap_known_non_trade_materialization_candidate(event_kind.as_str());
}
fn is_pump_swap_payload(payload: &serde_json::Value) -> bool {
let decoder = extract_first_string(payload, &["decoder", "protocolName", "protocol_name"]);
match decoder.as_deref() {
Some("pump_swap") => return true,
_ => return false,
}
}
fn transaction_has_effective_error(transaction: &crate::ChainTransactionDto) -> bool {
let err_json = match transaction.err_json.as_ref() {
Some(err_json) => err_json.trim(),
@@ -3627,6 +3801,81 @@ fn extract_first_number_as_string(
return None;
}
#[cfg(test)]
mod pump_swap_cleanup_tests {
fn make_decoded_event(event_kind: &str) -> crate::DexDecodedEventDto {
return crate::DexDecodedEventDto::new(
1,
Some(2),
"pump_swap".to_string(),
crate::PUMP_SWAP_PROGRAM_ID.to_string(),
event_kind.to_string(),
None,
None,
None,
None,
None,
"{}".to_string(),
);
}
#[test]
fn pump_swap_known_non_trade_candidates_bypass_legacy_trade_skip_reason() {
let decoded_event = make_decoded_event("pump_swap.update_fee_config");
let payload = serde_json::json!({
"decoder": "pump_swap",
"eventKind": "pump_swap.update_fee_config",
"skipAdminReason": "legacy_pump_swap_admin_skip",
"skipTradeReason": "pump_swap_non_trade_or_incomplete_trade_instruction",
"skipCandleReason": "pump_swap_non_trade_or_incomplete_trade_instruction"
});
assert!(super::should_skip_non_trade_event_due_to_explicit_reason(
&decoded_event,
&payload,
));
assert!(super::should_attempt_pump_swap_explicit_skip_materialization(
&decoded_event,
&payload,
));
}
#[test]
fn pump_swap_forced_admin_candidates_include_fee_config_event() {
assert!(super::is_pump_swap_forced_pool_admin_materialization_candidate(
"pump_swap.update_fee_config_event"
));
assert!(super::is_pump_swap_forced_pool_admin_materialization_candidate(
"pump_swap.migrate_pool_coin_creator_event"
));
assert!(!super::is_pump_swap_forced_pool_admin_materialization_candidate(
"pump_swap.buy_event"
));
}
#[test]
fn pump_swap_materializable_anchor_payload_is_not_audit_only() {
let payload = serde_json::json!({
"decoder": "pump_swap",
"eventKind": "pump_swap.buy_event",
"anchorEventAuditOnly": true
});
let decoded_event = crate::DexDecodedEventDto::new(
1,
Some(1),
"pump_swap".to_string(),
crate::PUMP_SWAP_PROGRAM_ID.to_string(),
"pump_swap.buy_event".to_string(),
Some("pool".to_string()),
None,
None,
None,
None,
payload.to_string(),
);
assert!(!super::is_anchor_event_audit_only(&decoded_event, &payload));
}
}
#[cfg(test)]
mod tests {

View File

@@ -122,7 +122,27 @@ pub(crate) async fn load_trade_aggregation_decoded_event_context(
};
let pool = match pool_option {
Some(pool) => pool,
None => return Ok(None),
None => {
let materialized =
crate::trade_aggregation_context::try_materialize_pump_swap_trade_pool(
database,
decoded_event,
)
.await;
match materialized {
Ok(true) => {
let refreshed_result =
crate::query_pools_get_by_address(database, pool_address.as_str()).await;
match refreshed_result {
Ok(Some(pool)) => pool,
Ok(None) => return Ok(None),
Err(error) => return Err(error),
}
},
Ok(false) => return Ok(None),
Err(error) => return Err(error),
}
},
};
let pool_id = match pool.id {
Some(pool_id) => pool_id,
@@ -140,7 +160,26 @@ pub(crate) async fn load_trade_aggregation_decoded_event_context(
};
let pair = match pair_option {
Some(pair) => pair,
None => return Ok(None),
None => {
let materialized =
crate::trade_aggregation_context::try_materialize_pump_swap_trade_pool(
database,
decoded_event,
)
.await;
match materialized {
Ok(true) => {
let refreshed_pair_result = crate::query_pairs_get_by_pool_id(database, pool_id).await;
match refreshed_pair_result {
Ok(Some(pair)) => pair,
Ok(None) => return Ok(None),
Err(error) => return Err(error),
}
},
Ok(false) => return Ok(None),
Err(error) => return Err(error),
}
},
};
let pair_id = match pair.id {
Some(pair_id) => pair_id,
@@ -195,6 +234,107 @@ pub(crate) async fn load_trade_aggregation_decoded_event_context(
}));
}
async fn try_materialize_pump_swap_trade_pool(
database: &crate::Database,
decoded_event: &crate::DexDecodedEventDto,
) -> Result<bool, crate::Error> {
if decoded_event.protocol_name != "pump_swap" {
return Ok(false);
}
if !crate::is_dex_trade_event_kind(decoded_event.event_kind.as_str()) {
return Ok(false);
}
if decoded_event.pool_account.is_none()
|| decoded_event.token_a_mint.is_none()
|| decoded_event.token_b_mint.is_none()
{
return Ok(false);
}
let dex_result = crate::query_dexs_get_by_code(database, decoded_event.protocol_name.as_str()).await;
let dex = match dex_result {
Ok(Some(dex)) => dex,
Ok(None) => return Ok(false),
Err(error) => return Err(error),
};
let dex_id = match dex.id {
Some(dex_id) => dex_id,
None => return Ok(false),
};
let payload_result = serde_json::from_str::<serde_json::Value>(decoded_event.payload_json.as_str());
let payload = match payload_result {
Ok(payload) => payload,
Err(_) => return Ok(false),
};
let token_a_vault_address = crate::trade_aggregation_context::extract_payload_string(
&payload,
&["poolBaseTokenAccount", "pool_base_token_account", "baseVault", "base_vault"],
);
let token_b_vault_address = crate::trade_aggregation_context::extract_payload_string(
&payload,
&["poolQuoteTokenAccount", "pool_quote_token_account", "quoteVault", "quote_vault"],
);
let materialization_input_result =
crate::dex_pool_materialization::DexPoolMaterializationInput::from_decoded_event(
decoded_event,
dex_id,
crate::PoolKind::Amm,
crate::PoolStatus::Active,
crate::dex_pool_materialization::DexPoolTokenOrder::AlreadyBaseQuote,
token_a_vault_address,
token_b_vault_address,
None,
);
let materialization_input = match materialization_input_result {
Ok(materialization_input) => materialization_input,
Err(_) => return Ok(false),
};
let materialization_result =
crate::dex_pool_materialization::materialize_dex_pool(database, &materialization_input).await;
match materialization_result {
Ok(_) => return Ok(true),
Err(error) => return Err(error),
}
}
fn extract_payload_string(
payload: &serde_json::Value,
candidate_keys: &[&str],
) -> std::option::Option<std::string::String> {
if let Some(object) = payload.as_object() {
for candidate_key in candidate_keys {
if let Some(value) = object.get(*candidate_key) {
if let Some(text) = value.as_str() {
let trimmed = text.trim();
if !trimmed.is_empty() {
return Some(trimmed.to_string());
}
}
}
}
for nested_value in object.values() {
let nested = crate::trade_aggregation_context::extract_payload_string(
nested_value,
candidate_keys,
);
if nested.is_some() {
return nested;
}
}
}
if let Some(array) = payload.as_array() {
for nested_value in array {
let nested = crate::trade_aggregation_context::extract_payload_string(
nested_value,
candidate_keys,
);
if nested.is_some() {
return nested;
}
}
}
return None;
}
fn find_pool_token_vault_address_by_token_id(
pool_tokens: &[crate::PoolTokenDto],
token_id: i64,

View File

@@ -81,6 +81,21 @@ pub(crate) async fn resolve_trade_amounts(
return Err(error);
}
}
if input.decoded_event.event_kind.starts_with("pump_swap.")
&& (base_amount_raw.is_none() || quote_amount_raw.is_none())
{
let resolution_result =
crate::trade_amount_resolution::apply_pump_swap_db_instruction_transfer_amount_fallback(
input,
&mut base_amount_raw,
&mut quote_amount_raw,
&mut resolved_trade_side,
)
.await;
if let Err(error) = resolution_result {
return Err(error);
}
}
if input.decoded_event.event_kind.starts_with("pump_fun.")
&& (base_amount_raw.is_none()
|| quote_amount_raw.is_none()
@@ -622,9 +637,240 @@ async fn apply_pump_swap_amount_fallbacks(
);
}
}
if price_quote_per_base.is_none() && base_amount_raw.is_some() && quote_amount_raw.is_some() {
*price_quote_per_base =
crate::trade_metric_update::compute_price_quote_per_base_from_raw_amounts_with_decimals(
base_amount_raw.as_ref().map(std::string::String::as_str),
quote_amount_raw.as_ref().map(std::string::String::as_str),
input.base_token_decimals,
input.quote_token_decimals,
);
tracing::debug!(
event_kind = %input.decoded_event.event_kind,
pool_account = ?input.decoded_event.pool_account,
decoded_event_id = ?input.decoded_event.id,
base_amount_raw = ?base_amount_raw,
quote_amount_raw = ?quote_amount_raw,
price_quote_per_base = ?price_quote_per_base,
"pump_swap trade price computed from resolved raw amounts"
);
}
return Ok(());
}
async fn apply_pump_swap_db_instruction_transfer_amount_fallback(
input: &crate::trade_amount_resolution::TradeAmountResolutionInput<'_>,
base_amount_raw: &mut std::option::Option<std::string::String>,
quote_amount_raw: &mut std::option::Option<std::string::String>,
resolved_trade_side: &mut std::option::Option<crate::SwapTradeSide>,
) -> Result<(), crate::Error> {
if base_amount_raw.is_some() && quote_amount_raw.is_some() {
return Ok(());
}
let decoded_instruction_result = crate::trade_amount_resolution::load_decoded_instruction(
input.database,
input.decoded_event,
)
.await;
let decoded_instruction = match decoded_instruction_result {
Ok(Some(decoded_instruction)) => decoded_instruction,
Ok(None) => return Ok(()),
Err(error) => return Err(error),
};
let instructions_result = crate::query_chain_instructions_list_by_transaction_id(
input.database,
input.decoded_event.transaction_id,
)
.await;
let instructions = match instructions_result {
Ok(instructions) => instructions,
Err(error) => return Err(error),
};
let payload_user_base_token_account =
crate::trade_amount_resolution::extract_string_by_candidate_keys(
input.payload,
&["userBaseTokenAccount", "user_base_token_account"],
);
let payload_user_quote_token_account =
crate::trade_amount_resolution::extract_string_by_candidate_keys(
input.payload,
&["userQuoteTokenAccount", "user_quote_token_account"],
);
let payload_pool_base_token_account =
crate::trade_amount_resolution::extract_string_by_candidate_keys(
input.payload,
&["poolBaseTokenAccount", "pool_base_token_account", "baseVault", "base_vault"],
);
let payload_pool_quote_token_account =
crate::trade_amount_resolution::extract_string_by_candidate_keys(
input.payload,
&["poolQuoteTokenAccount", "pool_quote_token_account", "quoteVault", "quote_vault"],
);
let pool_base_token_account = match input.base_vault_address {
Some(base_vault_address) => Some(base_vault_address),
None => payload_pool_base_token_account.as_deref(),
};
let pool_quote_token_account = match input.quote_vault_address {
Some(quote_vault_address) => Some(quote_vault_address),
None => payload_pool_quote_token_account.as_deref(),
};
let user_base_token_account = payload_user_base_token_account.as_deref();
let user_quote_token_account = payload_user_quote_token_account.as_deref();
let is_buy = input.decoded_event.event_kind.ends_with(".buy")
|| input.decoded_event.event_kind.ends_with(".buy_exact_quote_in");
let is_sell = input.decoded_event.event_kind.ends_with(".sell");
if !is_buy && !is_sell {
return Ok(());
}
let mut base_transfer_direction = None;
let mut quote_transfer_direction = None;
for instruction in &instructions {
if !crate::trade_amount_resolution::instruction_is_inside_same_pump_swap_instruction_window(
&decoded_instruction,
instruction,
) {
continue;
}
let parsed_transfer_result =
crate::trade_amount_resolution::parse_transfer_checked_instruction(instruction);
let parsed_transfer = match parsed_transfer_result {
Ok(Some(parsed_transfer)) => parsed_transfer,
Ok(None) => continue,
Err(error) => return Err(error),
};
if is_buy {
if base_amount_raw.is_none()
&& crate::trade_amount_resolution::string_option_equals(
input.base_token_mint,
parsed_transfer.mint.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
pool_base_token_account,
parsed_transfer.source.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
user_base_token_account,
parsed_transfer.destination.as_str(),
)
{
*base_amount_raw = Some(parsed_transfer.amount_raw.clone());
base_transfer_direction =
Some(crate::trade_amount_resolution::VaultTransferDirection::OutOfVault);
continue;
}
if quote_amount_raw.is_none()
&& crate::trade_amount_resolution::string_option_equals(
input.quote_token_mint,
parsed_transfer.mint.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
user_quote_token_account,
parsed_transfer.source.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
pool_quote_token_account,
parsed_transfer.destination.as_str(),
)
{
*quote_amount_raw = Some(parsed_transfer.amount_raw.clone());
quote_transfer_direction =
Some(crate::trade_amount_resolution::VaultTransferDirection::IntoVault);
continue;
}
}
if is_sell {
if base_amount_raw.is_none()
&& crate::trade_amount_resolution::string_option_equals(
input.base_token_mint,
parsed_transfer.mint.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
user_base_token_account,
parsed_transfer.source.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
pool_base_token_account,
parsed_transfer.destination.as_str(),
)
{
*base_amount_raw = Some(parsed_transfer.amount_raw.clone());
base_transfer_direction =
Some(crate::trade_amount_resolution::VaultTransferDirection::IntoVault);
continue;
}
if quote_amount_raw.is_none()
&& crate::trade_amount_resolution::string_option_equals(
input.quote_token_mint,
parsed_transfer.mint.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
pool_quote_token_account,
parsed_transfer.source.as_str(),
)
&& crate::trade_amount_resolution::account_option_equals(
user_quote_token_account,
parsed_transfer.destination.as_str(),
)
{
*quote_amount_raw = Some(parsed_transfer.amount_raw.clone());
quote_transfer_direction =
Some(crate::trade_amount_resolution::VaultTransferDirection::OutOfVault);
continue;
}
}
}
if resolved_trade_side.is_none() {
*resolved_trade_side =
crate::trade_amount_resolution::infer_trade_side_from_transfer_directions(
base_transfer_direction,
quote_transfer_direction,
);
}
if base_amount_raw.is_some() || quote_amount_raw.is_some() {
tracing::debug!(
event_kind = %input.decoded_event.event_kind,
decoded_event_id = ?input.decoded_event.id,
transaction_signature = %input.transaction.signature,
base_amount_raw = ?base_amount_raw,
quote_amount_raw = ?quote_amount_raw,
resolved_trade_side = ?resolved_trade_side,
"pump_swap trade amounts recovered from persisted instruction transfer window"
);
}
return Ok(());
}
fn instruction_is_inside_same_pump_swap_instruction_window(
decoded_instruction: &crate::ChainInstructionDto,
candidate_instruction: &crate::ChainInstructionDto,
) -> bool {
if candidate_instruction.transaction_id != decoded_instruction.transaction_id {
return false;
}
if candidate_instruction.instruction_index != decoded_instruction.instruction_index {
return false;
}
let candidate_inner_instruction_index = match candidate_instruction.inner_instruction_index {
Some(candidate_inner_instruction_index) => candidate_inner_instruction_index,
None => return false,
};
if let Some(decoded_inner_instruction_index) = decoded_instruction.inner_instruction_index {
if candidate_inner_instruction_index <= decoded_inner_instruction_index {
return false;
}
}
return true;
}
fn account_option_equals(left: std::option::Option<&str>, right: &str) -> bool {
let left = match left {
Some(left) => left,
None => return false,
};
return crate::trade_amount_resolution::account_equals(left, right);
}
fn apply_raydium_launchpad_amount_fallback(
input: &crate::trade_amount_resolution::TradeAmountResolutionInput<'_>,
base_amount_raw: &mut std::option::Option<std::string::String>,

View File

@@ -240,11 +240,7 @@ pub(crate) fn extract_trade_amounts_from_instruction_token_transfers(
Some(destination) => destination,
None => continue,
};
let amount_option =
crate::trade_solana_amounts::extract_scalar_as_string_by_candidate_keys(
info,
&["amount"],
);
let amount_option = crate::trade_solana_amounts::extract_spl_transfer_amount_raw(info);
let amount = match amount_option {
Some(amount) => amount,
None => continue,
@@ -418,6 +414,26 @@ pub(crate) fn compute_price_quote_per_base_with_decimals(
return inferred.2;
}
fn extract_spl_transfer_amount_raw(
info: &serde_json::Value,
) -> std::option::Option<std::string::String> {
let amount = crate::trade_solana_amounts::extract_scalar_as_string_by_candidate_keys(
info,
&["amount"],
);
if amount.is_some() {
return amount;
}
let token_amount = match info.get("tokenAmount") {
Some(token_amount) => token_amount,
None => return None,
};
return crate::trade_solana_amounts::extract_scalar_as_string_by_candidate_keys(
token_amount,
&["amount"],
);
}
fn is_spl_token_transfer_instruction(instruction: &serde_json::Value) -> bool {
let program_id_option = instruction.get("programId").and_then(|value| return value.as_str());
if let Some(program_id) = program_id_option {
@@ -854,6 +870,69 @@ mod tests {
assert_eq!(amounts.2, None);
}
#[test]
fn instruction_transfer_amounts_read_transfer_checked_token_amount() {
let meta_json = serde_json::json!({
"innerInstructions": [
{
"index": 7,
"instructions": [
{
"programId": crate::SPL_TOKEN_PROGRAM_ID,
"parsed": {
"type": "transferChecked",
"info": {
"source": "UserQuote111",
"destination": "QuoteVault111",
"mint": "QuoteMint111",
"tokenAmount": {
"amount": "199900249",
"decimals": 9,
"uiAmountString": "0.199900249"
}
}
}
},
{
"programId": crate::SPL_TOKEN_PROGRAM_ID,
"parsed": {
"type": "transferChecked",
"info": {
"source": "BaseVault111",
"destination": "UserBase111",
"mint": "BaseMint111",
"tokenAmount": {
"amount": "437138968699",
"decimals": 6,
"uiAmountString": "437138.968699"
}
}
}
}
]
}
]
});
let meta_json_text = meta_json.to_string();
let result = super::extract_trade_amounts_from_instruction_token_transfers(
Some(meta_json_text.as_str()),
Some(7),
Some("QuoteVault111"),
Some("BaseVault111"),
Some("UserQuote111"),
Some("UserBase111"),
Some("BaseVault111"),
Some("QuoteVault111"),
);
let amounts = match result {
Ok(amounts) => amounts,
Err(error) => panic!("transferChecked tokenAmount extraction should succeed: {}", error),
};
assert_eq!(amounts.0, Some("437138968699".to_string()));
assert_eq!(amounts.1, Some("199900249".to_string()));
assert_eq!(amounts.2, None);
}
#[test]
fn pump_fun_amounts_extract_token_delta_and_native_delta() {
let transaction_json = serde_json::json!({