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