subspace_networking/protocols/request_response/handlers/
generic_request_handler.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
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
//! Generic request-response handler, typically is used with a type implementing [`GenericRequest`]
//! to significantly reduce boilerplate when implementing [`RequestHandler`].

use crate::protocols::request_response::request_response_factory::{
    IncomingRequest, OutgoingResponse, ProtocolConfig, RequestHandler,
};
use async_trait::async_trait;
use futures::channel::mpsc;
use futures::prelude::*;
use libp2p::PeerId;
use parity_scale_codec::{Decode, Encode};
use std::pin::Pin;
use std::sync::Arc;
use tracing::{debug, trace};

/// Could be changed after the production feedback.
const REQUESTS_BUFFER_SIZE: usize = 50;

/// Generic request with associated response
pub trait GenericRequest: Encode + Decode + Send + Sync + 'static {
    /// Defines request-response protocol name.
    const PROTOCOL_NAME: &'static str;
    /// Specifies log-parameters for tracing.
    const LOG_TARGET: &'static str;
    /// Response type that corresponds to this request
    type Response: Encode + Decode + Send + Sync + 'static;
}

type RequestHandlerFn<Request> = Arc<
    dyn (Fn(
            PeerId,
            Request,
        )
            -> Pin<Box<dyn Future<Output = Option<<Request as GenericRequest>::Response>> + Send>>)
        + Send
        + Sync
        + 'static,
>;

/// Defines generic request-response protocol handler.
pub struct GenericRequestHandler<Request: GenericRequest> {
    request_receiver: mpsc::Receiver<IncomingRequest>,
    request_handler: RequestHandlerFn<Request>,
    protocol_config: ProtocolConfig,
}

impl<Request: GenericRequest> GenericRequestHandler<Request> {
    /// Creates new [`GenericRequestHandler`] by given handler.
    pub fn create<RH, Fut>(request_handler: RH) -> Box<dyn RequestHandler>
    where
        RH: (Fn(PeerId, Request) -> Fut) + Send + Sync + 'static,
        Fut: Future<Output = Option<Request::Response>> + Send + 'static,
    {
        let (request_sender, request_receiver) = mpsc::channel(REQUESTS_BUFFER_SIZE);

        let mut protocol_config = ProtocolConfig::new(Request::PROTOCOL_NAME);
        protocol_config.inbound_queue = Some(request_sender);

        Box::new(Self {
            request_receiver,
            request_handler: Arc::new(move |peer_id, request| {
                Box::pin(request_handler(peer_id, request))
            }),
            protocol_config,
        })
    }

    /// Invokes external protocol handler.
    async fn handle_request(
        &mut self,
        peer: PeerId,
        payload: Vec<u8>,
    ) -> Result<Vec<u8>, RequestHandlerError> {
        trace!(%peer, protocol=Request::LOG_TARGET, "Handling request...");
        let request = Request::decode(&mut payload.as_slice())
            .map_err(|_| RequestHandlerError::InvalidRequestFormat)?;
        let response = (self.request_handler)(peer, request).await;

        Ok(response.ok_or(RequestHandlerError::NoResponse)?.encode())
    }
}

#[async_trait]
impl<Request: GenericRequest> RequestHandler for GenericRequestHandler<Request> {
    /// Run [`RequestHandler`].
    async fn run(&mut self) {
        while let Some(request) = self.request_receiver.next().await {
            let IncomingRequest {
                peer,
                payload,
                pending_response,
            } = request;

            match self.handle_request(peer, payload).await {
                Ok(response_data) => {
                    let response = OutgoingResponse {
                        result: Ok(response_data),
                        sent_feedback: None,
                    };

                    match pending_response.send(response) {
                        Ok(()) => trace!(target = Request::LOG_TARGET, %peer, "Handled request",),
                        Err(_) => debug!(
                            target = Request::LOG_TARGET,
                            protocol = Request::PROTOCOL_NAME,
                            %peer,
                            "Failed to handle request: {}",
                            RequestHandlerError::SendResponse
                        ),
                    };
                }
                Err(e) => {
                    debug!(
                        target = Request::LOG_TARGET,
                        protocol = Request::PROTOCOL_NAME,
                        %e,
                        "Failed to handle request.",
                    );

                    let response = OutgoingResponse {
                        result: Err(()),
                        sent_feedback: None,
                    };

                    if pending_response.send(response).is_err() {
                        debug!(
                            target = Request::LOG_TARGET,
                            protocol = Request::PROTOCOL_NAME,
                            %peer,
                            "Failed to handle request: {}", RequestHandlerError::SendResponse
                        );
                    };
                }
            }
        }
    }

    fn protocol_config(&self) -> ProtocolConfig {
        self.protocol_config.clone()
    }

    fn protocol_name(&self) -> &'static str {
        Request::PROTOCOL_NAME
    }

    fn clone_box(&self) -> Box<dyn RequestHandler> {
        let (request_sender, request_receiver) = mpsc::channel(REQUESTS_BUFFER_SIZE);

        let mut protocol_config = ProtocolConfig::new(Request::PROTOCOL_NAME);
        protocol_config.inbound_queue = Some(request_sender);

        Box::new(Self {
            request_receiver,
            request_handler: Arc::clone(&self.request_handler),
            protocol_config,
        })
    }
}

#[derive(Debug, thiserror::Error)]
enum RequestHandlerError {
    #[error("Failed to send response.")]
    SendResponse,

    #[error("Incorrect request format.")]
    InvalidRequestFormat,

    #[error("No response.")]
    NoResponse,
}