tuwunel_service/fetcher/
mod.rs1mod error;
9mod inflight;
10mod opts;
11mod select;
12mod transport;
13mod validate;
14mod worker;
15
16#[cfg(test)]
17mod tests;
18
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use loole::{Receiver, Sender, unbounded};
23use tokio::sync::{
24 oneshot::{self, channel},
25 watch,
26};
27use tuwunel_core::{Result, implement};
28
29pub use self::opts::{EventWindow, FanoutGrowth, Op, Opts, Outcome};
34use self::{
35 error::Failure,
36 inflight::{Key, SharedResult, Subscription},
37 select::{RoomCandidates, Select},
38 transport::{FederationTransport, Transport},
39};
40use crate::services::OnceServices;
41
42const REQUESTS_MAX: usize = 256;
44
45pub struct Service {
51 services: Arc<OnceServices>,
52 channel: (Sender<Msg>, Receiver<Msg>),
53 transport: Arc<dyn Transport>,
54 select: Arc<dyn Select>,
55 capacity: usize,
56}
57
58struct Msg {
61 key: Key,
62 reply: oneshot::Sender<Subscription>,
63}
64
65#[async_trait]
66impl crate::Service for Service {
67 fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
68 let services = args.services.clone();
69 let transport: Arc<dyn Transport> =
70 Arc::new(FederationTransport { services: services.clone() });
71
72 let select: Arc<dyn Select> = Arc::new(RoomCandidates { services: services.clone() });
73
74 Ok(Arc::new(Self {
75 services,
76 channel: unbounded(),
77 transport,
78 select,
79 capacity: REQUESTS_MAX,
80 }))
81 }
82
83 async fn worker(self: Arc<Self>) -> Result {
84 self.run_worker().await;
85 Ok(())
86 }
87
88 async fn interrupt(&self) {
89 let (sender, _) = &self.channel;
90 if !sender.is_closed() {
91 sender.close();
92 }
93 }
94
95 fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
96}
97
98#[implement(Service)]
105#[tracing::instrument(
106 level = "debug",
107 skip_all,
108 fields(
109 op = ?opts.op,
110 room_id = ?opts.room_id,
111 event_id = ?opts.event_id,
112 ),
113)]
114pub async fn fetch(&self, opts: Opts) -> Result<Arc<Outcome>> {
115 let key = Key::new(opts);
116 let (reply, reply_rx) = channel();
117
118 self.channel
119 .0
120 .send(Msg { key, reply })
121 .map_err(|_| Failure::Cancelled)?;
122
123 let (rx, _interest) = reply_rx.await.map_err(|_| Failure::Cancelled)?;
125
126 await_result(rx).await.map_err(Into::into)
127}
128
129async fn await_result(mut rx: watch::Receiver<Option<SharedResult>>) -> SharedResult {
130 rx.wait_for(Option::is_some)
131 .await
132 .map_or(Err(Failure::Cancelled), |value| value.clone().expect("present by predicate"))
133}