-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathchunk_downloader.rs
136 lines (127 loc) · 4.83 KB
/
chunk_downloader.rs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
use crate::{
config::Config,
types::{ChunkDownloaderMessage, ManagerMessage, ManagerMessageKind, ShutdownSignal},
};
use near_jsonrpc_client::{methods, JsonRpcClient};
use near_jsonrpc_primitives::types::chunks::ChunkReference;
use near_primitives::{hash::CryptoHash, views::ChunkView};
use std::time::Duration;
use tokio::{
sync::mpsc::{self, Receiver, Sender},
task::JoinHandle,
};
/// An "actor" which represents a background task to poll the Near RPC
/// at regular intervals for new blocks.
pub struct ChunkDownloader {
id: String,
client: JsonRpcClient,
retry_frequency: Duration,
manager_channel: Sender<ManagerMessage>,
incoming_channel: Receiver<ChunkDownloaderMessage>,
max_retry_count: usize,
}
impl ChunkDownloader {
pub fn new(
config: &Config,
manager_channel: Sender<ManagerMessage>,
id_no: usize,
) -> (Self, Sender<ChunkDownloaderMessage>) {
let id = format!("ChunkDownloader_{id_no}");
let max_retry_count = config.max_download_retry.into();
let retry_frequency = Duration::from_millis(config.polling_frequency_ms);
let client = JsonRpcClient::new_client().connect(&config.near_rpc_url);
let (sender, incoming_channel) = mpsc::channel(100);
let this = Self {
id,
client,
retry_frequency,
manager_channel,
incoming_channel,
max_retry_count,
};
(this, sender)
}
pub fn start(mut self) -> JoinHandle<anyhow::Result<()>> {
tokio::task::spawn(async move {
while let Some(message) = self.incoming_channel.recv().await {
tracing::debug!("{} received a message from the Manager", self.id);
match message {
ChunkDownloaderMessage::Download {
chunk_hash,
next_block_hash: block_hash,
} => {
match download_chunk_with_retry(
&self.client,
chunk_hash,
self.retry_frequency,
self.max_retry_count,
)
.await
{
Ok(chunk) => {
if let Err(e) = self
.send_manager_message(ManagerMessageKind::NewChunk {
chunk: Box::new(chunk),
next_block_hash: block_hash,
})
.await
{
tracing::error!(
"ChunkDownloader failed to communicate with Manager."
);
return Err(e);
}
}
Err(e) => {
tracing::warn!("Failed to download chunk: {:?}", e);
self.send_manager_message(ManagerMessageKind::Shutdown(
ShutdownSignal,
))
.await
.ok();
return Err(anyhow::anyhow!("Failed to download chunks"));
}
}
}
ChunkDownloaderMessage::Shutdown(ShutdownSignal) => {
tracing::info!("ChunkDownloader received ShutdownSignal");
break;
}
}
}
Ok(())
})
}
async fn send_manager_message(&self, kind: ManagerMessageKind) -> anyhow::Result<()> {
let message = ManagerMessage {
worker_id: self.id.clone(),
kind,
};
self.manager_channel.send(message).await?;
Ok(())
}
}
async fn download_chunk_with_retry(
client: &JsonRpcClient,
chunk_hash: CryptoHash,
retry_frequency: Duration,
max_retries: usize,
) -> anyhow::Result<ChunkView> {
for _ in 0..max_retries {
match download_chunk(client, chunk_hash).await {
Ok(chunk) => return Ok(chunk),
Err(e) => {
tracing::warn!("Failed to download chunk: {:?}", e);
tokio::time::sleep(retry_frequency).await;
}
}
}
Err(anyhow::anyhow!("Failed to download chunk"))
}
async fn download_chunk(client: &JsonRpcClient, chunk_id: CryptoHash) -> anyhow::Result<ChunkView> {
let request = methods::chunk::RpcChunkRequest {
chunk_reference: ChunkReference::ChunkHash { chunk_id },
};
let chunk = client.call(request).await?;
Ok(chunk)
}