chore: Apply code formatting

This commit is contained in:
Santiago Carmuega 2022-07-17 16:42:53 -03:00
parent c377a6aaec
commit ecca0aa800
4 changed files with 142 additions and 137 deletions

View file

@ -54,10 +54,9 @@ impl VersionTable {
(PROTOCOL_V10, VersionData(network_magic)), (PROTOCOL_V10, VersionData(network_magic)),
(PROTOCOL_V11, VersionData(network_magic)), (PROTOCOL_V11, VersionData(network_magic)),
(PROTOCOL_V12, VersionData(network_magic)), (PROTOCOL_V12, VersionData(network_magic)),
]
] .into_iter()
.into_iter() .collect::<HashMap<u64, VersionData>>();
.collect::<HashMap<u64, VersionData>>();
VersionTable { values } VersionTable { values }
} }

View file

@ -5,8 +5,8 @@ pub mod blockfetch;
pub mod chainsync; pub mod chainsync;
pub mod handshake; pub mod handshake;
pub mod localstate; pub mod localstate;
pub mod txsubmission;
pub mod txmonitor; pub mod txmonitor;
pub mod txsubmission;
pub use common::*; pub use common::*;
pub use machines::*; pub use machines::*;

View file

@ -1,9 +1,7 @@
use super::{Message, MsgRequest, MsgResponse, MempoolSizeAndCapacity}; use super::{MempoolSizeAndCapacity, Message, MsgRequest, MsgResponse};
use pallas_codec::minicbor::{decode, encode, Decode, Encode, Encoder}; use pallas_codec::minicbor::{decode, encode, Decode, Encode, Encoder};
impl Encode<()> for Message impl Encode<()> for Message {
{
fn encode<W: encode::Write>( fn encode<W: encode::Write>(
&self, &self,
e: &mut Encoder<W>, e: &mut Encoder<W>,
@ -12,66 +10,59 @@ impl Encode<()> for Message
match self { match self {
Message::MsgDone => { Message::MsgDone => {
e.array(1)?.u16(0)?; e.array(1)?.u16(0)?;
}, }
Message::MsgAcquire => { Message::MsgAcquire => {
e.array(1)?.u16(1)?; e.array(1)?.u16(1)?;
}
},
Message::MsgAcquired(slot) => { Message::MsgAcquired(slot) => {
e.array(2)?.u16(2)?; e.array(2)?.u16(2)?;
e.encode(slot)?; e.encode(slot)?;
}, }
Message::MsgQuery(query) => { Message::MsgQuery(query) => {
query.encode(e, ctx)?; query.encode(e, ctx)?;
}, }
Message::MsgResponse(response) => { Message::MsgResponse(response) => {
response.encode(e, ctx)?; response.encode(e, ctx)?;
} }
} }
log::debug!("encode message: {:?}",self); log::debug!("encode message: {:?}", self);
Ok(()) Ok(())
} }
} }
impl<'b> Decode<'b,()> for Message { impl<'b> Decode<'b, ()> for Message {
fn decode(d: &mut pallas_codec::minicbor::Decoder<'b>, _ctx: &mut ()) -> Result<Self, decode::Error> { fn decode(
d.array()?; d: &mut pallas_codec::minicbor::Decoder<'b>,
let label = d.u16()?; _ctx: &mut (),
log::debug!("decode message: {:?}",label); ) -> Result<Self, decode::Error> {
match label { d.array()?;
0 => { let label = d.u16()?;
Ok(Message::MsgDone) log::debug!("decode message: {:?}", label);
}, match label {
1 => { 0 => Ok(Message::MsgDone),
Ok(Message::MsgAcquire) 1 => Ok(Message::MsgAcquire),
},
2 => { 2 => {
let slot = d.decode()?; let slot = d.decode()?;
Ok(Message::MsgAcquired(slot)) Ok(Message::MsgAcquired(slot))
}, }
3 => { 3 => Ok(Message::MsgQuery(MsgRequest::MsgRelease)),
Ok(Message::MsgQuery(MsgRequest::MsgRelease)) 5 => Ok(Message::MsgQuery(MsgRequest::MsgNextTx)),
},
5 => {
Ok(Message::MsgQuery(MsgRequest::MsgNextTx))
},
6 => { 6 => {
log::trace!("Decoding 6, 1. Array: {:?}",d); log::trace!("Decoding 6, 1. Array: {:?}", d);
let de : Result<Option<u64>,pallas_codec::minicbor::decode::Error> = d.array(); let de: Result<Option<u64>, pallas_codec::minicbor::decode::Error> = d.array();
log::trace!("Decoding 6, 2. Array: {:?}",de); log::trace!("Decoding 6, 2. Array: {:?}", de);
let tag : Result<u8,pallas_codec::minicbor::decode::Error> = d.u8(); let tag: Result<u8, pallas_codec::minicbor::decode::Error> = d.u8();
let mut tx = None; let mut tx = None;
if let Ok(_) = tag { if let Ok(_) = tag {
log::trace!("Decoding 6, Tag: {:?}",tag); log::trace!("Decoding 6, Tag: {:?}", tag);
let det = d.tag(); let det = d.tag();
log::trace!("Decoding 6, Bytes: {:?}",det); log::trace!("Decoding 6, Bytes: {:?}", det);
let cbor = d.bytes()?; let cbor = d.bytes()?;
tx = Some(hex::encode(cbor)); tx = Some(hex::encode(cbor));
log::trace!("Decoding 6, Tx: {:?}",tx); log::trace!("Decoding 6, Tx: {:?}", tx);
} }
Ok(Message::MsgResponse(MsgResponse::MsgReplyNextTx(tx))) Ok(Message::MsgResponse(MsgResponse::MsgReplyNextTx(tx)))
}, }
7 => { 7 => {
let txid = d.decode()?; let txid = d.decode()?;
Ok(Message::MsgQuery(MsgRequest::MsgHasTx(txid))) Ok(Message::MsgQuery(MsgRequest::MsgHasTx(txid)))
@ -80,29 +71,25 @@ fn decode(d: &mut pallas_codec::minicbor::Decoder<'b>, _ctx: &mut ()) -> Result<
let has = d.decode()?; let has = d.decode()?;
Ok(Message::MsgResponse(MsgResponse::MsgReplyHasTx(has))) Ok(Message::MsgResponse(MsgResponse::MsgReplyHasTx(has)))
} }
9 => { 9 => Ok(Message::MsgQuery(MsgRequest::MsgGetSizes)),
Ok(Message::MsgQuery(MsgRequest::MsgGetSizes))
}
10 => { 10 => {
d.array()?; d.array()?;
let capacity = d.decode()?; let capacity = d.decode()?;
let size_in_bytes = d.decode()?; let size_in_bytes = d.decode()?;
let number_of_tx = d.decode()?; let number_of_tx = d.decode()?;
Ok( Ok(Message::MsgResponse(MsgResponse::MsgReplyGetSizes(
Message::MsgResponse(MsgResponse::MsgReplyGetSizes(MempoolSizeAndCapacity { MempoolSizeAndCapacity {
capacity_in_bytes : capacity, capacity_in_bytes: capacity,
size_in_bytes : size_in_bytes, size_in_bytes: size_in_bytes,
number_of_txs : number_of_tx, number_of_txs: number_of_tx,
})) },
) )))
} }
_ => Err(decode::Error::message( _ => Err(decode::Error::message("can't decode Message")),
"can't decode Message", }
))
} }
}
fn nil() -> Option<Self> { fn nil() -> Option<Self> {
None None
} }
} }
@ -116,23 +103,22 @@ impl Encode<()> for MsgRequest {
match self { match self {
MsgRequest::MsgAwaitAcquire => { MsgRequest::MsgAwaitAcquire => {
e.array(1)?.u16(1)?; e.array(1)?.u16(1)?;
}, }
MsgRequest::MsgGetSizes => { MsgRequest::MsgGetSizes => {
e.array(1)?.u16(9)?; e.array(1)?.u16(9)?;
}, }
MsgRequest::MsgHasTx(tx) => { MsgRequest::MsgHasTx(tx) => {
e.array(2)?.u16(7)?; e.array(2)?.u16(7)?;
e.encode(tx)?; e.encode(tx)?;
}, }
MsgRequest::MsgNextTx => { MsgRequest::MsgNextTx => {
e.array(1)?.u16(5)?; e.array(1)?.u16(5)?;
}, }
MsgRequest::MsgRelease => { MsgRequest::MsgRelease => {
e.array(1)?.u16(3)?; e.array(1)?.u16(3)?;
}, }
} }
log::debug!("encode message: {:?}",self); log::debug!("encode message: {:?}", self);
Ok(()) Ok(())
} }
} }
@ -150,18 +136,18 @@ impl Encode<()> for MsgResponse {
e.encode(sz.capacity_in_bytes)?; e.encode(sz.capacity_in_bytes)?;
e.encode(sz.size_in_bytes)?; e.encode(sz.size_in_bytes)?;
e.encode(sz.number_of_txs)?; e.encode(sz.number_of_txs)?;
}, }
MsgResponse::MsgReplyHasTx(tx) => { MsgResponse::MsgReplyHasTx(tx) => {
e.array(2)?.u16(8)?; e.array(2)?.u16(8)?;
e.encode(tx)?; e.encode(tx)?;
}, }
MsgResponse::MsgReplyNextTx(None) => { MsgResponse::MsgReplyNextTx(None) => {
e.array(1)?.u16(6)?; e.array(1)?.u16(6)?;
}, }
MsgResponse::MsgReplyNextTx(Some(tx)) => { MsgResponse::MsgReplyNextTx(Some(tx)) => {
e.array(2)?.u16(6)?; e.array(2)?.u16(6)?;
e.encode(tx.to_string())?; e.encode(tx.to_string())?;
}, }
} }
Ok(()) Ok(())
} }

View file

@ -1,8 +1,8 @@
mod codec; mod codec;
use std::{fmt::{Debug}};
use pallas_codec::Fragment;
use crate::machines::{Agent, MachineError, Transition}; use crate::machines::{Agent, MachineError, Transition};
use pallas_codec::Fragment;
use std::fmt::Debug;
type Slot = u64; type Slot = u64;
type TxId = String; type TxId = String;
@ -26,9 +26,9 @@ pub enum State {
#[derive(Debug, PartialEq, Clone)] #[derive(Debug, PartialEq, Clone)]
pub struct MempoolSizeAndCapacity { pub struct MempoolSizeAndCapacity {
pub capacity_in_bytes : u32, pub capacity_in_bytes: u32,
pub size_in_bytes : u32, pub size_in_bytes: u32,
pub number_of_txs : u32, pub number_of_txs: u32,
} }
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
@ -40,7 +40,6 @@ pub enum Message {
MsgDone, MsgDone,
} }
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub enum MsgRequest { pub enum MsgRequest {
MsgAwaitAcquire, MsgAwaitAcquire,
@ -58,47 +57,45 @@ pub enum MsgResponse {
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct LocalTxMonitor { pub struct LocalTxMonitor {
pub state : State, pub state: State,
pub snapshot : Option<Slot>, pub snapshot: Option<Slot>,
pub request : Option<MsgRequest>, pub request: Option<MsgRequest>,
pub output : Option<MsgResponse>, pub output: Option<MsgResponse>,
} }
impl LocalTxMonitor impl LocalTxMonitor
where where
Message : Fragment, Message: Fragment,
{ {
pub fn initial(state: State) -> Self {
pub fn initial(state : State) -> Self {
Self { Self {
state : state, state: state,
snapshot : None, snapshot: None,
request : None, request: None,
output : None, output: None,
} }
} }
fn on_acquired(self, s: Slot) -> Transition<Self> { fn on_acquired(self, s: Slot) -> Transition<Self> {
log::debug!("acquired Slot: '{:?}' ",s); log::debug!("acquired Slot: '{:?}' ", s);
Ok(Self { Ok(Self {
state : State::StAcquired, state: State::StAcquired,
snapshot : Some(s), snapshot: Some(s),
output : None, output: None,
..self ..self
}) })
} }
fn on_reply_next_tx(self, tx: Option<Tx>) -> Transition<Self> { fn on_reply_next_tx(self, tx: Option<Tx>) -> Transition<Self> {
log::debug!("Next Transaction: {:?}", tx); log::debug!("Next Transaction: {:?}", tx);
Ok(Self { Ok(Self {
output: Some(MsgResponse::MsgReplyNextTx(tx)), output: Some(MsgResponse::MsgReplyNextTx(tx)),
..self ..self
}) })
} }
fn on_reply_has_tx(self, arg: bool) -> Transition<Self> { fn on_reply_has_tx(self, arg: bool) -> Transition<Self> {
log::debug!("Mempool has transaction: {:?}", arg); log::debug!("Mempool has transaction: {:?}", arg);
Ok(Self { Ok(Self {
output: Some(MsgResponse::MsgReplyHasTx(arg)), output: Some(MsgResponse::MsgReplyHasTx(arg)),
@ -106,7 +103,7 @@ impl LocalTxMonitor
}) })
} }
fn on_reply_get_size(self, msc: MempoolSizeAndCapacity) -> Transition<Self> { fn on_reply_get_size(self, msc: MempoolSizeAndCapacity) -> Transition<Self> {
log::debug!("Mempool Status: {:?}", msc); log::debug!("Mempool Status: {:?}", msc);
Ok(Self { Ok(Self {
@ -117,8 +114,9 @@ impl LocalTxMonitor
} }
impl Agent for LocalTxMonitor impl Agent for LocalTxMonitor
where where
Message: Fragment, { Message: Fragment,
{
type Message = Message; type Message = Message;
type State = State; type State = State;
@ -128,14 +126,18 @@ impl Agent for LocalTxMonitor
} }
fn is_done(&self) -> bool { fn is_done(&self) -> bool {
let done = self.state == State::StDone; let done = self.state == State::StDone;
log::debug!("is_done: {:?}",done); log::debug!("is_done: {:?}", done);
done done
} }
fn has_agency(&self) -> bool{ fn has_agency(&self) -> bool {
log::trace!("Hase Agency: State: {:?}, Request: {:?}, Response: {:?}",self.state,self.request,self.output); log::trace!(
"Hase Agency: State: {:?}, Request: {:?}, Response: {:?}",
self.state,
self.request,
self.output
);
match &self.state { match &self.state {
State::StIdle => true, State::StIdle => true,
State::StAcquiring => false, State::StAcquiring => false,
@ -146,17 +148,28 @@ impl Agent for LocalTxMonitor
} }
fn build_next(&self) -> Self::Message { fn build_next(&self) -> Self::Message {
log::debug!("build next; State: {:?}, request: {:?}, output: {:?}",&self.state, &self.request, &self.output); log::debug!(
"build next; State: {:?}, request: {:?}, output: {:?}",
&self.state,
&self.request,
&self.output
);
match (&self.state, &self.request, &self.output) { match (&self.state, &self.request, &self.output) {
(State::StIdle, None ,None) => Message::MsgAcquire, (State::StIdle, None, None) => Message::MsgAcquire,
(State::StAcquired, None , None) => Message::MsgAcquire, (State::StAcquired, None, None) => Message::MsgAcquire,
(State::StAcquired, Some(MsgRequest::MsgAwaitAcquire), None) => Message::MsgAcquire, (State::StAcquired, Some(MsgRequest::MsgAwaitAcquire), None) => Message::MsgAcquire,
(State::StAcquired, Some(MsgRequest::MsgNextTx),None) => Message::MsgQuery(MsgRequest::MsgNextTx), (State::StAcquired, Some(MsgRequest::MsgNextTx), None) => {
(State::StAcquired, Some(MsgRequest::MsgHasTx(tx)),None) => Message::MsgQuery(MsgRequest::MsgHasTx(tx.clone())), Message::MsgQuery(MsgRequest::MsgNextTx)
(State::StAcquired, Some(MsgRequest::MsgGetSizes),None) => Message::MsgQuery(MsgRequest::MsgGetSizes), }
(State::StAcquired, None, Some(_)) => Message::MsgAcquire, (State::StAcquired, Some(MsgRequest::MsgHasTx(tx)), None) => {
(State::StAcquired, Some(req), Some(_)) => Message::MsgQuery(req.to_owned()), Message::MsgQuery(MsgRequest::MsgHasTx(tx.clone()))
_ => panic!("I do not have agency, don't know what to do") }
(State::StAcquired, Some(MsgRequest::MsgGetSizes), None) => {
Message::MsgQuery(MsgRequest::MsgGetSizes)
}
(State::StAcquired, None, Some(_)) => Message::MsgAcquire,
(State::StAcquired, Some(req), Some(_)) => Message::MsgQuery(req.to_owned()),
_ => panic!("I do not have agency, don't know what to do"),
} }
} }
@ -168,13 +181,13 @@ impl Agent for LocalTxMonitor
fn apply_outbound(self, msg: Self::Message) -> Transition<Self> { fn apply_outbound(self, msg: Self::Message) -> Transition<Self> {
log::debug!("apply outbound"); log::debug!("apply outbound");
match (self.state, msg) { match (self.state, msg) {
(State::StIdle, Message::MsgAcquire) => { (State::StIdle, Message::MsgAcquire) => {
log::debug!("apply outbound : MsgAcquire"); log::debug!("apply outbound : MsgAcquire");
Ok(Self { Ok(Self {
state: State::StAcquiring, state: State::StAcquiring,
..self ..self
})}, })
}
(State::StAcquired, Message::MsgQuery(MsgRequest::MsgNextTx)) => Ok(Self { (State::StAcquired, Message::MsgQuery(MsgRequest::MsgNextTx)) => Ok(Self {
state: State::StBusy(StBusyKind::NextTx), state: State::StBusy(StBusyKind::NextTx),
..self ..self
@ -201,20 +214,27 @@ impl Agent for LocalTxMonitor
..self ..self
}), }),
_ => panic!("PANIC! Cannot match outbound") _ => panic!("PANIC! Cannot match outbound"),
} }
} }
fn apply_inbound(self, msg: Self::Message) -> Transition<Self> { fn apply_inbound(self, msg: Self::Message) -> Transition<Self> {
log::debug!("apply inbound"); log::debug!("apply inbound");
match (&self.state , msg) { match (&self.state, msg) {
(State::StAcquiring, Message::MsgAcquired(s)) => self.on_acquired(s), (State::StAcquiring, Message::MsgAcquired(s)) => self.on_acquired(s),
(State::StBusy(StBusyKind::NextTx), Message::MsgResponse(MsgResponse::MsgReplyNextTx(tx))) => self.on_reply_next_tx(tx), (
(State::StBusy(StBusyKind::HasTx), Message::MsgResponse(MsgResponse::MsgReplyHasTx(arg))) => self.on_reply_has_tx(arg), State::StBusy(StBusyKind::NextTx),
(State::StBusy(StBusyKind::GetSizes), Message::MsgResponse(MsgResponse::MsgReplyGetSizes(msc))) => self.on_reply_get_size(msc), Message::MsgResponse(MsgResponse::MsgReplyNextTx(tx)),
) => self.on_reply_next_tx(tx),
(
State::StBusy(StBusyKind::HasTx),
Message::MsgResponse(MsgResponse::MsgReplyHasTx(arg)),
) => self.on_reply_has_tx(arg),
(
State::StBusy(StBusyKind::GetSizes),
Message::MsgResponse(MsgResponse::MsgReplyGetSizes(msc)),
) => self.on_reply_get_size(msc),
(state, msg) => Err(MachineError::invalid_msg::<Self>(&state, &msg)), (state, msg) => Err(MachineError::invalid_msg::<Self>(&state, &msg)),
} }
} }
} }