Skip to main content

tuwunel_service/sending/sender/dispatch/
mod.rs

1mod appservice;
2mod federation;
3mod push;
4
5use futures::{FutureExt, future::BoxFuture};
6use tuwunel_core::{Error, Result, implement};
7
8use super::split::Split;
9use crate::sending::{
10	Destination, Service,
11	data::{Keys, QueueItem},
12};
13
14pub(super) type SendingError = (Destination, Error);
15pub(super) type SendingResult = Result<Destination, SendingError>;
16pub(super) type SendingFuture<'a> = BoxFuture<'a, Completion>;
17
18/// A finished transaction with the durable rows it carried.
19///
20/// Success acknowledges exactly these keys, leaving any other active row of
21/// the destination in place. A transaction sending one room of a rejected
22/// batch carries the split it advances. One whose rows all failed to load sent
23/// nothing, which proves nothing about the peer, so it ends the split instead.
24pub(super) struct Completion {
25	pub(super) result: SendingResult,
26	pub(super) keys: Keys,
27	pub(super) split: Option<Split>,
28}
29
30#[implement(Service)]
31pub(super) fn send_events(
32	&self,
33	dest: Destination,
34	items: Vec<QueueItem>,
35	split: Option<Split>,
36) -> SendingFuture<'_> {
37	debug_assert!(!items.is_empty(), "sending empty transaction");
38
39	let (keys, events): (Keys, Vec<_>) = items.into_iter().unzip();
40	let complete = |result, split| Completion { result, keys, split };
41
42	match dest {
43		| Destination::Federation(server) => self
44			.send_events_dest_federation(server, events)
45			.map(|(result, sent)| complete(result, split.filter(|_| sent)))
46			.boxed(), // heterogeneous SendingFutures
47		| Destination::Appservice(id) => self
48			.send_events_dest_appservice(id, events)
49			.map(|result| complete(result, split))
50			.boxed(),
51		| Destination::Push(user_id, pushkey) => self
52			.send_events_dest_push(user_id, pushkey, events)
53			.map(|result| complete(result, split))
54			.boxed(),
55	}
56}