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