Skip to content

Commit 1350fb5

Browse files
committed
problem: head may have different block on top and if multiple streams are archiving they may mix the data
solution: use block hash as part of the filename for stream archives (single block ranges)
1 parent 7b56427 commit 1350fb5

16 files changed

Lines changed: 281 additions & 228 deletions

File tree

src/archiver/archiver.rs

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -5,12 +5,11 @@ use chrono::Utc;
55
use tokio::sync::mpsc::Sender;
66
use crate::archiver::block::ArchiveBlock;
77
use crate::archiver::table::ArchiveTable;
8-
use crate::blockchain::{BlockchainData, BlockchainTypes, MultiBlockReference};
9-
use crate::blockchain::connection::Height;
8+
use crate::blockchain::{BlockchainData, BlockchainTypes};
109
use crate::archiver::datakind::{DataKind, DataOptions};
1110
use crate::notify::empty::EmptyNotifier;
1211
use crate::notify::{Maturity, Notification, Notifier, RunMode};
13-
use crate::archiver::range::Range;
12+
use crate::archiver::range::{Height, Range};
1413
use crate::global;
1514
use crate::storage::TargetStorage;
1615

@@ -71,8 +70,8 @@ impl<B: BlockchainTypes, TS: TargetStorage> ArchiveAll<Height> for Archiver<B, T
7170
location: "".to_string(),
7271
};
7372

74-
let blocks = self.process_blocks(MultiBlockReference::Single(what.clone()), notification.clone()).await?;
75-
let range = Range::Single(what.height);
73+
let blocks = self.process_blocks(Range::Single(what.clone()), notification.clone()).await?;
74+
let range = Range::Single(what.clone());
7675

7776
if let Some(tx_options) = &options.tx {
7877
self.process_table(range.clone(), notification.clone(), &blocks, tx_options).await?;
@@ -109,7 +108,7 @@ impl<B: BlockchainTypes, TS: TargetStorage> ArchiveAll<Range> for Archiver<B, TS
109108
location: "".to_string(),
110109
};
111110

112-
let blocks = self.process_blocks(MultiBlockReference::Range(what.clone()), notification.clone()).await?;
111+
let blocks = self.process_blocks(what.clone(), notification.clone()).await?;
113112

114113
if let Some(tx_options) = &options.tx {
115114
self.process_table(what.clone(), notification.clone(), &blocks, tx_options).await?;

src/archiver/block.rs

Lines changed: 8 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,12 @@
1-
use std::str::FromStr;
21
use anyhow::anyhow;
32
use async_trait::async_trait;
43
use chrono::Utc;
54
use crate::archiver::archiver::Archiver;
65
use crate::archiver::BlockTransactions;
7-
use crate::blockchain::{BlockReference, BlockchainData, BlockchainTypes, MultiBlockReference};
6+
use crate::blockchain::{BlockReference, BlockchainData, BlockchainTypes};
87
use crate::archiver::datakind::DataKind;
98
use crate::notify::Notification;
10-
use crate::archiver::range::Range;
9+
use crate::archiver::range::{Height, Range};
1110
use crate::global;
1211
use crate::storage::{TargetFile, TargetFileWriter, TargetStorage};
1312

@@ -17,46 +16,30 @@ use crate::storage::{TargetFile, TargetFileWriter, TargetStorage};
1716
pub trait ArchiveBlock<B: BlockchainTypes> {
1817
///
1918
/// @param blocks - which blocks to process (by height or hash)
20-
async fn process_blocks(&self, blocks: MultiBlockReference, notification: Notification) -> anyhow::Result<BlockTransactions<B>>;
19+
async fn process_blocks(&self, blocks: Range, notification: Notification) -> anyhow::Result<BlockTransactions<B>>;
2120
}
2221

2322
#[async_trait]
2423
impl<B: BlockchainTypes, TS: TargetStorage> ArchiveBlock<B> for Archiver<B, TS> {
25-
async fn process_blocks(&self, blocks: MultiBlockReference, notification: Notification) -> anyhow::Result<BlockTransactions<B>> {
24+
async fn process_blocks(&self, blocks: Range, notification: Notification) -> anyhow::Result<BlockTransactions<B>> {
2625
let shutdown = global::get_shutdown();
2726
if shutdown.is_signalled() {
2827
return Ok(vec![]);
2928
}
3029
let dry_run = global::is_dry_run();
31-
let range: Range = blocks.clone().into();
32-
let file = self.target.create(DataKind::Blocks, &range)
30+
let file = self.target.create(DataKind::Blocks, &blocks)
3331
.await
3432
.map_err(|e| anyhow!("Unable to create file: {}", e))?;
3533
let file_url = file.get_url();
3634

37-
let heights = match blocks {
38-
MultiBlockReference::Single(h) => {
39-
let r = if let Some(hash) = &h.hash {
40-
BlockReference::Hash(
41-
B::BlockHash::from_str(hash).map_err(|_| anyhow!("Not a valid hash"))?
42-
)
43-
} else {
44-
BlockReference::height(h.height)
45-
};
46-
vec![r]
47-
},
48-
MultiBlockReference::Range(range) => {
49-
range.iter().map(|h| BlockReference::height(h)).collect()
50-
}
51-
};
52-
5335
let mut results = Vec::new();
54-
for height in heights {
36+
for height in blocks.iter_height().collect::<Vec<Height>>() {
5537
if shutdown.is_signalled() {
5638
tracing::info!("Shutdown signalled, stopping");
5739
break;
5840
}
59-
let (record, block, txes) = self.data_provider.fetch_block(&height).await?;
41+
let block_ref = BlockReference::Height(height);
42+
let (record, block, txes) = self.data_provider.fetch_block(&block_ref).await?;
6043
if !dry_run {
6144
let _ = file.append(record).await?;
6245
}

src/archiver/filenames.rs

Lines changed: 92 additions & 50 deletions
Original file line numberDiff line numberDiff line change
@@ -2,10 +2,10 @@ use std::str::FromStr;
22
use lazy_static::lazy_static;
33
use regex::Regex;
44
use crate::archiver::datakind::DataKind;
5-
use crate::archiver::range::Range;
5+
use crate::archiver::range::{Height, Range};
66

77
lazy_static! {
8-
static ref RE_SINGLE: Regex = Regex::new(r"^(\d+)\.(\w+)\.(\w+\.)?avro$").unwrap();
8+
static ref RE_SINGLE: Regex = Regex::new(r"^(\d+)\.(([a-f0-9]{64})\.)?(\w+)\.(\w+\.)?avro$").unwrap();
99
static ref RE_RANGE: Regex = Regex::new(r"^range-(\d+)_(\d+)\.(\w+)\.(\w+\.)?avro$").unwrap();
1010
}
1111

@@ -28,67 +28,70 @@ impl Filenames {
2828

2929
pub fn parse(filename: String) -> Option<(DataKind, Range)> {
3030
if let Some(cap) = RE_SINGLE.captures(filename.as_str()) {
31-
let height = cap.get(1).unwrap().as_str().parse().unwrap();
32-
let kind = DataKind::from_str(cap.get(2).unwrap().as_str());
31+
let height: u64 = cap.get(1).unwrap().as_str().parse().unwrap();
32+
let hash = cap.get(3).map(|x| x.as_str().to_string());
33+
let kind = DataKind::from_str(cap.get(4).unwrap().as_str());
3334
if kind.is_err() {
3435
return None;
3536
}
36-
return Some((kind.unwrap(), Range::Single(height)));
37+
return Some((kind.unwrap(), Range::Single(Height::new(height, hash))));
3738
}
3839
if let Some(cap) = RE_RANGE.captures(filename.as_str()) {
39-
let start = cap.get(1).unwrap().as_str().parse().unwrap();
40-
let end = cap.get(2).unwrap().as_str().parse().unwrap();
40+
let start: u64 = cap.get(1).unwrap().as_str().parse().unwrap();
41+
let end: u64 = cap.get(2).unwrap().as_str().parse().unwrap();
4142
let kind = DataKind::from_str(cap.get(3).unwrap().as_str());
4243
if kind.is_err() {
4344
return None;
4445
}
45-
return Some((kind.unwrap(), Range::Multiple(start, end)));
46+
return Some((kind.unwrap(), Range::Multiple(start.into(), end.into())));
4647
}
4748
None
4849
}
4950

50-
pub fn filename(&self, kind: &DataKind, range: &Range) -> String {
51+
pub fn filename(&self, kind: &DataKind, blocks: &Range) -> String {
5152
let suffix = match kind {
52-
DataKind::Blocks => match range {
53-
Range::Single(_) => "block",
54-
Range::Multiple(_, _) => "blocks"
55-
},
53+
DataKind::Blocks => if blocks.len() == 1 { "block" } else { "blocks" },
5654
DataKind::Transactions => "txes",
5755
DataKind::TransactionTraces => "traces"
5856
};
5957

60-
match range {
61-
Range::Single(height) => {
62-
format!("{}.{}.avro", self.range_padded(*height), suffix)
63-
},
64-
Range::Multiple(start, end) => {
65-
format!("range-{}_{}.{}.avro", self.range_padded(*start), self.range_padded(*end), suffix)
58+
if blocks.len() == 1 {
59+
let height = blocks.first_height();
60+
if let Some(hash) = &height.hash {
61+
format!("{}.{}.{}.avro", self.range_padded(blocks.start()), hash, suffix)
62+
} else {
63+
format!("{}.{}.avro", self.range_padded(blocks.start()), suffix)
6664
}
65+
} else {
66+
format!("range-{}_{}.{}.avro", self.range_padded(blocks.start()), self.range_padded(blocks.end()), suffix)
6767
}
6868
}
6969

70-
pub fn relative_path(&self, kind: &DataKind, range: &Range) -> String {
71-
match range {
72-
Range::Single(start) => {
73-
format!("{}/{}/{}",
74-
self.level_1(*start),
75-
self.level_2(*start),
76-
self.filename(kind, range)
77-
)
78-
}
79-
Range::Multiple(start, _end) => {
80-
format!("{}/{}",
81-
self.level_1(*start),
82-
self.filename(kind, range)
83-
)
84-
}
70+
pub fn relative_path(&self, kind: &DataKind, blocks: &Range) -> String {
71+
if blocks.len() == 1 {
72+
format!("{}/{}/{}",
73+
self.level_1(blocks.start()),
74+
self.level_2(blocks.start()),
75+
self.filename(kind, blocks)
76+
)
77+
} else {
78+
format!("{}/{}",
79+
self.level_1(blocks.start()),
80+
self.filename(kind, blocks)
81+
)
8582
}
8683
}
8784

85+
///
86+
/// Offset is a position in a directory list. Since Dshackle Archive uses a S3-like storage, which sorts files by their name,
87+
/// this is where to start to look for files in a given range.
8888
pub fn offset(&self, range: &Range) -> String {
8989
match range {
90-
Range::Single(start) => self.range_padded(*start),
91-
Range::Multiple(start, _) => format!("range-{}", self.range_padded(*start))
90+
// For a single block, it's just the height
91+
Range::Single(start) => self.range_padded(start.height),
92+
// For a range we use a prefix `range-` only to distinguish them from the single files,
93+
// otherwise they would be mixed and single ranges listed twice, etc.
94+
Range::Multiple(start, _) => format!("range-{}", self.range_padded(start.height))
9295
}
9396
}
9497

@@ -226,56 +229,56 @@ mod tests {
226229
fn single_block_file() {
227230
let filenames = Filenames::default();
228231
let kind = DataKind::Blocks;
229-
let range = Range::Single(12000000);
232+
let range = Range::Single(12000000.into());
230233
assert_eq!(filenames.filename(&kind, &range), "012000000.block.avro");
231234
}
232235

233236
#[test]
234237
fn single_block_txes_file() {
235238
let filenames = Filenames::default();
236239
let kind = DataKind::Transactions;
237-
let range = Range::Single(12000000);
240+
let range = Range::Single(12000000.into());
238241
assert_eq!(filenames.filename(&kind, &range), "012000000.txes.avro");
239242
}
240243

241244
#[test]
242245
fn single_block_tx_traces_file() {
243246
let filenames = Filenames::default();
244247
let kind = DataKind::TransactionTraces;
245-
let range = Range::Single(12000000);
248+
let range = Range::Single(12000000.into());
246249
assert_eq!(filenames.filename(&kind, &range), "012000000.traces.avro");
247250
}
248251

249252
#[test]
250253
fn single_block_path() {
251254
let filenames = Filenames::default();
252255
let kind = DataKind::Blocks;
253-
assert_eq!(filenames.path(&kind, &Range::Single(12000005)), "012000000/012000000/012000005.block.avro");
254-
assert_eq!(filenames.path(&kind, &Range::Single(12004999)), "012000000/012004000/012004999.block.avro");
255-
assert_eq!(filenames.path(&kind, &Range::Single(12005000)), "012000000/012005000/012005000.block.avro");
256-
assert_eq!(filenames.path(&kind, &Range::Single(12005001)), "012000000/012005000/012005001.block.avro");
257-
assert_eq!(filenames.path(&kind, &Range::Single(12345678)), "012000000/012345000/012345678.block.avro");
256+
assert_eq!(filenames.path(&kind, &Range::Single(12000005.into())), "012000000/012000000/012000005.block.avro");
257+
assert_eq!(filenames.path(&kind, &Range::Single(12004999.into())), "012000000/012004000/012004999.block.avro");
258+
assert_eq!(filenames.path(&kind, &Range::Single(12005000.into())), "012000000/012005000/012005000.block.avro");
259+
assert_eq!(filenames.path(&kind, &Range::Single(12005001.into())), "012000000/012005000/012005001.block.avro");
260+
assert_eq!(filenames.path(&kind, &Range::Single(12345678.into())), "012000000/012345000/012345678.block.avro");
258261
}
259262

260263
#[test]
261264
fn multi_block_path() {
262265
let filenames = Filenames::default();
263266
let kind = DataKind::Blocks;
264-
assert_eq!(filenames.path(&kind, &Range::Multiple(12000000, 12000999)), "012000000/range-012000000_012000999.blocks.avro");
267+
assert_eq!(filenames.path(&kind, &Range::Multiple(12000000.into(), 12000999.into())), "012000000/range-012000000_012000999.blocks.avro");
265268
}
266269

267270
#[test]
268271
fn multi_tx_path() {
269272
let filenames = Filenames::default();
270273
let kind = DataKind::Transactions;
271-
assert_eq!(filenames.path(&kind, &Range::Multiple(12000000, 12000999)), "012000000/range-012000000_012000999.txes.avro");
274+
assert_eq!(filenames.path(&kind, &Range::Multiple(12000000.into(), 12000999.into())), "012000000/range-012000000_012000999.txes.avro");
272275
}
273276

274277
#[test]
275278
fn multi_tx_traces_path() {
276279
let filenames = Filenames::default();
277280
let kind = DataKind::TransactionTraces;
278-
assert_eq!(filenames.path(&kind, &Range::Multiple(12000000, 12000999)), "012000000/range-012000000_012000999.traces.avro");
281+
assert_eq!(filenames.path(&kind, &Range::Multiple(12000000.into(), 12000999.into())), "012000000/range-012000000_012000999.traces.avro");
279282
}
280283

281284
#[test]
@@ -343,20 +346,59 @@ mod tests {
343346
fn parse_single_block_file() {
344347
let (kind, range) = Filenames::parse("021625120.block.avro".to_string()).unwrap();
345348
assert_eq!(kind, DataKind::Blocks);
346-
assert_eq!(range, Range::Single(21625120));
349+
assert_eq!(range, Range::Single(21625120.into()));
350+
assert!(range.first_height().hash.is_none());
351+
}
352+
353+
#[test]
354+
fn parse_single_block_file_with_hash() {
355+
let (kind, range) = Filenames::parse("021625120.b4e72a78fd2cb0e75768401a92e5258618eec892b7e4edb49b2a592e9a5de5c4.block.avro".to_string()).unwrap();
356+
assert_eq!(kind, DataKind::Blocks);
357+
assert_eq!(range, Range::Single(Height::new(21625120, Some("b4e72a78fd2cb0e75768401a92e5258618eec892b7e4edb49b2a592e9a5de5c4".to_string()))));
347358
}
348359

349360
#[test]
350361
fn parse_single_tx_file() {
351362
let (kind, range) = Filenames::parse("021625139.txes.avro".to_string()).unwrap();
352363
assert_eq!(kind, DataKind::Transactions);
353-
assert_eq!(range, Range::Single(21625139));
364+
assert_eq!(range, Range::Single(21625139.into()));
365+
assert!(range.first_height().hash.is_none());
366+
}
367+
368+
#[test]
369+
fn parse_single_tx_file_with_hash() {
370+
let (kind, range) = Filenames::parse("021625139.b4e72a78fd2cb0e75768401a92e5258618eec892b7e4edb49b2a592e9a5de5c4.txes.avro".to_string()).unwrap();
371+
assert_eq!(kind, DataKind::Transactions);
372+
assert_eq!(range, Range::Single(Height::new(21625139, Some("b4e72a78fd2cb0e75768401a92e5258618eec892b7e4edb49b2a592e9a5de5c4".to_string()))));
354373
}
355374

356375
#[test]
357376
fn parse_single_tx_traces_file() {
358377
let (kind, range) = Filenames::parse("021625139.traces.avro".to_string()).unwrap();
359378
assert_eq!(kind, DataKind::TransactionTraces);
360-
assert_eq!(range, Range::Single(21625139));
379+
assert_eq!(range, Range::Single(21625139.into()));
380+
assert!(range.first_height().hash.is_none());
381+
}
382+
383+
#[test]
384+
fn offset_single_block() {
385+
let filenames = Filenames::default();
386+
let single = Range::Single(12000000.into());
387+
assert_eq!(filenames.offset(&single), "012000000");
388+
}
389+
390+
#[test]
391+
fn offset_multi_block() {
392+
let filenames = Filenames::default();
393+
let multi = Range::Multiple(12000000.into(), 12000999.into());
394+
assert_eq!(filenames.offset(&multi), "range-012000000");
395+
}
396+
397+
#[test]
398+
fn offset_different_for_ranges() {
399+
let filenames = Filenames::default();
400+
let multi = Range::Multiple(12000000.into(), 12000999.into());
401+
let single = multi.first();
402+
assert_ne!(filenames.offset(&single), filenames.offset(&multi));
361403
}
362404
}

src/archiver/mod.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,3 +16,5 @@ use crate::blockchain::BlockchainTypes;
1616

1717
#[allow(type_alias_bounds)]
1818
pub type BlockTransactions<B: BlockchainTypes> = Vec<(B::BlockParsed, Vec<B::TxId>)>;
19+
20+
pub type BlockHash = String;

0 commit comments

Comments
 (0)