polkadot_node_network_protocol/request_response/
outgoing.rs1use futures::{channel::oneshot, prelude::Future, FutureExt};
18
19use codec::{Decode, Encode, Error as DecodingError};
20use network::ProtocolName;
21
22use sc_network as network;
23use sc_network_types::PeerId;
24
25use polkadot_primitives::AuthorityDiscoveryId;
26
27use super::{v1, v2, v3, IsRequest, Protocol};
28
29#[derive(Debug)]
31pub enum Requests {
32 ChunkFetching(OutgoingRequest<v2::ChunkFetchingRequest, v1::ChunkFetchingRequest>),
34 CollationFetchingV1(OutgoingRequest<v1::CollationFetchingRequest>),
36 PoVFetchingV1(OutgoingRequest<v1::PoVFetchingRequest>),
38 AvailableDataFetchingV1(OutgoingRequest<v1::AvailableDataFetchingRequest>),
40 DisputeSendingV1(OutgoingRequest<v1::DisputeRequest>),
42
43 AttestedCandidateV2(OutgoingRequest<v2::AttestedCandidateRequest>),
45 CollationFetchingV2(OutgoingRequest<v2::CollationFetchingRequest>),
48 CollationFetchingV3(OutgoingRequest<v3::CollationFetchingRequest>),
52}
53
54impl Requests {
55 pub fn encode_request(self) -> (Protocol, OutgoingRequest<Vec<u8>>) {
63 match self {
64 Self::ChunkFetching(r) => r.encode_request(),
65 Self::CollationFetchingV1(r) => r.encode_request(),
66 Self::CollationFetchingV2(r) => r.encode_request(),
67 Self::CollationFetchingV3(r) => r.encode_request(),
68 Self::PoVFetchingV1(r) => r.encode_request(),
69 Self::AvailableDataFetchingV1(r) => r.encode_request(),
70 Self::DisputeSendingV1(r) => r.encode_request(),
71 Self::AttestedCandidateV2(r) => r.encode_request(),
72 }
73 }
74}
75
76pub type ResponseSender = oneshot::Sender<Result<(Vec<u8>, ProtocolName), network::RequestFailure>>;
78
79#[derive(Debug, thiserror::Error)]
81pub enum RequestError {
82 #[error("Response could not be decoded: {0}")]
84 InvalidResponse(#[from] DecodingError),
85
86 #[error("{0}")]
88 NetworkError(#[from] network::RequestFailure),
89
90 #[error("Response channel got canceled")]
92 Canceled(#[from] oneshot::Canceled),
93}
94
95impl RequestError {
96 pub fn is_timed_out(&self) -> bool {
98 match self {
99 Self::Canceled(_) |
100 Self::NetworkError(network::RequestFailure::Obsolete) |
101 Self::NetworkError(network::RequestFailure::Network(
102 network::OutboundFailure::Timeout,
103 )) => true,
104 _ => false,
105 }
106 }
107}
108
109#[derive(Debug)]
120pub struct OutgoingRequest<Req, FallbackReq = Req> {
121 pub peer: Recipient,
123 pub payload: Req,
125 pub fallback_request: Option<(FallbackReq, Protocol)>,
127 pub pending_response: ResponseSender,
129}
130
131#[derive(Debug, Eq, Hash, PartialEq, Clone)]
133pub enum Recipient {
134 Peer(PeerId),
136 Authority(AuthorityDiscoveryId),
138}
139
140pub type OutgoingResult<Res> = Result<Res, RequestError>;
142
143impl<Req, FallbackReq> OutgoingRequest<Req, FallbackReq>
144where
145 Req: IsRequest + Encode,
146 Req::Response: Decode,
147 FallbackReq: IsRequest + Encode,
148 FallbackReq::Response: Decode,
149{
150 pub fn new(
155 peer: Recipient,
156 payload: Req,
157 ) -> (Self, impl Future<Output = OutgoingResult<Req::Response>>) {
158 let (tx, rx) = oneshot::channel();
159 let r = Self { peer, payload, pending_response: tx, fallback_request: None };
160 (r, receive_response::<Req>(rx.map(|r| r.map(|r| r.map(|(resp, _)| resp)))))
161 }
162
163 pub fn new_with_fallback(
170 peer: Recipient,
171 payload: Req,
172 fallback_request: FallbackReq,
173 ) -> (Self, impl Future<Output = OutgoingResult<(Vec<u8>, ProtocolName)>>) {
174 let (tx, rx) = oneshot::channel();
175 let r = Self {
176 peer,
177 payload,
178 pending_response: tx,
179 fallback_request: Some((fallback_request, FallbackReq::PROTOCOL)),
180 };
181 (r, async { Ok(rx.await??) })
182 }
183
184 pub fn encode_request(self) -> (Protocol, OutgoingRequest<Vec<u8>>) {
189 let OutgoingRequest { peer, payload, pending_response, fallback_request } = self;
190 let encoded = OutgoingRequest {
191 peer,
192 payload: payload.encode(),
193 fallback_request: fallback_request.map(|(r, p)| (r.encode(), p)),
194 pending_response,
195 };
196 (Req::PROTOCOL, encoded)
197 }
198}
199
200async fn receive_response<Req>(
202 rec: impl Future<Output = Result<Result<Vec<u8>, network::RequestFailure>, oneshot::Canceled>>,
203) -> OutgoingResult<Req::Response>
204where
205 Req: IsRequest,
206 Req::Response: Decode,
207{
208 let raw = rec.await??;
209 Ok(Decode::decode(&mut raw.as_ref())?)
210}