Skip to main content

sol_parser_sdk/
rpc_parser.rs

1//! RPC Transaction Parser
2//!
3//! 提供独立的 RPC 交易解析功能,不依赖 gRPC streaming
4//! 可以用于测试验证和离线分析
5
6use crate::core::events::DexEvent;
7use crate::grpc::types::EventTypeFilter;
8use crate::instr::read_pubkey_fast;
9use crate::transaction_cost::{parse_yellowstone_transaction_cost, TransactionCost};
10use smallvec::SmallVec;
11use solana_client::rpc_client::RpcClient;
12use solana_client::rpc_config::{RpcTransactionConfig, UiTransactionEncoding};
13use solana_client::rpc_response::{
14    EncodedTransaction, UiInstruction, UiTransactionStatusMeta, UiTransactionTokenBalance,
15};
16use solana_sdk::pubkey::Pubkey;
17use solana_sdk::signature::Signature;
18use solana_transaction_status::{
19    option_serializer::OptionSerializer, EncodedConfirmedTransactionWithStatusMeta,
20};
21use std::collections::HashMap;
22use yellowstone_grpc_proto::prelude::{
23    CompiledInstruction, InnerInstruction, InnerInstructions, Message, MessageAddressTableLookup,
24    MessageHeader, TokenBalance, Transaction, TransactionError, TransactionStatusMeta,
25    UiTokenAmount,
26};
27
28/// Parse a transaction from RPC by signature
29///
30/// # Arguments
31/// * `rpc_client` - RPC client to fetch the transaction
32/// * `signature` - Transaction signature
33/// * `filter` - Optional event type filter
34///
35/// # Returns
36/// Vector of parsed DEX events
37///
38/// # Example
39/// ```no_run
40/// use solana_client::rpc_client::RpcClient;
41/// use solana_sdk::signature::Signature;
42/// use sol_parser_sdk::parse_transaction_from_rpc;
43/// use std::str::FromStr;
44///
45/// let client = RpcClient::new("https://api.mainnet-beta.solana.com".to_string());
46/// let sig = Signature::from_str("your-signature-here").unwrap();
47/// let events = parse_transaction_from_rpc(&client, &sig, None).unwrap();
48/// ```
49pub fn parse_transaction_from_rpc(
50    rpc_client: &RpcClient,
51    signature: &Signature,
52    filter: Option<&EventTypeFilter>,
53) -> Result<Vec<DexEvent>, ParseError> {
54    // Fetch transaction from RPC with V1 transaction support.
55    let config = RpcTransactionConfig {
56        encoding: Some(UiTransactionEncoding::Base64),
57        commitment: None,
58        max_supported_transaction_version: Some(1),
59    };
60
61    let rpc_tx = rpc_client.get_transaction_with_config(signature, config).map_err(|e| {
62        let msg = e.to_string();
63        if msg.contains("invalid type: null") && msg.contains("EncodedConfirmedTransactionWithStatusMeta") {
64            ParseError::RpcError(format!(
65                "Transaction not found (RPC returned null). Common causes: 1) Transaction is too old and pruned (use an archive RPC). 2) Wrong network or invalid signature. Try SOLANA_RPC_URL with an archive endpoint (e.g. Helius, QuickNode) or a more recent tx. Original: {}",
66                msg
67            ))
68        } else {
69            ParseError::RpcError(msg)
70        }
71    })?;
72
73    parse_rpc_transaction(&rpc_tx, filter)
74}
75
76/// Parse a RPC transaction structure
77///
78/// # Arguments
79/// * `rpc_tx` - RPC transaction to parse
80/// * `filter` - Optional event type filter
81///
82/// # Returns
83/// Vector of parsed DEX events
84///
85/// # Example
86/// ```no_run
87/// use sol_parser_sdk::parse_rpc_transaction;
88///
89/// // Assuming you have an rpc_tx from RPC
90/// // let events = parse_rpc_transaction(&rpc_tx, None).unwrap();
91/// ```
92pub fn parse_rpc_transaction(
93    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
94    filter: Option<&EventTypeFilter>,
95) -> Result<Vec<DexEvent>, ParseError> {
96    let (grpc_meta, grpc_tx) = convert_rpc_to_grpc_for_parsing(rpc_tx)?;
97    let signature = extract_grpc_signature(&grpc_tx)?;
98    parse_converted_rpc_transaction(rpc_tx, grpc_meta, grpc_tx, signature, filter)
99}
100
101/// Result of parsing RPC events and transaction costs from one shared decode.
102#[derive(Debug)]
103pub struct ParsedRpcTransaction {
104    pub events: Vec<DexEvent>,
105    pub cost: TransactionCost,
106    pub signature: Signature,
107}
108
109/// Parses events and transaction costs while decoding the RPC payload once.
110pub fn parse_rpc_transaction_with_cost(
111    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
112    filter: Option<&EventTypeFilter>,
113) -> Result<ParsedRpcTransaction, ParseError> {
114    let (grpc_meta, grpc_tx) = convert_rpc_to_grpc_for_parsing(rpc_tx)?;
115    let signature = extract_grpc_signature(&grpc_tx)?;
116    let cost = parse_yellowstone_transaction_cost(&grpc_tx, &grpc_meta)
117        .ok_or_else(|| ParseError::MissingField("transaction.message".to_string()))?;
118    let events = parse_converted_rpc_transaction(rpc_tx, grpc_meta, grpc_tx, signature, filter)?;
119    Ok(ParsedRpcTransaction { events, cost, signature })
120}
121
122/// Parses only transaction cost and signature from one shared RPC decode.
123pub fn parse_rpc_transaction_cost_with_signature(
124    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
125) -> Result<(TransactionCost, Signature), ParseError> {
126    let (grpc_meta, grpc_tx) = convert_rpc_to_grpc_for_parsing(rpc_tx)?;
127    let signature = extract_grpc_signature(&grpc_tx)?;
128    let cost = parse_yellowstone_transaction_cost(&grpc_tx, &grpc_meta)
129        .ok_or_else(|| ParseError::MissingField("transaction.message".to_string()))?;
130    Ok((cost, signature))
131}
132
133fn extract_grpc_signature(transaction: &Transaction) -> Result<Signature, ParseError> {
134    transaction
135        .signatures
136        .first()
137        .ok_or_else(|| ParseError::MissingField("transaction.signatures[0]".to_string()))
138        .and_then(|bytes| {
139            Signature::try_from(bytes.as_slice()).map_err(|error| {
140                ParseError::ConversionError(format!("Invalid transaction signature: {error}"))
141            })
142        })
143}
144
145fn parse_converted_rpc_transaction(
146    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
147    grpc_meta: TransactionStatusMeta,
148    grpc_tx: Transaction,
149    signature: Signature,
150    filter: Option<&EventTypeFilter>,
151) -> Result<Vec<DexEvent>, ParseError> {
152    // Extract metadata
153    let slot = rpc_tx.slot;
154    let block_time_us = rpc_tx.block_time.map(|t| t * 1_000_000);
155    let grpc_recv_us = std::time::SystemTime::now()
156        .duration_since(std::time::UNIX_EPOCH)
157        .unwrap_or_default()
158        .as_micros() as i64;
159
160    // Wrap grpc_tx in Option for reuse
161    let grpc_tx_opt = Some(grpc_tx);
162
163    let mut program_invokes: HashMap<Pubkey, Vec<(i32, i32)>> = HashMap::new();
164
165    if let Some(ref tx) = grpc_tx_opt {
166        if let Some(ref msg) = tx.message {
167            let keys_len = msg.account_keys.len();
168            let writable_len = grpc_meta.loaded_writable_addresses.len();
169            let get_key = |i: usize| -> Option<&Vec<u8>> {
170                if i < keys_len {
171                    msg.account_keys.get(i)
172                } else if i < keys_len + writable_len {
173                    grpc_meta.loaded_writable_addresses.get(i - keys_len)
174                } else {
175                    grpc_meta.loaded_readonly_addresses.get(i - keys_len - writable_len)
176                }
177            };
178
179            for (i, ix) in msg.instructions.iter().enumerate() {
180                let pid = get_key(ix.program_id_index as usize)
181                    .map_or(Pubkey::default(), |k| read_pubkey_fast(k));
182                if crate::grpc::program_ids::needs_invoke_context(&pid) {
183                    program_invokes.entry(pid).or_default().push((i as i32, -1));
184                }
185            }
186
187            for inner in &grpc_meta.inner_instructions {
188                let outer_idx = inner.index as usize;
189                for (j, inner_ix) in inner.instructions.iter().enumerate() {
190                    let pid = get_key(inner_ix.program_id_index as usize)
191                        .map_or(Pubkey::default(), |k| read_pubkey_fast(k));
192                    if crate::grpc::program_ids::needs_invoke_context(&pid) {
193                        program_invokes.entry(pid).or_default().push((outer_idx as i32, j as i32));
194                    }
195                }
196            }
197        }
198    }
199
200    let needs_pumpfun = filter.map(EventTypeFilter::includes_pumpfun).unwrap_or(true);
201    let log_messages = rpc_log_messages(rpc_tx);
202    let is_created_buy =
203        needs_pumpfun && crate::logs::optimized_matcher::detect_pumpfun_create(log_messages);
204
205    // Parse instructions
206    let instr_events =
207        crate::grpc::instruction_parser::parse_instructions_enhanced_with_created_buy(
208            &grpc_meta,
209            &grpc_tx_opt,
210            signature,
211            slot,
212            0, // tx_idx
213            block_time_us,
214            grpc_recv_us,
215            filter,
216            is_created_buy,
217        );
218
219    // Parse logs (for protocols like PumpFun that emit events in logs)
220    struct ActiveProgram<'a> {
221        encoded: &'a str,
222        pubkey: Pubkey,
223    }
224
225    let mut active_program_stack: SmallVec<[ActiveProgram<'_>; 8]> = SmallVec::new();
226    let mut log_events = Vec::new();
227
228    for log in log_messages {
229        if let Some((pid, depth)) = crate::logs::optimized_matcher::parse_invoke_info(log) {
230            let pk = crate::grpc::program_ids::known_program_id(pid).unwrap_or_default();
231            active_program_stack.truncate(depth - 1);
232            active_program_stack.push(ActiveProgram { encoded: pid, pubkey: pk });
233        }
234
235        if let Some(mut event) = crate::logs::parse_log_with_program_id(
236            log,
237            signature,
238            slot,
239            0, // tx_index
240            block_time_us,
241            grpc_recv_us,
242            filter,
243            is_created_buy,
244            None,
245            active_program_stack.last().map(|active| &active.pubkey),
246        ) {
247            // Fill account fields - use same function as gRPC parsing
248            crate::core::account_dispatcher::fill_accounts_with_owned_keys(
249                &mut event,
250                &grpc_meta,
251                &grpc_tx_opt,
252                &program_invokes,
253            );
254
255            // Fill additional data fields (e.g., PumpSwap is_pump_pool)
256            crate::core::common_filler::fill_data(
257                &mut event,
258                &grpc_meta,
259                &grpc_tx_opt,
260                &program_invokes,
261            );
262
263            log_events.push(event);
264        }
265
266        if let Some(pid) = crate::logs::optimized_matcher::parse_program_complete_info(log) {
267            if let Some(pos) = active_program_stack.iter().rposition(|active| active.encoded == pid)
268            {
269                active_program_stack.truncate(pos);
270            }
271        }
272    }
273
274    let mut events = merge_log_and_instruction_events(log_events, instr_events);
275    fill_rpc_event_metadata(&mut events, rpc_tx, &grpc_meta, &grpc_tx_opt);
276    Ok(events)
277}
278
279fn fill_rpc_event_metadata(
280    events: &mut [DexEvent],
281    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
282    meta: &TransactionStatusMeta,
283    transaction: &Option<Transaction>,
284) {
285    if let Some(rpc_meta) = rpc_tx.transaction.meta.as_ref() {
286        for event in events.iter_mut() {
287            fill_rpc_token_balances(event, rpc_meta, meta, transaction);
288        }
289    }
290    crate::grpc::transaction_meta::fill_recent_blockhash(events, transaction);
291}
292
293#[inline]
294fn fill_rpc_token_balances(
295    event: &mut DexEvent,
296    rpc_meta: &UiTransactionStatusMeta,
297    meta: &TransactionStatusMeta,
298    transaction: &Option<Transaction>,
299) {
300    let trade = match event {
301        DexEvent::PumpFunTrade(event)
302        | DexEvent::PumpFunBuy(event)
303        | DexEvent::PumpFunSell(event)
304        | DexEvent::PumpFunBuyExactSolIn(event) => event,
305        _ => return,
306    };
307
308    if let Some(user_index) = rpc_account_index(transaction, meta, &trade.user) {
309        trade.sol_balance = rpc_meta.post_balances.get(user_index).copied();
310    }
311
312    if trade.associated_user == Pubkey::default() {
313        return;
314    }
315
316    let matches_account = |balance: &UiTransactionTokenBalance| {
317        rpc_account_key(transaction, meta, balance.account_index as usize)
318            .is_some_and(|key| key.as_slice() == trade.associated_user.as_ref())
319    };
320
321    if let OptionSerializer::Some(balances) = &rpc_meta.post_token_balances {
322        if let Some(balance) = balances.iter().find(|balance| matches_account(balance)) {
323            trade.token_balance = balance.ui_token_amount.amount.parse().ok();
324            return;
325        }
326    }
327
328    if let OptionSerializer::Some(balances) = &rpc_meta.pre_token_balances {
329        if balances.iter().any(matches_account) {
330            trade.token_balance = Some(0);
331        }
332    }
333}
334
335#[inline]
336fn rpc_account_index(
337    transaction: &Option<Transaction>,
338    meta: &TransactionStatusMeta,
339    account: &Pubkey,
340) -> Option<usize> {
341    if *account == Pubkey::default() {
342        return None;
343    }
344
345    let message = transaction.as_ref()?.message.as_ref()?;
346    message
347        .account_keys
348        .iter()
349        .chain(&meta.loaded_writable_addresses)
350        .chain(&meta.loaded_readonly_addresses)
351        .position(|key| key.as_slice() == account.as_ref())
352}
353
354#[inline]
355fn rpc_account_key<'a>(
356    transaction: &'a Option<Transaction>,
357    meta: &'a TransactionStatusMeta,
358    index: usize,
359) -> Option<&'a Vec<u8>> {
360    let message = transaction.as_ref()?.message.as_ref()?;
361    let static_len = message.account_keys.len();
362    let writable_len = meta.loaded_writable_addresses.len();
363
364    if index < static_len {
365        message.account_keys.get(index)
366    } else if index < static_len + writable_len {
367        meta.loaded_writable_addresses.get(index - static_len)
368    } else {
369        meta.loaded_readonly_addresses.get(index - static_len - writable_len)
370    }
371}
372
373#[inline]
374fn rpc_log_messages(rpc_tx: &EncodedConfirmedTransactionWithStatusMeta) -> &[String] {
375    let Some(meta) = rpc_tx.transaction.meta.as_ref() else { return &[] };
376    match &meta.log_messages {
377        OptionSerializer::Some(messages) => messages,
378        _ => &[],
379    }
380}
381
382fn merge_log_and_instruction_events(
383    log_events: Vec<DexEvent>,
384    instr_events: Vec<DexEvent>,
385) -> Vec<DexEvent> {
386    crate::grpc::log_instr_dedup::dedupe_log_instruction_events(log_events, instr_events)
387}
388
389/// Parse error types
390#[derive(Debug)]
391pub enum ParseError {
392    RpcError(String),
393    ConversionError(String),
394    MissingField(String),
395}
396
397impl std::fmt::Display for ParseError {
398    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
399        match self {
400            ParseError::RpcError(msg) => write!(f, "RPC error: {}", msg),
401            ParseError::ConversionError(msg) => write!(f, "Conversion error: {}", msg),
402            ParseError::MissingField(msg) => write!(f, "Missing field: {}", msg),
403        }
404    }
405}
406
407impl std::error::Error for ParseError {}
408
409// ============================================================================
410// Internal conversion functions
411// ============================================================================
412
413pub fn convert_rpc_to_grpc(
414    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
415) -> Result<(TransactionStatusMeta, Transaction), ParseError> {
416    convert_rpc_to_grpc_impl(rpc_tx, true)
417}
418
419/// Converts only metadata that parsing and cost calculation must own; logs and balances stay
420/// borrowed from the original RPC response on these internal paths.
421#[inline]
422pub(crate) fn convert_rpc_to_grpc_for_parsing(
423    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
424) -> Result<(TransactionStatusMeta, Transaction), ParseError> {
425    convert_rpc_to_grpc_impl(rpc_tx, false)
426}
427
428fn convert_rpc_to_grpc_impl(
429    rpc_tx: &EncodedConfirmedTransactionWithStatusMeta,
430    include_borrowed_metadata: bool,
431) -> Result<(TransactionStatusMeta, Transaction), ParseError> {
432    let rpc_meta = rpc_tx
433        .transaction
434        .meta
435        .as_ref()
436        .ok_or_else(|| ParseError::MissingField("meta".to_string()))?;
437
438    let (loaded_writable_addresses, loaded_readonly_addresses) = rpc_meta
439        .loaded_addresses
440        .as_ref()
441        .map(|addresses| {
442            let writable = addresses
443                .writable
444                .iter()
445                .map(|address| parse_loaded_address(address))
446                .collect::<Result<Vec<_>, _>>()?;
447            let readonly = addresses
448                .readonly
449                .iter()
450                .map(|address| parse_loaded_address(address))
451                .collect::<Result<Vec<_>, _>>()?;
452            Ok((writable, readonly))
453        })
454        .transpose()?
455        .unwrap_or_default();
456
457    let err = rpc_meta
458        .err
459        .clone()
460        .map(|error| {
461            let error: solana_sdk::transaction::TransactionError = error.into();
462            wincode::serialize(&error).map(|err| TransactionError { err }).map_err(|error| {
463                ParseError::ConversionError(format!(
464                    "Failed to serialize transaction error: {error}"
465                ))
466            })
467        })
468        .transpose()?;
469
470    // Convert meta
471    let mut grpc_meta = TransactionStatusMeta {
472        err,
473        fee: rpc_meta.fee,
474        pre_balances: if include_borrowed_metadata {
475            rpc_meta.pre_balances.clone()
476        } else {
477            Vec::new()
478        },
479        post_balances: if include_borrowed_metadata {
480            rpc_meta.post_balances.clone()
481        } else {
482            Vec::new()
483        },
484        inner_instructions: Vec::new(),
485        log_messages: if include_borrowed_metadata {
486            rpc_meta.log_messages.as_ref().map(|messages| messages.clone()).unwrap_or_default()
487        } else {
488            Vec::new()
489        },
490        pre_token_balances: if include_borrowed_metadata {
491            rpc_meta
492                .pre_token_balances
493                .as_ref()
494                .map(|balances| convert_token_balances(balances))
495                .unwrap_or_default()
496        } else {
497            Vec::new()
498        },
499        post_token_balances: if include_borrowed_metadata {
500            rpc_meta
501                .post_token_balances
502                .as_ref()
503                .map(|balances| convert_token_balances(balances))
504                .unwrap_or_default()
505        } else {
506            Vec::new()
507        },
508        rewards: Vec::new(),
509        loaded_writable_addresses,
510        loaded_readonly_addresses,
511        return_data: None,
512        compute_units_consumed: rpc_meta.compute_units_consumed.clone().into(),
513
514        inner_instructions_none: !rpc_meta.inner_instructions.is_some(),
515        log_messages_none: !rpc_meta.log_messages.is_some(),
516        return_data_none: !rpc_meta.return_data.is_some(),
517        cost_units: rpc_meta.cost_units.clone().into(),
518    };
519
520    // Convert inner instructions
521    if let solana_transaction_status::option_serializer::OptionSerializer::Some(
522        inner_instructions,
523    ) = rpc_meta.inner_instructions.as_ref()
524    {
525        for inner in inner_instructions {
526            let mut grpc_inner =
527                InnerInstructions { index: inner.index as u32, instructions: Vec::new() };
528
529            for ix in &inner.instructions {
530                if let UiInstruction::Compiled(compiled) = ix {
531                    // Decode base58 data
532                    let data = base58_turbo::BITCOIN.decode(&compiled.data).map_err(|e| {
533                        ParseError::ConversionError(format!(
534                            "Failed to decode instruction data: {}",
535                            e
536                        ))
537                    })?;
538
539                    grpc_inner.instructions.push(InnerInstruction {
540                        program_id_index: compiled.program_id_index as u32,
541                        accounts: compiled.accounts.clone(),
542                        data,
543                        stack_height: compiled.stack_height,
544                    });
545                }
546            }
547
548            grpc_meta.inner_instructions.push(grpc_inner);
549        }
550    }
551
552    // Convert transaction
553    let ui_tx = &rpc_tx.transaction.transaction;
554
555    let (message, signatures) = match ui_tx {
556        EncodedTransaction::Binary(_, _) | EncodedTransaction::LegacyBinary(_) => {
557            // Solana's decoder handles Base58/Base64, wincode deserialization and sanitization.
558            let versioned_tx = ui_tx.decode().ok_or_else(|| {
559                ParseError::ConversionError(
560                    "Failed to decode or sanitize binary transaction".to_string(),
561                )
562            })?;
563
564            let sigs: Vec<Vec<u8>> =
565                versioned_tx.signatures.iter().map(|s| s.as_ref().to_vec()).collect();
566
567            let message = match versioned_tx.message {
568                solana_sdk::message::VersionedMessage::Legacy(legacy_msg) => {
569                    convert_legacy_message(legacy_msg)?
570                }
571                solana_sdk::message::VersionedMessage::V0(v0_msg) => convert_v0_message(v0_msg)?,
572                solana_sdk::message::VersionedMessage::V1(v1_msg) => convert_v1_message(v1_msg)?,
573            };
574
575            (message, sigs)
576        }
577        EncodedTransaction::Json(_) => {
578            return Err(ParseError::ConversionError(
579                "JSON encoded transactions not supported yet".to_string(),
580            ));
581        }
582        _ => {
583            return Err(ParseError::ConversionError(
584                "Unsupported transaction encoding".to_string(),
585            ));
586        }
587    };
588
589    let grpc_tx = Transaction { signatures, message: Some(message) };
590
591    Ok((grpc_meta, grpc_tx))
592}
593
594fn parse_loaded_address(address: &str) -> Result<Vec<u8>, ParseError> {
595    address.parse::<Pubkey>().map(|pubkey| pubkey.to_bytes().to_vec()).map_err(|error| {
596        ParseError::ConversionError(format!("Invalid loaded address {address}: {error}"))
597    })
598}
599
600fn convert_token_balances(balances: &[UiTransactionTokenBalance]) -> Vec<TokenBalance> {
601    balances
602        .iter()
603        .map(|balance| TokenBalance {
604            account_index: balance.account_index as u32,
605            mint: balance.mint.clone(),
606            ui_token_amount: Some(UiTokenAmount {
607                ui_amount: balance.ui_token_amount.ui_amount.unwrap_or_default(),
608                decimals: balance.ui_token_amount.decimals as u32,
609                amount: balance.ui_token_amount.amount.clone(),
610                ui_amount_string: balance.ui_token_amount.ui_amount_string.clone(),
611            }),
612            owner: balance.owner.as_ref().map(|owner| owner.clone()).unwrap_or_default(),
613            program_id: balance
614                .program_id
615                .as_ref()
616                .map(|program_id| program_id.clone())
617                .unwrap_or_default(),
618        })
619        .collect()
620}
621
622fn convert_legacy_message(
623    msg: solana_sdk::message::legacy::Message,
624) -> Result<Message, ParseError> {
625    let account_keys: Vec<Vec<u8>> =
626        msg.account_keys.iter().map(|k| k.to_bytes().to_vec()).collect();
627
628    let instructions: Vec<CompiledInstruction> = msg
629        .instructions
630        .into_iter()
631        .map(|ix| CompiledInstruction {
632            program_id_index: ix.program_id_index as u32,
633            accounts: ix.accounts,
634            data: ix.data,
635        })
636        .collect();
637
638    Ok(Message {
639        header: Some(MessageHeader {
640            num_required_signatures: msg.header.num_required_signatures as u32,
641            num_readonly_signed_accounts: msg.header.num_readonly_signed_accounts as u32,
642            num_readonly_unsigned_accounts: msg.header.num_readonly_unsigned_accounts as u32,
643        }),
644        account_keys,
645        recent_blockhash: msg.recent_blockhash.to_bytes().to_vec(),
646        instructions,
647        versioned: false,
648        address_table_lookups: Vec::new(),
649        config: None,
650    })
651}
652
653fn convert_v0_message(msg: solana_sdk::message::v0::Message) -> Result<Message, ParseError> {
654    let account_keys: Vec<Vec<u8>> =
655        msg.account_keys.iter().map(|k| k.to_bytes().to_vec()).collect();
656
657    let instructions: Vec<CompiledInstruction> = msg
658        .instructions
659        .into_iter()
660        .map(|ix| CompiledInstruction {
661            program_id_index: ix.program_id_index as u32,
662            accounts: ix.accounts,
663            data: ix.data,
664        })
665        .collect();
666
667    Ok(Message {
668        header: Some(MessageHeader {
669            num_required_signatures: msg.header.num_required_signatures as u32,
670            num_readonly_signed_accounts: msg.header.num_readonly_signed_accounts as u32,
671            num_readonly_unsigned_accounts: msg.header.num_readonly_unsigned_accounts as u32,
672        }),
673        account_keys,
674        recent_blockhash: msg.recent_blockhash.to_bytes().to_vec(),
675        instructions,
676        versioned: true,
677        address_table_lookups: msg
678            .address_table_lookups
679            .into_iter()
680            .map(|lookup| MessageAddressTableLookup {
681                account_key: lookup.account_key.to_bytes().to_vec(),
682                writable_indexes: lookup.writable_indexes,
683                readonly_indexes: lookup.readonly_indexes,
684            })
685            .collect(),
686        config: None,
687    })
688}
689
690fn convert_v1_message(msg: solana_sdk::message::v1::Message) -> Result<Message, ParseError> {
691    let account_keys = msg.account_keys.iter().map(|key| key.to_bytes().to_vec()).collect();
692    let instructions = msg
693        .instructions
694        .into_iter()
695        .map(|ix| CompiledInstruction {
696            program_id_index: ix.program_id_index as u32,
697            accounts: ix.accounts,
698            data: ix.data,
699        })
700        .collect();
701
702    Ok(Message {
703        header: Some(MessageHeader {
704            num_required_signatures: msg.header.num_required_signatures as u32,
705            num_readonly_signed_accounts: msg.header.num_readonly_signed_accounts as u32,
706            num_readonly_unsigned_accounts: msg.header.num_readonly_unsigned_accounts as u32,
707        }),
708        account_keys,
709        recent_blockhash: msg.lifetime_specifier.to_bytes().to_vec(),
710        instructions,
711        versioned: true,
712        address_table_lookups: Vec::new(),
713        config: Some(yellowstone_grpc_proto::prelude::TransactionConfig {
714            priority_fee: msg.config.priority_fee,
715            compute_unit_limit: msg.config.compute_unit_limit,
716            loaded_accounts_data_size_limit: msg.config.loaded_accounts_data_size_limit,
717            heap_size: msg.config.heap_size,
718        }),
719    })
720}
721
722#[cfg(test)]
723mod tests {
724    use super::*;
725    use crate::core::events::{
726        DexEvent, EventMetadata, PumpFunTradeEvent, PumpSwapCreatePoolEvent,
727    };
728    use base64::{engine::general_purpose, Engine as _};
729    use solana_client::rpc_response::{
730        UiLoadedAddresses, UiTokenAmount as RpcUiTokenAmount, UiTransactionStatusMeta,
731        UiTransactionTokenBalance,
732    };
733    use solana_sdk::{
734        hash::Hash,
735        message::{legacy, MessageHeader, VersionedMessage},
736        pubkey::Pubkey,
737        signature::Signature,
738        transaction::VersionedTransaction,
739    };
740    use solana_transaction_status::{
741        option_serializer::OptionSerializer, EncodedTransactionWithStatusMeta,
742        TransactionBinaryEncoding,
743    };
744
745    fn rpc_fixture(
746        user: Pubkey,
747        token_account: Pubkey,
748    ) -> EncodedConfirmedTransactionWithStatusMeta {
749        let transaction = VersionedTransaction {
750            signatures: vec![Signature::from([7; 64])],
751            message: VersionedMessage::Legacy(legacy::Message {
752                header: MessageHeader { num_required_signatures: 1, ..Default::default() },
753                account_keys: vec![user, token_account],
754                recent_blockhash: Hash::new_unique(),
755                instructions: Vec::new(),
756            }),
757        };
758        let bytes = wincode::serialize(&transaction).expect("serialize RPC fixture");
759        let token_balance = |amount: &str| UiTransactionTokenBalance {
760            account_index: 1,
761            mint: Pubkey::new_unique().to_string(),
762            ui_token_amount: RpcUiTokenAmount {
763                ui_amount: None,
764                decimals: 6,
765                amount: amount.to_string(),
766                ui_amount_string: amount.to_string(),
767            },
768            owner: OptionSerializer::Some(user.to_string()),
769            program_id: OptionSerializer::None,
770        };
771
772        EncodedConfirmedTransactionWithStatusMeta {
773            slot: 42,
774            transaction: EncodedTransactionWithStatusMeta {
775                transaction: EncodedTransaction::Binary(
776                    general_purpose::STANDARD.encode(bytes),
777                    TransactionBinaryEncoding::Base64,
778                ),
779                meta: Some(UiTransactionStatusMeta {
780                    err: None,
781                    status: Ok(()),
782                    fee: 5_000,
783                    pre_balances: vec![50_000, 2_039_280],
784                    post_balances: vec![40_000, 2_039_280],
785                    inner_instructions: OptionSerializer::None,
786                    log_messages: OptionSerializer::None,
787                    pre_token_balances: OptionSerializer::Some(vec![token_balance("10")]),
788                    post_token_balances: OptionSerializer::Some(vec![token_balance("35")]),
789                    rewards: OptionSerializer::None,
790                    loaded_addresses: OptionSerializer::None,
791                    return_data: OptionSerializer::None,
792                    compute_units_consumed: OptionSerializer::Some(123),
793                    cost_units: OptionSerializer::Some(456),
794                }),
795                version: None,
796            },
797            block_time: None,
798            transaction_index: None,
799        }
800    }
801
802    fn dummy_meta() -> EventMetadata {
803        EventMetadata {
804            signature: Signature::default(),
805            slot: 1,
806            tx_index: 0,
807            block_time_us: 0,
808            grpc_recv_us: 0,
809            recent_blockhash: None,
810        }
811    }
812
813    #[test]
814    fn rpc_merge_keeps_instruction_cashback_for_log_only_pumpswap_create_pool() {
815        let pool = Pubkey::new_unique();
816        let base_mint = Pubkey::new_unique();
817        let quote_mint = Pubkey::new_unique();
818
819        let log_create = PumpSwapCreatePoolEvent {
820            metadata: dummy_meta(),
821            pool,
822            base_mint,
823            quote_mint,
824            is_cashback_coin: false,
825            ..Default::default()
826        };
827        let ix_create = PumpSwapCreatePoolEvent {
828            metadata: dummy_meta(),
829            pool,
830            base_mint,
831            quote_mint,
832            is_cashback_coin: true,
833            ..Default::default()
834        };
835
836        let merged = merge_log_and_instruction_events(
837            vec![DexEvent::PumpSwapCreatePool(log_create)],
838            vec![DexEvent::PumpSwapCreatePool(ix_create)],
839        );
840
841        assert_eq!(merged.len(), 1);
842        match &merged[0] {
843            DexEvent::PumpSwapCreatePool(e) => assert!(e.is_cashback_coin),
844            other => panic!("expected PumpSwapCreatePool, got {other:?}"),
845        }
846    }
847
848    #[test]
849    fn optimized_rpc_parsing_skips_borrowed_metadata_clones_but_fills_pumpfun_trade() {
850        let user = Pubkey::new_unique();
851        let token_account = Pubkey::new_unique();
852        let rpc_tx = rpc_fixture(user, token_account);
853        let (public_meta, _) = convert_rpc_to_grpc(&rpc_tx).expect("convert RPC fixture");
854
855        assert_eq!(
856            public_meta.pre_token_balances[0].ui_token_amount.as_ref().unwrap().amount,
857            "10"
858        );
859        assert_eq!(
860            public_meta.post_token_balances[0].ui_token_amount.as_ref().unwrap().amount,
861            "35"
862        );
863        assert_eq!(public_meta.pre_balances, [50_000, 2_039_280]);
864        assert_eq!(public_meta.post_balances, [40_000, 2_039_280]);
865        let (meta, transaction) =
866            convert_rpc_to_grpc_for_parsing(&rpc_tx).expect("convert RPC fixture for parsing");
867        assert!(meta.pre_balances.is_empty());
868        assert!(meta.post_balances.is_empty());
869        assert!(meta.log_messages.is_empty());
870        assert!(meta.pre_token_balances.is_empty());
871        assert!(meta.post_token_balances.is_empty());
872        assert_eq!(meta.compute_units_consumed, Some(123));
873        assert_eq!(meta.cost_units, Some(456));
874
875        let mut events = vec![DexEvent::PumpFunTrade(PumpFunTradeEvent {
876            user,
877            associated_user: token_account,
878            ..Default::default()
879        })];
880        fill_rpc_event_metadata(&mut events, &rpc_tx, &meta, &Some(transaction));
881
882        let DexEvent::PumpFunTrade(trade) = &events[0] else {
883            panic!("expected PumpFun trade");
884        };
885        assert_eq!(trade.token_balance, Some(35));
886        assert_eq!(trade.sol_balance, Some(40_000));
887    }
888
889    #[test]
890    fn optimized_rpc_balance_fill_handles_closed_token_accounts() {
891        let user = Pubkey::new_unique();
892        let token_account = Pubkey::new_unique();
893        let mut rpc_tx = rpc_fixture(user, token_account);
894        rpc_tx.transaction.meta.as_mut().unwrap().post_token_balances = OptionSerializer::None;
895        let (meta, transaction) =
896            convert_rpc_to_grpc_for_parsing(&rpc_tx).expect("convert RPC fixture for parsing");
897        let mut events = vec![DexEvent::PumpFunSell(PumpFunTradeEvent {
898            user,
899            associated_user: token_account,
900            ..Default::default()
901        })];
902
903        fill_rpc_event_metadata(&mut events, &rpc_tx, &meta, &Some(transaction));
904
905        let DexEvent::PumpFunSell(trade) = &events[0] else {
906            panic!("expected PumpFun sell");
907        };
908        assert_eq!(trade.token_balance, Some(0));
909        assert_eq!(trade.sol_balance, Some(40_000));
910    }
911
912    #[test]
913    fn optimized_rpc_balance_fill_keeps_malformed_amount_unknown() {
914        let user = Pubkey::new_unique();
915        let token_account = Pubkey::new_unique();
916        let mut rpc_tx = rpc_fixture(user, token_account);
917        let rpc_meta = rpc_tx.transaction.meta.as_mut().unwrap();
918        let OptionSerializer::Some(balances) = &mut rpc_meta.post_token_balances else {
919            panic!("post token balances");
920        };
921        balances[0].ui_token_amount.amount = "invalid".to_string();
922        let (meta, transaction) =
923            convert_rpc_to_grpc_for_parsing(&rpc_tx).expect("convert RPC fixture for parsing");
924        let mut events = vec![DexEvent::PumpFunBuy(PumpFunTradeEvent {
925            user,
926            associated_user: token_account,
927            ..Default::default()
928        })];
929
930        fill_rpc_event_metadata(&mut events, &rpc_tx, &meta, &Some(transaction));
931
932        let DexEvent::PumpFunBuy(trade) = &events[0] else {
933            panic!("expected PumpFun buy");
934        };
935        assert_eq!(trade.token_balance, None);
936        assert_eq!(trade.sol_balance, Some(40_000));
937    }
938
939    #[test]
940    fn invalid_rpc_loaded_address_returns_error_instead_of_panicking() {
941        let mut rpc_tx = rpc_fixture(Pubkey::new_unique(), Pubkey::new_unique());
942        rpc_tx.transaction.meta.as_mut().unwrap().loaded_addresses =
943            OptionSerializer::Some(UiLoadedAddresses {
944                writable: vec!["not-a-pubkey".to_string()],
945                readonly: Vec::new(),
946            });
947
948        let error = convert_rpc_to_grpc(&rpc_tx).expect_err("invalid address must fail");
949        assert!(
950            matches!(error, ParseError::ConversionError(message) if message.contains("Invalid loaded address"))
951        );
952    }
953
954    #[test]
955    fn base58_rpc_transaction_is_supported() {
956        let mut rpc_tx = rpc_fixture(Pubkey::new_unique(), Pubkey::new_unique());
957        let EncodedTransaction::Binary(data, TransactionBinaryEncoding::Base64) =
958            &rpc_tx.transaction.transaction
959        else {
960            panic!("expected base64 fixture");
961        };
962        let bytes = general_purpose::STANDARD.decode(data).expect("decode fixture");
963        rpc_tx.transaction.transaction = EncodedTransaction::Binary(
964            bs58::encode(bytes).into_string(),
965            TransactionBinaryEncoding::Base58,
966        );
967
968        convert_rpc_to_grpc(&rpc_tx).expect("base58 binary transaction must decode");
969    }
970
971    #[test]
972    fn unsanitized_rpc_transaction_is_rejected() {
973        let mut rpc_tx = rpc_fixture(Pubkey::new_unique(), Pubkey::new_unique());
974        let invalid = VersionedTransaction {
975            signatures: vec![Signature::from([7; 64])],
976            message: VersionedMessage::Legacy(legacy::Message {
977                header: MessageHeader { num_required_signatures: 1, ..Default::default() },
978                account_keys: vec![Pubkey::new_unique()],
979                recent_blockhash: Hash::new_unique(),
980                instructions: vec![
981                    solana_sdk::message::compiled_instruction::CompiledInstruction {
982                        program_id_index: 9,
983                        accounts: Vec::new(),
984                        data: Vec::new(),
985                    },
986                ],
987            }),
988        };
989        let bytes = wincode::serialize(&invalid).expect("serialize invalid fixture");
990        rpc_tx.transaction.transaction = EncodedTransaction::Binary(
991            general_purpose::STANDARD.encode(bytes),
992            TransactionBinaryEncoding::Base64,
993        );
994
995        let error = convert_rpc_to_grpc(&rpc_tx).expect_err("unsanitized transaction must fail");
996        assert!(matches!(
997            error,
998            ParseError::ConversionError(message) if message.contains("decode or sanitize")
999        ));
1000    }
1001}