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