1use std::{
2 collections::{HashMap, HashSet},
3 sync::Arc,
4};
5
6use thiserror::Error;
7use tokio::sync::{mpsc::UnboundedReceiver, watch};
8use tycho_client::feed::{synchronizer::Snapshot, BlockHeader, FeedMessage};
9use tycho_common::{
10 models::{
11 blockchain::{Block, BlockAggregatedChanges, DCIUpdate, PendingBlock},
12 protocol::{ComponentBalance, ProtocolComponent, ProtocolComponentStateDelta},
13 Chain,
14 },
15 traits::TxDeltaIndexer,
16 Bytes,
17};
18
19use crate::{
20 evm::decoder::{StreamDecodeError, TychoStreamDecoder},
21 protocol::models::Update,
22};
23
24pub struct PendingUpdate {
30 pub label: String,
31 pub update: Update,
32}
33
34#[derive(Debug, Error)]
35pub enum PendingError {
36 #[error("parent block {needed} not yet confirmed (current: {current})")]
40 ParentNotYetConfirmed { needed: u64, current: u64 },
41 #[error("decoder error: {0}")]
42 Decoder(#[from] StreamDecodeError),
43 #[error("indexer error for extractor '{extractor}': {message}")]
44 Indexer { extractor: String, message: String },
45}
46
47pub struct PendingBlockProcessor {
82 indexers: HashMap<String, Box<dyn TxDeltaIndexer>>,
83 decoder: Arc<TychoStreamDecoder<BlockHeader>>,
84 chain: Chain,
85 current_confirmed_block: u64,
87 confirmed_block_tx: watch::Sender<u64>,
89 block_rx: UnboundedReceiver<FeedMessage<BlockHeader>>,
91}
92
93impl PendingBlockProcessor {
94 pub(crate) fn new(
95 indexers: HashMap<String, Box<dyn TxDeltaIndexer>>,
96 decoder: Arc<TychoStreamDecoder<BlockHeader>>,
97 chain: Chain,
98 block_rx: UnboundedReceiver<FeedMessage<BlockHeader>>,
99 ) -> Self {
100 let (confirmed_block_tx, _) = watch::channel(0u64);
101 Self { indexers, decoder, chain, current_confirmed_block: 0, confirmed_block_tx, block_rx }
102 }
103
104 pub fn subscribe_confirmed_block(&self) -> watch::Receiver<u64> {
110 self.confirmed_block_tx.subscribe()
111 }
112
113 pub fn current_confirmed_block(&self) -> u64 {
115 self.current_confirmed_block
116 }
117
118 pub fn advance(&mut self, msg: &FeedMessage<BlockHeader>) -> Result<(), PendingError> {
124 self.advance_inner(msg)
125 }
126
127 pub async fn generate_pending_update(
145 &mut self,
146 pending: &PendingBlock,
147 label: String,
148 ) -> Result<PendingUpdate, PendingError> {
149 while let Ok(msg) = self.block_rx.try_recv() {
151 self.advance_inner(&msg)?;
152 }
153
154 let target_block = pending.block();
155 let parent = target_block.number.saturating_sub(1);
156 if self.current_confirmed_block < parent {
157 return Err(PendingError::ParentNotYetConfirmed {
158 needed: parent,
159 current: self.current_confirmed_block,
160 });
161 }
162 let target_header = BlockHeader::from(target_block);
163
164 let mut pending_deltas: HashMap<String, BlockAggregatedChanges> = HashMap::new();
165 for (extractor, indexer) in &self.indexers {
166 let changes = indexer.generate_deltas(pending);
167 pending_deltas.insert(extractor.clone(), changes);
168 }
169
170 let update = self
171 .decoder
172 .apply_deltas_ephemeral(&pending_deltas, target_header)
173 .await?;
174 Ok(PendingUpdate { label, update })
175 }
176
177 fn advance_inner(&mut self, msg: &FeedMessage<BlockHeader>) -> Result<(), PendingError> {
178 let msg_block = msg
179 .state_msgs
180 .values()
181 .map(|s| s.header.number)
182 .max()
183 .unwrap_or(0);
184
185 for (extractor, state_msg) in &msg.state_msgs {
186 let Some(indexer) = self.indexers.get_mut(extractor) else {
187 continue;
188 };
189
190 if !state_msg.snapshots.states.is_empty() {
191 let block_changes = snapshot_to_block_changes(
192 extractor,
193 &state_msg.snapshots,
194 &state_msg.header,
195 self.chain,
196 );
197 indexer
198 .apply_block(&block_changes)
199 .map_err(|e| PendingError::Indexer {
200 extractor: extractor.clone(),
201 message: format!("{e:#}"),
202 })?;
203 }
204
205 if let Some(deltas) = &state_msg.deltas {
206 indexer
207 .apply_block(deltas)
208 .map_err(|e| PendingError::Indexer {
209 extractor: extractor.clone(),
210 message: format!("{e:#}"),
211 })?;
212 }
213 }
214
215 if msg_block > self.current_confirmed_block {
216 self.current_confirmed_block = msg_block;
217 let _ = self.confirmed_block_tx.send(msg_block);
219 }
220 Ok(())
221 }
222}
223
224fn snapshot_to_block_changes(
230 extractor: &str,
231 snapshot: &Snapshot,
232 header: &BlockHeader,
233 chain: Chain,
234) -> BlockAggregatedChanges {
235 let ts = chrono::DateTime::from_timestamp(header.timestamp as i64, 0)
236 .unwrap_or_default()
237 .naive_utc();
238 let block = Block {
239 number: header.number,
240 chain,
241 hash: header.hash.clone(),
242 parent_hash: header.parent_hash.clone(),
243 ts,
244 };
245
246 let mut new_protocol_components: HashMap<String, ProtocolComponent> = HashMap::new();
247 let mut state_deltas: HashMap<String, ProtocolComponentStateDelta> = HashMap::new();
248 let mut component_balances: HashMap<String, HashMap<Bytes, ComponentBalance>> = HashMap::new();
249 let mut dci_update = DCIUpdate::default();
250
251 for (id, comp_with_state) in &snapshot.states {
252 new_protocol_components.insert(id.clone(), comp_with_state.component.clone());
253
254 for (entrypoint, trace) in &comp_with_state.entrypoints {
255 let ep_id = entrypoint
256 .entry_point
257 .external_id
258 .clone();
259 dci_update
260 .new_entrypoints
261 .entry(id.clone())
262 .or_default()
263 .insert(entrypoint.entry_point.clone());
264 dci_update
265 .new_entrypoint_params
266 .entry(ep_id.clone())
267 .or_default()
268 .insert((entrypoint.params.clone(), id.clone()));
269 dci_update
270 .trace_results
271 .entry(ep_id)
272 .or_default()
273 .merge(trace.clone());
274 }
275
276 state_deltas.insert(
277 id.clone(),
278 ProtocolComponentStateDelta {
279 component_id: id.clone(),
280 updated_attributes: comp_with_state.state.attributes.clone(),
281 deleted_attributes: HashSet::new(),
282 created_attributes: HashSet::new(),
283 },
284 );
285
286 let token_balances: HashMap<Bytes, ComponentBalance> = comp_with_state
287 .state
288 .balances
289 .iter()
290 .map(|(token, balance)| {
291 (
292 token.clone(),
293 ComponentBalance {
294 token: token.clone(),
295 balance: balance.clone(),
296 balance_float: 0.0,
297 modify_tx: Bytes::default(),
298 component_id: id.clone(),
299 },
300 )
301 })
302 .collect();
303 component_balances.insert(id.clone(), token_balances);
304 }
305
306 BlockAggregatedChanges {
307 extractor: extractor.to_string(),
308 chain,
309 block,
310 finalized_block_height: header.number,
311 new_protocol_components,
312 state_deltas,
313 component_balances,
314 dci_update,
315 ..Default::default()
316 }
317}
318
319#[cfg(test)]
320mod tests {
321 use std::sync::Mutex;
322
323 use tycho_common::models::blockchain::{Block, PendingBlock};
324
325 use super::*;
326
327 struct RecordingIndexer {
329 seen: Arc<Mutex<Vec<Block>>>,
330 }
331
332 impl TxDeltaIndexer for RecordingIndexer {
333 fn apply_block(&mut self, _block: &BlockAggregatedChanges) -> anyhow::Result<()> {
334 Ok(())
335 }
336
337 fn generate_deltas(&self, pending: &PendingBlock) -> BlockAggregatedChanges {
338 self.seen
339 .lock()
340 .unwrap()
341 .push(pending.block().clone());
342 BlockAggregatedChanges::default()
343 }
344 }
345
346 fn target_block(number: u64, timestamp: i64) -> Block {
347 Block {
348 number,
349 chain: Chain::Ethereum,
350 hash: Bytes::from([1u8; 32]),
351 parent_hash: Bytes::from([2u8; 32]),
352 ts: chrono::DateTime::from_timestamp(timestamp, 0)
353 .unwrap()
354 .naive_utc(),
355 }
356 }
357
358 fn processor(seen: Arc<Mutex<Vec<Block>>>) -> PendingBlockProcessor {
359 let indexers: HashMap<String, Box<dyn TxDeltaIndexer>> =
360 HashMap::from([("fluid".to_string(), Box::new(RecordingIndexer { seen }) as _)]);
361 let (_tx, rx) = tokio::sync::mpsc::unbounded_channel();
362 PendingBlockProcessor::new(
363 indexers,
364 Arc::new(TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum)),
365 Chain::Ethereum,
366 rx,
367 )
368 }
369
370 #[test]
371 fn test_snapshot_entrypoints_land_in_dci_update() {
372 use tycho_client::feed::synchronizer::ComponentWithState;
373 use tycho_common::models::{
374 blockchain::{
375 EntryPoint, EntryPointWithTracingParams, RPCTracerParams, TracingParams,
376 TracingResult,
377 },
378 protocol::ProtocolComponentState,
379 };
380
381 let component_id = "0xpool".to_string();
382 let entry_point = EntryPoint {
383 external_id: "0xpool:get_virtual_price()".to_string(),
384 target: Bytes::from([0xaa; 20]),
385 signature: "get_virtual_price()".to_string(),
386 };
387 let params =
388 TracingParams::RPCTracer(RPCTracerParams::new(None, Bytes::from([0x12, 0x34])));
389 let oracle = Bytes::from([0xbb; 20]);
390 let trace = TracingResult::new(
391 HashSet::new(),
392 HashMap::from([(oracle.clone(), HashSet::from([Bytes::from([0u8; 32])]))]),
393 );
394 let snapshot = Snapshot {
395 states: HashMap::from([(
396 component_id.clone(),
397 ComponentWithState {
398 state: ProtocolComponentState::new(
399 &component_id,
400 HashMap::new(),
401 HashMap::new(),
402 ),
403 component: ProtocolComponent::default(),
404 component_tvl: None,
405 entrypoints: vec![(
406 EntryPointWithTracingParams::new(entry_point.clone(), params.clone()),
407 trace,
408 )],
409 },
410 )]),
411 vm_storage: HashMap::new(),
412 };
413 let header = BlockHeader { number: 7, ..Default::default() };
414
415 let changes = snapshot_to_block_changes("vm:curve", &snapshot, &header, Chain::Ethereum);
416
417 let dci = &changes.dci_update;
418 assert_eq!(dci.new_entrypoints[&component_id], HashSet::from([entry_point.clone()]));
419 assert_eq!(
420 dci.new_entrypoint_params[&entry_point.external_id],
421 HashSet::from([(params, component_id)])
422 );
423 assert!(dci.trace_results[&entry_point.external_id]
424 .accessed_slots
425 .contains_key(&oracle));
426 }
427
428 #[tokio::test]
429 async fn test_indexer_receives_the_callers_target_block() {
430 let seen = Arc::new(Mutex::new(Vec::new()));
431 let mut pending_processor = processor(seen.clone());
432 let block = target_block(1, 1_759_842_947);
434
435 pending_processor
436 .generate_pending_update(
437 &PendingBlock::new(block.clone(), vec![], HashMap::new()),
438 "bundle-1".to_string(),
439 )
440 .await
441 .expect("pending update failed");
442
443 let seen = seen.lock().unwrap();
444 assert_eq!(
445 seen.as_slice(),
446 [block],
447 "The indexer must be handed the caller's block, not one derived from the parent."
448 );
449 }
450
451 fn confirmed_at(number: u64) -> FeedMessage<BlockHeader> {
453 FeedMessage {
454 state_msgs: HashMap::from([(
455 "fluid".to_string(),
456 tycho_client::feed::synchronizer::StateSyncMessage {
457 header: BlockHeader { number, ..Default::default() },
458 ..Default::default()
459 },
460 )]),
461 sync_states: HashMap::new(),
462 }
463 }
464
465 #[tokio::test]
468 async fn test_stamped_header_comes_from_the_pending_block() {
469 let seen = Arc::new(Mutex::new(Vec::new()));
470 let mut pending_processor = processor(seen);
471 pending_processor
472 .advance(&confirmed_at(5))
473 .expect("advance failed");
474
475 let update = pending_processor
476 .generate_pending_update(
477 &PendingBlock::new(target_block(3, 1_759_842_947), vec![], HashMap::new()),
478 "bundle-1".to_string(),
479 )
480 .await
481 .expect("pending update failed");
482
483 assert_eq!(
484 update.update.block_number_or_timestamp, 3,
485 "The update must be stamped with the pending block, not the confirmed tip."
486 );
487 }
488
489 #[tokio::test]
490 async fn test_parent_guard_reads_the_pending_blocks_number() {
491 let seen = Arc::new(Mutex::new(Vec::new()));
492 let mut pending_processor = processor(seen.clone());
493 let block = target_block(23_526_115, 1_759_842_947);
494
495 let result = pending_processor
496 .generate_pending_update(
497 &PendingBlock::new(block, vec![], HashMap::new()),
498 "bundle-1".to_string(),
499 )
500 .await;
501
502 match result {
503 Err(PendingError::ParentNotYetConfirmed { needed, current }) => {
504 assert_eq!(needed, 23_526_114);
505 assert_eq!(current, 0);
506 }
507 Err(other) => panic!("expected ParentNotYetConfirmed, got {other:?}"),
508 Ok(_) => panic!("expected ParentNotYetConfirmed, got a successful update"),
509 }
510 assert!(
511 seen.lock().unwrap().is_empty(),
512 "No indexer should run when the parent is not confirmed."
513 );
514 }
515}