1use std::collections::HashSet;
7use std::num::NonZeroUsize;
8use std::str::FromStr;
9
10use async_trait::async_trait;
11use futures::future::join_all;
12use jsonrpc_core::types::error::Error;
13use jsonrpc_core::types::error::ErrorCode;
14use jsonrpc_core::Result;
15use rings_core::dht::entry::Entry;
16use rings_core::dht::Did;
17use rings_core::ecc::PublicKey;
18use rings_core::message::e2e;
19use rings_core::message::Decoder;
20use rings_core::message::Encoded;
21use rings_core::message::Encoder;
22use rings_core::message::MessagePayload;
23use rings_core::message::MessageVerificationExt;
24use rings_rpc::protos::rings_node::*;
25use rings_rpc::protos::rings_node_handler::HandleRpc;
26
27use crate::error::Error as ServerError;
28use crate::processor::Processor;
29use crate::seed::Seed;
30
31const DEFAULT_PEER_MEASUREMENT_PAGE_SIZE: u32 = 100;
32const MAX_PEER_MEASUREMENT_PAGE_SIZE: u32 = 1_000;
33
34#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
35#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
36impl HandleRpc<ConnectPeerViaHttpRequest, ConnectPeerViaHttpResponse> for Processor {
37 async fn handle_rpc(
38 &self,
39 req: ConnectPeerViaHttpRequest,
40 ) -> Result<ConnectPeerViaHttpResponse> {
41 let client = rings_rpc::jsonrpc::Client::new(&req.url);
42
43 let did = client
44 .node_did(&NodeDidRequest {})
45 .await
46 .map_err(|e| ServerError::RemoteRpcError(e.to_string()))?
47 .did;
48
49 let offer = self
50 .handle_rpc(CreateOfferRequest { did: did.clone() })
51 .await?
52 .offer;
53
54 let answer = client
55 .answer_offer(&AnswerOfferRequest { offer })
56 .await
57 .map_err(|e| ServerError::RemoteRpcError(e.to_string()))?
58 .answer;
59
60 self.handle_rpc(AcceptAnswerRequest { answer }).await?;
61
62 Ok(ConnectPeerViaHttpResponse { did })
63 }
64}
65
66#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
67#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
68impl HandleRpc<ConnectWithDidRequest, ConnectWithDidResponse> for Processor {
69 async fn handle_rpc(&self, req: ConnectWithDidRequest) -> Result<ConnectWithDidResponse> {
70 let did = s2d(&req.did)?;
71 self.connect_with_did(did).await.map_err(Error::from)?;
72 Ok(ConnectWithDidResponse {})
73 }
74}
75
76#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
77#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
78impl HandleRpc<ConnectWithSeedRequest, ConnectWithSeedResponse> for Processor {
79 async fn handle_rpc(&self, req: ConnectWithSeedRequest) -> Result<ConnectWithSeedResponse> {
80 let seed: Seed = Seed::try_from(req)?;
81
82 let mut connected: HashSet<String> =
83 HashSet::from_iter(self.swarm.peers().into_iter().map(|peer| peer.did));
84 connected.insert(self.swarm.did().to_string());
85
86 let tasks = seed
87 .peers
88 .iter()
89 .filter(|&x| !connected.contains(&x.did))
90 .map(|x| {
91 self.handle_rpc(ConnectPeerViaHttpRequest {
92 url: x.url.to_string(),
93 })
94 });
95
96 let results = join_all(tasks).await;
97
98 let first_err = results.into_iter().find(|x| x.is_err());
99 if let Some(err) = first_err {
100 err?;
101 }
102
103 Ok(ConnectWithSeedResponse {})
104 }
105}
106
107#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
108#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
109impl HandleRpc<ListPeersRequest, ListPeersResponse> for Processor {
110 async fn handle_rpc(&self, _req: ListPeersRequest) -> Result<ListPeersResponse> {
111 let peers = self
112 .swarm
113 .peers()
114 .into_iter()
115 .map(|peer| peer.into())
116 .collect();
117 Ok(ListPeersResponse { peers })
118 }
119}
120
121#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
122#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
123impl HandleRpc<CreateOfferRequest, CreateOfferResponse> for Processor {
124 async fn handle_rpc(&self, req: CreateOfferRequest) -> Result<CreateOfferResponse> {
125 let did = s2d(&req.did)?;
126 let offer_payload = self
127 .swarm
128 .create_offer(did)
129 .await
130 .map_err(ServerError::CreateOffer)
131 .map_err(Error::from)?;
132
133 let encoded = offer_payload
134 .encode()
135 .map_err(|_| ServerError::EncodeError)?;
136
137 Ok(CreateOfferResponse {
138 offer: encoded.to_string(),
139 })
140 }
141}
142
143#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
144#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
145impl HandleRpc<AnswerOfferRequest, AnswerOfferResponse> for Processor {
146 async fn handle_rpc(&self, req: AnswerOfferRequest) -> Result<AnswerOfferResponse> {
147 if req.offer.is_empty() {
148 return Err(Error::invalid_params("Offer is empty"));
149 }
150 let encoded: Encoded = <Encoded as From<String>>::from(req.offer);
151
152 let offer_payload =
153 MessagePayload::from_encoded(&encoded).map_err(|_| ServerError::DecodeError)?;
154
155 let answer_payload = self
156 .swarm
157 .answer_offer(offer_payload)
158 .await
159 .map_err(ServerError::AnswerOffer)
160 .map_err(Error::from)?;
161
162 tracing::debug!("connect_peer_via_ice response: {:?}", answer_payload);
163 let encoded = answer_payload
164 .encode()
165 .map_err(|_| ServerError::EncodeError)?;
166
167 Ok(AnswerOfferResponse {
168 answer: encoded.to_string(),
169 })
170 }
171}
172
173#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
174#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
175impl HandleRpc<AcceptAnswerRequest, AcceptAnswerResponse> for Processor {
176 async fn handle_rpc(&self, req: AcceptAnswerRequest) -> Result<AcceptAnswerResponse> {
177 if req.answer.is_empty() {
178 return Err(Error::invalid_params("Answer is empty"));
179 }
180 let encoded = Encoded::from(req.answer);
181
182 let answer_payload =
183 MessagePayload::from_encoded(&encoded).map_err(|_| ServerError::DecodeError)?;
184 answer_payload.transaction.signer();
185
186 self.swarm
187 .accept_answer(answer_payload)
188 .await
189 .map_err(ServerError::AcceptAnswer)?;
190
191 Ok(AcceptAnswerResponse {})
192 }
193}
194
195#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
196#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
197impl HandleRpc<DisconnectRequest, DisconnectResponse> for Processor {
198 async fn handle_rpc(&self, req: DisconnectRequest) -> Result<DisconnectResponse> {
199 let did = s2d(&req.did)?;
200 self.disconnect(did).await?;
201 Ok(DisconnectResponse {})
202 }
203}
204
205#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
206#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
207impl HandleRpc<SendBackendMessageRequest, SendBackendMessageResponse> for Processor {
208 async fn handle_rpc(
209 &self,
210 req: SendBackendMessageRequest,
211 ) -> Result<SendBackendMessageResponse> {
212 let destination = s2d(&req.destination_did)?;
213 let payload = base64::decode(req.data.as_str())
214 .map_err(|e| Error::invalid_params(format!("data is not valid base64: {e:?}")))?;
215 let envelope =
216 crate::extension::ext::Envelope::new(req.namespace, bytes::Bytes::from(payload));
217 self.send_envelope(destination, &envelope).await?;
218 Ok(SendBackendMessageResponse {})
219 }
220}
221
222#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
223#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
224impl HandleRpc<SendE2eHandshakeRequest, SendE2eHandshakeResponse> for Processor {
225 async fn handle_rpc(&self, req: SendE2eHandshakeRequest) -> Result<SendE2eHandshakeResponse> {
226 let destination = s2d(&req.destination_did)?;
227 let tx_id = self.send_e2e_handshake(destination).await?;
228 Ok(SendE2eHandshakeResponse {
229 tx_id: tx_id.to_string(),
230 })
231 }
232}
233
234#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
235#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
236impl HandleRpc<SendE2eMessageRequest, SendE2eMessageResponse> for Processor {
237 async fn handle_rpc(&self, req: SendE2eMessageRequest) -> Result<SendE2eMessageResponse> {
238 let destination = s2d(&req.destination_did)?;
239 let recipient_public_key = s2pk(&req.recipient_public_key)?;
240 let payload = base64::decode(req.data.as_str())
241 .map_err(|e| Error::invalid_params(format!("data is not valid base64: {e:?}")))?;
242 let frame_len = if req.max_plaintext_frame_len == 0 {
243 e2e::DEFAULT_E2E_PLAINTEXT_FRAME_LEN
244 } else {
245 req.max_plaintext_frame_len as usize
246 };
247
248 let stream_id = self
249 .send_e2e_message_with_frame_len(destination, recipient_public_key, &payload, frame_len)
250 .await?;
251 Ok(SendE2eMessageResponse {
252 stream_id: stream_id.to_string(),
253 })
254 }
255}
256
257#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
258#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
259impl HandleRpc<PublishMessageToTopicRequest, PublishMessageToTopicResponse> for Processor {
260 async fn handle_rpc(
261 &self,
262 req: PublishMessageToTopicRequest,
263 ) -> Result<PublishMessageToTopicResponse> {
264 let encoded = req
265 .data
266 .encode()
267 .map_err(|e| Error::invalid_params(format!("Failed to encode data: {e:?}")))?;
268 self.storage_append_data(&req.topic, encoded).await?;
269 Ok(PublishMessageToTopicResponse {})
270 }
271}
272
273#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
274#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
275impl HandleRpc<FetchTopicMessagesRequest, FetchTopicMessagesResponse> for Processor {
276 async fn handle_rpc(
277 &self,
278 req: FetchTopicMessagesRequest,
279 ) -> Result<FetchTopicMessagesResponse> {
280 let entry_key = Entry::gen_did(&req.topic)
281 .map_err(|_| Error::invalid_params("Failed to get id of topic"))?;
282
283 self.storage_fetch(entry_key).await?;
284 let result = self.storage_check_cache(entry_key).await;
285
286 let Some(entry) = result else {
287 return Ok(FetchTopicMessagesResponse { data: vec![] });
288 };
289
290 let data = entry
291 .data
292 .iter()
293 .skip(req.skip as usize)
294 .map(|v| v.decode())
295 .filter_map(|v| v.ok())
296 .collect::<Vec<String>>();
297
298 Ok(FetchTopicMessagesResponse { data })
299 }
300}
301
302#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
303#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
304impl HandleRpc<RegisterServiceRequest, RegisterServiceResponse> for Processor {
305 async fn handle_rpc(&self, req: RegisterServiceRequest) -> Result<RegisterServiceResponse> {
306 self.register_service(&req.name).await?;
307 Ok(RegisterServiceResponse {})
308 }
309}
310
311#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
312#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
313impl HandleRpc<LookupServiceRequest, LookupServiceResponse> for Processor {
314 async fn handle_rpc(&self, req: LookupServiceRequest) -> Result<LookupServiceResponse> {
315 let entry_key = Entry::gen_did(&req.name)
316 .map_err(|_| Error::invalid_params("Failed to get id of topic"))?;
317
318 self.storage_fetch(entry_key).await?;
319 let result = self.storage_check_cache(entry_key).await;
320
321 let Some(entry) = result else {
322 return Ok(LookupServiceResponse { dids: vec![] });
323 };
324
325 let dids = entry
326 .data
327 .iter()
328 .map(|v| v.decode())
329 .filter_map(|v| v.ok())
330 .collect::<Vec<String>>();
331
332 Ok(LookupServiceResponse { dids })
333 }
334}
335
336#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
337#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
338impl HandleRpc<LookupOnlineNodesRequest, LookupOnlineNodesResponse> for Processor {
339 async fn handle_rpc(&self, req: LookupOnlineNodesRequest) -> Result<LookupOnlineNodesResponse> {
340 let nodes = self
341 .lookup_online_nodes(req.include_expired)
342 .await
343 .map_err(Error::from)?;
344 Ok(LookupOnlineNodesResponse {
345 nodes: crate::rpc_dto::online_node_descriptor_infos(nodes)?,
346 })
347 }
348}
349
350#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
351#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
352impl HandleRpc<LookupOnionExitsRequest, LookupOnionExitsResponse> for Processor {
353 async fn handle_rpc(&self, req: LookupOnionExitsRequest) -> Result<LookupOnionExitsResponse> {
354 let exits = self
355 .lookup_onion_exits(&req.service, req.include_expired)
356 .await
357 .map_err(Error::from)?;
358 Ok(LookupOnionExitsResponse {
359 exits: crate::rpc_dto::onion_exit_descriptor_infos(exits)?,
360 })
361 }
362}
363
364#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
365#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
366impl HandleRpc<BuildOnionRouteRequest, BuildOnionRouteResponse> for Processor {
367 async fn handle_rpc(&self, req: BuildOnionRouteRequest) -> Result<BuildOnionRouteResponse> {
368 let route = self
369 .build_onion_route(req.service, req.hop_count as usize, req.allow_short_paths)
370 .await
371 .map_err(Error::from)?;
372 crate::rpc_dto::onion_route_response(route).map_err(Error::from)
373 }
374}
375
376#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
377#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
378impl HandleRpc<NodeInfoRequest, NodeInfoResponse> for Processor {
379 async fn handle_rpc(&self, _req: NodeInfoRequest) -> Result<NodeInfoResponse> {
380 self.get_node_info()
381 .await
382 .map_err(|_| Error::new(ErrorCode::InternalError))
383 }
384}
385
386#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
387#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
388impl HandleRpc<PeerMeasurementRequest, PeerMeasurementResponse> for Processor {
389 async fn handle_rpc(&self, req: PeerMeasurementRequest) -> Result<PeerMeasurementResponse> {
390 let did = s2d(&req.did)?;
391 Ok(PeerMeasurementResponse {
392 measurement: crate::rpc_dto::optional_peer_measurement_info(
393 self.peer_measurement(did).await,
394 )?,
395 })
396 }
397}
398
399#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
400#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
401impl HandleRpc<ListPeerMeasurementsRequest, ListPeerMeasurementsResponse> for Processor {
402 async fn handle_rpc(
403 &self,
404 req: ListPeerMeasurementsRequest,
405 ) -> Result<ListPeerMeasurementsResponse> {
406 let (after, limit) = peer_measurement_page_request(&req)?;
407 let page = self.peer_measurements_page(after, limit).await;
408 Ok(ListPeerMeasurementsResponse {
409 measurements: crate::rpc_dto::peer_measurement_infos(page.measurements)?,
410 next_cursor: page.next_cursor.map(|did| did.to_string()),
411 })
412 }
413}
414
415fn peer_measurement_page_request(
416 request: &ListPeerMeasurementsRequest,
417) -> Result<(Option<Did>, NonZeroUsize)> {
418 let requested = request.limit.unwrap_or(DEFAULT_PEER_MEASUREMENT_PAGE_SIZE);
419 let limit = usize::try_from(requested)
420 .ok()
421 .and_then(NonZeroUsize::new)
422 .filter(|limit| limit.get() <= MAX_PEER_MEASUREMENT_PAGE_SIZE as usize)
423 .ok_or_else(|| {
424 Error::invalid_params(format!(
425 "measurement page limit must be between 1 and {MAX_PEER_MEASUREMENT_PAGE_SIZE}"
426 ))
427 })?;
428 let after = request.cursor.as_deref().map(s2d).transpose()?;
429 Ok((after, limit))
430}
431
432#[cfg_attr(all(feature = "browser", target_family = "wasm"), async_trait(?Send))]
433#[cfg_attr(not(all(feature = "browser", target_family = "wasm")), async_trait)]
434impl HandleRpc<NodeDidRequest, NodeDidResponse> for Processor {
435 async fn handle_rpc(&self, _req: NodeDidRequest) -> Result<NodeDidResponse> {
436 let did = self.did();
437 Ok(NodeDidResponse {
438 did: did.to_string(),
439 })
440 }
441}
442
443fn s2d(s: &str) -> Result<Did> {
445 Did::from_str(s).map_err(|_| Error::invalid_params(format!("Invalid Did: {s}")))
446}
447
448fn s2pk(s: &str) -> Result<PublicKey<33>> {
449 PublicKey::try_from_b58m(s)
450 .or_else(|_| PublicKey::from_hex_string(s))
451 .map_err(|_| Error::invalid_params("Invalid secp256k1 public key"))
452}
453
454#[cfg(test)]
455mod tests {
456 use super::*;
457
458 #[test]
459 fn peer_measurement_page_request_applies_defaults_and_parses_cursor() {
460 let request = ListPeerMeasurementsRequest {
461 limit: None,
462 cursor: Some(Did::from(2_u32).to_string()),
463 };
464 let (after, limit) = peer_measurement_page_request(&request)
465 .unwrap_or_else(|error| panic!("bounded page request must parse: {error}"));
466 assert_eq!(after, Some(Did::from(2_u32)));
467 assert_eq!(limit.get(), DEFAULT_PEER_MEASUREMENT_PAGE_SIZE as usize);
468 }
469
470 #[test]
471 fn peer_measurement_page_rejects_unbounded_or_zero_limit() {
472 let too_large = ListPeerMeasurementsRequest {
473 limit: Some(MAX_PEER_MEASUREMENT_PAGE_SIZE + 1),
474 cursor: None,
475 };
476 let zero = ListPeerMeasurementsRequest {
477 limit: Some(0),
478 cursor: None,
479 };
480 assert!(peer_measurement_page_request(&too_large).is_err());
481 assert!(peer_measurement_page_request(&zero).is_err());
482 }
483}