From 06a9f284e891256f157ac1fd967c05f01b023d85 Mon Sep 17 00:00:00 2001 From: Matthieu Vachon Date: Mon, 5 Dec 2022 15:53:14 -0500 Subject: [PATCH] The `SubstreamsBlockStream` was using `request.clone()` side stepping latest cursor value The `request.clone()` does not correctly use the `latest_cursor` value which is the valid up to date in memory cursor to use on re-connection. This led to poisining error in `graph-node` where the same block was processed multiple time because the cursor was not correctly used. Fixed by moving the request creation directly where it's needed which will use the correct up to date `latest_cursor` value now. --- .../src/blockchain/substreams_block_stream.rs | 26 +++++++++---------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/graph/src/blockchain/substreams_block_stream.rs b/graph/src/blockchain/substreams_block_stream.rs index 7ff873c4035..e0a875f903b 100644 --- a/graph/src/blockchain/substreams_block_stream.rs +++ b/graph/src/blockchain/substreams_block_stream.rs @@ -178,18 +178,6 @@ fn stream_blocks>( let stop_block_num = manifest_end_block_num as u64; - let request = Request { - start_block_num, - start_cursor: latest_cursor.clone(), - stop_block_num, - fork_steps: vec![StepNew as i32, StepUndo as i32], - irreversibility_condition: "".to_string(), - modules, - output_modules: vec![module_name], - production_mode: true, - ..Default::default() - }; - // Back off exponentially whenever we encounter a connection error or a stream with bad data let mut backoff = ExponentialBackoff::new(Duration::from_millis(500), Duration::from_secs(45)); @@ -211,7 +199,19 @@ fn stream_blocks>( skip_backoff = false; let mut connect_start = Instant::now(); - let result = endpoint.clone().substreams(request.clone()).await; + let request = Request { + start_block_num, + start_cursor: latest_cursor.clone(), + stop_block_num, + fork_steps: vec![StepNew as i32, StepUndo as i32], + irreversibility_condition: "".to_string(), + modules: modules.clone(), + output_modules: vec![module_name.clone()], + production_mode: true, + ..Default::default() + }; + + let result = endpoint.clone().substreams(request).await; match result { Ok(stream) => {