hopr_protocol_hopr/codec/
encoder.rs1use bytes::{BufMut, BytesMut};
2use hopr_api::{
3 chain::*,
4 types::{
5 crypto::{crypto_traits::Randomizable, prelude::*},
6 internal::prelude::*,
7 primitive::prelude::*,
8 },
9};
10use hopr_crypto_packet::prelude::*;
11
12use crate::{HoprCodecConfig, OutgoingPacket, PacketEncoder, SurbStore, errors::HoprProtocolError};
13
14pub const MAX_ACKNOWLEDGEMENTS_BATCH_SIZE: usize =
19 (HoprPacket::PAYLOAD_SIZE - size_of::<u16>()) / Acknowledgement::SIZE;
20
21pub struct HoprEncoder<Chain, S, T> {
23 chain_api: Chain,
24 surb_store: S,
25 ticket_factory: T,
26 chain_key: ChainKeypair,
27 channels_dst: Hash,
28 cfg: HoprCodecConfig,
29}
30
31impl<Chain, S, T> HoprEncoder<Chain, S, T> {
32 pub fn new(
34 chain_key: ChainKeypair,
35 chain_api: Chain,
36 surb_store: S,
37 ticket_factory: T,
38 channels_dst: Hash,
39 cfg: HoprCodecConfig,
40 ) -> Self {
41 Self {
42 chain_api,
43 surb_store,
44 ticket_factory,
45 chain_key,
46 channels_dst,
47 cfg,
48 }
49 }
50}
51
52impl<Chain, S, T> HoprEncoder<Chain, S, T>
53where
54 Chain: ChainKeyOperations + ChainReadChannelOperations + ChainReadTicketOperations + ChainValues + Sync,
55 S: SurbStore,
56 T: hopr_api::tickets::TicketFactory + Sync,
57{
58 fn encode_packet_internal<D: AsRef<[u8]> + Send + 'static, Sig: Into<PacketSignals> + Send + 'static>(
59 &self,
60 next_peer: OffchainPublicKey,
61 data: D,
62 num_hops: usize,
63 signals: Sig,
64 routing: PacketRouting<ValidatedPath>,
65 pseudonym: HoprPseudonym,
66 ) -> Result<OutgoingPacket, HoprProtocolError> {
67 let next_peer = self
68 .chain_api
69 .packet_key_to_chain_key(&next_peer)
70 .map_err(HoprProtocolError::resolver)?
71 .ok_or(HoprProtocolError::KeyNotFound)?;
72
73 let next_ticket = if num_hops > 1 {
75 let channel = self
76 .chain_api
77 .channel_by_parties(self.chain_key.as_ref(), &next_peer)
78 .map_err(HoprProtocolError::resolver)?
79 .ok_or_else(|| HoprProtocolError::ChannelNotFound(*self.chain_key.as_ref(), next_peer))?;
80
81 let (outgoing_ticket_win_prob, outgoing_ticket_price) = self
82 .chain_api
83 .outgoing_ticket_values(self.cfg.outgoing_win_prob, self.cfg.outgoing_ticket_price)
84 .map_err(HoprProtocolError::resolver)?;
85
86 self.ticket_factory
87 .new_multihop_ticket(
88 &channel,
89 (num_hops as u8).try_into().expect("cannot fail due to num_hops > 1"),
90 outgoing_ticket_win_prob,
91 outgoing_ticket_price,
92 )
93 .map_err(HoprProtocolError::ticket_factory)?
94 } else {
95 TicketBuilder::zero_hop().counterparty(next_peer)
96 };
97
98 let (packet, openers) = HoprPacket::into_outgoing(
100 data.as_ref(),
101 &pseudonym,
102 routing,
103 &self.chain_key,
104 next_ticket,
105 self.chain_api.key_id_mapper_ref(),
106 &self.channels_dst,
107 signals,
108 )?;
109
110 let mut minted_surbs = Vec::with_capacity(openers.len());
113 openers.into_iter().for_each(|(surb_id, opener)| {
114 minted_surbs.push(surb_id);
115 self.surb_store
116 .insert_reply_opener(HoprSenderId::from_pseudonym_and_id(&pseudonym, surb_id), opener);
117 });
118
119 let out = packet.try_as_outgoing().ok_or(HoprProtocolError::InvalidState(
120 "cannot send out packet that is not outgoing",
121 ))?;
122
123 let mut transport_payload = BytesMut::with_capacity(HoprPacket::SIZE);
124 transport_payload.put_slice(out.packet.as_ref());
125 transport_payload.put_slice(&out.ticket.into_encoded());
126
127 Ok(OutgoingPacket {
128 next_hop: out.next_hop,
129 ack_challenge: out.ack_challenge,
130 data: transport_payload.freeze(),
131 minted_surbs,
132 })
133 }
134}
135
136impl<Chain, S, T> PacketEncoder for HoprEncoder<Chain, S, T>
137where
138 Chain: ChainKeyOperations + ChainReadChannelOperations + ChainReadTicketOperations + ChainValues + Send + Sync,
139 S: SurbStore + Send + Sync,
140 T: hopr_api::tickets::TicketFactory + Send + Sync,
141{
142 type Error = HoprProtocolError;
143
144 #[tracing::instrument(skip_all, level = "trace")]
145 fn encode_packet<D: AsRef<[u8]> + Send + 'static, Sig: Into<PacketSignals> + Send + 'static>(
146 &self,
147 data: D,
148 routing: ResolvedTransportRouting<HoprSurb>,
149 signals: Sig,
150 ) -> Result<OutgoingPacket, Self::Error> {
151 let (next_peer, num_hops, pseudonym, routing) = match routing {
153 ResolvedTransportRouting::Forward {
154 pseudonym,
155 forward_path,
156 return_paths,
157 } => (
158 forward_path[0],
159 forward_path.num_hops(),
160 pseudonym,
161 PacketRouting::ForwardPath {
162 forward_path,
163 return_paths,
164 },
165 ),
166 ResolvedTransportRouting::Return(sender_id, surb) => {
167 let next = self
168 .chain_api
169 .key_id_mapper_ref()
170 .map_id_to_public(&surb.first_relayer)
171 .ok_or(HoprProtocolError::KeyNotFound)?;
172
173 (
174 next,
175 surb.additional_data_receiver.proof_of_relay_values().chain_length() as usize,
176 sender_id.pseudonym(),
177 PacketRouting::Surb(sender_id.surb_id(), surb),
178 )
179 }
180 };
181
182 tracing::trace!(len = data.as_ref().len(), "encoding packet");
183 self.encode_packet_internal(next_peer, data, num_hops, signals, routing, pseudonym)
184 }
185
186 #[tracing::instrument(skip_all, level = "trace", fields(destination = destination.to_peerid_str()))]
187 fn encode_acknowledgements(
188 &self,
189 acks: &[VerifiedAcknowledgement],
190 destination: &OffchainPublicKey,
191 ) -> Result<OutgoingPacket, Self::Error> {
192 tracing::trace!(num_acks = acks.len(), "encoding acknowledgements");
193
194 let mut all_acks = Vec::<u8>::with_capacity(size_of::<u16>() + acks.len() * Acknowledgement::SIZE);
195 all_acks.extend((acks.len() as u16).to_be_bytes());
196 acks.iter().for_each(|ack| all_acks.extend(ack.leak().as_ref()));
197
198 self.encode_packet_internal(
199 *destination,
200 all_acks,
201 0,
202 None,
203 PacketRouting::NoAck(*destination),
204 HoprPseudonym::random(),
205 )
206 }
207}