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