0.7.24-pre.2
This commit is contained in:
@@ -105,6 +105,30 @@ impl KbDexDetectService {
|
||||
};
|
||||
detection_results.push(detect_result);
|
||||
}
|
||||
if decoded_event.protocol_name == "pump_fun"
|
||||
&& decoded_event.event_kind == "pump_fun.buy"
|
||||
{
|
||||
let detect_result = self
|
||||
.detect_pump_fun_trade(&transaction, decoded_event)
|
||||
.await;
|
||||
let detect_result = match detect_result {
|
||||
Ok(detect_result) => detect_result,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
detection_results.push(detect_result);
|
||||
}
|
||||
if decoded_event.protocol_name == "pump_fun"
|
||||
&& decoded_event.event_kind == "pump_fun.sell"
|
||||
{
|
||||
let detect_result = self
|
||||
.detect_pump_fun_trade(&transaction, decoded_event)
|
||||
.await;
|
||||
let detect_result = match detect_result {
|
||||
Ok(detect_result) => detect_result,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
detection_results.push(detect_result);
|
||||
}
|
||||
if decoded_event.protocol_name == "pump_swap"
|
||||
&& decoded_event.event_kind == "pump_swap.buy"
|
||||
{
|
||||
@@ -773,6 +797,243 @@ impl KbDexDetectService {
|
||||
})
|
||||
}
|
||||
|
||||
async fn detect_pump_fun_trade(
|
||||
&self,
|
||||
transaction: &crate::KbChainTransactionDto,
|
||||
decoded_event: &crate::KbDexDecodedEventDto,
|
||||
) -> Result<crate::KbDexPoolDetectionResult, crate::KbError> {
|
||||
let decoded_event_id_option = decoded_event.id;
|
||||
let decoded_event_id = match decoded_event_id_option {
|
||||
Some(decoded_event_id) => decoded_event_id,
|
||||
None => {
|
||||
return Err(crate::KbError::InvalidState(
|
||||
"decoded dex event has no internal id".to_string(),
|
||||
));
|
||||
}
|
||||
};
|
||||
let dex_id_result = self.ensure_pump_fun_dex().await;
|
||||
let dex_id = match dex_id_result {
|
||||
Ok(dex_id) => dex_id,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let pool_address_option = decoded_event.pool_account.clone();
|
||||
let pool_address = match pool_address_option {
|
||||
Some(pool_address) => pool_address,
|
||||
None => {
|
||||
return Err(crate::KbError::InvalidState(format!(
|
||||
"decoded event '{}' has no pool_account",
|
||||
decoded_event_id
|
||||
)));
|
||||
}
|
||||
};
|
||||
let token_a_mint_option = decoded_event.token_a_mint.clone();
|
||||
let token_a_mint = match token_a_mint_option {
|
||||
Some(token_a_mint) => token_a_mint,
|
||||
None => {
|
||||
return Err(crate::KbError::InvalidState(format!(
|
||||
"decoded event '{}' has no token_a_mint",
|
||||
decoded_event_id
|
||||
)));
|
||||
}
|
||||
};
|
||||
let token_b_mint_option = decoded_event.token_b_mint.clone();
|
||||
let token_b_mint = match token_b_mint_option {
|
||||
Some(token_b_mint) => token_b_mint,
|
||||
None => {
|
||||
return Err(crate::KbError::InvalidState(format!(
|
||||
"decoded event '{}' has no token_b_mint",
|
||||
decoded_event_id
|
||||
)));
|
||||
}
|
||||
};
|
||||
let base_is_token_a =
|
||||
kb_choose_base_quote_order(token_a_mint.as_str(), token_b_mint.as_str());
|
||||
let base_mint = if base_is_token_a {
|
||||
token_a_mint.clone()
|
||||
} else {
|
||||
token_b_mint.clone()
|
||||
};
|
||||
let quote_mint = if base_is_token_a {
|
||||
token_b_mint.clone()
|
||||
} else {
|
||||
token_a_mint.clone()
|
||||
};
|
||||
let base_token_id_result = self.ensure_token(base_mint.as_str()).await;
|
||||
let base_token_id = match base_token_id_result {
|
||||
Ok(base_token_id) => base_token_id,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let payload_value_result = kb_parse_payload_json(decoded_event.payload_json.as_str());
|
||||
let payload_value = match payload_value_result {
|
||||
Ok(payload_value) => payload_value,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let vault_addresses = kb_extract_pump_fun_vault_addresses(&payload_value);
|
||||
let token_a_vault_address = vault_addresses.0;
|
||||
let token_b_vault_address = vault_addresses.1;
|
||||
|
||||
let base_vault_address = if base_is_token_a {
|
||||
token_a_vault_address.clone()
|
||||
} else {
|
||||
token_b_vault_address.clone()
|
||||
};
|
||||
let quote_vault_address = if base_is_token_a {
|
||||
token_b_vault_address.clone()
|
||||
} else {
|
||||
token_a_vault_address.clone()
|
||||
};
|
||||
let quote_token_id_result = self.ensure_token(quote_mint.as_str()).await;
|
||||
let quote_token_id = match quote_token_id_result {
|
||||
Ok(quote_token_id) => quote_token_id,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let existing_pool_result =
|
||||
crate::get_pool_by_address(self.database.as_ref(), pool_address.as_str()).await;
|
||||
let existing_pool_option = match existing_pool_result {
|
||||
Ok(existing_pool_option) => existing_pool_option,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let created_pool = existing_pool_option.is_none();
|
||||
let pool_id = match existing_pool_option {
|
||||
Some(pool) => {
|
||||
let pool_id_option = pool.id;
|
||||
match pool_id_option {
|
||||
Some(pool_id) => pool_id,
|
||||
None => {
|
||||
return Err(crate::KbError::InvalidState(format!(
|
||||
"pool '{}' has no internal id",
|
||||
pool.address
|
||||
)));
|
||||
}
|
||||
}
|
||||
}
|
||||
None => {
|
||||
let pool_dto = crate::KbPoolDto::new(
|
||||
dex_id,
|
||||
pool_address.clone(),
|
||||
crate::KbPoolKind::BondingCurve,
|
||||
crate::KbPoolStatus::Active,
|
||||
);
|
||||
let upsert_result = crate::upsert_pool(self.database.as_ref(), &pool_dto).await;
|
||||
match upsert_result {
|
||||
Ok(pool_id) => pool_id,
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
};
|
||||
let existing_pair_result =
|
||||
crate::get_pair_by_pool_id(self.database.as_ref(), pool_id).await;
|
||||
let existing_pair_option = match existing_pair_result {
|
||||
Ok(existing_pair_option) => existing_pair_option,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let created_pair = existing_pair_option.is_none();
|
||||
let pair_symbol = kb_build_pair_symbol(base_mint.as_str(), quote_mint.as_str());
|
||||
let pair_dto =
|
||||
crate::KbPairDto::new(dex_id, pool_id, base_token_id, quote_token_id, pair_symbol);
|
||||
let pair_id_result = crate::upsert_pair(self.database.as_ref(), &pair_dto).await;
|
||||
let pair_id = match pair_id_result {
|
||||
Ok(pair_id) => pair_id,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let upsert_base_pool_token_result = crate::upsert_pool_token(
|
||||
self.database.as_ref(),
|
||||
&crate::KbPoolTokenDto::new(
|
||||
pool_id,
|
||||
base_token_id,
|
||||
crate::KbPoolTokenRole::Base,
|
||||
base_vault_address,
|
||||
Some(0),
|
||||
),
|
||||
)
|
||||
.await;
|
||||
if let Err(error) = upsert_base_pool_token_result {
|
||||
return Err(error);
|
||||
}
|
||||
let upsert_quote_pool_token_result = crate::upsert_pool_token(
|
||||
self.database.as_ref(),
|
||||
&crate::KbPoolTokenDto::new(
|
||||
pool_id,
|
||||
quote_token_id,
|
||||
crate::KbPoolTokenRole::Quote,
|
||||
quote_vault_address,
|
||||
Some(1),
|
||||
),
|
||||
)
|
||||
.await;
|
||||
if let Err(error) = upsert_quote_pool_token_result {
|
||||
return Err(error);
|
||||
}
|
||||
let existing_listing_result =
|
||||
crate::get_pool_listing_by_pool_id(self.database.as_ref(), pool_id).await;
|
||||
let existing_listing_option = match existing_listing_result {
|
||||
Ok(existing_listing_option) => existing_listing_option,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
let created_listing = existing_listing_option.is_none();
|
||||
let pool_listing_id = match existing_listing_option {
|
||||
Some(pool_listing) => pool_listing.id,
|
||||
None => {
|
||||
let listing_id_result = self
|
||||
.upsert_pool_listing_from_decoded_event(dex_id, pool_id, pair_id, transaction)
|
||||
.await;
|
||||
match listing_id_result {
|
||||
Ok(listing_id) => Some(listing_id),
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
};
|
||||
if created_pool {
|
||||
let signal_result = self
|
||||
.record_detection_signal(
|
||||
transaction,
|
||||
"signal.dex.pump_fun.new_pool",
|
||||
crate::KbAnalysisSignalSeverity::Low,
|
||||
payload_value.clone(),
|
||||
)
|
||||
.await;
|
||||
if let Err(error) = signal_result {
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
if created_pair {
|
||||
let signal_result = self
|
||||
.record_detection_signal(
|
||||
transaction,
|
||||
"signal.dex.pump_fun.new_pair",
|
||||
crate::KbAnalysisSignalSeverity::Low,
|
||||
payload_value.clone(),
|
||||
)
|
||||
.await;
|
||||
if let Err(error) = signal_result {
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
if created_listing {
|
||||
let signal_result = self
|
||||
.record_detection_signal(
|
||||
transaction,
|
||||
"signal.dex.pump_fun.first_listing_seen",
|
||||
crate::KbAnalysisSignalSeverity::Low,
|
||||
payload_value,
|
||||
)
|
||||
.await;
|
||||
if let Err(error) = signal_result {
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
Ok(crate::KbDexPoolDetectionResult {
|
||||
decoded_event_id,
|
||||
dex_id,
|
||||
pool_id,
|
||||
pair_id,
|
||||
pool_listing_id,
|
||||
created_pool,
|
||||
created_pair,
|
||||
created_listing,
|
||||
})
|
||||
}
|
||||
|
||||
async fn detect_pump_swap_trade(
|
||||
&self,
|
||||
transaction: &crate::KbChainTransactionDto,
|
||||
@@ -905,13 +1166,8 @@ impl KbDexDetectService {
|
||||
};
|
||||
let created_pair = existing_pair_option.is_none();
|
||||
let pair_symbol = kb_build_pair_symbol(base_mint.as_str(), quote_mint.as_str());
|
||||
let pair_dto = crate::KbPairDto::new(
|
||||
dex_id,
|
||||
pool_id,
|
||||
base_token_id,
|
||||
quote_token_id,
|
||||
pair_symbol,
|
||||
);
|
||||
let pair_dto =
|
||||
crate::KbPairDto::new(dex_id, pool_id, base_token_id, quote_token_id, pair_symbol);
|
||||
let pair_id_result = crate::upsert_pair(self.database.as_ref(), &pair_dto).await;
|
||||
let pair_id = match pair_id_result {
|
||||
Ok(pair_id) => pair_id,
|
||||
@@ -2862,6 +3118,27 @@ fn kb_extract_string_from_array_index(
|
||||
Some(text.to_string())
|
||||
}
|
||||
|
||||
fn kb_extract_pump_fun_vault_addresses(
|
||||
payload_value: &serde_json::Value,
|
||||
) -> (
|
||||
std::option::Option<std::string::String>,
|
||||
std::option::Option<std::string::String>,
|
||||
) {
|
||||
let accounts_option = payload_value.get("accounts");
|
||||
let accounts = match accounts_option {
|
||||
Some(accounts) => accounts,
|
||||
None => return (None, None),
|
||||
};
|
||||
let accounts_array_option = accounts.as_array();
|
||||
let accounts_array = match accounts_array_option {
|
||||
Some(accounts_array) => accounts_array,
|
||||
None => return (None, None),
|
||||
};
|
||||
let token_a_vault_address = kb_extract_string_from_array_index(accounts_array, 4);
|
||||
let token_b_native_address = kb_extract_string_from_array_index(accounts_array, 3);
|
||||
(token_a_vault_address, token_b_native_address)
|
||||
}
|
||||
|
||||
fn kb_extract_pump_swap_vault_addresses(
|
||||
payload_value: &serde_json::Value,
|
||||
) -> (
|
||||
|
||||
Reference in New Issue
Block a user