| 413 | } |
| 414 | |
| 415 | fn destroy_session(&mut self, ctx: u64, response_channel: oneshot::Sender<Result<()>>) { |
| 416 | if let Some(bs) = self.swarm.behaviour().bitswap.as_ref() { |
| 417 | let workers = self.bitswap_sessions.remove(&ctx); |
| 418 | let client = bs.client().clone(); |
| 419 | tokio::task::spawn(async move { |
| 420 | debug!("stopping session {}", ctx); |
| 421 | if let Some(workers) = workers { |
| 422 | debug!("stopping workers {} for session {}", workers.len(), ctx); |
| 423 | // first shutdown workers |
| 424 | for (closer, worker) in workers { |
| 425 | if closer.send(()).is_ok() { |
| 426 | worker.await.ok(); |
| 427 | } |
| 428 | } |
| 429 | debug!("all workers stopped for session {}", ctx); |
| 430 | } |
| 431 | if let Err(err) = client.stop_session(ctx).await { |
| 432 | warn!("failed to stop session {}: {:?}", ctx, err); |
| 433 | } |
| 434 | // Ignore error if the otherside already hung up. |
| 435 | let _ = response_channel.send(Ok(())); |
| 436 | debug!("session {} stopped", ctx); |
| 437 | }); |
| 438 | } else { |
| 439 | let _ = response_channel.send(Err(anyhow!("no bitswap available"))); |
| 440 | } |
| 441 | } |
| 442 | |
| 443 | /// Send a request for data over bitswap |
| 444 | fn want_block( |