use std::collections::HashSet;
use std::sync::Arc;
use crate::{instr::read_pubkey_fast, DexEvent};
use solana_sdk::pubkey::Pubkey;
use solana_sdk::signature::Signature;
use yellowstone_grpc_proto::prelude::{TokenBalance, Transaction, TransactionStatusMeta};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum YellowstoneMessageVersion {
Legacy,
V0,
V1,
}
#[inline]
pub fn yellowstone_message_version(
message: &yellowstone_grpc_proto::prelude::Message,
) -> YellowstoneMessageVersion {
if message.config.is_some() {
YellowstoneMessageVersion::V1
} else if message.versioned {
YellowstoneMessageVersion::V0
} else {
YellowstoneMessageVersion::Legacy
}
}
#[inline]
pub(crate) fn fill_recent_blockhash(events: &mut [DexEvent], transaction: &Option<Transaction>) {
let Some(blockhash) = transaction
.as_ref()
.and_then(|tx| tx.message.as_ref())
.map(|message| message.recent_blockhash.as_slice())
.filter(|blockhash| !blockhash.is_empty())
else {
return;
};
let mut metadata = events.iter_mut().filter_map(DexEvent::metadata_mut);
let Some(first) = metadata.next() else { return };
let encoded = bs58::encode(blockhash).into_string();
for item in metadata {
item.recent_blockhash = Some(encoded.clone());
}
first.recent_blockhash = Some(encoded);
}
#[inline]
pub fn pubkey_bytes_to_bs58(bytes: &[u8]) -> Option<String> {
let a: [u8; 32] = bytes.try_into().ok()?;
Some(solana_sdk::pubkey::Pubkey::from(a).to_string())
}
pub fn collect_account_keys_bs58(
tx: &Transaction,
meta: &TransactionStatusMeta,
) -> Option<Vec<String>> {
let msg = tx.message.as_ref()?;
let mut keys: Vec<String> =
msg.account_keys.iter().filter_map(|b| pubkey_bytes_to_bs58(b.as_slice())).collect();
for b in &meta.loaded_writable_addresses {
keys.push(pubkey_bytes_to_bs58(b)?);
}
for b in &meta.loaded_readonly_addresses {
keys.push(pubkey_bytes_to_bs58(b)?);
}
Some(keys)
}
#[inline]
pub fn lamport_balance_deltas(meta: &TransactionStatusMeta) -> Vec<i128> {
meta.pre_balances
.iter()
.zip(meta.post_balances.iter())
.map(|(pre, post)| *post as i128 - *pre as i128)
.collect()
}
pub fn heuristic_sol_counterparties_for_watched_keys(
account_keys_bs58: &[String],
lamport_deltas: &[i128],
watched_bs58: &HashSet<&str>,
min_outflow_lamports: u64,
) -> Vec<(String, String)> {
let min_l = min_outflow_lamports as i128;
let mut pairs = Vec::new();
for (i, key) in account_keys_bs58.iter().enumerate() {
if !watched_bs58.contains(key.as_str()) {
continue;
}
let d = lamport_deltas.get(i).copied().unwrap_or(0);
if d >= -min_l {
continue;
}
for (j, dj) in lamport_deltas.iter().enumerate() {
if i == j || *dj <= min_l / 2 {
continue;
}
pairs.push((key.clone(), account_keys_bs58[j].clone()));
}
}
pairs
}
pub fn collect_watch_transfer_counterparty_pairs(
tx: &Transaction,
meta: &TransactionStatusMeta,
watched_bs58: &[String],
min_native_outflow_lamports: u64,
spl_min_watch_decrease_raw: u64,
) -> Option<Vec<(String, String)>> {
let keys = collect_account_keys_bs58(tx, meta)?;
let n = keys.len();
if meta.pre_balances.len() != n || meta.post_balances.len() != n {
return None;
}
let deltas = lamport_balance_deltas(meta);
let watched_h: HashSet<&str> = watched_bs58.iter().map(|s| s.as_str()).collect();
let mut pairs = heuristic_sol_counterparties_for_watched_keys(
&keys,
&deltas,
&watched_h,
min_native_outflow_lamports,
);
for w in watched_bs58 {
pairs.extend(spl_token_counterparty_by_owner(meta, w, spl_min_watch_decrease_raw));
}
pairs.sort_by(|a, b| a.1.cmp(&b.1));
pairs.dedup_by(|a, b| a.0 == b.0 && a.1 == b.1);
Some(pairs)
}
#[inline]
pub fn token_balance_raw_amount(t: &TokenBalance) -> u64 {
try_token_balance_raw_amount(t).unwrap_or(0)
}
#[inline]
pub fn try_token_balance_raw_amount(t: &TokenBalance) -> Option<u64> {
t.ui_token_amount.as_ref()?.amount.parse().ok()
}
pub fn spl_token_counterparty_by_owner(
meta: &TransactionStatusMeta,
watch_owner_bs58: &str,
min_watch_decrease_raw: u64,
) -> Vec<(String, String)> {
use std::collections::{HashMap, HashSet};
let pre = meta.pre_token_balances.as_slice();
let post = meta.post_token_balances.as_slice();
let mut pre_m: HashMap<(String, String), u64> = HashMap::new();
for b in pre {
if b.owner.is_empty() {
continue;
}
let k = (b.mint.clone(), b.owner.clone());
*pre_m.entry(k).or_insert(0) += token_balance_raw_amount(b);
}
let mut post_m: HashMap<(String, String), u64> = HashMap::new();
for b in post {
if b.owner.is_empty() {
continue;
}
let k = (b.mint.clone(), b.owner.clone());
*post_m.entry(k).or_insert(0) += token_balance_raw_amount(b);
}
let mut mints = HashSet::new();
for (m, o) in pre_m.keys() {
if o == watch_owner_bs58 {
mints.insert(m.clone());
}
}
for (m, o) in post_m.keys() {
if o == watch_owner_bs58 {
mints.insert(m.clone());
}
}
let mut out = Vec::new();
let min_l = min_watch_decrease_raw;
for mint in mints {
let w_pre = pre_m.get(&(mint.clone(), watch_owner_bs58.to_string())).copied().unwrap_or(0);
let w_post =
post_m.get(&(mint.clone(), watch_owner_bs58.to_string())).copied().unwrap_or(0);
let lost = w_pre.saturating_sub(w_post);
if lost < min_l.max(1) {
continue;
}
for ((m, owner), po) in &post_m {
if m != &mint || owner == watch_owner_bs58 {
continue;
}
let pr = pre_m.get(&(mint.clone(), owner.clone())).copied().unwrap_or(0);
if *po > pr {
out.push((watch_owner_bs58.to_string(), owner.clone()));
}
}
}
out.sort_by(|a, b| a.1.cmp(&b.1));
out.dedup_by(|a, b| a.0 == b.0 && a.1 == b.1);
out
}
#[inline]
pub fn yellowstone_static_account_keys_arc(tx: &Option<Transaction>) -> Arc<[Pubkey]> {
let Some(t) = tx.as_ref() else {
return Arc::from(Vec::<Pubkey>::new().into_boxed_slice());
};
let Some(msg) = t.message.as_ref() else {
return Arc::from(Vec::<Pubkey>::new().into_boxed_slice());
};
let keys: Vec<Pubkey> =
msg.account_keys.iter().map(|bytes| read_pubkey_fast(bytes.as_slice())).collect();
Arc::from(keys.into_boxed_slice())
}
#[inline]
pub fn try_yellowstone_signature(sig: &[u8]) -> Option<Signature> {
if sig.len() != 64 {
return None;
}
let a: [u8; 64] = sig.try_into().ok()?;
Some(Signature::from(a))
}
#[cfg(test)]
mod tests {
use super::*;
use yellowstone_grpc_proto::prelude::{Message, TransactionConfig};
#[test]
fn message_version_uses_config_before_versioned_flag() {
let legacy = Message::default();
assert_eq!(yellowstone_message_version(&legacy), YellowstoneMessageVersion::Legacy);
let v0 = Message { versioned: true, ..Message::default() };
assert_eq!(yellowstone_message_version(&v0), YellowstoneMessageVersion::V0);
let v1 = Message {
versioned: true,
config: Some(TransactionConfig::default()),
..Message::default()
};
assert_eq!(yellowstone_message_version(&v1), YellowstoneMessageVersion::V1);
let v1_with_inconsistent_legacy_flag = Message {
versioned: false,
config: Some(TransactionConfig::default()),
..Message::default()
};
assert_eq!(
yellowstone_message_version(&v1_with_inconsistent_legacy_flag),
YellowstoneMessageVersion::V1
);
}
}