sol_parser_sdk/grpc/
transaction_meta.rs1use std::collections::HashSet;
6use std::sync::Arc;
7
8use crate::{instr::read_pubkey_fast, DexEvent};
9use solana_sdk::pubkey::Pubkey;
10use solana_sdk::signature::Signature;
11use yellowstone_grpc_proto::prelude::{TokenBalance, Transaction, TransactionStatusMeta};
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
18pub enum YellowstoneMessageVersion {
19 Legacy,
20 V0,
21 V1,
22}
23
24#[inline]
28pub fn yellowstone_message_version(
29 message: &yellowstone_grpc_proto::prelude::Message,
30) -> YellowstoneMessageVersion {
31 if message.config.is_some() {
32 YellowstoneMessageVersion::V1
33 } else if message.versioned {
34 YellowstoneMessageVersion::V0
35 } else {
36 YellowstoneMessageVersion::Legacy
37 }
38}
39
40#[inline]
45pub(crate) fn fill_recent_blockhash(events: &mut [DexEvent], transaction: &Option<Transaction>) {
46 let Some(blockhash) = transaction
47 .as_ref()
48 .and_then(|tx| tx.message.as_ref())
49 .map(|message| message.recent_blockhash.as_slice())
50 .filter(|blockhash| !blockhash.is_empty())
51 else {
52 return;
53 };
54
55 let mut metadata = events.iter_mut().filter_map(DexEvent::metadata_mut);
56 let Some(first) = metadata.next() else { return };
57 let encoded = bs58::encode(blockhash).into_string();
58 for item in metadata {
59 item.recent_blockhash = Some(encoded.clone());
60 }
61 first.recent_blockhash = Some(encoded);
62}
63
64#[inline]
66pub fn pubkey_bytes_to_bs58(bytes: &[u8]) -> Option<String> {
67 let a: [u8; 32] = bytes.try_into().ok()?;
68 Some(solana_sdk::pubkey::Pubkey::from(a).to_string())
69}
70
71pub fn collect_account_keys_bs58(
74 tx: &Transaction,
75 meta: &TransactionStatusMeta,
76) -> Option<Vec<String>> {
77 let msg = tx.message.as_ref()?;
78 let mut keys: Vec<String> =
79 msg.account_keys.iter().filter_map(|b| pubkey_bytes_to_bs58(b.as_slice())).collect();
80 for b in &meta.loaded_writable_addresses {
81 keys.push(pubkey_bytes_to_bs58(b)?);
82 }
83 for b in &meta.loaded_readonly_addresses {
84 keys.push(pubkey_bytes_to_bs58(b)?);
85 }
86 Some(keys)
87}
88
89#[inline]
91pub fn lamport_balance_deltas(meta: &TransactionStatusMeta) -> Vec<i128> {
92 meta.pre_balances
93 .iter()
94 .zip(meta.post_balances.iter())
95 .map(|(pre, post)| *post as i128 - *pre as i128)
96 .collect()
97}
98
99pub fn heuristic_sol_counterparties_for_watched_keys(
102 account_keys_bs58: &[String],
103 lamport_deltas: &[i128],
104 watched_bs58: &HashSet<&str>,
105 min_outflow_lamports: u64,
106) -> Vec<(String, String)> {
107 let min_l = min_outflow_lamports as i128;
108 let mut pairs = Vec::new();
109 for (i, key) in account_keys_bs58.iter().enumerate() {
110 if !watched_bs58.contains(key.as_str()) {
111 continue;
112 }
113 let d = lamport_deltas.get(i).copied().unwrap_or(0);
114 if d >= -min_l {
115 continue;
116 }
117 for (j, dj) in lamport_deltas.iter().enumerate() {
118 if i == j || *dj <= min_l / 2 {
119 continue;
120 }
121 pairs.push((key.clone(), account_keys_bs58[j].clone()));
122 }
123 }
124 pairs
125}
126
127pub fn collect_watch_transfer_counterparty_pairs(
131 tx: &Transaction,
132 meta: &TransactionStatusMeta,
133 watched_bs58: &[String],
134 min_native_outflow_lamports: u64,
135 spl_min_watch_decrease_raw: u64,
136) -> Option<Vec<(String, String)>> {
137 let keys = collect_account_keys_bs58(tx, meta)?;
138 let n = keys.len();
139 if meta.pre_balances.len() != n || meta.post_balances.len() != n {
140 return None;
141 }
142 let deltas = lamport_balance_deltas(meta);
143 let watched_h: HashSet<&str> = watched_bs58.iter().map(|s| s.as_str()).collect();
144
145 let mut pairs = heuristic_sol_counterparties_for_watched_keys(
146 &keys,
147 &deltas,
148 &watched_h,
149 min_native_outflow_lamports,
150 );
151 for w in watched_bs58 {
152 pairs.extend(spl_token_counterparty_by_owner(meta, w, spl_min_watch_decrease_raw));
153 }
154 pairs.sort_by(|a, b| a.1.cmp(&b.1));
155 pairs.dedup_by(|a, b| a.0 == b.0 && a.1 == b.1);
156 Some(pairs)
157}
158
159#[inline]
161pub fn token_balance_raw_amount(t: &TokenBalance) -> u64 {
162 try_token_balance_raw_amount(t).unwrap_or(0)
163}
164
165#[inline]
167pub fn try_token_balance_raw_amount(t: &TokenBalance) -> Option<u64> {
168 t.ui_token_amount.as_ref()?.amount.parse().ok()
169}
170
171pub fn spl_token_counterparty_by_owner(
176 meta: &TransactionStatusMeta,
177 watch_owner_bs58: &str,
178 min_watch_decrease_raw: u64,
179) -> Vec<(String, String)> {
180 use std::collections::{HashMap, HashSet};
181
182 let pre = meta.pre_token_balances.as_slice();
183 let post = meta.post_token_balances.as_slice();
184
185 let mut pre_m: HashMap<(String, String), u64> = HashMap::new();
186 for b in pre {
187 if b.owner.is_empty() {
188 continue;
189 }
190 let k = (b.mint.clone(), b.owner.clone());
191 *pre_m.entry(k).or_insert(0) += token_balance_raw_amount(b);
192 }
193 let mut post_m: HashMap<(String, String), u64> = HashMap::new();
194 for b in post {
195 if b.owner.is_empty() {
196 continue;
197 }
198 let k = (b.mint.clone(), b.owner.clone());
199 *post_m.entry(k).or_insert(0) += token_balance_raw_amount(b);
200 }
201
202 let mut mints = HashSet::new();
203 for (m, o) in pre_m.keys() {
204 if o == watch_owner_bs58 {
205 mints.insert(m.clone());
206 }
207 }
208 for (m, o) in post_m.keys() {
209 if o == watch_owner_bs58 {
210 mints.insert(m.clone());
211 }
212 }
213
214 let mut out = Vec::new();
215 let min_l = min_watch_decrease_raw;
216 for mint in mints {
217 let w_pre = pre_m.get(&(mint.clone(), watch_owner_bs58.to_string())).copied().unwrap_or(0);
218 let w_post =
219 post_m.get(&(mint.clone(), watch_owner_bs58.to_string())).copied().unwrap_or(0);
220 let lost = w_pre.saturating_sub(w_post);
221 if lost < min_l.max(1) {
222 continue;
223 }
224 for ((m, owner), po) in &post_m {
225 if m != &mint || owner == watch_owner_bs58 {
226 continue;
227 }
228 let pr = pre_m.get(&(mint.clone(), owner.clone())).copied().unwrap_or(0);
229 if *po > pr {
230 out.push((watch_owner_bs58.to_string(), owner.clone()));
231 }
232 }
233 }
234 out.sort_by(|a, b| a.1.cmp(&b.1));
235 out.dedup_by(|a, b| a.0 == b.0 && a.1 == b.1);
236 out
237}
238
239#[inline]
242pub fn yellowstone_static_account_keys_arc(tx: &Option<Transaction>) -> Arc<[Pubkey]> {
243 let Some(t) = tx.as_ref() else {
244 return Arc::from(Vec::<Pubkey>::new().into_boxed_slice());
245 };
246 let Some(msg) = t.message.as_ref() else {
247 return Arc::from(Vec::<Pubkey>::new().into_boxed_slice());
248 };
249 let keys: Vec<Pubkey> =
250 msg.account_keys.iter().map(|bytes| read_pubkey_fast(bytes.as_slice())).collect();
251 Arc::from(keys.into_boxed_slice())
252}
253
254#[inline]
256pub fn try_yellowstone_signature(sig: &[u8]) -> Option<Signature> {
257 if sig.len() != 64 {
258 return None;
259 }
260 let a: [u8; 64] = sig.try_into().ok()?;
261 Some(Signature::from(a))
262}
263
264#[cfg(test)]
265mod tests {
266 use super::*;
267 use yellowstone_grpc_proto::prelude::{Message, TransactionConfig};
268
269 #[test]
270 fn message_version_uses_config_before_versioned_flag() {
271 let legacy = Message::default();
272 assert_eq!(yellowstone_message_version(&legacy), YellowstoneMessageVersion::Legacy);
273
274 let v0 = Message { versioned: true, ..Message::default() };
275 assert_eq!(yellowstone_message_version(&v0), YellowstoneMessageVersion::V0);
276
277 let v1 = Message {
278 versioned: true,
279 config: Some(TransactionConfig::default()),
280 ..Message::default()
281 };
282 assert_eq!(yellowstone_message_version(&v1), YellowstoneMessageVersion::V1);
283
284 let v1_with_inconsistent_legacy_flag = Message {
285 versioned: false,
286 config: Some(TransactionConfig::default()),
287 ..Message::default()
288 };
289 assert_eq!(
290 yellowstone_message_version(&v1_with_inconsistent_legacy_flag),
291 YellowstoneMessageVersion::V1
292 );
293 }
294}