Skip to main content

sol_parser_sdk/grpc/
subscribe_builder.rs

1//! Yellowstone [`SubscribeRequest`] 构造(DEX 订阅、钱包 mentions 转账监控等共用)。
2
3use std::collections::HashMap;
4
5use yellowstone_grpc_proto::prelude::{
6    CommitmentLevel, SubscribeRequest, SubscribeRequestFilterAccounts,
7    SubscribeRequestFilterBlocksMeta, SubscribeRequestFilterTransactions,
8};
9
10use super::types::{AccountFilter, EventTypeFilter, TransactionFilter};
11
12#[inline]
13fn tx_filter_to_proto(f: &TransactionFilter) -> SubscribeRequestFilterTransactions {
14    SubscribeRequestFilterTransactions {
15        vote: Some(false),
16        failed: Some(false),
17        signature: None,
18        account_include: f.account_include.clone(),
19        account_exclude: f.account_exclude.clone(),
20        account_required: f.account_required.clone(),
21        ..Default::default()
22    }
23}
24
25#[inline]
26fn acc_filter_to_proto(f: &AccountFilter) -> SubscribeRequestFilterAccounts {
27    SubscribeRequestFilterAccounts {
28        account: f.account.clone(),
29        owner: f.owner.clone(),
30        filters: f.filters.clone(),
31        nonempty_txn_signature: None,
32        cuckoo_accounts_filter: None,
33    }
34}
35
36fn finalize(
37    transactions: HashMap<String, SubscribeRequestFilterTransactions>,
38    accounts: HashMap<String, SubscribeRequestFilterAccounts>,
39    blocks_meta: HashMap<String, SubscribeRequestFilterBlocksMeta>,
40    commitment: CommitmentLevel,
41) -> SubscribeRequest {
42    SubscribeRequest {
43        slots: HashMap::new(),
44        accounts,
45        transactions,
46        transactions_status: HashMap::new(),
47        blocks: HashMap::new(),
48        blocks_meta,
49        entry: HashMap::new(),
50        commitment: Some(commitment as i32),
51        accounts_data_slice: Vec::new(),
52        ping: None,
53        from_slot: None,
54    }
55}
56
57/// 构建订阅请求:`tx_0`…`tx_n`、`acc_0`…、**commitment = Processed**(与历史行为一致)。
58pub fn build_subscribe_request(
59    tx_filters: &[TransactionFilter],
60    acc_filters: &[AccountFilter],
61) -> SubscribeRequest {
62    build_subscribe_request_with_commitment(tx_filters, acc_filters, CommitmentLevel::Processed)
63}
64
65/// 与 [`build_subscribe_request`] 相同,可指定 commitment(例如 Confirmed)。
66pub fn build_subscribe_request_with_commitment(
67    tx_filters: &[TransactionFilter],
68    acc_filters: &[AccountFilter],
69    commitment: CommitmentLevel,
70) -> SubscribeRequest {
71    build_subscribe_request_with_event_filter(tx_filters, acc_filters, None, commitment)
72}
73
74/// 与 [`build_subscribe_request_with_commitment`] 相同,但会按事件过滤器订阅区块元数据。
75pub fn build_subscribe_request_with_event_filter(
76    tx_filters: &[TransactionFilter],
77    acc_filters: &[AccountFilter],
78    event_type_filter: Option<&EventTypeFilter>,
79    commitment: CommitmentLevel,
80) -> SubscribeRequest {
81    let transactions = tx_filters
82        .iter()
83        .enumerate()
84        .map(|(i, f)| (format!("tx_{}", i), tx_filter_to_proto(f)))
85        .collect();
86    let accounts = acc_filters
87        .iter()
88        .enumerate()
89        .map(|(i, f)| (format!("acc_{}", i), acc_filter_to_proto(f)))
90        .collect();
91    let blocks_meta = if event_type_filter.map(|f| f.includes_block_meta()).unwrap_or(false) {
92        HashMap::from([("block_meta".to_string(), SubscribeRequestFilterBlocksMeta {})])
93    } else {
94        HashMap::new()
95    };
96    finalize(transactions, accounts, blocks_meta, commitment)
97}
98
99/// 自定义交易订阅在 `SubscribeRequest.transactions` 中的 key(便于日志区分多条订阅)。
100pub fn build_subscribe_transaction_filters_named<N: AsRef<str>>(
101    named_tx_filters: &[(N, TransactionFilter)],
102    acc_filters: &[AccountFilter],
103    commitment: CommitmentLevel,
104) -> SubscribeRequest {
105    let transactions = named_tx_filters
106        .iter()
107        .map(|(name, f)| (name.as_ref().to_string(), tx_filter_to_proto(f)))
108        .collect();
109    let accounts = acc_filters
110        .iter()
111        .enumerate()
112        .map(|(i, f)| (format!("acc_{}", i), acc_filter_to_proto(f)))
113        .collect();
114    finalize(transactions, accounts, HashMap::new(), commitment)
115}
116
117#[cfg(test)]
118mod tests {
119    use super::*;
120    use crate::grpc::types::EventType;
121
122    #[test]
123    fn event_filter_controls_block_meta_subscription() {
124        let without_block_meta = build_subscribe_request_with_event_filter(
125            &[],
126            &[],
127            Some(&EventTypeFilter::include_only(vec![EventType::PumpFunTrade])),
128            CommitmentLevel::Processed,
129        );
130        assert!(without_block_meta.blocks_meta.is_empty());
131
132        let with_block_meta = build_subscribe_request_with_event_filter(
133            &[],
134            &[],
135            Some(&EventTypeFilter::include_only(vec![EventType::BlockMeta])),
136            CommitmentLevel::Processed,
137        );
138        assert_eq!(with_block_meta.blocks_meta.len(), 1);
139    }
140}