diff --git a/src/main.rs b/src/main.rs index dd110cd..1f6fabc 100755 --- a/src/main.rs +++ b/src/main.rs @@ -188,7 +188,13 @@ async fn test_shreds() -> Result<(), Box> { fn create_event_callback() -> impl Fn(Box) { |event: Box| { - println!("🎉 Event received! Type: {:?}, ID: {}", event.event_type(), event.id()); + println!( + "🎉 Event received! Type: {:?}, ID: {}, Slot: {}, Transaction Index: {:?}", + event.event_type(), + event.id(), + event.slot(), + event.transaction_index(), + ); match_event!(event, { // -------------------------- block meta ----------------------- BlockMetaEvent => |e: BlockMetaEvent| { diff --git a/src/streaming/common/event_processor.rs b/src/streaming/common/event_processor.rs index e69d699..7bbacd2 100644 --- a/src/streaming/common/event_processor.rs +++ b/src/streaming/common/event_processor.rs @@ -113,7 +113,7 @@ impl EventProcessor { let signature = transaction_pretty.signature; // 使用缓存获取解析器 let parser = self.get_parser(); - let all_events = parser + let mut all_events = parser .parse_transaction( transaction_pretty.tx.clone(), signature, @@ -125,6 +125,11 @@ impl EventProcessor { .await .unwrap_or_else(|_e| vec![]); + // 为所有事件设置交易索引 + for event in &mut all_events { + event.set_transaction_index(transaction_pretty.transaction_index); + } + let max_time_consuming_us = all_events .iter() .map(|event| event.program_handle_time_consuming_us()) diff --git a/src/streaming/event_parser/common/mod.rs b/src/streaming/event_parser/common/mod.rs index d421559..3562a92 100755 --- a/src/streaming/event_parser/common/mod.rs +++ b/src/streaming/event_parser/common/mod.rs @@ -69,6 +69,14 @@ macro_rules! impl_unified_event { fn index(&self) -> String { self.metadata.index.clone() } + + fn transaction_index(&self) -> Option { + self.metadata.transaction_index + } + + fn set_transaction_index(&mut self, transaction_index: Option) { + self.metadata.set_transaction_index(transaction_index); + } } }; } diff --git a/src/streaming/event_parser/common/types.rs b/src/streaming/event_parser/common/types.rs index 37e75e4..e1ab15e 100755 --- a/src/streaming/event_parser/common/types.rs +++ b/src/streaming/event_parser/common/types.rs @@ -337,6 +337,7 @@ pub struct EventMetadata { pub id: String, pub signature: String, pub slot: u64, + pub transaction_index: Option, // 新增:交易在slot中的索引 pub block_time: i64, pub block_time_ms: i64, pub program_received_time_us: i64, @@ -347,7 +348,7 @@ pub struct EventMetadata { #[deprecated(note = "Please use swap_data instead")] pub transfer_datas: Vec, pub swap_data: Option, - pub index: String, + pub index: String, // 保留原有的指令索引 } impl EventMetadata { @@ -368,6 +369,7 @@ impl EventMetadata { id, signature, slot, + transaction_index: None, // 默认为None,后续设置 block_time, block_time_ms, program_received_time_us, @@ -389,6 +391,11 @@ impl EventMetadata { self.swap_data = Some(swap_data); } + /// 设置交易索引 + pub fn set_transaction_index(&mut self, transaction_index: Option) { + self.transaction_index = transaction_index; + } + /// Recycle EventMetadata to object pool pub fn recycle(self) { EVENT_METADATA_POOL.release(self); diff --git a/src/streaming/event_parser/core/traits.rs b/src/streaming/event_parser/core/traits.rs index d29e52e..67d6af8 100755 --- a/src/streaming/event_parser/core/traits.rs +++ b/src/streaming/event_parser/core/traits.rs @@ -64,6 +64,12 @@ pub trait UnifiedEvent: Debug + Send + Sync { /// Get index fn index(&self) -> String; + + /// Get transaction index in slot + fn transaction_index(&self) -> Option; + + /// Set transaction index in slot + fn set_transaction_index(&mut self, transaction_index: Option); } /// 事件解析器trait - 定义了事件解析的核心方法 diff --git a/src/streaming/grpc/types.rs b/src/streaming/grpc/types.rs index b0ef05f..9238f9f 100644 --- a/src/streaming/grpc/types.rs +++ b/src/streaming/grpc/types.rs @@ -69,6 +69,7 @@ impl fmt::Debug for BlockMetaPretty { #[derive(Clone)] pub struct TransactionPretty { pub slot: u64, + pub transaction_index: Option, // 新增:交易在slot中的索引 pub block_hash: String, pub block_time: Option, pub signature: Signature, @@ -81,6 +82,7 @@ impl fmt::Debug for TransactionPretty { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("TransactionPretty") .field("slot", &self.slot) + .field("transaction_index", &self.transaction_index) .field("signature", &self.signature) .field("is_vote", &self.is_vote) .field("program_received_time_us", &self.program_received_time_us) @@ -132,8 +134,11 @@ impl From<(SubscribeUpdateTransaction, Option)> for TransactionPretty ), ) -> Self { let tx = transaction.expect("should be defined"); + // 根据用户说明,交易索引在 transaction.index 中 + let transaction_index = tx.index; Self { slot, + transaction_index: Some(transaction_index), // 提取交易索引 block_time, block_hash: "".to_string(), signature: Signature::try_from(tx.signature.as_slice()).expect("valid signature"),