feat(blockfetch): add more observer events
This commit is contained in:
parent
a7fbbe4a28
commit
123cc68a60
1 changed files with 35 additions and 13 deletions
|
|
@ -1,6 +1,6 @@
|
||||||
use std::sync::mpsc::Receiver;
|
use std::sync::mpsc::Receiver;
|
||||||
|
|
||||||
use log::info;
|
use log::{debug, info};
|
||||||
use pallas_machines::{
|
use pallas_machines::{
|
||||||
primitives::Point, Agent, CodecError, DecodePayload, EncodePayload, MachineOutput,
|
primitives::Point, Agent, CodecError, DecodePayload, EncodePayload, MachineOutput,
|
||||||
PayloadDecoder, PayloadEncoder, Transition,
|
PayloadDecoder, PayloadEncoder, Transition,
|
||||||
|
|
@ -88,8 +88,20 @@ impl DecodePayload for Message {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub trait Observer {
|
pub trait Observer {
|
||||||
fn on_block(&self, body: Vec<u8>) -> Result<(), Box<dyn std::error::Error>> {
|
fn on_block_received(&self, body: Vec<u8>) -> Result<(), Box<dyn std::error::Error>> {
|
||||||
log::debug!("block fetched {:?}", body);
|
log::debug!("block received, sice: {}", body.len());
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn on_block_range_requested(
|
||||||
|
&self,
|
||||||
|
range: &(Point, Point),
|
||||||
|
) -> Result<(), Box<dyn std::error::Error>> {
|
||||||
|
log::debug!(
|
||||||
|
"block range requested, from: {:?}, to: {:?}",
|
||||||
|
range.0,
|
||||||
|
range.1
|
||||||
|
);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -128,11 +140,21 @@ where
|
||||||
|
|
||||||
tx.send_msg(&msg)?;
|
tx.send_msg(&msg)?;
|
||||||
|
|
||||||
|
self.observer.on_block_range_requested(&self.range)?;
|
||||||
|
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
state: State::Busy,
|
state: State::Busy,
|
||||||
..self
|
..self
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn on_block(self, body: Vec<u8>) -> Transition<Self> {
|
||||||
|
debug!("received block body, size {}", body.len());
|
||||||
|
|
||||||
|
self.observer.on_block_received(body)?;
|
||||||
|
|
||||||
|
Ok(self)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<O> Agent for BatchClient<O>
|
impl<O> Agent for BatchClient<O>
|
||||||
|
|
@ -171,11 +193,7 @@ where
|
||||||
state: State::Done,
|
state: State::Done,
|
||||||
..self
|
..self
|
||||||
}),
|
}),
|
||||||
(State::Streaming, Message::Block { body }) => {
|
(State::Streaming, Message::Block { body }) => self.on_block(body),
|
||||||
info!("received block body of size {}", body.len());
|
|
||||||
self.observer.on_block(body)?;
|
|
||||||
Ok(self)
|
|
||||||
}
|
|
||||||
(State::Streaming, Message::BatchDone) => Ok(Self {
|
(State::Streaming, Message::BatchDone) => Ok(Self {
|
||||||
state: State::Done,
|
state: State::Done,
|
||||||
..self
|
..self
|
||||||
|
|
@ -221,6 +239,14 @@ where
|
||||||
..self
|
..self
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn on_block(self, body: Vec<u8>) -> Transition<Self> {
|
||||||
|
debug!("received block body, size {}", body.len());
|
||||||
|
|
||||||
|
self.observer.on_block_received(body)?;
|
||||||
|
|
||||||
|
Ok(self)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<O> Agent for OnDemandClient<O>
|
impl<O> Agent for OnDemandClient<O>
|
||||||
|
|
@ -261,11 +287,7 @@ where
|
||||||
state: State::Idle,
|
state: State::Idle,
|
||||||
..self
|
..self
|
||||||
}),
|
}),
|
||||||
(State::Streaming, Message::Block { body }) => {
|
(State::Streaming, Message::Block { body }) => self.on_block(body),
|
||||||
info!("received block body of size {}", body.len());
|
|
||||||
self.observer.on_block(body)?;
|
|
||||||
Ok(self)
|
|
||||||
}
|
|
||||||
(State::Streaming, Message::BatchDone) => Ok(Self {
|
(State::Streaming, Message::BatchDone) => Ok(Self {
|
||||||
state: State::Idle,
|
state: State::Idle,
|
||||||
..self
|
..self
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue