Send a request for data over bitswap
(
&mut self,
ctx: u64,
cid: Cid,
providers: HashSet<PeerId>,
mut chan: OneShotSender<Result<Block, String>>,
)
| 442 | |
| 443 | /// Send a request for data over bitswap |
| 444 | fn want_block( |
| 445 | &mut self, |
| 446 | ctx: u64, |
| 447 | cid: Cid, |
| 448 | providers: HashSet<PeerId>, |
| 449 | mut chan: OneShotSender<Result<Block, String>>, |
| 450 | ) -> Result<()> { |
| 451 | if let Some(bs) = self.swarm.behaviour().bitswap.as_ref() { |
| 452 | let client = bs.client().clone(); |
| 453 | let (closer_s, closer_r) = oneshot::channel(); |
| 454 | |
| 455 | let entry = self.bitswap_sessions.entry(ctx).or_default(); |
| 456 | |
| 457 | let providers: Vec<_> = providers.into_iter().collect(); |
| 458 | let worker = tokio::task::spawn(async move { |
| 459 | tokio::select! { |
| 460 | _ = closer_r => { |
| 461 | // Explicit session stop. |
| 462 | debug!("session {}: stopped: closed", ctx); |
| 463 | } |
| 464 | _ = chan.closed() => { |
| 465 | // RPC dropped |
| 466 | debug!("session {}: stopped: request canceled", ctx); |
| 467 | } |
| 468 | block = client.get_block_with_session_id(ctx, &cid, &providers) => match block { |
| 469 | Ok(block) => { |
| 470 | if let Err(e) = chan.send(Ok(block)) { |
| 471 | warn!("failed to send block response: {:?}", e); |
| 472 | } |
| 473 | } |
| 474 | Err(err) => { |
| 475 | chan.send(Err(err.to_string())).ok(); |
| 476 | } |
| 477 | }, |
| 478 | } |
| 479 | }); |
| 480 | entry.push((closer_s, worker)); |
| 481 | |
| 482 | Ok(()) |
| 483 | } else { |
| 484 | bail!("no bitswap available"); |
| 485 | } |
| 486 | } |
| 487 | |
| 488 | // TODO fix skip_all |
| 489 | #[tracing::instrument(skip_all)] |