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