1- use crate :: cli:: { Cli , NetworkKind } ;
2- use crate :: fs:: create_fs;
3- use crate :: ingest:: ingest_from_service;
4- use crate :: layout:: Layout ;
5- use crate :: metrics;
6- use crate :: proc:: Proc ;
7- use crate :: server:: run_server;
8- use crate :: writer:: Writer ;
1+ use std:: time:: Duration ;
2+
93use anyhow:: { ensure, Context } ;
104use prometheus_client:: registry:: Registry ;
11- use sqd_data:: bitcoin :: tables :: BitcoinChunkBuilder ;
12- use sqd_data :: evm:: tables:: EvmChunkBuilder ;
13- use sqd_data :: hyperliquid_fills:: tables:: HyperliquidFillsChunkBuilder ;
14- use sqd_data :: hyperliquid_replica_cmds:: tables:: HyperliquidReplicaCmdsChunkBuilder ;
15- use sqd_data :: solana :: tables:: SolanaChunkBuilder ;
16- use sqd_data :: tron :: tables :: TronChunkBuilder ;
5+ use sqd_data:: {
6+ bitcoin :: tables :: BitcoinChunkBuilder , evm:: tables:: EvmChunkBuilder ,
7+ hyperliquid_fills:: tables:: HyperliquidFillsChunkBuilder ,
8+ hyperliquid_replica_cmds:: tables:: HyperliquidReplicaCmdsChunkBuilder , solana :: tables :: SolanaChunkBuilder ,
9+ tron :: tables:: TronChunkBuilder
10+ } ;
1711use sqd_primitives:: BlockNumber ;
18- use std:: time:: Duration ;
1912
13+ use crate :: {
14+ cli:: { Cli , NetworkKind } ,
15+ fs:: create_fs,
16+ ingest:: ingest_from_service,
17+ layout:: Layout ,
18+ metrics,
19+ proc:: Proc ,
20+ server:: run_server,
21+ writer:: Writer
22+ } ;
2023
2124pub async fn run ( args : Cli ) -> anyhow:: Result < ( ) > {
2225 ensure ! (
@@ -27,12 +30,9 @@ pub async fn run(args: Cli) -> anyhow::Result<()> {
2730 let fs = create_fs ( & args. dest ) . await ?;
2831 let layout = Layout :: new ( fs. clone ( ) ) ;
2932
30- let chunk_tracker = layout. create_chunk_tracker (
31- & chunk_check,
32- args. top_dir_size ,
33- args. first_block ,
34- args. last_block
35- ) . await ?;
33+ let chunk_tracker = layout
34+ . create_chunk_tracker ( & chunk_check, args. top_dir_size , args. first_block , args. last_block )
35+ . await ?;
3636
3737 if let Some ( last_block) = args. last_block {
3838 if chunk_tracker. next_block ( ) > last_block {
@@ -78,32 +78,31 @@ pub async fn run(args: Cli) -> anyhow::Result<()> {
7878 let attach_idx_field = args. attach_idx_field ;
7979 let write_task = tokio:: spawn ( async move {
8080 let mut writer = Writer :: new ( fs, chunk_receiver, attach_idx_field) ;
81- writer. start ( ) . await
81+ writer. start ( ) . await
8282 } ) ;
83-
83+
8484 match write_task. await . context ( "write task panicked" ) {
8585 Ok ( Ok ( _) ) => {
8686 proc_task. await . context ( "processing task panicked" ) ??;
87- } ,
87+ }
8888 Ok ( Err ( err) ) => {
8989 proc_task. abort ( ) ;
90- return Err ( err)
91- } ,
90+ return Err ( err) ;
91+ }
9292 Err ( err) => {
9393 proc_task. abort ( ) ;
94- return Err ( err)
94+ return Err ( err) ;
9595 }
9696 }
9797
9898 Ok ( ( ) )
9999}
100100
101-
102101fn chunk_check ( filelist : & [ String ] ) -> bool {
103102 for file in filelist {
104103 if file. starts_with ( "blocks.parquet" ) {
105104 return true ;
106105 }
107106 }
108107 false
109- }
108+ }
0 commit comments