Skip to main content

sol_parser_sdk/grpc/
transaction_meta.rs

1//! Yellowstone [`Transaction`] / [`TransactionStatusMeta`] 通用工具。
2//!
3//! 不依赖 DEX 日志或指令解析,适用于:mentions 订阅后的 SOL/SPL 转账分析、审计、风控等。
4
5use 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/// Transaction message version represented by Yellowstone's normalized protobuf.
14///
15/// Yellowstone uses `Message.config` to identify V1. The older `versioned` flag
16/// alone cannot distinguish V0 from V1.
17#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
18pub enum YellowstoneMessageVersion {
19    Legacy,
20    V0,
21    V1,
22}
23
24/// Returns the transaction version encoded by a Yellowstone protobuf message.
25///
26/// `config` must be checked first because both V0 and V1 set `versioned`.
27#[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/// Fill the public recent blockhash field once final events are known.
41///
42/// Encoding after log/instruction deduplication avoids doing the same base58
43/// allocation independently on both parser branches.
44#[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/// 32 字节公钥 → base58 地址字符串。
65#[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
71/// 消息静态 `account_keys` + meta 中 `loaded_writable_addresses` / `loaded_readonly_addresses`,
72/// 顺序与 `pre_balances` / `post_balances` 对齐。
73pub 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/// 每个账户索引的 lamports 变化(post - pre)。
90#[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
99/// 启发式原生 SOL:对 `watched_bs58` 中出现的账户,若 lamports 净减少 ≥ `min_outflow_lamports`,
100/// 再与其它索引配对,要求对方 delta ≥ `min_outflow_lamports/2`(与常见 mentions 转账监控一致)。
101pub 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
127/// 汇总「监控地址」在一笔交易中的转出对手方(原生 SOL 启发式 + SPL token balance 启发式)。
128///
129/// 返回 `None` 当账户 key 与 balance 数组长度不一致。
130pub 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/// `TokenBalance.ui_token_amount.amount` 解析为原始整数;失败为 0。
160#[inline]
161pub fn token_balance_raw_amount(t: &TokenBalance) -> u64 {
162    try_token_balance_raw_amount(t).unwrap_or(0)
163}
164
165/// Parse `TokenBalance.ui_token_amount.amount` without treating malformed metadata as zero.
166#[inline]
167pub fn try_token_balance_raw_amount(t: &TokenBalance) -> Option<u64> {
168    t.ui_token_amount.as_ref()?.amount.parse().ok()
169}
170
171/// SPL:对给定 owner(TokenBalance.owner,base58),当其某 mint 上余额净减少 ≥ `min_watch_decrease_raw` 时,
172/// 找出同 mint 下余额增加的其它 owner,返回 `(watch_owner, counterparty_owner)`。
173///
174/// 用于启发式「谁转给谁」配对(非链上 Transfer 事件级精确解析)。
175pub 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/// 仅消息头里的静态 `account_keys`(与 ShredStream `VersionedTransaction::static_account_keys()` 语义对齐;
240/// 不含 ALT 加载地址)。
241#[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/// Yellowstone 交易签名原始字节(64)→ `solana_sdk::signature::Signature`。
255#[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}